A service inside an asyncio application (FastAPI 36,700 , aiohttp) must not block its event loop in poll(). confluent-kafka 2.12 (October 2025) added AIOProducer and AIOConsumer as a beta: they run the blocking calls in a small thread pool and return awaitables. BookNest's shipping service turns each paid order into a request for the warehouse, keyed by order:
"""shipping.py: an asyncio service that turns paid orders into shipment requests."""
import asyncio, json
from confluent_kafka.aio import AIOConsumer, AIOProducer
B = "localhost:32092"
def shipment(event): # what the warehouse needs
return json.dumps({"order_id": event["order_id"], "paid_at": event["ts"]})
async def main():
consumer = AIOConsumer({"bootstrap.servers": B, "group.id": "booknest-shipping",
"auto.offset.reset": "earliest", "enable.auto.commit": False})
await consumer.subscribe(["booknest.order-events"])
sent = idle = 0
async with AIOProducer({"bootstrap.servers": B, "enable.idempotence": True,
"partitioner": "murmur2_random"}) as producer:
while idle < 3: # stop after 3 empty seconds
batch = await consumer.consume(5000, timeout=1)
idle = 0 if batch else idle + 1
paid = [m for m in batch
if dict(m.headers() or []).get("event-type") == b"order_paid"]
pending = [await producer.produce("booknest.shipments",
shipment(json.loads(m.value())), m.key())
for m in paid]
await producer.flush() # hand the buffer to librdkafka, wait
await asyncio.gather(*pending) # raises if a request was not delivered
if batch:
await consumer.commit(asynchronous=False) # only after the requests landed
sent += len(pending)
await consumer.close()
print(f"{sent:,} shipment requests sent")
asyncio.run(main())94,106 shipment requests sent
One request per paid order. Offsets are committed only after a batch's requests are acknowledged: at-least-once, so the warehouse deduplicates by order_id (Transactions and Exactly-Once's transactions make the pair atomic). AIOProducer buffers 1,000 messages or one second before handing them to librdkafka, and its futures resolve only after that, hence flush(). In 2.15.1 it refuses headers (NotImplementedError), so the shop producer's event-type header still needs the plain Producer.