Bootcamp backendontwikkeling met FastAPI · Les

Kafka-events asynchroon produceren en consumeren

Integreer aiokafka met FastAPI om events te publiceren en te consumeren zonder de eventloop te blokkeren.

Les 1 van 413 stappen

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 send en getone zijn coroutines waarop je await gebruikt.
  • 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 RecordMetadata vertelt 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_serializer ontvangt je object en moet bytes retourneren.
  • 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 fut

De 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.value op dezelfde manier als waarop je deze aan de producerzijde hebt gecodeerd.
  • Je kunt ook een value_deserializer meegeven aan de constructor van de consumer.
  • msg.timestamp bevat 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 offset

Gelijktijdigheid 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 message

Korte 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 en stop() bij het afsluiten — nooit per aanvraag.
  • send_and_wait bevestigt de aflevering; send maximaliseert 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.
Gratis beginnen

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

  1. Kafka-events asynchroon produceren en consumeren
  2. Schema Registry en Avro-contractevolutie
  3. Het transactional-outboxpatroon
  4. Idempotente consumers en exactly-once-semantiek
← Terug naar Bootcamp backendontwikkeling met FastAPI