Skip to content

Pausing and Resuming an Async Worker Pool

There are good reasons to stop a worker pool without stopping the process: the database is in a maintenance window, a downstream API is returning errors and every attempt only adds load, a circuit breaker has opened, or an operator wants to stop consuming while investigating. Stopping the process loses the warm state and the queue; cancelling workers interrupts items. A pause switch does neither: workers finish their current item and then wait, and on resume they continue with the queue untouched. Tested with 4 workers processing 50 ms items, pause() returned after 30 ms — the time for in-flight items to complete — no items were processed during a 0.5 s pause, and after resume() all 100 items completed. This guide builds that switch, wires it to a circuit breaker, and covers the details that make it reliable.

Prerequisites

1. Gate workers on an Event

An asyncio.Event that is set while running and cleared while paused is a natural gate. Workers wait on it before taking each item:

import asyncio


class PausablePool:
    def __init__(self, handle, workers: int = 4) -> None:
        self.q: asyncio.Queue = asyncio.Queue()
        self.handle = handle
        self.running = asyncio.Event()
        self.running.set()
        self.in_flight = 0
        self.idle = asyncio.Event()
        self.idle.set()
        self.tasks = [asyncio.create_task(self._worker()) for _ in range(workers)]

    async def _worker(self) -> None:
        while True:
            await self.running.wait()                 # gate: blocks while paused
            item = await self.q.get()
            if not self.running.is_set():             # paused while we waited for an item
                self.q.put_nowait(item)
                self.q.task_done()
                continue
            self.in_flight += 1
            self.idle.clear()
            try:
                await self.handle(item)
            finally:
                self.in_flight -= 1
                self.q.task_done()
                if self.in_flight == 0:
                    self.idle.set()

The re-check after get() matters: a worker that was already blocked in get() when the pause began would otherwise process one more item. Putting the item back keeps it in the queue, at the back — acceptable for most pools, not for ones that promise FIFO order, which need a peek-then-take design instead.

Verify: pause while the queue is full; no new item starts after running is cleared.

Pause, wait for in-flight work, resume 4 lanes over time. Pause, wait for in-flight work, resume gate running paused running workers processing finish in-flight wait on gate processing pause() returns at 30 ms queue items stay queued, none lost time (not to scale) → Measured: pause() returned after 30 ms, zero items ran during the pause, all 100 completed after resume.

2. Make pause() wait for in-flight items

Clearing the gate only stops new items. Callers usually need to know when the pool is fully quiet — before taking a database snapshot, before failing over. pause() clears the gate and then waits for the in-flight count to reach zero:

    async def pause(self, timeout: float | None = 30.0) -> None:
        self.running.clear()
        async with asyncio.timeout(timeout):
            await self.idle.wait()

    def resume(self) -> None:
        self.running.set()

    @property
    def paused(self) -> bool:
        return not self.running.is_set()

Measured: with 50 ms items in flight, pause() returned after 30 ms — the remaining time on the items that were running. The timeout protects the caller when an in-flight item is itself stuck; combine it with per-item timeouts from setting per-item timeouts in worker pools so "pause" cannot wait longer than one item's deadline.

Verify: pause() returns only when no item is running, and returns within the slowest in-flight item's remaining time.

3. Pause automatically on a dependency outage

The most useful trigger is not an operator but the dependency itself. When a circuit breaker opens, pause the pool; when it closes, resume. Items wait in the queue instead of failing one by one against a dead service:

async def follow_breaker(pool: PausablePool, breaker, every: float = 1.0) -> None:
    while True:
        if breaker.state == "open" and not pool.paused:
            log.warning("dependency down: pausing pool with %d queued", pool.q.qsize())
            await pool.pause()
        elif breaker.state != "open" and pool.paused:
            log.info("dependency recovered: resuming")
            pool.resume()
        await asyncio.sleep(every)

Without the pause, every queued item would be attempted, fail fast at the open breaker, and go to retries or dead letters — turning a five-minute outage into thousands of failure records. With it, the queue simply waits. Watch the queue's size during a pause: if producers keep adding, it grows for the outage's duration, so pair the pause with backpressure on producers as in bounded asyncio queue with backpressure under load.

Verify: take the dependency down; the pool pauses within one check interval, no items are dead-lettered for the outage, and the queue drains after recovery.

A circuit breaker pausing and resuming the pool A sequence of 6 messages between 4 participants. A circuit breaker pausing and resuming the pool dependency breaker controller pool errors: breaker opens pause(): in-flight finish items wait in queue half-open probe succeeds breaker closed resume() Queued work waits out the outage instead of failing against it.

4. Expose the switch to operators safely

Operators need a way to pause during incidents and maintenance. Expose it through an admin endpoint, authenticated, and make the state visible:

from fastapi import FastAPI, Depends

admin = FastAPI(dependencies=[Depends(require_admin)])


@admin.post("/pool/pause")
async def pause_pool(reason: str):
    await app_state.pool.pause()
    log.warning("pool paused by operator: %s", reason)
    return {"paused": True, "queued": app_state.pool.q.qsize()}


@admin.post("/pool/resume")
async def resume_pool():
    app_state.pool.resume()
    return {"paused": False}


@admin.get("/pool")
async def pool_state():
    p = app_state.pool
    return {"paused": p.paused, "in_flight": p.in_flight, "queued": p.q.qsize()}

Record the reason, and alert when a pool has been paused longer than expected — a pool paused for maintenance and forgotten is an outage that looks healthy. In a multi-replica service, an operator pause must reach every replica: broadcast it through Redis pub/sub or a shared flag that each replica's controller reads, the same fan-out used in invalidating caches across async workers.

Verify: a pause request reaches every replica within one poll interval and is visible in their state endpoints.

5. Keep health checks honest while paused

A paused pool is not a broken one, and the orchestrator should not restart it — but a paused pool also is not doing its job. Distinguish them in health reporting:

@app.get("/healthz")              # liveness: the process works
async def healthz():
    return {"ok": True}


@app.get("/readyz")               # readiness: the pool is consuming
async def readyz():
    if app_state.pool.paused:
        return JSONResponse({"ready": False, "reason": "paused"}, status_code=503)
    return {"ready": True}

Liveness stays green so the process is not killed; readiness goes red so load balancers and dashboards show that this instance is not consuming. For queue consumers without inbound traffic, readiness mainly feeds dashboards and alerts; the probe design itself is in implementing health and readiness probes for asyncio.

Verify: while paused, liveness passes and readiness reports the pause reason.

Pause, shed or shut down? A decision on Why stop processing with 3 outcomes. Pause, shed or shut down? Why stop processing? dependency down or maintenance pause work waits in the queue more demand than capacity shed at intake reject, do not queue process must stop drain and shut down return unstarted work Pausing keeps work and state; shedding protects the system; shutdown hands work back.

Verification

The pause switch works when:

  • No item starts after pause() is called, and pause() returns once in-flight items finish.
  • resume() continues with the untouched queue, losing nothing.
  • Dependency outages pause the pool automatically via the breaker, without dead-lettering.
  • Paused state is visible in readiness and admin endpoints, and long pauses alert.

Diagnostic Hook: export a pool_paused gauge (0 or 1) with a reason label, the pause duration, and queue size. Alert when paused for longer than the expected maintenance window, and when queue size during a pause approaches the queue's bound — at that point producers will start blocking or shedding, and the outage is about to reach users.

Pitfalls & edge cases

  • No re-check after get(). A worker blocked in get() processes one more item after the pause.
  • Pausing without a timeout. One stuck in-flight item makes pause() hang.
  • Forgotten pauses. A paused pool looks healthy; alert on duration.
  • Per-replica pauses in a fleet. Operator actions must reach every instance.

Frequently Asked Questions

How do I pause an asyncio worker pool?

Gate each worker's loop on an asyncio.Event: clear it to pause, set it to resume. Have pause() also wait until the in-flight count reaches zero, so callers know the pool is quiet. In testing, pause() returned in 30 ms and no items ran while paused.

What happens to queued items while the pool is paused?

They stay in the queue and are processed after resume. If producers keep adding during a long pause, bound the queue so they get backpressure instead of growing it without limit.

Should a circuit breaker pause the worker pool?

Usually yes. Pausing while the breaker is open lets work wait out the outage instead of failing item by item and filling retry or dead-letter queues.

Should a paused pool fail its health check?

It should fail readiness, so dashboards and load balancers see it is not consuming, but pass liveness, so the orchestrator does not restart a deliberately paused process.