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 2026: Stream Processing, Event Time en Sollicitatievragen

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 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.

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())
    );

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.

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());
    }
}

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.

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());
    }
}

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.

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);
    }
}

Deze stateful processor volgt transactiepatronen per account. De state blijft behouden over checkpoints heen en overleeft storingen zonder verlies van fraudedetectie context.

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.

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.

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.

AspectFlinkSpark Structured Streaming
VerwerkingsmodelEchte streamingMicro-batch
LatentieMillisecondenSeconden (batch interval)
State BackendRocksDB, HashMapsIn-memory, HDFS
Exactly-OnceNatief met checkpointsVereist idempotente sinks
Event TimeEersteklas ondersteuningOndersteund sinds 2.1
SQL OndersteuningVolledige streaming SQLBeperkte 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.

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.

De Flink Kubernetes Operator 1.15 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.

Productie deployments vereisen aandacht voor parallellisme, geheugen en state backend configuratie. Zie de ETL en data pipeline patronen voor integratie overwegingen.

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
        ));
    }
}

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.

  • 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
Dagelijkse challenge

Zie jij de bug in Data Engineering?

Een echt codefragment, een verborgen bug, één poging per dag. Zonder account uit te proberen.

Anthony Fillion-Maillet

Geschreven door

Anthony Fillion-Maillet

Oprichter 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

#apache flink
#stream processing
#event time
#data engineering
#flink vs spark

Delen

Gerelateerde artikelen