skip to content

What does queue.Queue.join() wait for, and how does task_done() feed it?

level: middleimportance: must knowfreq 48%

answer

  1. Two things are being counted, not one
  2. Taking an item is not finishing it
  3. One call in, one call out
  4. Put it in a finally block
  5. Zero is what unblocks the waiter

basics

~20 s

The queue keeps a count of unfinished items: every put increments it and every task_done call decrements it. join blocks until that count reaches zero, so it waits for work to be finished, not merely for the queue to be drained.

solid answer

~40 s

`queue.Queue` maintains an internal unfinished-task counter. Each `put()` adds one; each `task_done()` from a consumer subtracts one; `join()` blocks until the counter hits zero and then returns. The distinction that matters is that `get()` alone does **not** decrement it — draining the queue is not the same as finishing the work, which is precisely why the protocol exists. So a consumer must call `task_done()` exactly once per successful `get()`, and it belongs in a `finally` block: if the handler raises and skips it, the counter never reaches zero and `join()` hangs forever. An extra call raises `ValueError: task_done() called too many times`. Note it is unrelated to `threading.Thread.join()`, which waits for a thread to exit; `queue.SimpleQueue` has neither method.

code

python · 22 lines
python
import queue
import threading

work = queue.Queue()
done = []

def worker():
    while True:
        item = work.get()
        try:
            done.append(item * 2)
        finally:
            work.task_done()          # even if the body raised

for _ in range(3):
    threading.Thread(target=worker, daemon=True).start()

for i in range(10):
    work.put(i)

work.join()                            # returns when all 10 are marked done
print(len(done), sorted(done))

go deeper

for a junior

Recall the pairing: one task_done() for each item you took with get(), and queue.Queue.join() returns once every put item has been reported done. Know that it is a different method from threading.Thread.join().

for a middle

Explain the counter mechanics — put() increments, task_done() decrements, get() does neither — and why the call belongs in a finally. Be able to say what ValueError: task_done() called too many times is telling you.

for a senior

Show the full coordinator sequence: drain with join(), then shut consumers down with sentinels or Queue.shutdown(), then join the threads. Explain why join() is a counting barrier that carries no results or exceptions, and when an executor with futures is the better tool.

for a principal

Decide when a raw queue barrier is the right primitive for the team at all. Weigh it against submitting callables and reading futures, and set the convention for how worker failures are recorded so a completion signal never silently means 'everything failed'.

## A counter, not a queue length `queue.Queue` tracks two different things. One is how many items are sitting in the buffer, which `qsize()` reports. The other is how many items have been *put* but not yet declared *finished* — an internal unfinished-task counter that no public method exposes. `put()` increments it. `task_done()` decrements it and, when it reaches zero, notifies everyone waiting in `join()`. `get()` does not touch it at all. That asymmetry is the whole design. An item that a worker has pulled off the queue is no longer in the buffer but is still in flight; if `join()` returned as soon as the buffer emptied, it would return while workers were still halfway through processing. The counter lets a coordinator wait for the *work* rather than for the *buffer*. ## The worker loop shape ```python def worker(): while True: item = q.get() try: handle(item) finally: q.task_done() ``` The `try/finally` is not decoration. If `handle(item)` raises — a malformed record, a timeout, anything — and `task_done()` is skipped, the counter is permanently short by one and the coordinator's `join()` waits forever for an item that no longer exists anywhere. This hang is one of the classic threaded-Python bugs, and the symptom is maddening: the queue is empty, the workers are idle, the CPU is at zero, and the program never exits. Equally, `task_done()` belongs *inside* the loop, once per `get()`. Calling it once after the loop, or twice for one item, raises `ValueError: task_done() called too many times`, because the counter would go negative. ## Which join is this? Python has two joins with nothing in common, and confusing them is a real interview tell. * `threading.Thread.join()` waits for **that thread** to terminate. * `queue.Queue.join()` waits for **all queued work** to be marked done; it does not know or care which threads exist. The practical consequence is that `Queue.join()` returning is not permission to assume your workers have stopped. They are typically still blocked in `get()` waiting for more. Shutting them down is a separate step: put one sentinel item per worker, or on 3.13+ call `Queue.shutdown()`, which makes every blocked `get()` raise `queue.ShutDown` so each worker can return. Then, if you need to be sure they are gone, `Thread.join()` each of them. ## Ordering the two waits The usual coordinator sequence is: put all the work, `q.join()` to wait for it to be processed, then shut the workers down, then `Thread.join()` them. Reversing the first two — shutting down before the work is finished — loses items. Note that `Queue.shutdown(immediate=True)` marks a task done for every item still in the buffer, which deliberately unblocks a waiting `join()` while discarding that work; the default `shutdown()` does not, so a `join()` that is already waiting still waits for the in-flight items to be reported finished. ## When not to use it The `task_done`/`join` protocol is a *counting* barrier: it tells you that N items were finished, and nothing else. It carries no results and no errors. An exception inside a worker vanishes unless the worker catches and records it, and `join()` will still return happily. When you want results or failures back, a queue is the wrong tool — submit callables to an executor and read the futures, which propagate return values and exceptions to the caller. Reach for `task_done`/`join` when the work is genuinely fire-and-forget and you only need a completion barrier. ## Variants `LifoQueue` and `PriorityQueue` inherit the protocol unchanged. `queue.SimpleQueue` has neither `task_done()` nor `join()` — it is a minimal channel with no bookkeeping — so if you need a completion barrier, `SimpleQueue` is not the class to pick. ## The counter is global to the queue The unfinished-task count belongs to the queue, not to a producer or a batch. If three producer threads all `put()` onto the same queue, one thread's `join()` waits for *all* of their items, including work it never submitted. That is usually what you want in a fan-in pipeline and exactly what you do not want if two independent jobs share a queue: there is no way to wait for only your own items. Use a separate queue per job, or track completion yourself with your own counter and a condition, when the barrier needs to be narrower than the queue. `join()` is also safe to call from more than one thread at once — each waiter is notified when the count reaches zero — and it is re-usable: after it returns, further `put()` calls raise the count above zero again and a later `join()` waits afresh. That makes the protocol workable for batch loops, where each round puts a batch, joins, and then puts the next. One more subtlety worth stating in an interview: because the counter is incremented by `put()` and not by `get()`, an item that is never taken still keeps `join()` waiting. A coordinator that stops its consumers and then calls `join()` will hang on whatever is left in the buffer. Drain first, then stop the workers — or use `shutdown(immediate=True)`, which marks the remaining items done on your behalf and lets the waiter through.

  • A worker raises an exception before reaching task_done() and the program hangs. What exactly is stuck?
    The coordinator's `queue.Queue.join()`. The failed item incremented the unfinished-task counter at `put()` time and nothing ever decremented it, so the counter never reaches zero and `join()` waits forever — with an empty queue and idle workers, which is what makes it confusing. The fix is structural: call `task_done()` in a `finally` block so it runs on every path, and record the failure separately rather than letting it escape the loop.
  • How is queue.Queue.join() different from threading.Thread.join()?
    They share a name and nothing else. `Thread.join()` waits for one thread to terminate. `Queue.join()` waits for the queue's unfinished-task counter to reach zero and knows nothing about threads. `Queue.join()` returning means the work is finished, not that the workers have exited — they are usually still blocked in `get()`. Shut them down separately with sentinels or `Queue.shutdown()`, then `Thread.join()` if you need to be certain they are gone.
  • Can you call queue.Queue.task_done() more times than you called get()?
    No — it raises `ValueError: task_done() called too many times`. The counter is not allowed to go negative, because that would let `join()` return while work was still outstanding. In practice the error surfaces a real bug: `task_done()` placed outside the consumer loop, called in both a `try` and a `finally`, or invoked for an item that was never taken from the queue.

saying these in an interview costs you the question

  • Thinks get() itself decrements the unfinished count
  • Says join() waits for the worker threads to exit
  • Calls task_done() once after the loop ends
  • Skips task_done() when the handler raises
  • Confuses Queue.join() with Thread.join()
  • Expects join() to surface worker exceptions

context