Collecting Errors from Async Worker Pools¶
A worker pool built on asyncio.TaskGroup has fail-fast semantics by default: the first exception cancels every worker. For a batch where some items are expected to fail — malformed records, missing files, a remote 404 — that is usually wrong, and the alternative, catching everything, has costs of its own. Measured on Python 3.14 with a pool of 20 workers processing 2,000 items, each failing with 2% probability: the fail-fast pool stopped after 19 successful items, with 1,961 never started and 19 more cancelled mid-flight. A pool that caught and collected errors per item finished 1,954 successes and 46 failures. Keeping the exception objects retained 784 KiB for those 46 failures, because each traceback kept its frame's local variables alive; keeping a tuple of item, type and message retained 21 KiB. An error budget of 10% let the 2% run finish, and stopped a run with a 50% failure rate after 68 items instead of 2,000. This guide builds the three policies and the reporting around them.
Prerequisites¶
- A TaskGroup worker pool, from building an async worker pool with TaskGroup.
- Exception groups, from handling specific errors with except*.
- The topic overview, Worker Pool Implementations.
1. Measure what fail-fast does¶
The plain pool lets an item's exception escape its worker:
async def run_fail_fast(items, workers=20):
q = asyncio.Queue()
for item in items:
q.put_nowait(item)
async def worker():
while not q.empty():
await process(q.get_nowait())
async with asyncio.TaskGroup() as tg:
for _ in range(workers):
tg.create_task(worker())
Measured with 2,000 items and a 2% failure rate: the first failure arrived among the first 20 items. The TaskGroup cancelled the other 19 workers, 19 items had already succeeded, 19 were cancelled part-way through, and 1,961 were never started. The caller received an ExceptionGroup containing one ItemError. That is the right behaviour when any failure invalidates the batch — a migration that must apply completely or not at all — and the wrong one for a batch of independent items.
One trap when handling the result: return, break and continue are not allowed inside an except* block — Python rejects the file with a SyntaxError. Record what you need in a variable and return after the block.
Verify: decide, per pool, whether one item's failure should stop the batch, and write the policy down next to the pool.
2. Catch per item and keep going¶
Move the try inside the worker loop, so an item's failure ends that item, not the worker:
async def run_collecting(items, workers=20):
q = asyncio.Queue()
for item in items:
q.put_nowait(item)
ok, failures = 0, []
async def worker():
nonlocal ok
while not q.empty():
item = q.get_nowait()
try:
await process(item)
ok += 1
except ItemError as e: # expected per-item failures only
failures.append((item, type(e).__name__, str(e)))
async with asyncio.TaskGroup() as tg:
for _ in range(workers):
tg.create_task(worker())
return ok, failures
Measured: 1,954 items succeeded and 46 failed, in 0.53 s. Catch the exception types you expect from bad items, not Exception: a TypeError from a code bug should still stop the pool, and CancelledError — a subclass of BaseException — must never be swallowed, or the pool cannot be shut down.
Verify: a run with a known number of bad items reports exactly that many failures and processes every other item.
3. Store errors without keeping their frames¶
An exception object holds its traceback, and the traceback holds every frame it passed through, with all their local variables. Collecting exceptions keeps all of that alive until the batch ends:
tracemalloc.start()
ok, failures = await run_collecting(items) # variant appending `e` itself
current, _ = tracemalloc.get_traced_memory()
Measured with a process function holding a 2,000-element list as a local: storing the 46 exception objects retained 784 KiB, about 17 KiB per failure; storing (item, type name, message) tuples retained 21 KiB. With large buffers in the failing frames — a request body, a parsed document — or tens of thousands of failures, collected exceptions become the pool's largest memory use. If the traceback is needed, format it at the point of failure with traceback.format_exception(e) and keep the string, or log it immediately and keep only the summary.
Verify: compare tracemalloc usage after a run with many failures against one with none; the difference should be small per failure.
4. Stop a broken run with an error budget¶
Collecting every error is wrong when everything is failing — a bad credential or a down dependency makes each item fail, and the pool spends its whole run producing 2,000 copies of one error. An error budget aborts when the failure rate crosses a threshold, after a minimum sample:
MIN_SAMPLE, BUDGET = 20, 0.10
async def run_with_budget(items, workers=20):
q = asyncio.Queue()
for item in items:
q.put_nowait(item)
ok, failures, aborted = 0, [], False
async def worker():
nonlocal ok, aborted
while not q.empty() and not aborted:
item = q.get_nowait()
try:
await process(item)
ok += 1
except ItemError as e:
failures.append((item, type(e).__name__, str(e)))
n = ok + len(failures)
if len(failures) > MIN_SAMPLE and len(failures) / n > BUDGET:
aborted = True # workers stop taking new items
async with asyncio.TaskGroup() as tg:
for _ in range(workers):
tg.create_task(worker())
return ok, failures, q.qsize(), aborted
Measured with a 10% budget: at a 2% failure rate the run finished normally, 1,962 succeeded and 38 failed. At a 50% failure rate it stopped after 39 successes and 29 failures, in 0.02 s, leaving 1,932 items unstarted for a later rerun. The minimum sample matters — without it, the first item failing is a 100% failure rate. Workers finish their current item rather than being cancelled, so no item is left half-processed.
Verify: a run against a deliberately broken dependency aborts within the first few dozen items and reports how many were not started.
5. Report failures so they can be acted on¶
At the end of the run, the caller needs three things: whether the batch as a whole succeeded, which items failed, and why — grouped, since 46 failures are usually two or three distinct causes:
def report(ok, failures, not_started, aborted):
by_cause = collections.Counter((kind, msg.split(":")[0]) for _, kind, msg in failures)
log.info("batch done", extra={"ok": ok, "failed": len(failures),
"not_started": not_started, "aborted": aborted})
for (kind, msg), count in by_cause.most_common(5):
log.warning("failure cause", extra={"type": kind, "message": msg, "count": count})
if aborted:
raise BatchAborted(f"{len(failures)} failures, {not_started} items not started")
return [item for item, _, _ in failures] # for a retry run
Return the failed items' identifiers so a rerun can process only those, and raise when the budget aborted the run, so a scheduler marks the job failed instead of successful. When the caller wants all failures as exceptions — a library API, for example — raise one ExceptionGroup built from the collected errors at the end, which the caller can split by type with except*, as described in exception groups and TaskGroups.
Verify: the batch report names the top failure causes with counts, and the failed item list can drive a rerun.
Verification¶
Error collection is correct when:
- Each pool has a stated policy: fail-fast, collect, or collect with a budget.
- Only expected exception types are caught, and
CancelledErroris never swallowed. - Collected errors are summaries, not exception objects, unless tracebacks are needed and memory has been measured.
- Aborted runs fail visibly and report the items not started.
Diagnostic Hook: when a batch job's memory grows with its failure count, check what the error list holds. Exception objects keep their frames' locals alive — 17 KiB per failure in this test, and far more when frames hold request bodies.
Pitfalls & edge cases¶
- Fail-fast by accident. Measured: 1,961 of 2,000 items never started.
returninsideexcept*. It is aSyntaxError.- Keeping exception objects. Measured: 784 KiB for 46 failures against 21 KiB.
- Collecting through a total outage. A budget stopped a 50%-failure run after 68 items.
Frequently Asked Questions¶
How do I stop one failing task from cancelling an asyncio worker pool?
Catch the expected exception types inside the worker loop, around each item, and record the failure. A TaskGroup only cancels siblings when an exception escapes a task.
Why does storing exceptions use so much memory?
Each exception keeps its traceback, which keeps every frame's local variables alive. 46 exceptions retained 784 KiB; storing type and message tuples retained 21 KiB.
What is an error budget in a batch job?
A failure-rate threshold that stops new work once exceeded after a minimum sample. A 10% budget let a 2% run finish and stopped a 50% run after 68 items.
Can I return from an except* block?
No. return, break and continue inside except* are a SyntaxError. Set a variable in the block and return after it.
Related¶
- Worker Pool Implementations — up to the topic overview.
- Setting per-item timeouts in worker pools — timeouts as another per-item failure.
- Concurrent Execution & Worker Patterns — the section overview.