Prise en main - Apache Kafka¶
2025 Python
Kafka est comme un hub central pour les données en temps réel, avec les propriétés principales suivantes : plateforme distribuée, tolérante aux pannes, à haut débit et de streaming. Différentes applications peuvent publier des données vers Kafka, et d’autres applications peuvent s’abonner aux données de Kafka. Kafka agit comme un tampon et garantit la livraison fiable des données, même si les applications d’envoi et de réception ne s’exécutent pas en même temps ou ont des vitesses de traitement différentes. Il découple les producteurs et les consommateurs de données, rendant les systèmes plus flexibles et évolutifs.
Première implémentation¶
Démarrer un conteneur avec le service kafka
docker run --rm --name kafka-server --hostname kafka-server \
-e KAFKA_CFG_NODE_ID=0 \
-e KAFKA_CFG_PROCESS_ROLES=controller,broker \
-e KAFKA_CFG_LISTENERS=PLAINTEXT://:9092,CONTROLLER://:9093,EXTERNAL://:9094 \
-e KAFKA_CFG_ADVERTISED_LISTENERS=PLAINTEXT://kafka:9092,EXTERNAL://localhost:9094 \
-e KAFKA_CFG_LISTENER_SECURITY_PROTOCOL_MAP=CONTROLLER:PLAINTEXT,EXTERNAL:PLAINTEXT,PLAINTEXT:PLAINTEXT \
-e KAFKA_CFG_CONTROLLER_QUORUM_VOTERS=0@localhost:9093 \
-e KAFKA_CFG_CONTROLLER_LISTENER_NAMES=CONTROLLER \
-p 9092:9092 -p 9093:9093 -p 9094:9094 bitnami/kafka:latest
Créer un topic
Dans le répertoire /opt/bitnami/kafka/bin, se trouve un ensemble de scripts pour interagir avec Kafka.
Se connecter au conteneur et utiliser le script kafka-topics.sh pour créer un topic.
Le topic nommé duck-topic sera utilisé par le producteur et le consommateur. »
# Connect to the container
docker exec -it kafka-server /bin/bash
# Creation of topic "duck-topic"
kafka-topics.sh --bootstrap-server localhost:9094 --topic duck-topic --create
# Check if the topic was successfully created
kafka-topics.sh --bootstrap-server localhost:9092 --list
Premier script producteur/consommateur en Python
Nous utilisons uv pour développer les deux scripts en python.
pyproject.toml
[project]
name = "Kafka"
requires-python = ">=3.12"
dependencies = [
"kafka-python-ng>=2.2.3",
]
productor.py — uv run .\productor.py
from kafka import KafkaProducer
from kafka.errors import KafkaError
import json
import time
try:
producer = KafkaProducer(
bootstrap_servers=["localhost:9094"],
value_serializer=lambda v: json.dumps(v).encode("utf-8"),
)
for i in range(10): # Produce 10 messages
message = {"message_id": i, "data": f"Message {i}"}
producer.send("duck-topic", message)
print(f"Produced message: {message}")
time.sleep(1) # Send a message every second
producer.flush() # Ensure all messages are sent
print("Finished producing messages.")
except KafkaError as e:
print(f"Error producing messages: {e}")
finally:
if producer is not None:
producer.close()
consumer.py — uv run .\consumers.py
from kafka import KafkaConsumer
from json import loads
consumer = KafkaConsumer(
"duck-topic",
bootstrap_servers=["localhost:9094"],
group_id="my_consumer_group",
value_deserializer=lambda x: (loads(x.decode("utf-8")) if x else None),
key_deserializer=lambda x: x.decode("utf-8") if x else None,
auto_offset_reset="earliest", # or 'latest' or 'none'
enable_auto_commit=False, # Important for manual commits
)
try:
for message in consumer:
print(f"Received message: Key={message.key}, Value={message.value}")
consumer.commit()
except KeyboardInterrupt:
pass # Allow Ctrl+C to exit gracefully
finally:
consumer.close()