Shutting Down asyncio Queues with Queue.shutdown in Python 3.13¶
Stopping consumers of an asyncio.Queue used to mean sentinels: put one None per consumer, make every consumer check for it, and hope nobody counted wrong — one sentinel too few leaves a consumer blocked forever, one too many leaves a stray None for the next run. Producers blocked on a full queue had no way to be told to stop at all. Python 3.13 added Queue.shutdown(), which makes the queue itself carry the "no more" signal. In a test with two consumers and four queued items, a graceful shutdown() let the consumers drain all four and then exit via QueueShutDown; shutdown(immediate=True) with six queued items let the two in-progress items finish and dropped the other four, leaving qsize() at 0. A producer blocked on a full queue was released with QueueShutDown, and join() returned at once after an immediate shutdown.
Prerequisites¶
- Python 3.13+ for
Queue.shutdown()andasyncio.QueueShutDown; a sentinel fallback for earlier versions is in step 5. - Queue basics, from Async Queue Management.
- Shutdown ordering, from Graceful Shutdown & Signal Handling.
1. Write consumers that exit on QueueShutDown¶
After shutdown(), get() keeps returning items while any remain (in graceful mode) and raises QueueShutDown once the queue is empty. Consumers catch it as their exit condition:
import asyncio
async def consumer(name: str, q: asyncio.Queue) -> None:
while True:
try:
item = await q.get()
except asyncio.QueueShutDown:
log.info("%s: queue shut down, exiting", name)
return
try:
await handle(item)
finally:
q.task_done()
No sentinel counting, and no special values in the item type: the consumer count can change freely, and every consumer blocked in get() is woken when the queue is shut down and empty.
Verify: with N consumers and M items, every item is handled exactly once and every consumer returns after shutdown.
2. Choose graceful or immediate¶
The two modes answer different shutdown questions:
# graceful: stop accepting, finish everything already queued
q.shutdown()
# immediate: stop accepting, discard everything not yet taken
q.shutdown(immediate=True)
In the graceful case, put() raises QueueShutDown immediately, but get() keeps returning queued items until the queue is empty. In the immediate case, the queue is emptied at once: queued items are discarded, and every waiting get() raises. Items a consumer has already taken keep being processed — measured, the two in-flight items completed — because the queue cannot reach into a consumer's hands.
Use graceful for work that must not be lost and fits in the shutdown budget: emails, writes to a store. Use immediate for work that is safe to drop or will be re-derived — cache refreshes, metrics flushes superseded by the next run — or as the second stage when a graceful drain runs out of time.
Verify: in each mode, count handled versus discarded items; graceful handles all, immediate handles only those already taken.
3. Release blocked producers¶
A bounded queue applies backpressure by blocking put() when full. Before 3.13, a producer blocked there could only be cancelled. Now shutdown() releases it with QueueShutDown:
async def producer(q: asyncio.Queue, source) -> None:
async for item in source:
try:
await q.put(item) # may block on a full queue
except asyncio.QueueShutDown:
log.info("queue closed; stopping producer with %r unsent", item)
return
Verified: a producer blocked on put() into a full maxsize=1 queue returned through this path when the queue was shut down. The producer can then record where it stopped — an offset, a cursor — for the next run, instead of being cancelled mid-statement. Backpressure design for bounded queues is in bounded asyncio queue with backpressure under load.
Verify: fill the queue, block a producer, shut down; the producer exits through its QueueShutDown handler.
4. Combine shutdown with a deadline¶
A production shutdown is usually two-stage: drain gracefully for a bounded time, then drop the rest. join() waits until every queued item has had task_done() called:
async def stop_pipeline(q: asyncio.Queue, consumers: list[asyncio.Task], grace: float = 10.0) -> None:
q.shutdown() # stage 1: no new work, drain
try:
async with asyncio.timeout(grace):
await q.join()
except TimeoutError:
log.warning("drain timed out with %d items left; dropping", q.qsize())
q.shutdown(immediate=True) # stage 2: drop the rest
await asyncio.gather(*consumers, return_exceptions=True)
After an immediate shutdown, join() returns at once — measured at effectively zero — because the discarded items are counted as done. That makes the two-stage pattern safe: the second stage cannot hang waiting for items nobody will process. The surrounding process-level shutdown, including signal handling and in-flight requests, is in draining in-flight requests before shutdown.
Verify: with slow consumers and a short grace period, the pipeline stops within grace plus the time for in-flight items, and logs how many were dropped.
5. Fall back to sentinels before 3.13¶
Libraries supporting 3.11 and 3.12 need the old pattern behind a capability check:
_STOP = object()
HAS_SHUTDOWN = hasattr(asyncio.Queue, "shutdown")
def close_queue(q: asyncio.Queue, consumers: int) -> None:
if HAS_SHUTDOWN:
q.shutdown()
else:
for _ in range(consumers):
q.put_nowait(_STOP) # requires spare capacity or an unbounded queue
async def next_item(q: asyncio.Queue):
if HAS_SHUTDOWN:
try:
return await q.get()
except asyncio.QueueShutDown:
return _STOP
return await q.get()
The sentinel version has the old limitations — it needs free capacity to enqueue the sentinels, and it cannot release blocked producers — which is the strongest argument for raising the floor to 3.13 in queue-heavy code. The capability-check approach to version differences is covered in Asyncio Across Python Versions.
Verify: the same consumer code runs and exits correctly on 3.12 and 3.13+.
Verification¶
Queue shutdown is correct when:
- Every consumer exits via
QueueShutDown(or the sentinel on older versions), with no sentinel counting. - Graceful mode handles every queued item; immediate mode drops only what was not taken.
- Blocked producers are released, and record where they stopped.
- The two-stage shutdown completes within its deadline.
Diagnostic Hook: at shutdown, log three numbers per queue: items handled during the drain, items dropped by the immediate stage, and time spent draining. A non-zero dropped count on routine deploys means the grace period is shorter than the backlog takes to clear, which is either a capacity problem or a grace period that needs raising.
Pitfalls & edge cases¶
- Catching
QueueShutDownaroundtask_done(). It is raised bygetandput, not bytask_done; keep thefinally: task_done()structure. - Calling
shutdown()from another thread. Like other asyncio methods it must run on the loop; useloop.call_soon_threadsafe(q.shutdown). - Expecting immediate mode to cancel handlers. It only drops queued items; in-flight work finishes unless you cancel the consumer tasks.
- Mixing sentinels and shutdown. Pick one per queue, or consumers may see both signals.
Frequently Asked Questions¶
How do I stop asyncio queue consumers in Python 3.13?
Call queue.shutdown(). Consumers receive the remaining items and then QueueShutDown from get(), which they catch as their exit condition. No sentinel values are needed.
What is the difference between shutdown() and shutdown(immediate=True)?
Graceful shutdown stops new puts but keeps delivering queued items. Immediate shutdown also discards queued items at once, so waiting consumers get QueueShutDown straight away and join() returns.
What happens to a producer blocked on a full queue when the queue is shut down?
Its put() raises QueueShutDown, so it can stop cleanly instead of being cancelled. This was not possible before Python 3.13.
Does Queue.join() hang after shutdown?
After an immediate shutdown it returns at once, because discarded items count as done. After a graceful shutdown it waits for task_done() on every remaining item, so bound it with a timeout.
Related¶
- Async Queue Management — up to the topic overview.
- Tracking unfinished work with task_done and join — the counter that join waits on.
- Concurrent Execution & Worker Patterns — the section overview.