A Python Producer and Consumer

A Python Producer and Consumer for BookNest

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:

clients/shipping.py: an asyncio consumer and producer in one servicePython
"""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())
Output
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.