Kafka-events asynchroon produceren en consumeren
Integreer aiokafka met FastAPI om events te publiceren en te consumeren zonder de eventloop te blokkeren.
Kafka-events asynchroon produceren en consumeren is een gratis Bootcamp backendontwikkeling met FastAPI-les op CoddyKit. Dit is les 1 van 4. Je kunt 3 lessen uit dit leerpad gratis volledig lezen — daarna ontgrendelt CoddyKit PRO alle lessen, plus praktische oefeningen met een ingebouwde code-editor en een AI-tutor die 24/7 beschikbaar is. Deze les maakt deel uit van het leertraject Bootcamp backendontwikkeling met FastAPI. Je voortgang wordt gesynchroniseerd op het web en in de CoddyKit-app. De cursus Bootcamp backendontwikkeling met FastAPI bevat in totaal 4 lessen.
Waarom asynchrone Kafka in FastAPI
FastAPI draait op een asyncio-eventlus. Als je Kafka publiceert of bevraagt met een blokkerende client (zoals de standaardbibliotheek kafka-python), bevriest elke netwerkoproep de volledige lus en worden alle gelijktijdige verzoeken opgehouden.
- aiokafka is een native asyncio-Kafka-client die de lus nooit blokkeert.
- De bewerkingen
sendengetonezijn coroutines waarop jeawaitgebruikt. - Daardoor kan één worker duizenden verzoeken tegelijk verwerken terwijl Kafka-I/O wacht.
In deze les koppel je een AIOKafkaProducer en een AIOKafkaConsumer op de juiste manier aan een FastAPI-app.
De levenscyclus van een producer met lifespan
Een producer onderhoudt TCP-verbindingen en een achtergrondtaak voor het verzenden. Je moet deze eenmaal starten bij het opstarten van de app met start() en stoppen bij het afsluiten met stop() — nooit per verzoek. De moderne FastAPI-manier is de contextmanager lifespan.
await producer.start()opent de verbindingen en de verzendlus.await producer.stop()verzendt wachtende batches en sluit netjes af.- Sla de producer op in
app.state, zodat routes erbij kunnen.
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)Een gebeurtenis publiceren vanuit een route
Haal binnen een route de gedeelde producer op en gebruik await producer.send_and_wait(...). De aanroep send_and_wait retourneert zodra de broker het record heeft bevestigd, zodat je tegendruk en bevestiging van de aflevering krijgt.
- Kafka-sleutels en -waarden zijn bytes — codeer JSON zelf of geef een serializer door.
- De geretourneerde
RecordMetadatavertelt je de partitie en offset. - Gebruik de sleutel van het bericht om de volgorde per entiteit over partities heen te garanderen.
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 tegenover handmatige codering
In plaats van bij elke verzending json.dumps(...).encode() aan te roepen, kun je aiokafka een value_serializer en key_serializer meegeven. De producer past deze automatisch toe, zodat routes gewone Python-objecten doorgeven.
value_serializerontvangt je object en moetbytesretourneren.- Hiermee centraliseer je de codering en voorkom je herhaling in routes.
- In productieteams wordt JSON vaak vervangen door Avro of Protobuf via een schemaregister.
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 tegenover send_and_wait
Twee methoden van de producer, twee afwegingen:
send()retourneert onmiddellijk een future en laat aiokafka records op de achtergrond bundelen — hoge verwerkingscapaciteit, maar je weet nog niet of de aflevering is gelukt.send_and_wait()wacht op de bevestiging van de broker — langzamer per aanroep, maar je krijgt een resultaat en fouten komen meteen aan het licht.
Geef voor verzoekafhandelaars waarbij de client bevestiging nodig heeft de voorkeur aan send_and_wait. Gebruik voor bulkverzendingen zonder verdere actie send en wacht eventueel later op de futures.
# 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 futDe blokkeringsvalkuil die je moet vermijden
De meest voorkomende fout is een synchrone client mengen met de asynchrone app. Hieronder staat een zelfstandig uitvoerbaar voorbeeld van waarom het blokkeren van de lus schadelijk is: een blokkerende time.sleep in een coroutine houdt alles tegen, terwijl asyncio.sleep de besturing overdraagt.
Voer het uit en merk op dat de versie met await beide taken laat overlappen — precies het gedrag dat aiokafka biedt voor echte netwerk-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())De consumer als achtergrondtaak
Een FastAPI-app serveert HTTP, maar een Kafka-consumer moet continu pollen. Het nette patroon is: start de consumer in lifespan en voer de poll-lus uit als een asyncio-achtergrondtaak; annuleer deze bij het afsluiten.
asyncio.create_task(...)start de lus zonder het opstarten te blokkeren.- Annuleer bij het afsluiten de taak en gebruik daarna
await consumer.stop(). - Omwikkel altijd de body van de lus, zodat één ongeldig bericht de consumer niet beëindigt.
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)Berichten doorlopen en decoderen
Een AIOKafkaConsumer is een asynchrone iterator: async for msg in consumer wacht op elk nieuw record. Elk msg bevat topic, partition, offset, key en value als bytes.
- Decodeer
msg.valueop dezelfde manier als waarop je deze aan de producerzijde hebt gecodeerd. - Je kunt ook een
value_deserializermeegeven aan de constructor van de consumer. msg.timestampbevat het tijdstip van de broker of producer voor latentiemetingen.
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)Offsets vastleggen: aflevering minstens één keer
Automatisch vastleggen (de standaardinstelling) legt offsets op een timer vast. Daardoor kunnen berichten verloren gaan als de worker crasht nadat de offset is vastgelegd maar voordat het bericht is verwerkt. Stel voor betrouwbare verwerking enable_auto_commit=False in en leg de offset na het verwerken van een bericht vast.
- Handmatig vastleggen geeft de semantiek minstens één keer — bij een crash wordt het laatste niet-vastgelegde bericht opnieuw afgespeeld.
- Omdat berichten kunnen worden herhaald, moeten je afhandelaars idempotent zijn.
- Leg offsets vast in kleine batches om verwerkingscapaciteit en kosten van opnieuw afspelen tegen elkaar af te wegen.
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 offsetGelijktijdigheid en volgorde binnen partities
Binnen één partitie bewaart Kafka de volgorde en verwerkt een async-for-lus deze records sequentieel. Om op te schalen heb je twee mogelijkheden:
- Meer partities en meer consumers in dezelfde
group_id— Kafka wijst partities automatisch toe aan de instanties. - Begrensde gelijktijdigheid per worker met een semafoor, wanneer het werk per bericht veel I/O vereist en een strikte volgorde niet nodig is.
Als de volgorde per entiteit belangrijk is, houd je die entiteit op één partitie via een stabiele berichtsleutel en verwerk je de partitie serieel.
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))Gecontroleerd afsluiten en foutafhandeling
Een robuuste implementatie ruimt alles op, zodat gegevens die nog worden verwerkt niet verloren gaan:
- Annuleer bij het afsluiten de consumer-taak en gebruik
await consumer.stop()— dit legt offsets vast (bij automatisch vastleggen) en verlaat de groep netjes, zodat herverdeling snel verloopt. await producer.stop()verzendt gebufferde records voordat de verbinding wordt gesloten.- Omwikkel de verwerking per bericht met try/except en stuur onherstelbare berichten naar een dead-letter-topic in plaats van de lus te laten crashen.
Roep start()/stop() nooit aan in verzoekafhandelaars — daardoor worden verbindingen onnodig steeds opnieuw gestart en wordt het lidmaatschap van de consumergroep verbroken.
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 messageKorte controle
Je hebt betrouwbare verwerking nodig waarbij een crash van een worker nooit stilzwijgend een gebeurtenis over een bestelling mag laten vallen. Welke configuratie van de consumer ondersteunt dit het best?
Samenvatting
Je hebt Kafka in FastAPI geïntegreerd zonder de gebeurtenislus te blokkeren:
- aiokafka biedt awaitable producer- en consumerclients die native zijn voor asyncio.
- Beheer de producer en consumer in de levensduurcontext:
start()bij het opstarten enstop()bij het afsluiten — nooit per aanvraag. send_and_waitbevestigt de aflevering;sendmaximaliseert de doorvoer.- Voer de poll-lus van de consumer uit als een asyncio-achtergrondtaak en iterereer met
async for. - Schakel automatisch committen uit en commit na de verwerking voor ten minste één keer-aflevering, en houd handlers idempotent.
- Schaal met meer partities en consumers in een groep, behoud de volgorde per entiteit via de berichtsleutel en sluit netjes af met een dead-lettertopic voor onbruikbare berichten.
Leer Bootcamp backendontwikkeling met FastAPI met een AI-tutor — gratis
Schrijf echte code en voer die uit in je browser, krijg direct hulp van een AI-tutor die 24/7 beschikbaar is en ga verder waar je gebleven bent op het web of in de app.
- Cursussen
- 21
- Lessen
- 84
Veelgestelde vragen
Is de les “Kafka-events asynchroon produceren en consumeren” gratis?
Ja — je kunt hier op het web alle 3 lessen van het leerpad Bootcamp backendontwikkeling met FastAPI, waaronder “Kafka-events asynchroon produceren en consumeren”, gratis volledig lezen. Daarna ontgrendelt CoddyKit PRO alle lessen, plus interactieve oefeningen met een ingebouwde code-editor en een AI-tutor die 24/7 beschikbaar is. De cursus Bootcamp backendontwikkeling met FastAPI bevat in totaal 4 lessen.
Wat leer ik in “Kafka-events asynchroon produceren en consumeren”?
Integreer aiokafka met FastAPI om events te publiceren en te consumeren zonder de eventloop te blokkeren. Je oefent met Bootcamp backendontwikkeling met FastAPI door code rechtstreeks in de browser uit te voeren. Een AI-begeleider die 24/7 beschikbaar is beantwoordt je vragen terwijl je de les doorwerkt.
Heb ik ervaring nodig om met Bootcamp backendontwikkeling met FastAPI te beginnen?
Ervaring vooraf is niet nodig. Bootcamp backendontwikkeling met FastAPI op CoddyKit is opgebouwd voor beginners tot gevorderden, zodat je hier of bij het begin kunt starten en in je eigen tempo kunt leren. Dit is les 1 van 4.
Hoe lang duurt de les “Kafka-events asynchroon produceren en consumeren”?
De meeste lessen van CoddyKit duren ongeveer 5–10 minuten. Elke les is kort en interactief, zodat je gestaag vooruitgaat en op het web en in de app precies verdergaat waar je was gebleven.
Kan ik code schrijven en uitvoeren in deze les over Bootcamp backendontwikkeling met FastAPI?
Ja. Elke les over Bootcamp backendontwikkeling met FastAPI bevat een ingebouwde code-editor, zodat je rechtstreeks in je browser echte code kunt schrijven en uitvoeren en direct feedback van AI krijgt — lokale installatie is niet nodig.
Alle lessen in deze cursus
- Kafka-events asynchroon produceren en consumeren
- Schema Registry en Avro-contractevolutie
- Het transactional-outboxpatroon
- Idempotente consumers en exactly-once-semantiek