An asyncio.Queue coordinates producer and consumer coroutines, with maxsize limiting the number of queued entries.
Python asyncio.Queue: backpressure and completion accounting
Operation contract
The producer awaits each put into a capacity-two queue, so it pauses when buffered work has not been removed. The consumer acknowledges each retrieved item in finally, including the sentinel used to stop it. join waits for acknowledgements, not for proof of a business success. The program records completed receipt IDs in one consumer, which preserves their order.
Failure and ownership boundary
This is a finite local pipeline. A failing consumer needs an explicit supervisor and producer cancellation policy; otherwise a producer or join can wait forever. maxsize bounds queued entries, not the bytes in each entry or already-running work. Python asyncio TaskGroup: cancel sibling work and retain failure evidence and Python asyncio timeout: cancellation and cleanup ownership complete those ownership decisions.
Working program
import asyncio
async def pipeline():
queue = asyncio.Queue(maxsize=2)
completed = []
async def consume():
while True:
receipt_id = await queue.get()
try:
if receipt_id is None:
return
completed.append(receipt_id)
finally:
queue.task_done()
worker = asyncio.create_task(consume())
for receipt_id in (41, 42, 43):
await queue.put(receipt_id)
await queue.put(None)
await queue.join()
await worker
return completed
print(asyncio.run(pipeline()))Output
[41, 42, 43]Costs and limits
For n fixed-size receipts the pipeline does O(n) queue work. Buffered entries stay within capacity, but the displayed completion list retains O(n) IDs.
Common Mistakes
- task_done acknowledges queue bookkeeping, not a successful durable write.
- Bound entry bytes and running work separately from queue length.
Connected lessons
Python asyncio TaskGroup: cancel sibling work and retain failure evidence, Python asyncio timeout: cancellation and cleanup ownership, Python Counter and deque: counts, queues and bounded history.
