Instalar com Codex ou Claude Copie este prompt, cole no Codex, Claude ou outro assistente e deixe que ele revise a página da skill e instale para você.
Um comando direto ignora o prompt de revisão. Verifique a origem antes de executá-lo.
Apache Flink est le moteur de traitement de flux le plus avancé, avec une philosophie radicale : le batch est un cas particulier du streaming à temps borné. Ses caractéristiques clés :
Event Time processing : traitement basé sur l'horloge de l'événement, pas du système.
State Backend : état distribué scalable (RocksDB, Heap, ForSt).
Savepoints / Checkpoints : reprise exactement-une-fois (exactly-once) sans perte.
Flink SQL : requêtes SQL stream avec des tables dynamiques.
CEP (Complex Event Processing) : détection de patterns temporels complexes.
Cette compétence couvre : PyFlink (Python), Flink SQL, CEP, gestion d'état (RocksDB), sauvegarde et restauration (savepoints), déploiement sur YARN/Kubernetes, intégration Kafka, watermarking, et fenêtrage avancé.
Quand l'utiliser
Activez cette compétence lorsque l'utilisateur :
A besoin de streaming temps réel avec exactly-once garanti.
Doit manipuler des états distribués (machine learning online, sessions, agrégations).
Veut détecter des séquences complexes d'événements (fraude, maintenance prédictive).
Demande du batch unifié avec le même code que le streaming.
Pose des questions sur les watermarks, les triggers, les state backends.
Prérequis
pip install apache-flink pyflink
Vérification :
from pyflink.datastream import StreamExecutionEnvironment
from pyflink.table import StreamTableEnvironment
env = StreamExecutionEnvironment.get_execution_environment()
env.set_parallelism(4)
print(f"Flink {env}"[:50])
1. Concepts Fondamentaux
1.1 Event Time vs Processing Time vs Ingestion Time
from pyflink.datastream import StreamExecutionEnvironment
from pyflink.datastream.time_characteristic TimeCharacteristic
env = StreamExecutionEnvironment.get_execution_environment()
env.set_stream_time_characteristic(TimeCharacteristic.EventTime)
import
# Event time (recommandé — basé sur l'horloge des capteurs)
# Processing time (basé sur l'horloge de l'exécuteur — le plus simple)
from pyflink.datastream.window import TriggerResult, Trigger
from pyflink.datastream.window import TimeWindow
classEarlyFireTrigger(Trigger):
"""Déclenche un résultat après 5 événements ET au bout de 30s max"""defon_element(self, element, timestamp, window, ctx):
count_ctx = ctx.get_partitioned_state(
ValueStateDescriptor('count', Types.INT()))
count = (count_ctx.value() or0) + 1
count_ctx.update(count)
if count >= 5:
count_ctx.clear()
return TriggerResult.FIRE_AND_PURGE
return TriggerResult.CONTINUE
defon_processing_time(self, time, window, ctx):
return TriggerResult.FIRE
defon_event_time(self, time, window, ctx):
return TriggerResult.FIRE_AND_PURGE
defclear(self, window, ctx):
ctx.get_partitioned_state(
ValueStateDescriptor('count', Types.INT())).clear()
4. Flink SQL : Tables Dynamiques
4.1 Source Kafka en SQL
CREATE TABLE sensor_readings (
machine_id STRING,
temperature DOUBLE,
pression DOUBLE,
`timestamp` TIMESTAMP(3) METADATA FROM'timestamp',
WATERMARK FOR `timestamp` AS `timestamp` -INTERVAL'10'SECOND
) WITH (
'connector'='kafka',
'topic'='sensors-raw',
'properties.bootstrap.servers'='kafka:9092',
'properties.group.id'='flink-sql-consumer',
'scan.startup.mode'='earliest-offset',
'format'='json',
'json.fail-on-missing-field'='true',
'json.ignore-parse-errors'='false'
);
4.2 Fenêtrage SQL (Tumbling / Hop / Session)
-- Fenêtre tumbling de 5 minutesSELECT
machine_id,
TUMBLE_END(`timestamp`, INTERVAL'5'MINUTE) AS fenetre_fin,
AVG(temperature) AS temp_moyenne,
MAX(temperature) AS temp_max,
COUNT(*) AS nb_mesures
FROM sensor_readings
GROUPBY
machine_id,
TUMBLE(`timestamp`, INTERVAL'5'MINUTE);
-- Fenêtre glissante (Hop)SELECT
machine_id,
HOP_END(`timestamp`, INTERVAL'5'MINUTE, INTERVAL'15'MINUTE) AS fenetre_fin,
AVG(temperature) AS temp_moyenne
FROM sensor_readings
GROUPBY
machine_id,
HOP(`timestamp`, INTERVAL'5'MINUTE, INTERVAL'15'MINUTE);
4.3 Jointure Flux - Table de Référence
-- Création d'une table de référence (machines)CREATE TABLE machines (
machine_id STRING PRIMARY KEYNOT ENFORCED,
site STRING,
zone STRING
) WITH (
'connector'='jdbc',
'url'='jdbc:postgresql://postgres:5432/sensors',
'table-name'='machines',
'lookup.cache.max-rows'='5000',
'lookup.cache.ttl'='10min'
);
-- Jointure stream-tableSELECT
s.machine_id,
s.temperature,
m.site,
m.zone
FROM sensor_readings AS s
LEFTJOIN machines FORSYSTEM_TIMEASOF s.proctime AS m
ON s.machine_id = m.machine_id;
4.4 Sink vers Delta Lake
CREATE TABLE temp_agg_delta (
machine_id STRING,
fenetre_fin TIMESTAMP(3),
temp_moyenne DOUBLE,
temp_max DOUBLE,
nb_mesures BIGINT
) WITH (
'connector'='delta',
'table-path'='s3://data-lake/curated/temp_agg/',
'delta.auto-optimize.auto-compact'='true'
);
INSERT INTO temp_agg_delta
SELECT ... FROM sensor_readings ...;
5. Complex Event Processing (CEP)
5.1 Détection de Patterns Temporels
from pyflink.datastream.cep.pattern import Pattern, AfterMatchSkipStrategy
from pyflink.datastream.cep import CEP
from pyflink.datastream.cep.functions import PatternProcessFunction
from pyflink.common.time import Time
# Pattern : augmentation rapide + pic critique dans les 2 min
pattern = Pattern.begin("mesure_stable") \
.where(lambda v: 80 <= v['temperature'] <= 90) \
.next("hausse_rapide") \
.where(lambda v: v['temperature'] > 100) \
.times_or_more(3) \
.within(Time.minutes(2))
# Sauter les patterns qui se chevauchent
skip_strategy = AfterMatchSkipStrategy.skip_past_last_event()
pattern = pattern.with_skip_strategy(skip_strategy)
# Appliquer le pattern
keyed_stream = parsed.key_by(lambda x: x['machine_id'])
pattern_stream = CEP.pattern(keyed_stream, pattern)
classAlertePatternFunction(PatternProcessFunction):
defprocess_match(self, match, ctx, collector):
hausse = match.get('hausse_rapide')
collector.collect({
"machine_id": hausse[-1]['machine_id'],
"type": "SURECHUFFEMENT_RAPIDE",
"derniere_temp": hausse[-1]['temperature'],
"nb_pics": len(hausse),
"timestamp": hausse[-1]['timestamp']
})
alerts = pattern_stream.process(AlertePatternFunction())
alerts.add_sink(...)
from pyflink.datastream import StreamExecutionEnvironment
from pyflink.datastream.state_backend import RocksDBStateBackend
env = StreamExecutionEnvironment.get_execution_environment()
# RocksDB (recommandé pour les états > 1 Go)
rocks_backend = RocksDBStateBackend(
checkpoint_data_uri='file:///data/flink/checkpoints',
enable_incremental=True
)
env.set_state_backend(rocks_backend)
# HeapStateBackend (pour petits états, en mémoire)# env.set_state_backend(StateBackend())# Configuration RocksDB fine
config = rocks_backend.get_db_options()
config.set_option("rocksdb.block.cache-size", "256mb")
config.set_option("rocksdb.writebuffer.size", "64mb")
config.set_option("rocksdb.max.write.buffer.number", "4")
6.2 Checkpointing de Production
from pyflink.common.time import Duration
env.enable_checkpointing(60000) # Checkpoint toutes les 60s
env.get_checkpoint_config().set_checkpoint_storage_dir(
's3://flink-checkpoints/sensors-job/')
env.get_checkpoint_config().set_min_pause_between_checkpoints(30000)
env.get_checkpoint_config().set_checkpoint_timeout(Duration.of_minutes(10))
env.get_checkpoint_config().set_max_concurrent_checkpoints(1)
env.get_checkpoint_config().enable_externalized_checkpoints(
'RETAIN_ON_CANCELLATION') # Garde les checkpoints après arrêt# Exactly-once (vs at-least-once)
env.get_checkpoint_config().set_checkpointing_mode('EXACTLY_ONCE')
6.3 Savepoints : Mise à jour à Chaud
# Déclencher un savepoint
bin/flink savepoint <job_id> s3://flink-savepoints/
# Redémarrer avec un nouveau JAR depuis un savepoint
bin/flink run -s s3://flink-savepoints/savepoint-xxxxx/ \
-c com.eva.SensorJob \
eva-flink-job-2.0.jar
# Arrêt avec savepoint
bin/flink stop --savepointPath s3://flink-savepoints/ <job_id>
# Vérifier le backpressure# Flink Web UI → http://flink-dashboard:8081# Jobs → Job actuel → Backpressure# Ou via REST API
curl http://flink-jobmanager:8081/jobs/<job_id>/backpressure
8. Batch Unifié (Même Code que le Streaming)
from pyflink.datastream import StreamExecutionEnvironment
from pyflink.table import StreamTableEnvironment
# Batch = stream borné
env = StreamExecutionEnvironment.get_execution_environment()
t_env = StreamTableEnvironment.create(env)
# Lecture batch (fichier borné)
t_env.execute_sql("""
CREATE TABLE batch_source (
machine_id STRING,
temperature DOUBLE,
`timestamp` TIMESTAMP(3)
) WITH (
'connector' = 'filesystem',
'path' = 's3://data-lake/raw/sensors/date=2026-07-22/',
'format' = 'parquet'
)
""")
# Même SQL que le streaming
t_env.execute_sql("""
INSERT INTO temp_agg
SELECT machine_id, AVG(temperature), MAX(temperature)
FROM batch_source
GROUP BY machine_id
""").wait()
Pièges Courants (Pitfalls)
Idle sources sans watermark.
Erreur : Une partition Kafka vide bloque toutes les fenêtres car le watermark n'avance pas.
Correction : Augmenter rocksdb.block.cache-size (256-512 MB par TM). Utiliser enable_incremental=True.
Checkpoints trop fréquents.
Erreur : Intervalle de checkpoint trop court (< 30s) → overhead I/O sur RocksDB et réseau.
Correction : Minimum 60s entre checkpoints. Utiliser set_min_pause_between_checkpoints.
Perte d'état après modification de code.
Erreur : Changer le nom d'un state descriptor invalide l'état existant.
Correction : Utiliser le paramètre uid() sur les opérateurs stateful et setUidHash() pour garantir la compatibilité des savepoints. Nommer explicitement valueState avec des UIDs stables.
Pas de gestion des late events.
Erreur : Les événements arrivant après le watermark sont ignorés silencieusement.
Correction : Utiliser sideOutputLateData() pour capturer les événements tardifs dans un flux secondaire.
Liste de Vérification (Checklist)
Watermark configuré (event time processing).
with_idleness() sur les sources potentiellement inactives.
State backend RocksDB pour états > 1 Go.
Checkpoints toutes les 60s minimum avec répertoire externe (S3/HDFS).
Savepoints avant toute mise à jour de code.
uid() défini sur tous les opérateurs stateful.
Side output pour les late events.
Backpressure monitoré dans l'interface Web.
Exactly-once activé (Kafka source + checkpoint).
Tests de reprise avec arrêt/reprise sur savepoint.