Apache Flink en 2026 : Traitement de Flux, Event Time et Questions d'Entretien

Guide complet sur Apache Flink 2.3 pour le traitement de flux en temps réel. Watermarks, fenêtrage, gestion d'état et préparation aux entretiens data engineering.

Architecture de traitement de flux Apache Flink avec event time et fenêtrage

Apache Flink 2.3 représente la référence du traitement de flux distribué avec une latence inférieure à la seconde et des garanties exactly-once. Contrairement aux systèmes orientés batch, Flink traite les données en continu dès leur arrivée, ce qui en fait le framework privilégié pour l'analytique temps réel, la détection de fraude et les architectures événementielles.

Point Clé Entretien

Flink se distingue de Spark Streaming par un véritable traitement de flux : Flink traite les événements un par un avec une sémantique event time, tandis que Spark Streaming traite des micro-batches avec le processing time par défaut.

Flink s'exécute sur une architecture distribuée avec un JobManager qui coordonne le travail entre plusieurs TaskManagers. Chaque TaskManager exécute des task slots qui traitent des portions des opérateurs parallèles du job. Cette séparation permet à Flink de scaler horizontalement tout en maintenant la tolérance aux pannes via des checkpoints distribués.

Le modèle dataflow dans Flink représente les calculs sous forme de graphes acycliques dirigés (DAG). Les données circulent des sources vers les sinks en passant par des transformations, chaque opérateur pouvant s'exécuter sur plusieurs instances parallèles.

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

La configuration de checkpoint ci-dessus stocke des instantanés de l'état distribué toutes les 10 secondes. En cas de panne, Flink restaure depuis le dernier checkpoint complété et rejoue les enregistrements depuis Kafka.

Event Time vs Processing Time : Comprendre les Sémantiques Temporelles

L'event time désigne le moment où un événement s'est réellement produit, intégré dans les données elles-mêmes. Le processing time correspond au moment où Flink traite l'enregistrement. Cette distinction est cruciale car les délais réseau, les livraisons désordonnées et les retards de traitement rendent le processing time peu fiable pour les opérations basées sur le temps.

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

La stratégie forBoundedOutOfOrderness indique à Flink que les événements peuvent arriver avec jusqu'à 2 minutes de retard. Les watermarks progressent lorsque Flink détermine qu'aucun événement avec un timestamp antérieur au watermark n'arrivera plus.

Question d'Entretien Fréquente

Que se passe-t-il pour les événements en retard dans Flink ? Par défaut, les événements arrivant après que le watermark a dépassé la fin de la fenêtre sont supprimés. Il est possible de configurer une latence autorisée avec .allowedLateness(Time.minutes(10)) pour traiter les arrivées tardives, ou d'utiliser des side outputs pour les capturer et les traiter séparément.

Stratégies de Fenêtrage pour l'Analytique Temps Réel

Flink propose quatre types de fenêtres : tumbling, sliding, session et global. Chacun répond à des besoins analytiques différents.

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

Les fenêtres de session se ferment après une période d'inactivité configurable. Ce pattern fonctionne parfaitement pour l'analyse du comportement utilisateur où la durée des sessions varie selon l'engagement.

Prêt à réussir tes entretiens Data Engineering ?

Entraîne-toi avec nos simulateurs interactifs, fiches express et tests techniques.

Gestion de l'État et Checkpointing

Flink maintient l'état des opérateurs et l'état keyed tout au long du traitement. L'état keyed partitionne les données par clé, permettant un traitement parallèle tout en gardant les enregistrements liés ensemble. L'état opérateur s'applique à l'instance entière de l'opérateur.

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

Ce processeur avec état suit les patterns de transaction par compte. L'état persiste à travers les checkpoints, survivant aux pannes sans perdre le contexte de détection de fraude.

Flink 2.3 étend ses capacités SQL avec les Materialized Tables pour la maintenance incrémentale des vues. La Table API fournit une interface unifiée pour le traitement batch et streaming.

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

L'approche SQL simplifie le développement pour les analystes familiers avec SQL tandis que Flink gère la complexité du traitement de flux en arrière-plan.

Les deux frameworks traitent des données en streaming, mais leurs architectures diffèrent fondamentalement. Flink traite les enregistrements individuellement avec un véritable streaming, tandis que Spark traite des micro-batches. Pour les comparaisons avec Apache Spark, les compromis latence/cohérence comptent en production.

AspectFlinkSpark Structured Streaming
Modèle de traitementVrai streamingMicro-batch
LatenceMillisecondesSecondes (intervalle batch)
State BackendRocksDB, HashMapsEn mémoire, HDFS
Exactly-OnceNatif avec checkpointsNécessite des sinks idempotents
Event TimeSupport natifSupporté depuis 2.1
Support SQLSQL streaming completFenêtrage limité
Conseil Entretien

Lorsqu'on pose la question Flink vs Spark pour le streaming, il convient de se concentrer sur l'adéquation au cas d'usage. Flink excelle pour le traitement d'événements à faible latence et les patterns d'événements complexes. Spark Streaming convient aux organisations utilisant déjà Spark pour le batch qui ont besoin d'un traitement unifié batch-stream.

Comment Flink atteint-il la sémantique exactly-once ?

Flink combine le checkpointing avec le two-phase commit pour les sinks qui supportent les transactions. Pendant un checkpoint, Flink capture un instantané de l'état des opérateurs et enregistre les offsets des sources. Pour les sinks Kafka, Flink pré-commit les enregistrements vers Kafka, complète le checkpoint, puis commit la transaction. Si une panne survient avant la complétion du checkpoint, les enregistrements non committés sont supprimés et le traitement reprend depuis le dernier checkpoint.

Comment expliquer la propagation des watermarks dans une topologie multi-sources ?

Quand un job lit depuis plusieurs partitions ou sources, chacune génère ses propres watermarks basés sur les événements entrants. Le watermark de l'opérateur en aval est égal au watermark minimum de tous les canaux d'entrée. Cela garantit qu'aucune fenêtre ne se ferme prématurément à cause d'une partition rapide qui avance avant les plus lentes. La configuration de withIdleness() permet d'avancer les watermarks quand certaines partitions cessent d'envoyer des données.

Quelles sont les causes de la backpressure dans Flink et comment la diagnostiquer ?

La backpressure survient quand les opérateurs en aval ne peuvent pas suivre le débit des données en amont. L'interface Web Flink affiche le statut de backpressure par opérateur. Les causes courantes incluent :

  • Les appels lents vers des systèmes externes (requêtes base de données, appels API)
  • Les calculs coûteux dans les fonctions map/process
  • Un parallélisme insuffisant pour le volume de données
  • Les opérations d'état volumineuses bloquant le traitement

Pour y remédier, il faut augmenter le parallélisme, optimiser les opérations lentes ou utiliser l'I/O asynchrone pour les appels externes.

Quelle est la différence entre savepoints et checkpoints ?

Les checkpoints sont automatiques, incrémentaux et optimisés pour la récupération après panne. Flink gère leur cycle de vie et supprime automatiquement les anciens. Les savepoints sont déclenchés par l'utilisateur, sont des instantanés complets destinés aux tâches opérationnelles : déployer du nouveau code, redimensionner le job ou migrer entre clusters. Les savepoints persistent jusqu'à suppression explicite et supportent l'évolution de schéma.

Le Flink Kubernetes Operator 1.15 simplifie le déploiement avec les ressources personnalisées FlinkDeployment. Il gère le cycle de vie des jobs, les mises à jour et le scaling.

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

Le paramètre upgradeMode: savepoint garantit que l'opérateur prend un savepoint avant la mise à jour, préservant l'état entre les déploiements.

Les déploiements en production nécessitent une attention particulière au parallélisme, à la mémoire et à la configuration du state backend. Consultez les patterns ETL et pipelines de données pour les considérations d'intégration.

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

Les checkpoints incrémentaux réduisent la taille des checkpoints en n'écrivant que l'état modifié depuis le dernier checkpoint. Cette optimisation devient critique lors de la gestion de gigaoctets d'état keyed.

Passe à la pratique !

Teste tes connaissances avec nos simulateurs d'entretien et tests techniques.

  • Flink 2.3 traite les événements individuellement avec une latence de l'ordre de la milliseconde, contrairement aux systèmes micro-batch
  • La sémantique event time avec les watermarks gère correctement les données désordonnées, répondant à la question d'entretien courante sur les événements en retard
  • L'état keyed partitionne les données pour un traitement parallèle tout en gardant les enregistrements liés ensemble
  • Les checkpoints fournissent des garanties exactly-once via des instantanés distribués et le two-phase commit
  • Le Kubernetes Operator automatise le déploiement, le scaling et les mises à jour avec la préservation de l'état basée sur les savepoints
  • Flink est préférable à Spark Streaming quand la latence sub-seconde ou les patterns d'événements complexes sont importants
  • La configuration de RocksDB avec des checkpoints incrémentaux est recommandée pour les workloads de production avec un état volumineux
Défi du jour

Tu saurais repérer le bug en Data Engineering ?

Un vrai bout de code, un bug caché, une tentative par jour. Sans compte pour essayer.

Anthony Fillion-Maillet

Écrit par

Anthony Fillion-Maillet

Fondateur de SharpSkill

Développeur fullstack depuis plus de 10 ans. Il dirige SharpSkill et répond de tout ce qui y est publié.

Mis à jour le 28 août 2026

Tags

#apache-flink
#traitement-flux
#data-engineering
#analytique-temps-reel
#event-time

Partager

Articles similaires