# 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. - Published: 2026-08-28 - Updated: 2026-08-28 - Author: Anthony Fillion-Maillet - Tags: apache flink, stream processing, event time, data engineering, flink vs spark - Reading time: 10 min --- 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. > **Essentieel voor Sollicitaties** > > 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. ```java // FlinkStreamJob.java // 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 rawStream = env.addSource( new FlinkKafkaConsumer<>("events", new SimpleStringSchema(), kafkaProps) ); // Parse and transform the stream DataStream events = rawStream .map(json -> objectMapper.readValue(json, Event.class)) .assignTimestampsAndWatermarks( WatermarkStrategy.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. ```java // EventTimeExample.java // Configuring event time with watermarks public class EventTimeProcessor { public DataStream processWithEventTime( DataStream readings) { return readings // Extract timestamp from the event payload .assignTimestampsAndWatermarks( WatermarkStrategy .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. > **Veelgestelde Sollicitatievraag** > > 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. ```java // WindowingStrategies.java // Different windowing approaches for stream processing public class WindowingStrategies { // Tumbling windows: fixed-size, non-overlapping // Use case: hourly aggregations, daily summaries public DataStream tumblingAggregation(DataStream 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 slidingAverage(DataStream 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 sessionAnalysis(DataStream 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. ## 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. ```java // StatefulProcessor.java // Managing state in a Flink KeyedProcessFunction public class FraudDetector extends KeyedProcessFunction { // Keyed state: one value per key (account) private ValueState lastAmountState; private ValueState lastTransactionTimeState; private MapState 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 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](https://flink.apache.org/downloads/) voor incrementeel view onderhoud. De Table API biedt een uniforme interface voor batch en stream processing. ```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); ``` 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](/blog/data-engineering/apache-spark-4-new-features-structured-streaming-interview) 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 | > **Sollicitatie Inzicht** > > 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](https://flink.apache.org/2026/05/26/apache-flink-kubernetes-operator-1.15.0-release-announcement/) vereenvoudigt deployment met FlinkDeployment custom resources. Het handelt job lifecycle, upgrades en scaling af. ```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 ``` De `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](/technologies/data-engineering/interview-questions/etl-elt-patterns) voor integratie overwegingen. ```java // ProductionConfig.java // 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. ## 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 --- Source: SharpSkill (https://sharpskill.dev), tech interview preparation for your real stack. HTML version of this page: https://sharpskill.dev/nl/blog/data-engineering/apache-flink-stream-processing-event-time-interview-2026