Skip to content

Draining an Async Worker Pool on Shutdown

A worker pool that stops abruptly loses work: items in flight are interrupted halfway, items still queued vanish with the process. A pool that waits politely for everything can hold a deploy for minutes. The goal of a drain is narrower and achievable: stop taking new work, let the work that has started finish within a deadline, and give back everything that did not start, so another worker can pick it up. Tested with 300 queued items of 50 ms and 10 workers, a drain with a 0.2 s grace period completed 60 items, interrupted 10 that were in flight when the deadline hit, and returned 230 unstarted items — all 300 accounted for; with a 2 s grace, all 300 completed in 1.39 s. This guide builds the drain step by step.

Prerequisites

1. Stop intake first

The first action on shutdown is to stop new work entering the pool. With an asyncio.Queue, shutdown() makes further put() calls raise QueueShutDown while consumers keep receiving what is already queued:

import asyncio


async def worker(q: asyncio.Queue, handle) -> None:
    while True:
        try:
            item = await q.get()
        except asyncio.QueueShutDown:
            return                                  # nothing left, and nothing more will come
        try:
            await handle(item)
        finally:
            q.task_done()


def begin_shutdown(q: asyncio.Queue) -> None:
    q.shutdown()                                    # producers stop; consumers drain

Upstream of the queue, stop whatever feeds it: unsubscribe from the broker, stop polling the database, close the listening socket. If the queue is fed by an HTTP endpoint, the endpoint should start returning 503 at the same moment, which is the readiness change described in draining in-flight requests before shutdown.

Verify: after begin_shutdown(), producers receive QueueShutDown and the queue length only goes down.

2. Drain within a deadline

Wait for the queue to empty, bounded by a grace period shorter than whatever the platform gives you before a hard kill:

async def drain(q: asyncio.Queue, workers: list[asyncio.Task], grace: float) -> list:
    q.shutdown()
    returned: list = []
    try:
        async with asyncio.timeout(grace):
            await q.join()                          # every queued item processed
    except TimeoutError:
        while not q.empty():                        # unstarted work: hand it back
            returned.append(q.get_nowait())
            q.task_done()
        for w in workers:
            w.cancel()                              # in-flight work: interrupt
    await asyncio.gather(*workers, return_exceptions=True)
    return returned

join() returns when every item has had task_done() called, so it covers in-flight items as well as queued ones. When the deadline expires, queued items are pulled out and returned to the caller to persist or re-enqueue elsewhere, and only then are workers cancelled. Measured: with 0.2 s, 60 completed, 10 interrupted, 230 returned; with 2 s, all 300 completed in 1.39 s.

Verify: the sum of completed, interrupted and returned items equals what was in the pool when shutdown began.

300 items at shutdown, 0.2 s grace 3 horizontal bars comparing returned unstarted with the others. 300 items at shutdown, 0.2 s grace returned unstarted 230 completed 60 interrupted in flight 10 10 workers, 50 ms per item; with a 2 s grace all 300 completed in 1.39 s. Every item ends in exactly one bucket; the drain's job is to make that true.

3. Make interrupted items safe to redo

Items interrupted at the deadline were partly processed. They must either be safe to run again — idempotent — or be rolled back by their own cleanup:

async def handle(item) -> None:
    try:
        async with db.transaction():
            await apply(item)                       # all-or-nothing in the database
    except asyncio.CancelledError:
        await requeue_later(item)                   # put it back for another worker
        raise

A transaction makes the database side atomic: cancellation rolls it back. Anything else — an external API call already made — needs an idempotency key so the retry does not repeat the effect, as in making background jobs idempotent. If an item cannot tolerate interruption at all, the grace period must be longer than the longest item, and the deadline should be applied before starting items rather than to their completion.

Verify: interrupt a handler mid-item in a test; the item is re-queued and its partial effects were rolled back.

A drain with a deadline 4 lanes over time. A drain with a deadline intake accepting shut down workers processing draining the queue cancel in-flight remaining items returned for re-queue platform grace period before SIGKILL time → The drain deadline must end comfortably before the platform's own kill deadline.

4. Put the returned items somewhere durable

Returning unstarted items is only useful if they survive the process. Where they go depends on where they came from:

async def shutdown_pool(q, workers, source, grace: float = 20.0) -> None:
    leftovers = await drain(q, workers, grace)
    if not leftovers:
        return
    if source.kind == "broker":
        await source.nack(leftovers)                # RabbitMQ/SQS: return to the queue
    elif source.kind == "database":
        await source.release_claims(leftovers)      # clear claimed_by so others pick them up
    else:
        await source.persist(leftovers)             # last resort: write to a recovery table
    log.info("returned %d unstarted items to %s", len(leftovers), source.kind)

With a message broker, unacknowledged messages are redelivered when the consumer disconnects, so "returning" can be as simple as not acking them — but explicit negative acknowledgement makes redelivery immediate rather than waiting for a timeout. With a database-backed queue, release the claims, as in building a durable job queue on Postgres with asyncio. An in-memory queue fed by nothing durable cannot return work anywhere — which is the argument for not holding important work only in memory.

Verify: after a shutdown with leftovers, those items are processed by another instance.

5. Wire it to signals and deadlines

Connect the drain to the process's shutdown signal and pick the grace from the platform's limit:

import signal


async def main() -> None:
    q: asyncio.Queue = asyncio.Queue(maxsize=1000)
    workers = [asyncio.create_task(worker(q, handle)) for _ in range(16)]
    feeder = asyncio.create_task(feed_from_broker(q))
    stop = asyncio.Event()
    loop = asyncio.get_running_loop()
    loop.add_signal_handler(signal.SIGTERM, stop.set)

    await stop.wait()
    feeder.cancel()                                  # stop pulling from the broker
    await asyncio.gather(feeder, return_exceptions=True)
    await shutdown_pool(q, workers, broker_source, grace=20.0)   # Kubernetes default grace: 30 s

Kubernetes sends SIGTERM and waits terminationGracePeriodSeconds (30 s by default) before SIGKILL; leave several seconds of margin for the rest of shutdown — closing pools, flushing telemetry. The full pod-level sequence is in shutting down asyncio pods in Kubernetes.

Verify: kubectl delete pod during load results in no lost items and an exit before the grace period ends.

The four steps of a pool drain A flow of 4 stages. The four steps of a pool drain stop intake feeder + queue.shutdown() drain with a deadline await q.join() return the rest nack or release claims cancel stragglers exit before SIGKILL Each step bounds the next; the deadline is what keeps a slow item from holding the deploy.

Verification

The drain is correct when:

  • Intake stops first, at the source and at the queue.
  • Every item is accounted for: completed, interrupted and re-queued, or returned unstarted.
  • Interrupted items are idempotent or rolled back.
  • The drain finishes before the platform's kill deadline.

Diagnostic Hook: log the drain's accounting — completed, interrupted, returned, and duration — on every shutdown and export it as metrics. Interrupted and returned counts that are routinely non-zero mean the grace period is shorter than the backlog takes to clear; a duration that creeps toward the platform's limit is an incident waiting for a slow day.

Pitfalls & edge cases

  • Cancelling workers before stopping intake. New items arrive with nobody to process them.
  • No deadline on join(). One stuck item holds shutdown until the hard kill.
  • Returning items to an in-memory queue. They die with the process; return them to a durable source.
  • Grace longer than the platform's. SIGKILL arrives mid-drain and nothing is accounted for.

Frequently Asked Questions

How do I shut down an asyncio worker pool without losing work?

Stop intake, wait for the queue to drain with a deadline, return unstarted items to a durable source, then cancel any workers still running. In testing, a 0.2 s drain completed 60 items, interrupted 10 and returned 230, with none lost.

How long should the drain grace period be?

Long enough for typical backlogs to finish, and several seconds shorter than the platform's hard-kill deadline — Kubernetes allows 30 seconds by default.

What happens to items that are in progress at the deadline?

They are cancelled. Make handlers idempotent or transactional so cancelled items can be safely retried, and re-queue them from the cancellation handler.

Do I need Queue.shutdown to drain a pool?

No, but it helps on Python 3.13+: it stops producers and wakes idle consumers. On earlier versions, stop producers yourself and send one sentinel per worker.