What does asyncio.Queue.join() wait for, and what must consumers call?
answer
- Two different notions of “empty”
- A counter of unfinished items
- get() does not decrement it
- Only the consumer can signal completion
- task_done() in a finally, then await join()
basics
~10 sjoin() waits until the queue's unfinished-items counter reaches zero. Every put increments it, and only a consumer calling task_done() decrements it. Consumers that get items but never call task_done() leave join() waiting forever.
solid answer
~40 s`asyncio.Queue` keeps a counter of unfinished items. `put()` and `put_nowait()` increment it; `get()` deliberately does **not** decrement it, because handing an item over is not the same as finishing it. The consumer signals completion with `task_done()`, and `join()` is a coroutine that returns once the counter is back to zero. That split is what makes the canonical worker pool work: start N long-lived worker tasks that loop on `get()`, feed the queue, `await queue.join()` to know all work is genuinely done, then cancel the idle workers. Call `task_done()` in a `finally` so a failing item still counts — otherwise one exception hangs the join permanently. Calling it more times than items were received raises `ValueError: task_done() called too many times`.
code
python · 22 linesimport asyncio
async def worker(q):
while True:
item = await q.get()
try:
await asyncio.sleep(0) # process the item
finally:
q.task_done() # even if processing raised
async def main():
q = asyncio.Queue()
workers = [asyncio.create_task(worker(q)) for _ in range(3)]
for i in range(10):
q.put_nowait(i)
await q.join() # every item accounted for
for w in workers:
w.cancel()
await asyncio.gather(*workers, return_exceptions=True)
print("done")
asyncio.run(main())go deeper
Remember the pairing: every item put on an asyncio.Queue must eventually be answered by one task_done() from whoever processed it, and join() is what waits for that. You will mostly see this in worker-pool examples.
Explain why get() deliberately does not decrement the unfinished counter, and place task_done() in a finally so a failed item still counts. Know that join() waits for completion, not for the buffer to empty.
Show the full lifecycle: bounded queue, fixed worker pool, join for completion, then cancellation or Queue.shutdown() for teardown, plus what happens to in-flight items when the process is asked to stop.
Weigh the pattern itself: a worker pool with explicit completion accounting versus a TaskGroup per item, and whether completion should be tracked in process memory at all when the work must survive a restart.
**The counter.** Every `asyncio.Queue` carries an internal count of *unfinished* items and an internal event that is set while that count is zero. `put()` / `put_nowait()` increment the count. `get()` / `get_nowait()` remove the item from the buffer but leave the count alone — a deliberate design choice, because a consumer that has merely *received* an item has not *processed* it. `task_done()` decrements the count and, when it reaches zero, sets the event that `join()` is awaiting. `join()` is a coroutine: it returns immediately when nothing is outstanding, and otherwise suspends the calling task until the last `task_done()` lands. **Why the split matters.** It is the difference between "the buffer is empty" and "the work is finished". `qsize() == 0` only says nothing is queued right now; workers may still be halfway through items they already pulled. `join()` is the only built-in way to ask the stronger question, and it is why the pattern is worth knowing even though many people reach for a list of tasks and `asyncio.gather` instead. **The canonical worker pool.** Start a fixed number of worker tasks, each looping forever on `await queue.get()`. Feed the queue from a producer. `await queue.join()` to wait for completion. Then cancel the workers — they are parked in `get()` and will never exit on their own — and await them with `return_exceptions=True` to swallow the resulting `CancelledError`. This shape gives you a fixed concurrency level (the worker count) that is independent of the number of items, unlike a fan-out where each item is its own task. **The failure modes, all of which hang.** - *Forgetting `task_done()` entirely.* The count never falls, `join()` never returns, and the program sits there with an empty queue and idle workers. This is the single most common asyncio.Queue bug and it looks like a deadlock with no obvious culprit. - *Calling it only on the success path.* One item raises, the worker's `except` logs and continues, and the count is permanently short by one. Always structure the loop body as `try: ... finally: queue.task_done()`. - *Calling it in the producer.* Only the consumer knows when the item is done; a `task_done()` next to the `put()` makes `join()` return the instant the queue drains. - *Calling it too often.* `task_done()` beyond the number of received items raises `ValueError("task_done() called too many times")` — a loud failure, which is the one merciful case. **Cancellation cuts both ways.** Putting `task_done()` in a `finally` is right, but understand what it means under cancellation: if a worker is cancelled while awaiting inside the processing block, the `finally` still runs and the item is counted as finished even though it was not. A waiting `join()` therefore completes, and the item is silently lost. That is usually the behaviour you want during shutdown — you are tearing down, not reconciling — but if items must not be dropped, the queue is the wrong place to track them: the durable record has to live outside the process, and the queue merely dispatches work that has already been recorded. **When to use something else.** If every item's result is needed, an `asyncio.TaskGroup` (3.11+) or `asyncio.gather` over per-item coroutines is simpler: completion is implicit and results come back directly. `join()` earns its place when workers are long-lived, when items are fed continuously rather than known up front, and when you want a fixed pool rather than one task per item. A sentinel value (`None` pushed once per worker) is the older idiom for shutting workers down; it terminates workers but tells you nothing about completion, so the two are complementary rather than alternatives. Since 3.13, `asyncio.Queue.shutdown()` offers a first-class alternative to sentinels: pending and subsequent `get()`/`put()` calls raise `asyncio.QueueShutDown`, and `shutdown(immediate=True)` also drains outstanding items and marks them done so a pending `join()` is released. **Related detail.** `join()` respects the *counter*, not the consumers: if you `put()` items and never start a consumer at all, `join()` waits forever with a full queue rather than raising. And because `asyncio.Queue` is a plain single-loop object with no locking, none of this is safe to drive from another thread — that role belongs to `queue.Queue`, which offers the same `join()` / `task_done()` contract for threads.
- A worker catches an exception per item and keeps going, and now join() never returns. What is wrong?The `task_done()` call sits on the success path, so any item that raised was never counted as finished and the unfinished-items counter can never reach zero. Move the call into a `finally` around the processing block, so it runs whether the item succeeded, failed or the worker was cancelled mid-item. Failure handling and completion accounting are separate concerns and must not share a code path.
- After `await queue.join()` returns, why do the worker tasks still need cancelling?Because they are infinite loops parked in `await queue.get()`; join() says nothing about them. Cancel each task and await them with `return_exceptions=True` so the resulting CancelledError does not propagate. Since 3.13 the alternative is `queue.shutdown()`, which makes pending and future get() calls raise asyncio.QueueShutDown so the workers can exit on their own.
- Why is `qsize() == 0` not equivalent to “all the work is done”?qsize() reports only what is buffered. Items already handed to workers have left the buffer but are still being processed, so an empty queue routinely coexists with several items in flight. The unfinished-items counter that join() watches includes those, which is exactly why get() does not decrement it and task_done() exists as a separate call.
A restaurant pass is empty when the last plate leaves it, but service is not finished until every waiter reports the table served — task_done() is that report.
saying these in an interview costs you the question
- "get() marks the item as done"
- Calling task_done() in the producer next to the put
- task_done() only on the success path, so one error hangs join() forever
- Treating qsize() == 0 as proof the work finished
- Expecting join() to stop or cancel the worker tasks
- Confusing asyncio.Queue.join() with Thread.join()