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¶
- Python 3.11+, stdlib only.
- Events as gates, from signalling with asyncio.Event set() and clear().
- Circuit breakers, from implementing an async circuit breaker.
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.
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.
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.
Verification¶
The pause switch works when:
- No item starts after
pause()is called, andpause()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 inget()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.
Related¶
- Worker Pool Implementations — up to the topic overview.
- Draining a worker pool on shutdown — the permanent version of stopping a pool.
- Concurrent Execution & Worker Patterns — the section overview.