Skip to content

Setting Per-Item Timeouts in Worker Pools

A worker pool's throughput is set by its workers, and a worker stuck on one item is a worker lost. A small fraction of items that hang — a slow upstream, a pathological input, a lock that is never released — can occupy every worker in turn. Tested with 8 workers and 500 items, where 2% of items hung for 5 seconds and the rest took 10 ms: without per-item timeouts the batch took 10.41 s, because the 10 hanging items each pinned a worker for 5 s; with a 0.5 s timeout around each item, the same batch finished in 1.54 s, with 490 items succeeding and the 10 hanging ones reported as timed out. This guide adds per-item timeouts, picks their value from data, and handles what happens to timed-out items.

Prerequisites

1. Wrap each item, not the whole worker

The timeout belongs around the handling of a single item, inside the worker loop, so that a timeout ends that item and the worker moves on:

import asyncio


async def worker(q: asyncio.Queue, handle, item_timeout: float, stats: dict) -> None:
    while True:
        item = await q.get()
        try:
            async with asyncio.timeout(item_timeout):
                await handle(item)
            stats["ok"] += 1
        except TimeoutError:
            stats["timed_out"] += 1
            await on_timeout(item)
        except Exception as exc:
            stats["failed"] += 1
            await on_error(item, exc)
        finally:
            q.task_done()

asyncio.timeout() cancels the handler's current await when the deadline passes and converts the cancellation into TimeoutError at the block's exit, so the worker continues with the next item. Measured: 10.41 s without the timeout, 1.54 s with 0.5 s — the hanging items now cost half a second each instead of five. A timeout around the worker would instead end the worker after the deadline regardless of how many items it had processed.

Verify: inject a handler that hangs forever for one item; that item is reported as timed out and the batch completes.

500 items, 2% of which hang for 5 s, on 8 workers 2 horizontal bars comparing no per-item timeout with the others. 500 items, 2% of which hang for 5 s, on 8 workers no per-item timeout 10.41 s 0.5 s per-item timeout 1.54 s, 10 timed out Healthy items take 10 ms; Python 3.14. A timeout bounds what one bad item can cost the whole pool.

2. Choose the timeout from the latency distribution

A timeout should cut off the pathological tail, not normal slow items. Measure item duration and set the timeout well above the high percentile:

import statistics


def suggest_timeout(durations_s: list[float], multiplier: float = 3.0, floor: float = 0.1) -> float:
    durations_s = sorted(durations_s)
    p99 = durations_s[int(len(durations_s) * 0.99) - 1]
    return max(floor, p99 * multiplier)

Three times p99 is a reasonable starting point: it times out only items far outside normal behaviour, and it bounds the worst case for a worker. Different item types often deserve different timeouts — a thumbnail versus a full video transcode — so key the timeout by type rather than using one global value. And the timeout must stay below anything the item's caller is waiting on: a pool serving a request with a 2 s deadline cannot use a 5 s item timeout, which is the deadline-propagation problem in propagating deadlines across async service calls.

Verify: the timed-out fraction in production is small — well under 1% — and timed-out items are genuinely abnormal when inspected.

3. Decide what happens to timed-out items

A timed-out item is unfinished work. Choices, by item semantics:

async def on_timeout(item) -> None:
    item.attempts += 1
    if item.attempts < 3 and item.retryable:
        await retry_queue.put(item, delay=backoff(item.attempts))     # transient: try later
    else:
        await dead_letters.put(item, reason="timeout")                # persistent: park it
        log.warning("item %s timed out %d times", item.id, item.attempts)

Retrying makes sense when the hang was environmental — an upstream blip. Parking makes sense when the item itself is pathological: the same input will hang again, and retrying it forever wastes workers. Count attempts per item so a poison item reaches the dead-letter path, as in implementing a dead letter queue with asyncio. Whatever the handler did before the timeout must be safe to repeat, since a retry will repeat it.

Verify: a deliberately pathological item ends in the dead-letter queue after its attempts, not in an endless retry loop.

What should happen to a timed-out item? A decision on Why did it time out with 3 outcomes. What should happen to a timed-out item? Why did it time out? dependency blip, likely transient retry with backoff count attempts same item times out again dead-letter queue inspect the input result no longer useful drop and log e.g. expired request Track attempts per item so a poison input cannot loop forever.

4. Handle work that ignores cancellation

asyncio.timeout() only works if the handler is cancellable at an await. Two kinds of work defeat it:

  • Blocking calls inside the coroutine. A synchronous call that hangs never reaches an await; the timeout fires only after it returns. Move such calls to asyncio.to_thread, which makes the await cancellable — though the thread keeps running, as covered in timing out blocking calls in threads.
  • Handlers that swallow CancelledError. A try/except Exception is fine, but except BaseException: or bare except: will catch the cancellation, and the timeout then does nothing.
async def handle(item):
    try:
        return await fetch_and_process(item)
    except Exception:                          # fine: CancelledError is not an Exception
        log.exception("item failed")
        raise
    # never: except BaseException / bare except — it would swallow the timeout's cancellation

A timed-out item whose work keeps running in a thread still consumes resources. Track those threads, or bound them with a separate executor, so a run of timeouts cannot exhaust the thread pool behind the scenes.

Verify: in a test with a handler that swallows cancellation, the timeout is not respected — then fix the handler and confirm it is.

5. Detect workers that are stuck anyway

Even with timeouts, a worker can get stuck: a bug outside the timed block, or a timeout that was never applied to some code path. A watchdog that tracks when each worker last made progress catches it:

import time


class Watched:
    def __init__(self, n: int) -> None:
        self.last_progress = [time.monotonic()] * n

    def tick(self, i: int) -> None:
        self.last_progress[i] = time.monotonic()

    async def watchdog(self, limit: float, every: float = 5.0) -> None:
        while True:
            await asyncio.sleep(every)
            now = time.monotonic()
            for i, t in enumerate(self.last_progress):
                if now - t > limit:
                    log.error("worker %d made no progress for %.0fs", i, now - t)

Each worker calls tick(i) after every item and while idle; a worker silent for longer than the item timeout plus margin is stuck. Pair the alert with a task dump, as in dumping stacks of a hung asyncio program, to see where.

Verify: a worker stuck in an untimed code path triggers the watchdog within one interval.

An item's path through a timed worker A flow of 4 stages. An item's path through a timed worker take item q.get() handle under timeout asyncio.timeout(t) ok, or timed out retry or dead-letter tick progress next item The worker always returns to the queue; only the item pays for the timeout.

Verification

Per-item timeouts work when:

  • Every item's handling is inside a timeout, and the worker loop continues after it fires.
  • The timeout is derived from measured latency and respects callers' deadlines.
  • Timed-out items are retried with limits or dead-lettered, never lost or retried forever.
  • Stuck workers are detected by a progress watchdog.

Diagnostic Hook: export the timeout rate per item type and the duration histogram per type. A timeout rate that jumps for one type while its median is unchanged points at a subset of pathological inputs; a rate that rises together with the median points at a slow dependency — the fix in that case is upstream, not a longer timeout.

Pitfalls & edge cases

  • Timeout around the worker loop. It kills the worker, not the item.
  • Blocking calls in handlers. The timeout cannot fire until they return.
  • Swallowed cancellations. except BaseException makes timeouts ineffective.
  • Retrying poison items forever. Count attempts and dead-letter them.

Frequently Asked Questions

How do I add a timeout to each item in an asyncio worker pool?

Wrap the handling of each item in async with asyncio.timeout(seconds) inside the worker loop, catch TimeoutError, record or re-queue the item, and continue to the next one. In testing this cut a batch from 10.4 s to 1.5 s when 2% of items hung.

What should a per-item timeout be?

Well above normal latency — around three times the p99 is a reasonable start — and below any deadline the caller is waiting on. Use different values for item types with very different durations.

Why doesn't asyncio.timeout stop my handler?

Either the handler is blocked in a synchronous call with no await to cancel, or it catches BaseException and swallows the cancellation. Offload blocking calls to threads and only catch Exception.

What should happen to items that time out?

Retry them with backoff if the cause is likely transient, and send them to a dead-letter queue after a few attempts so a pathological item cannot occupy workers forever.