Tracking Unfinished Work with task_done and join¶
asyncio.Queue.join() waits until every item that was ever put into the queue has been processed — not merely taken. It works through a counter: put() increments it, task_done() decrements it, and join() returns when it reaches zero. The contract is simple and unforgiving. A consumer that skips task_done() on one code path leaves the counter above zero forever: in a test with five items, a consumer that forgot task_done() when it skipped one item left one unfinished task, and join() hung until its 0.5 s timeout. Calling it once too often raises ValueError: task_done() called too many times. Putting task_done() in a finally block made the same consumer — including a path that raised and was handled — join cleanly. This guide covers the counter, the patterns that keep it right, and how to wait on it safely.
Prerequisites¶
- Python 3.11+, stdlib only.
- Queue basics, from Async Queue Management.
- Worker pools, from building an async worker pool with TaskGroup.
1. Understand the counter¶
q = asyncio.Queue()
q.put_nowait("a") # unfinished = 1
q.put_nowait("b") # unfinished = 2
item = await q.get() # unfinished = 2 — taking does not count as finishing
q.task_done() # unfinished = 1
item = await q.get()
q.task_done() # unfinished = 0 → join() returns
await q.join()
get() and task_done() are separate on purpose: an item is "unfinished" from the moment it is put until a consumer says it is done, so join() waits for in-flight processing, not just for an empty queue. qsize() and the unfinished count are different numbers — a queue can be empty with items still being processed — which is why "wait until qsize() is zero" is the wrong way to wait for completion.
The counter is not tied to items: task_done() decrements it no matter which item you just finished. One call per get() keeps it balanced; anything else breaks it.
Verify: after processing a batch, q._unfinished_tasks (private, for inspection in tests) is zero and join() returns immediately.
2. Call task_done in a finally block¶
Every way out of processing an item must call task_done() exactly once. Structure the consumer so there is only one place it happens:
async def consumer(q: asyncio.Queue) -> None:
while True:
item = await q.get()
try:
if should_skip(item):
continue # finally still runs
await process(item)
except Exception:
log.exception("item %r failed", item)
finally:
q.task_done() # exactly once per get()
A continue inside try runs the finally before continuing, so the skip path is covered. Cancellation is covered too: if the consumer is cancelled while processing, the finally runs and the item is counted as done — which is what you want for join() to return during shutdown, though it means the item was not actually processed. If cancelled items should be retried, re-queue them before task_done().
The bug from the introduction was a continue placed before the try, skipping the decrement. That is the whole class of bug: any return path outside the try/finally.
Verify: add a skip path and an exception path to a test; join() still returns.
3. Wait on join with a deadline¶
join() has no timeout and no progress reporting, and a stuck consumer makes it wait forever. Wrap it:
async def wait_until_done(q: asyncio.Queue, timeout: float, report_every: float = 5.0) -> None:
joiner = asyncio.ensure_future(q.join())
deadline = asyncio.get_running_loop().time() + timeout
try:
while not joiner.done():
remaining = deadline - asyncio.get_running_loop().time()
if remaining <= 0:
raise TimeoutError(f"{q._unfinished_tasks} items still unfinished")
done, _ = await asyncio.wait({joiner}, timeout=min(report_every, remaining))
if not done:
log.info("waiting: %d queued, %d unfinished", q.qsize(), q._unfinished_tasks)
finally:
joiner.cancel()
Cancelling the join() waiter does not affect the queue; it only stops this caller waiting. Logging both qsize() and the unfinished count every few seconds distinguishes "consumers are slow" (both falling) from "something is stuck" (queue empty, unfinished not falling) — the second is nearly always a missed task_done() or a consumer hung on an await without a timeout.
Verify: with a consumer deliberately stuck, the function raises after timeout and the logs show unfinished stuck above zero with an empty queue.
4. Keep consumers alive, or join waits for nobody¶
join() cannot know whether anyone is still consuming. If every consumer task crashed on an unexpected exception, the queue still holds unfinished items and join() waits forever. Two defences: catch per-item exceptions inside the consumer loop (step 2), and run consumers inside a TaskGroup alongside the join(), so a consumer crash cancels the wait:
async def run_batch(items, workers: int = 8) -> None:
q: asyncio.Queue = asyncio.Queue()
for it in items:
q.put_nowait(it)
async with asyncio.TaskGroup() as tg:
consumers = [tg.create_task(consumer(q)) for _ in range(workers)]
await q.join() # all items processed
for c in consumers:
c.cancel() # consumers loop forever; stop them
If a consumer raises outside its per-item handler, the TaskGroup cancels the join() and the batch fails visibly instead of hanging. On Python 3.13+, replace the cancellation loop with q.shutdown() so consumers exit through their own code path, as described in shutting down queues with Queue.shutdown.
Verify: inject a bug that crashes a consumer outside the per-item try; run_batch raises instead of hanging.
5. Use join for batches, not for continuous pipelines¶
join() answers "is everything that was ever put now done?" That is the right question for a batch job and the wrong one for a service whose producers never stop — the counter only reaches zero momentarily between bursts. For continuous pipelines, measure progress instead:
class TrackedQueue(asyncio.Queue):
def __init__(self, *a, **kw) -> None:
super().__init__(*a, **kw)
self.done_total = 0
def task_done(self) -> None:
super().task_done()
self.done_total += 1
Exporting done_total as a counter and the unfinished count as a gauge gives throughput and in-flight work directly, which is what the monitoring in monitoring queue depth and item age builds on. Reserve join() for drain-at-shutdown and batch completion.
Verify: in a continuous service, dashboards use throughput and in-flight counts, and join() appears only in shutdown and batch code.
Verification¶
task_done/join usage is correct when:
- Every
get()is matched by exactly onetask_done(), in afinallyblock. join()is always bounded by a timeout or a TaskGroup that can cancel it.- Stuck waits are diagnosable from logged
qsize()and unfinished counts. - Continuous pipelines use throughput metrics rather than
join().
Diagnostic Hook: expose the unfinished count as a gauge next to qsize(). The difference between them is items in flight; if it equals the number of consumers and stays there, every consumer is stuck on one item — typically an await with no timeout. If the queue is empty and the difference is non-zero with idle consumers, a code path skipped task_done().
Pitfalls & edge cases¶
continueorreturnbefore thetry. The classic skippedtask_done().- Calling
task_done()for items obtained withoutget(). RaisesValueErroronce the counter would go negative. - Re-queueing failed items after
task_done(). Put the item back first, then calltask_done(), orjoin()can return with work outstanding. - Waiting on
qsize() == 0. Returns while items are still being processed.
Frequently Asked Questions¶
What does asyncio.Queue.task_done do?
It tells the queue that one item previously taken with get() has been fully processed, decrementing the count of unfinished items that join() waits on.
Why does asyncio.Queue.join hang forever?
Usually a consumer took an item and never called task_done() for it, often on a skip or error path outside a try/finally. It also hangs if every consumer has crashed while items remain.
What does 'task_done() called too many times' mean?
task_done() was called more often than get() returned items, so the unfinished counter would go below zero. Call it exactly once per get(), in a finally block.
How do I add a timeout to Queue.join?
Wrap it: async with asyncio.timeout(seconds): await q.join(). Cancelling the wait does not affect the queue, and logging qsize and the unfinished count while waiting shows why it is slow.
Related¶
- Async Queue Management — up to the topic overview.
- Implementing a dead letter queue with asyncio — where failed items go instead of being retried forever.
- Concurrent Execution & Worker Patterns — the section overview.