2026'da Apache Flink: Akis Isleme, Event Time ve Mulakat Sorulari

Apache Flink 2.3 ile event time semantigi, watermark'lar ve pencereleme konularinda kapsamli rehber. Veri muhendisligi mulakatlarina hazirlik.

Event time ve pencereleme ile Apache Flink akis isleme mimarisi

Apache Flink 2.3, exactly-once garantisi ve saniyenin altinda gecikme ile buyuk olcekli akis islemesini yonetir. Toplu isleme odakli sistemlerin aksine, Flink verileri geldikce surekli olarak isler ve bu da onu gercek zamanli analitik, dolandiricilik tespiti ve olay gudumlu mimariler icin tercih edilen framework yapar.

Mulakat Icin Kritik

Flink, gercek akis isleme ile Spark Streaming'den ayrilir: Flink olaylari event time semantigi ile tek tek islerken, Spark Streaming varsayilan olarak processing time ile mikro-batch'ler isler.

Flink, birden fazla TaskManager arasinda isi koordine eden bir JobManager ile dagitik mimari uzerinde calisir. Her TaskManager, gorev paralel operatorlerinin bolumlerni calistiran gorev yuvalarini calistirir. Bu ayrim, Flink'in dagitik kontrol noktalari araciligiyla hata toleransini korurken yatay olarak olceklenmesini saglar.

Flink'teki veri akisi modeli, hesaplamalari yonlu asiklik graflar (DAG) olarak temsil eder. Veriler kaynaklardan donusumler araciligiyla havuzlara akar ve her operator potansiyel olarak birden fazla paralel ornek uzerinde calisabilir.

FlinkStreamJob.javajava
// Basic Flink streaming application setup
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

// Enable checkpointing every 10 seconds for fault tolerance
env.enableCheckpointing(10000, CheckpointingMode.EXACTLY_ONCE);

// Configure state backend for large state
env.setStateBackend(new EmbeddedRocksDBStateBackend());

// Define the data source - Kafka in production scenarios
DataStream<String> rawStream = env.addSource(
    new FlinkKafkaConsumer<>("events", new SimpleStringSchema(), kafkaProps)
);

// Parse and transform the stream
DataStream<Event> events = rawStream
    .map(json -> objectMapper.readValue(json, Event.class))
    .assignTimestampsAndWatermarks(
        WatermarkStrategy.<Event>forBoundedOutOfOrderness(Duration.ofSeconds(5))
            .withTimestampAssigner((event, timestamp) -> event.getTimestamp())
    );

Yukaridaki kontrol noktasi yapilandirmasi, dagitik durumun anlık goruntusunu her 10 saniyede bir depolar. Bir hata olusursa, Flink son tamamlanan kontrol noktasindan geri yukler ve Kafka'dan kayitlari yeniden oynatir.

Event Time ve Processing Time Semantigi

Event time, bir olayin gercekte ne zaman gerceklestigini ifade eder ve verinin icine gomulmustur. Processing time, Flink'in kaydi isleme zamandir. Bu ayrım onemlidir cunku ag gecikmeleri, sira disi teslim ve isleme birikimleri, zamana dayali islemler icin processing time'i guvenilmez kilar.

EventTimeExample.javajava
// Configuring event time with watermarks
public class EventTimeProcessor {
    
    public DataStream<AggregatedMetric> processWithEventTime(
            DataStream<SensorReading> readings) {
        
        return readings
            // Extract timestamp from the event payload
            .assignTimestampsAndWatermarks(
                WatermarkStrategy
                    .<SensorReading>forBoundedOutOfOrderness(Duration.ofMinutes(2))
                    .withTimestampAssigner((reading, ts) -> reading.getEventTime())
                    .withIdleness(Duration.ofMinutes(5)) // Handle idle partitions
            )
            // Window by event time, not wall clock
            .keyBy(SensorReading::getSensorId)
            .window(TumblingEventTimeWindows.of(Time.minutes(1)))
            .aggregate(new AverageAggregator());
    }
}

forBoundedOutOfOrderness stratejisi, Flink'e olaylarin 2 dakikaya kadar gec gelebilecegini soyler. Watermark'lar, Flink watermark'tan onceki zaman damgalarina sahip daha fazla olayin gelmeyecegini belirlediginde ilerler.

Sik Sorulan Mulakat Sorusu

Flink'te gec gelen olaylara ne olur? Varsayilan olarak, watermark pencerenin bitis zamanini gectikten sonra gelen olaylar atilir. Gec gelenleri islemek icin .allowedLateness(Time.minutes(10)) ile izin verilen gecikmeyi yapilandirin veya ayri isleme icin side output'lari kullanin.

Gercek Zamanli Analitik icin Pencereleme Stratejileri

Flink dort pencere turu saglar: tumbling, sliding, session ve global pencereler. Her biri farkli analitik ihtiyaclara hizmet eder.

WindowingStrategies.javajava
// Different windowing approaches for stream processing
public class WindowingStrategies {

    // Tumbling windows: fixed-size, non-overlapping
    // Use case: hourly aggregations, daily summaries
    public DataStream<Summary> tumblingAggregation(DataStream<Transaction> txns) {
        return txns
            .keyBy(Transaction::getAccountId)
            .window(TumblingEventTimeWindows.of(Time.hours(1)))
            .sum("amount");
    }

    // Sliding windows: fixed-size, overlapping
    // Use case: moving averages, rolling metrics
    public DataStream<Double> slidingAverage(DataStream<Metric> metrics) {
        return metrics
            .keyBy(Metric::getCategory)
            .window(SlidingEventTimeWindows.of(Time.minutes(10), Time.minutes(1)))
            .aggregate(new AverageAggregator());
    }

    // Session windows: activity-based, variable size
    // Use case: user sessions, conversation analysis
    public DataStream<Session> sessionAnalysis(DataStream<ClickEvent> clicks) {
        return clicks
            .keyBy(ClickEvent::getUserId)
            .window(EventTimeSessionWindows.withGap(Time.minutes(30)))
            .process(new SessionBuilder());
    }
}

Oturum pencereleri, yapilandiriabilir bir etkisizlik suresinden sonra kapanir. Bu desen, oturum uzunlugunun katilima gore degistigi kullanici davranisi analizi icin calisir.

Data Engineering mülakatlarında başarılı olmaya hazır mısın?

İnteraktif simülatörler, flashcards ve teknik testlerle pratik yap.

Durum Yonetimi ve Checkpointing

Flink, isleme boyunca operator durumunu ve anahtarli durumu korur. Anahtarli durum verileri anahtara gore bolumlendirerek, iliskili kayitlari bir arada tutarken paralel islemeyi mumkun kilar. Operator durumu tum operator ornegine uygulanir.

StatefulProcessor.javajava
// Managing state in a Flink KeyedProcessFunction
public class FraudDetector extends KeyedProcessFunction<String, Transaction, Alert> {
    
    // Keyed state: one value per key (account)
    private ValueState<Double> lastAmountState;
    private ValueState<Long> lastTransactionTimeState;
    private MapState<String, Integer> merchantCountState;
    
    @Override
    public void open(Configuration parameters) {
        // Initialize state descriptors
        lastAmountState = getRuntimeContext().getState(
            new ValueStateDescriptor<>("lastAmount", Double.class));
        lastTransactionTimeState = getRuntimeContext().getState(
            new ValueStateDescriptor<>("lastTime", Long.class));
        merchantCountState = getRuntimeContext().getMapState(
            new MapStateDescriptor<>("merchantCounts", String.class, Integer.class));
    }
    
    @Override
    public void processElement(Transaction txn, Context ctx, Collector<Alert> out) 
            throws Exception {
        Double lastAmount = lastAmountState.value();
        Long lastTime = lastTransactionTimeState.value();
        
        // Detect suspicious patterns
        if (lastAmount != null && lastTime != null) {
            long timeDelta = txn.getTimestamp() - lastTime;
            // Flag transactions 10x larger than previous within 1 minute
            if (txn.getAmount() > lastAmount * 10 && timeDelta < 60000) {
                out.collect(new Alert(txn.getAccountId(), "SUSPICIOUS_SPIKE", txn));
            }
        }
        
        // Update state for next transaction
        lastAmountState.update(txn.getAmount());
        lastTransactionTimeState.update(txn.getTimestamp());
        
        // Track merchant frequency
        Integer count = merchantCountState.get(txn.getMerchantId());
        merchantCountState.put(txn.getMerchantId(), (count == null ? 0 : count) + 1);
    }
}

Bu durumlu islemci, her hesap icin islem desenlerini izler. Durum kontrol noktalari boyunca kalici olarak korunur ve dolandiricilik tespit baglamini kaybetmeden hatalari atlatir.

Flink 2.3, artimli gorunum bakimi icin Materialized Tables ile SQL yeteneklerini genisletir. Table API, toplu ve akis isleme icin birlesik bir arayuz saglar.

sql
-- flink_sql_streaming.sql
-- Create a streaming source table from Kafka
CREATE TABLE orders (
    order_id STRING,
    customer_id STRING,
    product_id STRING,
    amount DECIMAL(10, 2),
    order_time TIMESTAMP(3),
    -- Define watermark for event time processing
    WATERMARK FOR order_time AS order_time - INTERVAL '5' SECOND
) WITH (
    'connector' = 'kafka',
    'topic' = 'orders',
    'properties.bootstrap.servers' = 'kafka:9092',
    'format' = 'json',
    'scan.startup.mode' = 'earliest-offset'
);

-- Streaming aggregation with tumbling window
SELECT 
    customer_id,
    TUMBLE_START(order_time, INTERVAL '1' HOUR) AS window_start,
    COUNT(*) AS order_count,
    SUM(amount) AS total_amount
FROM orders
GROUP BY 
    customer_id,
    TUMBLE(order_time, INTERVAL '1' HOUR);

SQL yaklasimi, SQL'e asina analistler icin gelistirmeyi basitlestirirken Flink altta akis isleme karmasikligini yonetir.

Her iki framework de akis verilerini isler, ancak mimarileri temelden farklidir. Flink kayitlari gercek akis ile tek tek islerken, Spark mikro-batch'ler isler. Apache Spark karsilastirmalari icin, gecikme ve tutarlilik odunlesimleri uretimde onemlidir.

OzellikFlinkSpark Structured Streaming
Isleme ModeliGercek akisMikro-batch
GecikmeMilisaniyeSaniye (batch araligi)
Durum BackendRocksDB, HashMapsBellek ici, HDFS
Exactly-OnceCheckpoint'lerle yerelIdempotent havuz gerektirir
Event TimeBirinci sinif destek2.1'den beri destekleniyor
SQL DestegiTam akis SQLSinirli pencereleme
Mulakat Ipucu

Akis icin Flink vs Spark soruldugunda, kullanim senaryosu uyumuna odaklanin. Flink dusuk gecikmeli olay isleme ve karmasik olay desenleri icin mukemmeldir. Spark Streaming, toplu islem icin zaten Spark kullanan ve birlesik batch-stream islemesine ihtiyac duyan organizasyonlara uygundur.

Flink exactly-once semantigini nasil saglar?

Flink, islem destekleyen havuzlar icin checkpointing'i iki fazli commit ile birlestirir. Bir kontrol noktasi sirasinda, Flink operator durumunun anlik goruntusunu alir ve kaynak offset'lerini kaydeder. Kafka havuzlari icin, Flink kayitlari Kafka'ya on-commit eder, checkpoint'i tamamlar, sonra islemi commit eder. Checkpoint tamamlanmadan once hata olusursa, commit edilmemis kayitlar atilir ve isleme son checkpoint'ten devam eder.

Coklu kaynak topolojisinde watermark yayilimini aciklayin.

Bir is birden fazla bolum veya kaynaktan okudugunda, her biri gelen olaylara dayali kendi watermark'larini uretir. Downstream operatorun watermark'i, tum giris kanallarindaki minimum watermark'a esittir. Bu, hizli bir bolum nedeniyle hicbir pencerenin erken kapanmamasini saglar. Bazi bolumler veri gondermeyi biraktiginda watermark'lari ilerletmek icin withIdleness() yapilandirin.

Flink'te backpressure'a ne sebep olur ve nasil teshis edilir?

Backpressure, downstream operatorler upstream veri hizlarina yetisemediklerinde olusur. Flink Web UI, operator basina backpressure durumunu gosterir. Yaygin nedenler sunlardir:

  • Yavas harici sistem cagrilari (veritabani sorgulari, API cagrilari)
  • Map/process fonksiyonlarinda pahali hesaplamalar
  • Veri hacmi icin yetersiz paralellik
  • Islemeyi engelleyen buyuk durum islemleri

Paraleligi artirarak, yavas islemleri optimize ederek veya harici cagrilar icin async I/O kullanarak cozun.

Savepoint'ler checkpoint'lerden nasil farklidir?

Checkpoint'ler otomatik, artimli ve hata kurtarma icin optimize edilmistir. Flink yasam dongulerini yoneterek eskileri otomatik olarak siler. Savepoint'ler kullanici tarafindan tetiklenen, operasyonel gorevler icin tasarlanmis tam anlik goruntulerdir: yeni kod dagitimi, isi yeniden olcekleme veya kumeler arasi gecis. Savepoint'ler acikca silinene kadar kalici olur ve sema evrimini destekler.

Flink Kubernetes Operator 1.15, FlinkDeployment ozel kaynaklari ile dagitimi basitlestiriri. Is yasam dongusu, yukseltmeler ve olceklemeyi yonetir.

yaml
# flink-deployment.yaml
# Kubernetes deployment for a Flink application
apiVersion: flink.apache.org/v1beta1
kind: FlinkDeployment
metadata:
  name: fraud-detection-job
spec:
  image: flink:1.20
  flinkVersion: v1_20
  flinkConfiguration:
    taskmanager.numberOfTaskSlots: "4"
    state.backend.type: rocksdb
    state.checkpoints.dir: s3://flink-checkpoints/fraud-detection
    execution.checkpointing.interval: "30s"
  serviceAccount: flink
  jobManager:
    resource:
      memory: "2048m"
      cpu: 1
  taskManager:
    resource:
      memory: "4096m"
      cpu: 2
    replicas: 3
  job:
    jarURI: s3://flink-artifacts/fraud-detection-1.0.jar
    entryClass: com.example.FraudDetectionJob
    parallelism: 12
    upgradeMode: savepoint

upgradeMode: savepoint ayari, operatorun yukseltmeden once savepoint almasi ve dagitimlar arasinda durumun korunmasini saglar.

Uretim dagitimlari paralellik, bellek ve durum backend yapilandirmasina dikkat gerektirir.

ProductionConfig.javajava
// Production-ready Flink configuration
public class ProductionConfig {
    
    public static void configureForProduction(StreamExecutionEnvironment env) {
        // Checkpoint configuration
        CheckpointConfig checkpointConfig = env.getCheckpointConfig();
        checkpointConfig.setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE);
        checkpointConfig.setMinPauseBetweenCheckpoints(30000); // 30 seconds
        checkpointConfig.setCheckpointTimeout(600000); // 10 minutes
        checkpointConfig.setMaxConcurrentCheckpoints(1);
        checkpointConfig.setExternalizedCheckpointRetention(
            ExternalizedCheckpointRetention.RETAIN_ON_CANCELLATION);
        
        // State backend with incremental checkpoints
        EmbeddedRocksDBStateBackend rocksDB = new EmbeddedRocksDBStateBackend(true);
        rocksDB.setDbStoragePath("/tmp/rocksdb");
        env.setStateBackend(rocksDB);
        
        // Restart strategy with exponential backoff
        env.setRestartStrategy(RestartStrategies.exponentialDelayRestart(
            Duration.ofSeconds(1),  // Initial delay
            Duration.ofMinutes(5),  // Max delay
            2.0,                    // Backoff multiplier
            Duration.ofHours(1),    // Reset backoff after
            0.1                     // Jitter
        ));
    }
}

Artimli checkpoint'ler, yalnizca son checkpoint'ten bu yana degisen durumu yazarak checkpoint boyutunu azaltir. Bu optimizasyon, gigabayt'larca anahtarli durumu yonetirken kritik hale gelir.

Pratik yapmaya başla!

Mülakat simülatörleri ve teknik testlerle bilgini test et.

  • Flink 2.3 olaylari milisaniye gecikmesiyle tek tek isler, mikro-batch sistemlerinden farkli olarak
  • Watermark'larla event time semantigi, gec gelen olaylar hakkindaki yaygin mulakat sorusuna cevap vererek sira disi verileri dogru sekilde isler
  • Anahtarli durum, iliskili kayitlari bir arada tutarken paralel isleme icin verileri bolumlendirir
  • Checkpoint'ler dagitik anlik goruntuler ve iki fazli commit araciligiyla exactly-once garantileri saglar
  • Kubernetes Operator, savepoint tabanli durum koruma ile dagitimi, olceklemeyi ve yukseltmeleri otomatiklestirir
  • Saniyenin altinda gecikme veya karmasik olay isleme desenleri onemli oldugunda Spark Streaming yerine Flink'i secin
  • Buyuk durumlu uretim is yukleri icin artimli checkpoint'lerle RocksDB yapilandirin
Günün meydan okuması

Data Engineering kodundaki hatayı bulabilir misin?

Gerçek bir kod parçası, gizli bir hata, günde bir deneme. Denemek için hesap gerekmez.

Anthony Fillion-Maillet

Yazan:

Anthony Fillion-Maillet

SharpSkill kurucusu

10 yılı aşkın süredir fullstack geliştirici. SharpSkill’i yönetiyor ve burada yayımlanan her şeyden sorumlu.

28 Ağustos 2026 tarihinde güncellendi

Etiketler

#apache-flink
#stream-processing
#data-engineering
#real-time-analytics
#event-time

Paylaş

İlgili makaleler