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.

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.
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.
Akis Isleme icin Flink 2.3 Mimarisi
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.
// 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.
// 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.
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.
// 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.
// 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.
Akis Isleme icin Flink SQL ve Table API
Flink 2.3, artimli gorunum bakimi icin Materialized Tables ile SQL yeteneklerini genisletir. Table API, toplu ve akis isleme icin birlesik bir arayuz saglar.
-- 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.
Flink vs Spark Structured Streaming
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.
| Ozellik | Flink | Spark Structured Streaming |
|---|---|---|
| Isleme Modeli | Gercek akis | Mikro-batch |
| Gecikme | Milisaniye | Saniye (batch araligi) |
| Durum Backend | RocksDB, HashMaps | Bellek ici, HDFS |
| Exactly-Once | Checkpoint'lerle yerel | Idempotent havuz gerektirir |
| Event Time | Birinci sinif destek | 2.1'den beri destekleniyor |
| SQL Destegi | Tam akis SQL | Sinirli pencereleme |
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.
Yaygin Apache Flink Mulakat Sorulari ve Cevaplari
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.
Kubernetes'te Flink Dagitimi
Flink Kubernetes Operator 1.15, FlinkDeployment ozel kaynaklari ile dagitimi basitlestiriri. Is yasam dongusu, yukseltmeler ve olceklemeyi yonetir.
# 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: savepointupgradeMode: savepoint ayari, operatorun yukseltmeden once savepoint almasi ve dagitimlar arasinda durumun korunmasini saglar.
Uretim icin Flink Uygulamalarini Optimize Etme
Uretim dagitimlari paralellik, bellek ve durum backend yapilandirmasina dikkat gerektirir.
// 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.
Apache Flink Akis Isleme icin Temel Cikarimlar
- 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
Data Engineering kodundaki hatayı bulabilir misin?
Gerçek bir kod parçası, gizli bir hata, günde bir deneme. Denemek için hesap gerekmez.

Yazan:
Anthony Fillion-MailletSharpSkill 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
Paylaş
İlgili makaleler

Apache Airflow 2026: Pipeline Orkestrasyonu, DAG Mimarisi ve Mülakat Soruları
Apache Airflow 3.2 ile Task SDK kullanarak DAG yazımı, dinamik görev eşlemesi, asset partition desteği ve veri mühendisliği mülakatlarında karşılaşılan soruları kapsayan uygulamalı rehber.

Apache Spark 4: Yeni Ozellikler, Structured Streaming ve Mulakat Sorulari (2026)
Apache Spark 4 hakkinda kapsamli teknik rehber. ANSI SQL modu, VARIANT veri tipi, Real-Time Mode Streaming, Spark Connect ve Data Engineering mulakat sorulari detayli kod ornekleriyle inceleniyor.

2026'da En Çok Sorulan 25 Veri Mühendisliği Mülakat Sorusu
2026 yılında veri mühendisliği mülakatlarında en sık karşılaşılan 25 soru ve detaylı cevapları. SQL optimizasyonu, veri hatları, ETL/ELT, Apache Spark, Kafka, veri modelleme, orkestrasyon ve sistem tasarımı konularını kapsar.