Apache Flink 2026: Stream Processing, Event Time en Sollicitatievragen
Apache Flink 2.3 verwerkt streaming data met exactly-once semantiek en sub-seconde latentie. Deze handleiding behandelt event-time verwerking, windowing strategieen, state management en veelgestelde sollicitatievragen voor Data Engineers.

Apache Flink 2.3 verwerkt streamingdata in realtime met exactly-once semantiek en sub-seconde latentie. In tegenstelling tot batch-georienteerde systemen verwerkt Flink data continu bij aankomst, waardoor het de voorkeurskeuze is voor realtime analytics, fraudedetectie en event-driven architecturen.
Flink onderscheidt zich van Spark Streaming door echte stream processing: Flink verwerkt events een voor een met event-time semantiek, terwijl Spark Streaming micro-batches verwerkt met processing time als standaard.
Flink 2.3 Architectuur voor Stream Processing
Flink draait op een gedistribueerde architectuur met een JobManager die het werk coordineert over meerdere TaskManagers. Elke TaskManager voert task slots uit die delen van de parallelle operators van de job uitvoeren. Deze scheiding maakt horizontale schaalbaarheid mogelijk terwijl fouttolerantie behouden blijft door gedistribueerde checkpoints.
Het dataflow model in Flink representeert berekeningen als gerichte acyclische grafen (DAGs). Data stroomt van bronnen door transformaties naar sinks, waarbij elke operator potentieel op meerdere parallelle instanties draait.
// 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())
);De checkpoint configuratie slaat elke 10 seconden snapshots op van de gedistribueerde state. Bij een storing herstelt Flink vanaf het laatste voltooide checkpoint en speelt records opnieuw af vanaf Kafka.
Event Time vs Processing Time Semantiek
Event time verwijst naar het moment waarop een event daadwerkelijk plaatsvond, ingebed in de data zelf. Processing time is het moment waarop Flink de record verwerkt. Het onderscheid is belangrijk omdat netwerkvertragingen, out-of-order levering en verwerkingsachterstanden processing time onbetrouwbaar maken voor tijdgebaseerde operaties.
// 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());
}
}De forBoundedOutOfOrderness strategie vertelt Flink dat events tot 2 minuten te laat kunnen aankomen. Watermarks schuiven op wanneer Flink bepaalt dat er geen events meer zullen aankomen met timestamps voor de watermark.
Wat gebeurt er met late events in Flink? Standaard worden events die aankomen nadat de watermark de eindtijd van het venster heeft gepasseerd, verworpen. Configureer .allowedLateness(Time.minutes(10)) om late aankomsten te verwerken, of gebruik side outputs om ze op te vangen voor aparte verwerking.
Windowing Strategieen voor Realtime Analytics
Flink biedt vier venstertypes: tumbling, sliding, session en global windows. Elk dient verschillende analytische behoeften.
// 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());
}
}Session windows sluiten na een configureerbare periode van inactiviteit. Dit patroon werkt voor analyse van gebruikersgedrag waarbij sessieduur varieert op basis van engagement.
Klaar om je Data Engineering gesprekken te halen?
Oefen met onze interactieve simulatoren, flashcards en technische tests.
State Management en Checkpointing
Flink onderhoudt operator state en keyed state tijdens de verwerking. Keyed state partitioneert data per sleutel, wat parallelle verwerking mogelijk maakt terwijl gerelateerde records bij elkaar blijven. Operator state is van toepassing op de gehele operator instantie.
// 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);
}
}Deze stateful processor volgt transactiepatronen per account. De state blijft behouden over checkpoints heen en overleeft storingen zonder verlies van fraudedetectie context.
Flink SQL en Table API voor Stream Processing
Flink 2.3 breidt zijn SQL mogelijkheden uit met Materialized Tables voor incrementeel view onderhoud. De Table API biedt een uniforme interface voor batch en stream processing.
-- 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);De SQL aanpak vereenvoudigt ontwikkeling voor analisten die bekend zijn met SQL, terwijl Flink de stream processing complexiteit op de achtergrond afhandelt.
Flink vs Spark Structured Streaming
Beide frameworks verwerken streaming data, maar hun architecturen verschillen fundamenteel. Flink verwerkt records individueel met echte streaming, terwijl Spark micro-batches verwerkt. Voor Apache Spark vergelijkingen zijn de latentie en consistentie trade-offs relevant in productie.
| Aspect | Flink | Spark Structured Streaming |
|---|---|---|
| Verwerkingsmodel | Echte streaming | Micro-batch |
| Latentie | Milliseconden | Seconden (batch interval) |
| State Backend | RocksDB, HashMaps | In-memory, HDFS |
| Exactly-Once | Natief met checkpoints | Vereist idempotente sinks |
| Event Time | Eersteklas ondersteuning | Ondersteund sinds 2.1 |
| SQL Ondersteuning | Volledige streaming SQL | Beperkte windowing |
Bij vragen over Flink vs Spark voor streaming, focus op use case fit. Flink excelleert bij low-latency event processing en complexe event patronen. Spark Streaming past bij organisaties die al Spark draaien voor batch en uniforme batch-stream verwerking nodig hebben.
Veelgestelde Flink Sollicitatievragen en Antwoorden
Hoe bereikt Flink exactly-once semantiek?
Flink combineert checkpointing met two-phase commit voor sinks die transacties ondersteunen. Tijdens een checkpoint maakt Flink een snapshot van de operator state en registreert source offsets. Voor Kafka sinks pre-commit Flink records naar Kafka, voltooit de checkpoint, en commit vervolgens de transactie. Als een storing optreedt voor checkpoint voltooiing, worden uncommitted records verworpen en wordt de verwerking hervat vanaf de laatste checkpoint.
Leg watermark propagatie uit in een multi-source topologie.
Wanneer een job leest van meerdere partities of bronnen, genereert elk zijn eigen watermarks gebaseerd op binnenkomende events. De watermark van de downstream operator is gelijk aan de minimale watermark over alle input kanalen. Dit zorgt ervoor dat geen venster voortijdig sluit doordat een snelle partitie vooruitloopt op tragere. Configureer withIdleness() om watermarks te laten voortschrijden wanneer sommige partities stoppen met data verzenden.
Wat veroorzaakt backpressure in Flink en hoe diagnosticeer je het?
Backpressure treedt op wanneer downstream operators de dataratio van upstream niet kunnen bijhouden. De Flink Web UI toont backpressure status per operator. Veelvoorkomende oorzaken zijn:
- Trage externe systeemaanroepen (database queries, API calls)
- Dure berekeningen in map/process functies
- Onvoldoende parallellisme voor het datavolume
- Grote state operaties die verwerking blokkeren
Oplossingen: verhoog parallellisme, optimaliseer trage operaties, of gebruik async I/O voor externe aanroepen.
Wat is het verschil tussen savepoints en checkpoints?
Checkpoints zijn automatisch, incrementeel en geoptimaliseerd voor failure recovery. Flink beheert hun levenscyclus en verwijdert oude automatisch. Savepoints zijn door gebruikers getriggerd, complete snapshots bedoeld voor operationele taken: deployen van nieuwe code, herschalen van de job, of migreren tussen clusters. Savepoints blijven bestaan tot ze expliciet worden verwijderd en ondersteunen schema evolutie.
Flink Deployen op Kubernetes
De Flink Kubernetes Operator 1.15 vereenvoudigt deployment met FlinkDeployment custom resources. Het handelt job lifecycle, upgrades en scaling af.
# 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: savepointDe upgradeMode: savepoint instelling zorgt ervoor dat de operator een savepoint maakt voor een upgrade, waardoor state behouden blijft over deployments heen.
Flink Applicaties Optimaliseren voor Productie
Productie deployments vereisen aandacht voor parallellisme, geheugen en state backend configuratie. Zie de ETL en data pipeline patronen voor integratie overwegingen.
// 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
));
}
}Incrementele checkpoints reduceren checkpoint grootte door alleen gewijzigde state sinds de laatste checkpoint te schrijven. Deze optimalisatie wordt kritiek bij het beheren van gigabytes aan keyed state.
Begin met oefenen!
Test je kennis met onze gespreksimulatoren en technische tests.
Kernpunten voor Apache Flink Stream Processing
- Flink 2.3 verwerkt events individueel met milliseconde latentie, in tegenstelling tot micro-batch systemen
- Event-time semantiek met watermarks handelt out-of-order data correct af en beantwoordt de veelgestelde sollicitatievraag over late events
- Keyed state partitioneert data voor parallelle verwerking terwijl gerelateerde records bij elkaar blijven
- Checkpoints bieden exactly-once garanties door gedistribueerde snapshots en two-phase commit
- De Kubernetes Operator automatiseert deployment, scaling en upgrades met savepoint-gebaseerde state preservatie
- Kies Flink boven Spark Streaming wanneer sub-seconde latentie of complexe event processing patronen belangrijk zijn
- Configureer RocksDB met incrementele checkpoints voor productie workloads met grote state
Zie jij de bug in Data Engineering?
Een echt codefragment, een verborgen bug, één poging per dag. Zonder account uit te proberen.

Geschreven door
Anthony Fillion-MailletOprichter van SharpSkill
Al meer dan 10 jaar fullstack-ontwikkelaar. Hij leidt SharpSkill en staat in voor alles wat hier verschijnt.
Bijgewerkt op 28 augustus 2026
Tags
Delen
Gerelateerde artikelen

Apache Beam vs Spark 2026: Unified Pipelines en Sollicitatievragen
Vergelijking van Apache Beam 2.76 en Spark 4.2 voor data pipelines. Portabiliteit, performance, sollicitatievragen en beslissingscriteria voor Data Engineers.

Apache Spark 4.2 vs Databricks in 2026: Architectuur, Prestaties en Sollicitatievragen
Vergelijking van Apache Spark 4.2 en Databricks in 2026. Architectuurverschillen, prestatiekenmerken en veelgestelde sollicitatievragen voor data engineering functies.

Delta Lake vs Apache Iceberg 2026: Lakehouse Architectuur en Sollicitatievragen
Vergelijking van Delta Lake en Apache Iceberg voor Data Lakehouse architecturen. Transactiemodellen, partitie-evolutie, engine-compatibiliteit en veelvoorkomende sollicitatievragen voor Data Engineers.