Apache Kafka est la plateforme de streaming d'événements décentralisée la plus déployée en production. Son architecture log distribuée (« commit log ») permet :
Ingestion en temps réel de millions d'événements/s.
Stockage durable et répliqué avec rétention configurable.
Connecteurs Kafka Connect pour l'intégration avec 200+ systèmes.
Exactly-once semantics entre producers, brokers et consumers.
Cette compétence couvre l'architecture KRaft (ZooKeeper déprécié), le Schema Registry (Avro/Protobuf/JSON Schema), les patterns de production (idempotence, transactions, répartition des partitions), Kafka Connect (source/sink), Kafka Streams (DSL bas niveau et haut niveau), ksqlDB, et l'observabilité (Cruise Control, Kafka Exporter, Burrow).
Quand l'utiliser
Activez cette compétence lorsque l'utilisateur :
Veut configurer un cluster Kafka (KRaft mode) productif.
Demande d'écrire des producers/consumers Python Java ou Go.
Besoin de streaming d'événements entre microservices.
Veut utiliser Kafka Connect pour ingérer/exporter des données.
Pose des questions sur l'ordonnancement, la répartition des partitions, le rebalancing.
A besoin de exactly-once, idempotence ou transactions.
Prérequis
# Client Python complet
pip install kafka-python confluent-kafka faust-streaming avro-python3 fastavro jsonschema
# Confluent CLI (optionnel — gestion de cluster local)
curl -sL https://cnfl.io/cli | sh -s -- -b /usr/local/bin
from confluent_kafka import Producer
import json
conf = {
'bootstrap.servers': 'kafka-1:9092,kafka-2:9092,kafka-3:9092',
'client.id': 'eva-sensor-producer',
'acks': 'all', # Tous les ISR accusent'enable.idempotence': True, # Exactly-once garanti'compression.type': 'snappy',
'linger.ms': 5, # Batch jusqu'à 5ms'batch.size': 65536, # 64KB par batch'max.in.flight.requests.per.connection': 5,
}
producer = Producer(conf)
defdelivery_report(err, msg):
if err:
print(f"ÉCHEC delivery : {err}")
else:
print(f"OK → partition {msg.partition()} | offset {msg.offset()}")
# Envoi avec clé = machine_id (garantit l'ordre par machine)
machine_id = "MACHINE-042"
payload = json.dumps({
"machine_id": machine_id,
"temperature": 87.3,
"pression": 2.15,
"timestamp": "2026-07-22T10:30:00Z"
}).encode('utf-8')
producer.produce(
topic="sensors-raw",
key=machine_id.encode('utf-8'),
value=payload,
callback=delivery_report
)
producer.flush()
2.2 Partitionnement Personnalisé
# Round-robin sans clé# sticky partitioner (défaut depuis Kafka 2.4) : groupe les messages en batch# Avec clé : hash par défaut (murmur2) garantit même partition pour même clé# Utile pour préserver l'ordre des événements d'une même entité# Partitionneur personnalisé (Python)from confluent_kafka import Producer
import hashlib
classSensorPartitioner:
defpartition(topic, key, partitions):
# Envoyer les alertes critiques sur les partitions 0-1ifb"CRITIQUE"in key:
return0# Hash standard pour les autresreturn hashlib.sha256(key).hexdigest().__hash__() % partitions
3. Consumers : Lecture d'Événements
3.1 Consumer fiable avec Gestion des Offsets
from confluent_kafka import Consumer, KafkaError
import json
conf = {
'bootstrap.servers': 'kafka-1:9092,kafka-2:9092',
'group.id': 'eva-sensor-processor',
'auto.offset.reset': 'earliest',
'enable.auto.commit': False, # Commit manuel'max.poll.interval.ms': 300000, # 5 min max'session.timeout.ms': 45000,
'heartbeat.interval.ms': 15000,
'fetch.min.bytes': 1024, # Au moins 1KB par fetch'fetch.max.wait.ms': 500,
'isolation.level': 'read_committed', # Évite les messages transactionnels non commités
}
consumer = Consumer(conf)
consumer.subscribe(['sensors-raw'])
try:
whileTrue:
msg = consumer.poll(1.0) # timeout 1sif msg isNone:
continueif msg.error():
if msg.error().code() == KafkaError._PARTITION_EOF:
continueelse:
print(f"Erreur : {msg.error()}")
break
data = json.loads(msg.value().decode('utf-8'))
print(f"{msg.key()} | p{msg.partition()}@{msg.offset()} | {data['temperature']}°C")
# Commit manuel après traitement réussi
consumer.commit(asynchronous=False)
finally:
consumer.close()
3.2 Reprise sur Erreur (Dead Letter Topic)
defprocess_message(msg):
try:
data = json.loads(msg.value())
# Traitement métierif data.get("temperature", 0) > 200:
raise ValueError("Température hors échelle")
return data
except Exception as e:
# Envoi vers Dead Letter Queue
dlq_producer.produce(
topic="sensors-dlq",
key=msg.key(),
value=msg.value(),
headers={"error": str(e), "original_topic": msg.topic()}
)
dlq_producer.flush()
returnNone
# KTable : table partitionnée (joindre avec un flux partitionné de la même clé)
machine_table = app.Table('machine-metadata', default=dict)
# GlobalKTable : copie complète sur tous les nœuds (joindre avec n'importe quelle clé)
global_config = app.GlobalTable('global-config', default=dict)
# Détection de déséquilibre
bin/kafka-cruise-control-start.sh \
--config config/cruisecontrol.properties \
--port 9090
# Proposer un plan de rebalancement
curl "http://localhost:9090/kafkacruisecontrol/rebalance?dryRun=true&verbose=true"# Exécuter
curl -X POST "http://localhost:9090/kafkacruisecontrol/rebalance?dryRun=false"
8.3 Burrow (Consumer Lag)
# Burrow surveille le lag de tous les consumers# Configuration :
[consumer.eva-sensor-processor]
servers = "kafka-1:9092,kafka-2:9092"
group = "eva-sensor-processor"# Endpoint HTTP : /v3/kafka/{cluster}/consumer/{group}/lag
curl http://burrow:8000/v3/kafka/local/consumer/eva-sensor-processor/lag
Erreur :session.timeout.ms trop bas ou max.poll.interval.ms trop court → le consumer est exclu du groupe.
Correction : Augmenter session.timeout.ms (45s-60s) et max.poll.interval.ms (5min+). Vérifier le temps de traitement des messages.
Messages dupliqués sans idempotence.
Erreur : Le producer retente un message après timeout réseau, et Kafka l'accepte deux fois.
Correction : Toujours activer enable.idempotence=true sur le producer.
Topics sans réplication (RF=1).
Erreur : Un broker tombe → perte totale de données.
Correction :default.replication.factor=3 et min.insync.replicas=2.
Ordre des messages non garanti.
Erreur : Plusieurs partitions sans clé → ordre non préservé entre partitions.
Correction : Utiliser une clé pour garantir l'ordre par entité (ex : machine_id). Une seule partition peut aussi garantir l'ordre global, mais limite le débit.
Lag de consommation non surveillé.
Erreur : Le consumer lag augmente silencieusement jusqu'à dépasser la rétention → perte de messages.
Correction : Déployer Burrow ou Kafka Lag Exporter, alerting Prometheus.
Évolution de schéma sans compatibilité.
Erreur : Ajout d'un champ requis → les consumers anciens crash.
Correction : Toujours définir default pour les nouveaux champs. Respecter la règle de compatibilité (BACKWARD par défaut).
Liste de Vérification (Checklist)
KRaft configuré (plus de ZooKeeper).
enable.idempotence=true sur tous les producers.
replication.factor >= 3 pour les topics critiques.
min.insync.replicas=2 pour garantir la durabilité.