Production et consommation asynchrones d’événements Kafka
Intégrez aiokafka à FastAPI pour publier et consommer des événements sans bloquer la boucle d’événements.
Production et consommation asynchrones d’événements Kafka est une leçon FastAPI Backend Development Bootcamp gratuite sur CoddyKit. Ceci est la leçon 1 sur 4. Tu peux lire la leçon complète ci-dessous gratuitement — puis la pratiquer en direct dans le navigateur avec un éditeur de code intégré et un tuteur IA 24/7. Elle fait partie du parcours d'apprentissage FastAPI Backend Development Bootcamp, et ta progression se synchronise sur le web et l'application CoddyKit. Le cours FastAPI Backend Development Bootcamp comprend 4 leçons au total.
Certaines parties de cette leçon n'ont pas encore été traduites et s'affichent en anglais.
Why async Kafka in FastAPI
FastAPI runs on an asyncio event loop. If you publish or poll Kafka with a blocking client (like the standard kafka-python library), every network call freezes the entire loop, stalling all concurrent requests.
- aiokafka is a native asyncio Kafka client that never blocks the loop.
- Its
sendandgetoneoperations are coroutines youawait. - This lets one worker handle thousands of in-flight requests while Kafka I/O is pending.
In this lesson you will wire an AIOKafkaProducer and an AIOKafkaConsumer into a FastAPI app the right way.
Producer lifecycle with lifespan
A producer maintains TCP connections and a background sender task. You must start() it once at app boot and stop() it on shutdown — never per request. The modern FastAPI way is the lifespan context manager.
await producer.start()opens connections and the sender loop.await producer.stop()flushes pending batches and closes cleanly.- Store the producer on
app.stateso routes can reach it.
from contextlib import asynccontextmanager
from fastapi import FastAPI
from aiokafka import AIOKafkaProducer
@asynccontextmanager
async def lifespan(app: FastAPI):
producer = AIOKafkaProducer(
bootstrap_servers="localhost:9092",
enable_idempotence=True,
)
await producer.start()
app.state.producer = producer
try:
yield
finally:
await producer.stop()
app = FastAPI(lifespan=lifespan)Publishing an event from a route
Inside a route, grab the shared producer and await producer.send_and_wait(...). The send_and_wait call returns once the broker has acknowledged the record, giving you back-pressure and delivery confirmation.
- Kafka keys and values are bytes — encode JSON yourself or pass a serializer.
- The returned
RecordMetadatatells you the partition and offset. - Use the message key to guarantee per-entity ordering across partitions.
import json
from fastapi import FastAPI, Request
app = FastAPI()
@app.post("/orders")
async def create_order(payload: dict, request: Request):
producer = request.app.state.producer
value = json.dumps(payload).encode("utf-8")
key = str(payload["order_id"]).encode("utf-8")
meta = await producer.send_and_wait(
"orders", value=value, key=key
)
return {"partition": meta.partition, "offset": meta.offset}Serializers vs manual encoding
Instead of calling json.dumps(...).encode() on every send, you can hand aiokafka a value_serializer and key_serializer. The producer applies them automatically, so routes pass plain Python objects.
value_serializerreceives your object and must returnbytes.- This centralizes encoding and avoids repetition across routes.
- In production teams often swap JSON for Avro or Protobuf via a schema registry.
import json
from aiokafka import AIOKafkaProducer
producer = AIOKafkaProducer(
bootstrap_servers="localhost:9092",
value_serializer=lambda v: json.dumps(v).encode("utf-8"),
key_serializer=lambda k: str(k).encode("utf-8"),
)
# Now routes can send native objects:
# await producer.send_and_wait("orders", value={"id": 7}, key=7)send vs send_and_wait
Two producer methods, two trade-offs:
send()returns a future immediately and lets aiokafka batch records in the background — high throughput, but you do not yet know if delivery succeeded.send_and_wait()awaits the broker ack — slower per call, but gives you a result and surfaces errors right away.
For request handlers where the client needs confirmation, prefer send_and_wait. For fire-and-forget bulk emits, use send and optionally await the futures later.
# Fire many records fast, then wait once for all of them
futures = []
for item in batch:
fut = await producer.send("events", value=item)
futures.append(fut)
# Awaiting the futures surfaces any delivery errors
for fut in futures:
record_meta = await futThe blocking trap to avoid
The single most common mistake is mixing a synchronous client into the async app. Below is a standalone demo of why blocking the loop hurts: a blocking time.sleep inside a coroutine stalls everything, while asyncio.sleep yields control.
Run it and notice the awaited version lets both tasks overlap — exactly the behavior aiokafka gives you for real network I/O.
import asyncio
import time
async def good_io(name):
await asyncio.sleep(0.2) # yields the loop
print(f"{name} done at {time.strftime('%X')}")
async def main():
start = time.perf_counter()
await asyncio.gather(good_io("A"), good_io("B"))
print(f"both finished in {time.perf_counter() - start:.2f}s")
asyncio.run(main())Consumer as a background task
A FastAPI app serves HTTP, but a Kafka consumer must poll continuously. The clean pattern: start the consumer in lifespan and run its poll loop as an asyncio background task, cancelling it on shutdown.
asyncio.create_task(...)launches the loop without blocking startup.- On shutdown, cancel the task, then
await consumer.stop(). - Always wrap the loop body so one bad message does not kill the consumer.
import asyncio
from contextlib import asynccontextmanager
from fastapi import FastAPI
from aiokafka import AIOKafkaConsumer
async def consume(consumer: AIOKafkaConsumer):
async for msg in consumer:
try:
handle(msg.value)
except Exception as exc:
print("handler failed:", exc)
@asynccontextmanager
async def lifespan(app: FastAPI):
consumer = AIOKafkaConsumer(
"orders",
bootstrap_servers="localhost:9092",
group_id="order-workers",
)
await consumer.start()
task = asyncio.create_task(consume(consumer))
try:
yield
finally:
task.cancel()
await consumer.stop()
app = FastAPI(lifespan=lifespan)Iterating messages and decoding
An AIOKafkaConsumer is an async iterator: async for msg in consumer awaits each new record. Each msg exposes topic, partition, offset, key, and value as bytes.
- Decode
msg.valuethe same way you encoded it on the producer side. - You can also pass a
value_deserializerto the consumer constructor. msg.timestampcarries the broker or producer timestamp for latency metrics.
import json
from aiokafka import AIOKafkaConsumer
consumer = AIOKafkaConsumer(
"orders",
bootstrap_servers="localhost:9092",
group_id="order-workers",
value_deserializer=lambda b: json.loads(b.decode("utf-8")),
)
async def run():
async for msg in consumer:
event = msg.value # already a dict
print(event["order_id"], "at offset", msg.offset)Offset commits: at-least-once delivery
Auto-commit (the default) commits offsets on a timer, which can lose messages if the worker crashes after committing but before processing. For reliable processing, set enable_auto_commit=False and commit after you finish handling a message.
- Manual commit gives at-least-once semantics — a crash replays the last uncommitted message.
- Because messages can repeat, your handlers must be idempotent.
- Commit in small batches to balance throughput against replay cost.
consumer = AIOKafkaConsumer(
"orders",
bootstrap_servers="localhost:9092",
group_id="order-workers",
enable_auto_commit=False,
auto_offset_reset="earliest",
)
async def run():
async for msg in consumer:
await process(msg.value) # do the work first
await consumer.commit() # then advance the offsetConcurrency and partition ordering
Within a single partition Kafka preserves order, and an async-for loop processes those records sequentially. To scale, you have two levers:
- More partitions + more consumers in the same
group_id— Kafka assigns partitions across instances automatically. - Bounded concurrency per worker with a semaphore, when per-message work is I/O-heavy and strict ordering is not required.
If ordering per entity matters, keep that entity on one partition via a stable message key and process its partition serially.
import asyncio
sem = asyncio.Semaphore(10)
async def handle_bounded(value):
async with sem:
await do_async_work(value)
async def run(consumer):
async for msg in consumer:
# Schedule work without blocking the poll loop
asyncio.create_task(handle_bounded(msg.value))Graceful shutdown and error handling
A robust deployment cleans up so in-flight data is not lost:
- On shutdown, cancel the consumer task and
await consumer.stop()— this commits offsets (if auto-commit) and leaves the group cleanly so rebalancing is fast. await producer.stop()flushes buffered records before closing.- Wrap per-message handling in try/except and route poison messages to a dead-letter topic instead of crashing the loop.
Never call start()/stop() inside request handlers — that thrashes connections and breaks consumer group membership.
async def consume(consumer, dlq_producer):
async for msg in consumer:
try:
await process(msg.value)
await consumer.commit()
except PermanentError:
await dlq_producer.send_and_wait(
"orders.dlq", value=msg.value, key=msg.key
)
await consumer.commit() # skip the poison messageQuick Check
You need reliable processing where a worker crash must never silently drop an order event. Which consumer configuration best supports this?
Recap
You integrated Kafka into FastAPI without blocking the event loop:
- aiokafka provides awaitable producer and consumer clients native to asyncio.
- Manage the producer and consumer in the lifespan context:
start()at boot,stop()at shutdown — never per request. send_and_waitconfirms delivery;sendmaximizes throughput.- Run the consumer poll loop as an asyncio background task and iterate with
async for. - Disable auto-commit and commit after processing for at-least-once delivery, keeping handlers idempotent.
- Scale with more partitions and consumers in a group, preserve per-entity order via the message key, and shut down gracefully with a dead-letter topic for poison messages.
Questions Fréquemment Posées
La leçon « Production et consommation asynchrones d’événements Kafka » est-elle gratuite ?
Oui — le texte complet de « Production et consommation asynchrones d’événements Kafka » est gratuit à lire ici sur le web. Pour la pratiquer de manière interactive (un éditeur de code intégré et un tuteur IA 24/7) et déverrouiller le reste du cours FastAPI Backend Development Bootcamp, passe à CoddyKit PRO. Le cours FastAPI Backend Development Bootcamp comprend 4 leçons au total.
Qu'est-ce que j'apprendrai dans « Production et consommation asynchrones d’événements Kafka » ?
Intégrez aiokafka à FastAPI pour publier et consommer des événements sans bloquer la boucle d’événements. Tu pratiques FastAPI Backend Development Bootcamp avec du code pratique que tu exécutes directement dans le navigateur, et un tuteur IA 24/7 répond à tes questions au fur et à mesure que tu avances dans la leçon.
Dois-je avoir de l'expérience pour commencer FastAPI Backend Development Bootcamp ?
Aucune expérience préalable n'est requise. FastAPI Backend Development Bootcamp sur CoddyKit est structuré pour les débutants jusqu'aux apprenants avancés, donc tu peux commencer ici ou depuis le début et avancer à ton rythme. Ceci est la leçon 1 sur 4.
Combien de temps prend la leçon « Production et consommation asynchrones d’événements Kafka » ?
La plupart des leçons CoddyKit prennent environ 5–10 minutes. Chacune est courte et interactive, tu progresses régulièrement et tu repiques exactement où tu t'es arrêté sur le web et l'app.
Peux-tu écrire et exécuter du code dans cette leçon FastAPI Backend Development Bootcamp ?
Oui. Chaque leçon FastAPI Backend Development Bootcamp inclut un éditeur de code intégré, tu écris et exécutes du vrai code directement dans ton navigateur et tu reçois des retours IA instantanés — aucune configuration locale requise.
Toutes les leçons de ce cours
- Production et consommation asynchrones d’événements Kafka
- Registre de schémas et évolution des contrats Avro
- Modèle de boîte de sortie transactionnelle
- Consommateurs idempotents et sémantique de livraison unique