Structured Concurrency in Python 3.11+: Replacing Ad-Hoc Task Management with asyncio.TaskGroup
Python 3.11's asyncio.TaskGroup enforces structured concurrency: child tasks are awaited on exit, exceptions cancel siblings and re-raise as ExceptionGroup. Covers design, trust boundaries, failure modes, and when to use a supervised nursery instead.
28 Jun 2026, 14:13 UTC

The Problem: Orphaned Tasks and Silent Failures
Long-running Python services often manage concurrent work by spawning asyncio.Task objects directly — storing them in lists, awaiting them with asyncio.gather, or forgetting them entirely. This approach creates three persistent problems: tasks can outlive their logical scope (orphaned tasks), exceptions in one task may be swallowed while siblings continue running (silent failures), and cleanup logic scattered across finally blocks becomes unreliable when cancellation arrives mid-operation.
Python 3.11 introduces asyncio.TaskGroup, a context manager that enforces structured concurrency: all child tasks are awaited on exit, any exception cancels siblings and re-raises as an ExceptionGroup, and the group’s lifetime becomes the explicit boundary for the unit of work. This article describes the smallest suitable design, trust and data boundaries, operational checks, failure modes, and the conditions that would lead you to choose a different pattern.
Requirements Driving the Design
- Deterministic cleanup: Every spawned task must either complete or be cancelled when the logical operation ends — no background leftovers.
- Explicit error propagation: The first exception must surface to the caller; siblings must be cancelled, not left running in a degraded state.
- Scope ownership: Only the code that owns the logical operation may spawn tasks within it; untrusted callers receive a pre-configured group or factory.
- Observability: Group entry/exit, active task count, cancellation rate, and
ExceptionGroupfrequency must be measurable without intrusive instrumentation.
Smallest Suitable Design: One TaskGroup per Logical Unit
Wrap a single logical operation — handling one HTTP request, running one pipeline stage, processing one message batch — in one TaskGroup. Nested groups model sub-operations with independent failure domains while preserving parent cancellation semantics.
async def handle_request(request: Request) -> Response:
async with asyncio.TaskGroup() as tg:
# Sub-operation with its own failure domain
async with asyncio.TaskGroup() as validation_tg:
validation_tg.create_task(validate_headers(request))
validation_tg.create_task(validate_payload(request))
# Main work — siblings cancelled if validation raised
tg.create_task(fetch_user_profile(request.user_id))
tg.create_task(fetch_permissions(request.user_id))
tg.create_task(audit_log_request(request))
# All tasks completed or cancelled; first exception re-raised as ExceptionGroup
return compose_response(...)
If validation fails, validation_tg exits with an ExceptionGroup, cancelling its children. The outer tg never starts the main work because the exception propagates out of the nested async with. If the outer group is cancelled (e.g., client disconnect), validation_tg and all main-work tasks receive CancelledError simultaneously.
Trust and Data Boundaries
TaskGroup enforces that only the group owner can spawn tasks inside its scope. Pass a pre-configured group or a factory to untrusted code rather than exposing create_task directly:
# Trusted orchestrator
async def run_pipeline(stage_configs: list[StageConfig]) -> PipelineResult:
async with asyncio.TaskGroup() as tg:
for config in stage_configs:
# Untrusted plugin receives a factory bound to this group
stage_factory = lambda coro: tg.create_task(coro)
await config.plugin.execute(stage_factory)
return collect_results()
Data boundaries stay explicit: results flow through awaited coroutines or explicit asyncio.Queue instances, not shared mutable state. Each task returns a value or raises; the group collects outcomes in ExceptionGroup.exceptions or via the task’s result() method after __aexit__.
Operational Checks
Instrument the group’s lifetime at the boundaries:
import time
import asyncio
from contextlib import asynccontextmanager
@asynccontextmanager
async def observed_task_group(name: str):
start = time.monotonic()
active = 0
async with asyncio.TaskGroup() as tg:
original_create = tg.create_task
def tracked_create(coro, *, name=None):
nonlocal active
active += 1
task = original_create(coro, name=name)
task.add_done_callback(lambda _: active - 1)
return task
tg.create_task = tracked_create # type: ignore[method-assign]
try:
yield tg
finally:
duration = time.monotonic() - start
log.info("task_group_complete", name=name, duration=duration, active_at_exit=active)
Expose metrics for:
active_task_count(gauge)cancellation_rate(counter, incremented inexcept asyncio.CancelledErrorblocks)exception_group_frequency(counter, incremented whenExceptionGroupis caught)
Correlate traces with asyncio.current_task().get_name() — set meaningful names when creating tasks: tg.create_task(work(), name="fetch-user-123").
Failure Modes and Mitigations
1. Blocking Synchronous Code Stalls Cancellation
A task executing a tight CPU loop or a blocking C extension ignores CancelledError because cancellation is cooperative. The group’s __aexit__ waits indefinitely.
Mitigation: Offload blocking work to an executor with a timeout:
async def safe_blocking_call(func, *args, timeout=5.0):
loop = asyncio.get_running_loop()
try:
return await asyncio.wait_for(
loop.run_in_executor(None, func, *args),
timeout=timeout
)
except asyncio.TimeoutError:
# Executor future cannot be cancelled; log and raise
raise RuntimeError(f"Blocking call {func.__name__} exceeded {timeout}s")
For CPU-bound work that dominates the workload, move to process pools and treat TaskGroup as orchestration only (see Conditions That Change the Design).
2. ExceptionGroup Masking
An ExceptionGroup may contain multiple exception types. Code that only catches ValueError will miss a KeyError sibling.
Mitigation (Python 3.11+): Use except* to handle matching exceptions while letting others propagate:
try:
async with asyncio.TaskGroup() as tg:
tg.create_task(risky_operation())
tg.create_task(another_risky())
except* ValueError as eg:
for exc in eg.exceptions:
handle_validation_error(exc)
except* ConnectionError as eg:
for exc in eg.exceptions:
handle_transient_failure(exc)
# Other exception types propagate outward
On older versions (or when backporting with exceptiongroup), iterate group.exceptions manually and re-raise unhandled ones.
3. Cancellation During __aexit__ Cleanup
If a parent cancels the group while a task’s finally block is running, the cleanup may be interrupted, leaving resources in an inconsistent state.
Mitigation: Shield critical finally sections:
async def write_with_checkpoint(data: bytes, path: Path):
tmp = path.with_suffix(".tmp")
try:
async with aiofiles.open(tmp, "wb") as f:
await f.write(data)
await asyncio.shield(tmp.replace(path)) # atomic rename protected
finally:
# Best-effort cleanup; shield if tmp must be removed
try:
await asyncio.shield(asyncio.to_thread(tmp.unlink, missing_ok=True))
except OSError:
pass
4. Cascade Cancellation Storms
Nested TaskGroups multiply cancellation signals: a parent cancel cancels all descendants simultaneously. With thousands of tasks, the storm of CancelledError deliveries can spike latency.
Mitigation: Rate-limit cancellation propagation with a semaphore in the supervisor, or stagger child group creation. For very wide fans, consider a dedicated supervisor task that drains children in batches (see Conditions That Change the Design).
Conditions That Change the Design
Tasks Must Outlive the Request Scope
If background jobs (periodic cleanup, webhook retries, metrics flush) must survive request cancellation, TaskGroup is the wrong primitive — its lifetime is bound to the async with block. Replace it with a supervised nursery pattern:
class SupervisedNursery:
def __init__(self):
self._tasks: set[asyncio.Task] = set()
self._shutdown = asyncio.Event()
async def spawn(self, coro, *, name=None):
if self._shutdown.is_set():
raise RuntimeError("Nursery is shutting down")
task = asyncio.create_task(coro, name=name)
self._tasks.add(task)
task.add_done_callback(self._tasks.discard)
return task
async def shutdown(self, timeout=30.0):
self._shutdown.set()
if self._tasks:
await asyncio.wait(self._tasks, timeout=timeout)
for t in self._tasks:
t.cancel()
await asyncio.gather(*self._tasks, return_exceptions=True)
The supervisor task owns the nursery and implements an explicit shutdown protocol. Request-scoped work still uses TaskGroup; only truly long-lived work uses the nursery.
CPU-Bound Work Dominates
When the workload is primarily numerical computation, image processing, or other CPU-heavy tasks, the GIL limits throughput. TaskGroup becomes an orchestration layer over concurrent.futures.ProcessPoolExecutor:
async def process_batch(items: list[Item]) -> list[Result]:
loop = asyncio.get_running_loop()
with concurrent.futures.ProcessPoolExecutor() as pool:
async with asyncio.TaskGroup() as tg:
for chunk in chunked(items, 100):
tg.create_task(
loop.run_in_executor(pool, cpu_intensive_transform, chunk)
)
# All futures completed or cancelled; collect results from tasks
return [task.result() for task in tg._tasks if not task.cancelled()]
The TaskGroup manages the lifecycle of the executor submissions; the actual work runs in separate processes.
Verification Checklist
Before adopting TaskGroup in production, run these checks:
- Version guard:
sys.version_info >= (3, 11)andhasattr(asyncio, 'TaskGroup'). Backports (exceptiongroup,asyncio-taskgroup) differ inExceptionGroupsemantics and cancellation timing — test explicitly if you must support 3.10. - Minimal behavior test: Spawn a task that raises
ValueErrorand another that sleeps for 10 seconds. Verify the group exits with anExceptionGroupcontaining theValueErrorand the sleeping task’sfinallyblock logs cancellation. - Cancellation latency profile: Spawn 1,000 tasks that
await asyncio.sleep(10), cancel the group after 100 ms, measure time to__aexit__completion. Expect sub-100 ms on idle systems; investigate if it exceeds 500 ms. except*syntax validation: Confirmtry/except* ValueError as eg:catches only matching exceptions from the group on your target Python version.
Limitations
- Requires Python 3.11+. Backports exist but are not drop-in replacements for
ExceptionGrouphandling. - Cancellation remains cooperative; tasks with blocking C extensions or tight loops without
awaitpoints will not respond to cancellation. - Nested groups increase cancellation signal fan-out; design for cascade storms at scale.
- No built-in supervision for tasks that must outlive the group — you must build or adopt a nursery pattern.
Practical Result Check
After migration, verify that:
- No
Taskobjects are stored in long-lived collections without an owningTaskGroupor nursery. - Every logical operation has a single
async with asyncio.TaskGroup()at its top level. - Metrics show
active_task_countdropping to zero at group exit under normal and error conditions. ExceptionGroupfrequency correlates with actual failure rates, not silent task deaths.
If these hold, you have replaced ad-hoc task management with a deterministic, observable concurrency boundary.
0 replies
A thoughtful contribution can make all the difference. Be the first to share one.