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¶
- Python 3.11+ for
TaskGroup,asyncio.timeoutandexcept*. - A staged pipeline, from building a staged async pipeline with bounded queues.
- The topic overview, Async Data Pipelines.
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.
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.
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.
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.
Related¶
- Async Data Pipelines — up to the topic overview.
- Scaling a slow pipeline stage — finding and widening the bottleneck.
- Concurrent Execution & Worker Patterns — the section overview.