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.

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.
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.
Architecture Flink 2.3 pour le Traitement de Flux
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.
// 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.
// 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.
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.
// 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.
// 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 SQL et Table API pour le Traitement de Flux
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.
-- 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.
Flink vs Spark Structured Streaming : Comparaison Technique
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.
| Aspect | Flink | Spark Structured Streaming |
|---|---|---|
| Modèle de traitement | Vrai streaming | Micro-batch |
| Latence | Millisecondes | Secondes (intervalle batch) |
| State Backend | RocksDB, HashMaps | En mémoire, HDFS |
| Exactly-Once | Natif avec checkpoints | Nécessite des sinks idempotents |
| Event Time | Support natif | Supporté depuis 2.1 |
| Support SQL | SQL streaming complet | Fenêtrage limité |
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.
Questions d'Entretien Courantes sur Flink avec Réponses
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.
Déploiement de Flink sur Kubernetes
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.
# 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: savepointLe paramètre upgradeMode: savepoint garantit que l'opérateur prend un savepoint avant la mise à jour, préservant l'état entre les déploiements.
Optimisation des Applications Flink pour la Production
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.
// 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.
Points Clés pour le Traitement de Flux avec Apache Flink
- 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
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.

Écrit par
Anthony Fillion-MailletFondateur 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
Partager
Articles similaires

Apache Beam vs Spark en 2026 : Pipelines Unifiés et Questions d'Entretien
Comparaison détaillée entre Apache Beam et Spark pour les pipelines de données en 2026. Portabilité, windowing, performances et questions d'entretien techniques.

Snowflake en 2026 : architecture, SQL et questions d'entretien data engineer
Un guide 2026 de l'architecture Snowflake pour data engineers : comment le stockage et le calcul se séparent, comment fonctionnent les virtual warehouses et les micro-partitions, et les questions d'entretien qui testent l'expérience en production.

Top 25 Questions d'Entretien Data Engineering en 2026
Les questions d'entretien data engineering les plus fréquentes en 2026 : SQL avancé, pipelines temps réel, architecture lakehouse, Spark, Airflow et optimisation des coûts cloud.