Standardmäßig ist der Prompt ausgewählt, der zuerst die Quelle prüft. Sie können zu einem direkten Befehl wechseln oder eine lokale Kopie herunterladen.
Quelldateien prüfen
Lesen Sie SKILL.md und alle von SkillsMP angezeigten Begleitdateien, bevor Sie sich für eine Installation entscheiden.
Mit Codex oder Claude installieren Kopieren Sie diesen Prompt, fügen Sie ihn in Codex, Claude oder einen anderen Assistant ein und lassen Sie die Skill-Seite prüfen und installieren.
Ein direkter Befehl überspringt den Prüf-Prompt. Prüfen Sie die Quelle, bevor Sie ihn ausführen.
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é.