Skip to content

Handling Per-Item Errors in Async Pipelines

A pipeline processes thousands of items, and some of them will fail: malformed records, upstream resets, an item that makes a call hang. How a pipeline reacts to the first failure decides whether a 1% error rate costs 1% of the output or all of it. Measured on Python 3.14 with a three-stage pipeline — a source, 20 enrichment workers and a sink connected by bounded queues — over 10,000 items, where 100 items raised a permanent ValueError, 200 failed once with a transient ConnectionError and 5 hung forever: run inside a TaskGroup with no per-item handling, the pipeline aborted after 19 items on the first transient error. Run with asyncio.gather(..., return_exceptions=True), it raised nothing and was still waiting after 5 s with 842 items processed — the workers that hit errors had died without sending their end-of-stream markers. With per-item error handling, a 1 s per-item timeout and a dead-letter list, it finished 9,695 items in 1.63 s and dead-lettered 305; adding up to three retries for transient errors finished 9,895 in 1.74 s and dead-lettered exactly the 105 items that could never succeed. This guide builds that last version.

Prerequisites

1. Decide what one item's failure should do

Errors in a pipeline belong to one of two kinds. An item error affects only that item — a bad record, a 404, a timeout — and the right response is to set it aside and carry on. A pipeline error means nothing else can succeed — the database is down, credentials are wrong, the output disk is full — and the right response is to stop everything quickly. A pipeline without per-item handling treats every error as a pipeline error:

async def worker():
    while (item := await q_in.get()) is not None:
        await q_out.put(await enrich(item))        # any exception kills this worker
    await q_out.put(None)

async with asyncio.TaskGroup() as tg:              # ...and the TaskGroup cancels the rest
    tg.create_task(source())
    tg.create_task(sink())
    for _ in range(20):
        tg.create_task(worker())

Measured: the first transient error, on item 3, aborted the whole run after 19 items had reached the sink. That is the correct behaviour for a pipeline error and the wrong one for 99.7% of this run's failures. Write down which exceptions are which for each stage before writing handlers; the classification is the design.

Verify: for each stage, you have a list of exceptions that are per-item and a list that should stop the pipeline.

2. Never let a worker die silently

The tempting fix is to stop errors from propagating, and the common way to do it hides them entirely:

await asyncio.gather(source(), sink(), *workers, return_exceptions=True)

Measured: no exception surfaced, and the run was still waiting after 5 s with 842 items processed and the input queue full. Each worker that raised had exited without putting its None end-of-stream marker on the output queue, so the sink waited for 20 markers that would never all arrive, and the source blocked on a full queue with nobody left to drain it. return_exceptions=True turned crashes into return values that nobody looked at until everything finished — which it never did. Keep the TaskGroup, so a genuine pipeline error still cancels everything and raises, and handle item errors inside the worker, around the item:

async def worker():
    try:
        while (item := await q_in.get()) is not None:
            await process_item(item)                # handles its own item errors
    finally:
        await q_out.put(None)                       # end-of-stream even on the way out

Verify: inject an exception into one item and confirm that the run either completes or raises — never waits.

10,000 items with 1% bad, 2% flaky and 5 hanging A grid of 4 rows by 5 columns. 10,000 items with 1% bad, 2% flaky and 5 hanging error handling outcome processed dead-lettered time none, TaskGroup aborted on first error 19 - 0.00 s none, gather(return_exceptions=True) hung, no error raised 842 - > 5 s per-item + 1 s timeout + dead-letter completed 9,695 305 1.63 s ... + 3 tries for transient errors completed 9,895 105 1.74 s 100 permanent errors, 200 one-off ConnectionErrors, 5 items that hang; 20 workers.

3. Bound every item with a timeout

An item that hangs does not raise; it just occupies a worker forever. Five hanging items out of 10,000 would eventually take five of twenty workers, and a pipeline whose upstream hangs more often loses all of them. Put a deadline around the processing of each item:

ITEM_TIMEOUT = 1.0

async def process_item(item):
    try:
        async with asyncio.timeout(ITEM_TIMEOUT):
            result = await enrich(item)
    except TimeoutError:
        dead_letter(item, stage="enrich", error="TimeoutError")
        return
    await q_out.put(result)

Measured: the five hanging items were dead-lettered as TimeoutError after 1 s each, while the other workers carried on, and the whole run finished in 1.63 s — the hangs added little because they overlapped with other work. Note that the await q_out.put(...) is outside the timeout: waiting for space in the next stage's queue is backpressure, not a stuck item, and must not be counted against the item's deadline. Choose the timeout from the stage's measured latency distribution, well above its p99.

Verify: a test item that never completes is dead-lettered within the timeout, and throughput of the other items is unaffected.

4. Retry transient errors, dead-letter the rest

Some item errors are worth retrying — a connection reset, a 503, a lock timeout — and some will fail the same way every time. Retry only the first kind, a bounded number of times, and record the rest with enough context to investigate or replay:

@dataclass
class DeadLetter:
    item: object
    stage: str
    error: str
    attempts: int

dead: list[DeadLetter] = []

async def process_item(item, max_tries=3):
    for attempt in range(1, max_tries + 1):
        try:
            async with asyncio.timeout(ITEM_TIMEOUT):
                result = await enrich(item)
            break
        except ConnectionError as exc:
            if attempt == max_tries:
                dead.append(DeadLetter(item, "enrich", repr(exc), attempt))
                return
            await asyncio.sleep(0.01 * attempt)
        except (ValueError, TimeoutError) as exc:   # permanent, or not worth waiting again
            dead.append(DeadLetter(item, "enrich", repr(exc), attempt))
            return
    await q_out.put(result)

Measured: without retries, the 200 transient failures went to the dead-letter list alongside the 100 genuinely bad records and 5 hangs — 305 in all. With up to three tries, all 200 succeeded on the second attempt, and the dead-letter list held exactly the 105 items that could never succeed, in 1.74 s against 1.63 s. Hangs are deliberately not retried here: an item that hung once is likely to hang again, and each retry would cost another full timeout. Backoff and jitter for retries are covered in Retry & Backoff Strategies.

Verify: after a run, the dead-letter list contains only permanent failures, each with its stage, error and attempt count.

What happens to one item A flow of 4 stages. What happens to one item Process under a 1 s timeout Transient error retry, up to 3 tries Permanent / timeout dead-letter with context Success put on next queue, outside timeout Item errors stay with the item; pipeline errors still stop the TaskGroup.

5. Stop the pipeline when the error rate says something is wrong

Per-item handling has a failure mode of its own: if the database goes down, every item fails, and a pipeline that dead-letters everything "completes" having done nothing. Add a circuit on the error rate, raising a pipeline error that the TaskGroup turns into a clean stop:

class PipelineAborted(Exception):
    pass

class ErrorBudget:
    def __init__(self, max_rate=0.2, min_items=200):
        self.max_rate, self.min_items = max_rate, min_items
        self.seen = self.failed = 0

    def record(self, ok: bool) -> None:
        self.seen += 1
        self.failed += not ok
        if self.seen >= self.min_items and self.failed / self.seen > self.max_rate:
            raise PipelineAborted(f"{self.failed}/{self.seen} items failed")

Call budget.record(...) after each item; the exception propagates out of the worker, the TaskGroup cancels the other stages, and the caller gets an ExceptionGroup containing PipelineAborted. In this run the error rate was 1.05% after retries, well under a 20% budget; in a second run where every item failed, the pipeline stopped with "PipelineAborted: 200/200 items failed" after 0.012 s instead of working through all 10,000. Persist the dead-letter list somewhere durable — a JSONL file or a table — so items can be replayed after a fix, and make the replay feed the same process_item, as in checkpointing progress in long-running async jobs.

Verify: a run in which every item fails stops with PipelineAborted after the minimum sample, instead of dead-lettering everything.

What should this error do? A decision on What does the error affect with 4 outcomes. What should this error do? What does the error affect? one bad item dead-letter with context 100 kept for review one item, transient retry up to 3 times 200 recovered one item that hangs per-item timeout 5 cut off at 1 s every item / over budget raise: TaskGroup stops all not 'completed' Classify first; the handler follows from the class.

Verification

A pipeline handles per-item errors correctly when:

  • Item errors are caught inside the worker, around the item, and pipeline errors propagate.
  • Workers always emit their end-of-stream marker, from a finally.
  • Each item has a timeout that excludes waiting on the next queue.
  • Transient errors are retried; the rest are dead-lettered with stage, error and attempts, and an error budget stops runs that fail wholesale.

Diagnostic Hook: export counters per stage for processed, retried and dead-lettered items, labelled by exception type. A healthy pipeline shows a steady trickle of dead letters of a few known types; a new exception type, or retries climbing while dead letters stay flat, is an upstream degrading before it fails.

Pitfalls & edge cases

  • Item errors escaping the worker. Measured: one transient error aborted 10,000 items after 19.
  • gather(return_exceptions=True) around workers. Measured: hung silently with 842 done.
  • Timeouts around queue.put. Backpressure would count as item failure.
  • Dead-lettering everything. Without an error budget, an outage looks like a successful run.

Frequently Asked Questions

How do I stop one bad item from crashing an asyncio pipeline?

Catch per-item exceptions inside the worker around each item, record the item in a dead-letter list, and continue; let only pipeline-wide errors escape to the TaskGroup. In testing this completed 9,895 of 10,000 items instead of aborting after 19.

Why does my asyncio pipeline hang when a worker fails?

A worker that dies never sends its end-of-stream marker, so downstream stages wait forever. gather with return_exceptions=True hid the error and hung with 842 of 10,000 done. Send the marker from a finally block and keep errors visible.

Should pipeline items be retried?

Only transient errors, a bounded number of times. Retrying 200 one-off connection errors recovered all of them; permanent errors and hangs were dead-lettered instead.

What is a dead-letter queue in an async pipeline?

A durable list of items that failed permanently, with the stage, error and attempt count, kept for investigation and replay after a fix.