Skip to main content

Concurrency and Asyncio: Ownership, Waiting, and Cancellation

A concurrent program can advance one operation while another is waiting. A request might be waiting for its response while the program processes data that has already arrived. Concurrency means overlapping the progress of several pieces of work, either by interleaving execution or by running in parallel; parallelism means executing work at the same instant, for example on different CPU cores. Overlapping waits can improve throughput without making an individual computation faster.

Before introducing concurrency, decide who starts each piece of work, receives its result, and stops it after a failure. The example below gives a whole batch one shared lifetime boundary.

Choosing coroutines, threads, or processes​

Python’s coroutines and tasks documentation distinguishes a coroutine object from a scheduled task. Calling an async def function creates a coroutine object; awaiting it directly waits for it, while creating a task schedules it to progress concurrently with other work. Two consecutive coroutine awaits can still execute sequentially.

Thread and process executors offer the same interface for submitting work and retrieving results. Start with what the work waits for and how its data must cross boundaries:

WorkStarting pointCost to account for
Network or other I/O through native asynchronous APIsasyncio coroutinesCalls must cooperate with the event loop and avoid long blocking operations
Blocking I/O in an existing synchronous libraryThreadPoolExecutor with a fixed worker limit; asyncio.to_thread inside an async programThreads share memory, so shared state and stopping need coordination
Substantial pure Python computation that needs several coresProcessPoolExecutorStartup and data transfer have costs; functions, arguments, and results must be picklable, and worker processes must be able to import the main module

The CPython condition in the threading documentation matters: within one interpreter with the GIL enabled, threads cannot execute Python bytecode simultaneously. Extensions that release the GIL may compute in parallel. Since Python 3.13, free-threaded builds can disable the GIL, but simultaneous Python bytecode execution across threads in one interpreter requires it to be actually disabled at runtime. Runtime options or incompatible C extensions can re-enable it. The choice still depends on the runtime, libraries, and workload.

Python 3.14 also provides InterpreterPoolExecutor: each worker thread has its own interpreter and GIL, allowing work to use multiple cores. Interpreters cannot directly share mutable objects, and submitted functions, arguments, and return values are transferred through pickle serialization.

Blocking calls and cooperative scheduling​

An event loop runs one task at a time. When a task suspends to wait for an incomplete operation, the loop can advance other work. await is not a guaranteed scheduling break: an operation that completes immediately may let the task continue without suspending.

Synchronous file reads, time.sleep, and long computations inside a coroutine occupy the event loop. Adding async def does not change what those calls do. asyncio.sleep suspends the current task; asyncio.to_thread moves a blocking function to a thread. For a CPU-heavy loop, consider processes or smaller computation batches with await asyncio.sleep(0) between them to yield control. Splitting the computation or creating more tasks alone does not resolve the blocking.

When running external programs, argument handling, pipes, and child cleanup still follow the subprocess boundaries. Async scheduling overlaps waits; it does not manage the child’s entire lifecycle for you.

Task ownership, lifetimes, and cancellation​

A TaskGroup gives related work a lifetime boundary: it waits for every member before leaving its context. A non-cancellation failure cancels the other members; after they finish, the group raises an exception group. KeyboardInterrupt and SystemExit are re-raised as the original exception after the tasks are brought to an end. Register children through the group’s create_task; extra background tasks created with ordinary asyncio.create_task do not automatically join it. TaskGroup and asyncio.timeout are available from Python 3.11.

With ordinary create_task, the caller must retain references, retrieve exceptions, and arrange to wait for tasks during shutdown. Default gather propagates the first exception without cancelling the other work. Receiving an exception therefore does not establish that the whole group has stopped.

A cancellation request raises CancelledError at the next opportunity. Release resources in finally; if you catch cancellation explicitly, normally propagate it after cleanup. Swallowing it can break task-group or timeout control flow. After calling cancel(), still wait for completion; cancellation does not undo a completed write. Multi-Agent Coordination extends similar ownership questions to files, messages, and multiple execution environments.

Threads have a different stopping boundary. An executor’s Future.cancel() cannot cancel a running call, and result(timeout=...) only limits the wait to retrieve its result. Consequently, cancelling an await of to_thread does not guarantee that the underlying thread stops. A blocking library needs its own timeout or cooperative stopping mechanism.

Shared-state races​

Even on one event loop, another task can intervene if an update suspends between reading state and writing it back. asyncio.Lock provides exclusive access between tasks on the same loop. It does not synchronize operating-system threads; use threading synchronization for those.

Save this complete program as race.py and run python3 race.py. It deliberately suspends after reading the counter:

import asyncio

async def count(use_lock):
value = 0
lock = asyncio.Lock()

async def update():
nonlocal value
before = value
await asyncio.sleep(0)
value = before + 1

async def worker():
if use_lock:
async with lock:
await update()
else:
await update()

async with asyncio.TaskGroup() as group:
group.create_task(worker())
group.create_task(worker())
return value

async def main():
print(f"unlocked={await count(False)}")
print(f"locked={await count(True)}")

asyncio.run(main())

With Python 3.13.15 and the default task factory, the output was:

unlocked=1
locked=2

Without the lock, both workers read 0 and then write 1, losing an update. With the lock, the second worker reads 1 after the first finishes and writes 2. Protect the entire read–modify–write operation; locking only the final assignment leaves the earlier read unprotected.

The suspension inside the lock exposes the race in this example. In ordinary work, keep critical sections short; compute independent results first and let one task combine them. An Event signals that a condition has occurred, while a Semaphore limits how many tasks enter a section of work at once. Creating a task for every input and making those tasks wait on a semaphore limits active work, but not the number of allocated tasks.

Bounded queues, backpressure, and timeouts​

See Queues for FIFO ordering and implementations. A concurrent pipeline must also decide what happens when its producer is too fast. A positive maxsize on asyncio.Queue limits waiting items; await put waits for space when the queue is full. Letting downstream capacity slow the producer is backpressure. With maxsize=0, queue length is unbounded. This queue also cannot be shared directly across threads.

Queue capacity, worker count, and result storage are separate bounds. The batch below holds at most 2 queued items and 2 worker-held items: 4 admitted work items. The producer may hold 1 more while waiting to enqueue it, making at most 5 items in handoff or processing. Stop markers are not work items. range produces input incrementally, but the result dictionary grows with input count. An infinite stream needs a downstream result sink with capacity control too, and item sizes need limits.

asyncio.timeout cancels the current task when its deadline expires and converts the cancellation it caused to TimeoutError outside the context; external cancellation remains CancelledError. A per-item deadline inside the worker measures processing after dequeueing; a batch deadline outside the task group also covers admission and draining waits. Event-loop blocking or cleanup can delay the actual return beyond the deadline. A timeout is not forced termination.

A bounded batch with failure handling​

This program simulates waiting with asynchronous sleeps, then squares the numbers 1 through 6. It uses no network or external processes. Save it as bounded.py and run python3 bounded.py with Python 3.11 or later.

ModeOperation and deadlinesBatch policy
normalEach item waits 0.01 seconds; item deadline 0.05 seconds, batch deadline 2 secondsPrint all results
item-timeoutItem 3 waits 0.2 seconds; same deadlinesRecord its timeout and continue with other items
failureItem 1 raises ValueError after 0.01 seconds; other items wait 0.2 seconds; same deadlinesCancel the group without printing partial results
batch-timeoutEach item waits 0.2 seconds; item deadline 0.5 seconds, batch deadline 0.02 secondsStop the batch without printing partial results
import asyncio

WORKERS = 2
CAPACITY = 2
STOP = object()

async def operation(number, mode):
delay = 0.01
if mode in {"failure", "batch-timeout"}:
delay = 0.2
if mode == "failure" and number == 1:
delay = 0.01
if mode == "item-timeout" and number == 3:
delay = 0.2
await asyncio.sleep(delay)
if mode == "failure" and number == 1:
raise ValueError("invalid item 1")
return number * number

async def run_batch(mode):
queue = asyncio.Queue(maxsize=CAPACITY)
results = {}
tasks = []
closed = 0
status = "ok"
item_limit = 0.5 if mode == "batch-timeout" else 0.05
batch_limit = 0.02 if mode == "batch-timeout" else 2.0

async def produce():
for number in range(1, 7):
await queue.put(number)
await queue.join()
for _ in range(WORKERS):
await queue.put(STOP)

async def consume():
nonlocal closed
try:
while True:
number = await queue.get()
try:
if number is STOP:
return
try:
async with asyncio.timeout(item_limit) as item_deadline:
value = await operation(number, mode)
except TimeoutError:
if not item_deadline.expired():
raise
results[number] = "timeout"
else:
results[number] = value
finally:
queue.task_done()
finally:
closed += 1

try:
async with asyncio.timeout(batch_limit):
try:
async with asyncio.TaskGroup() as group:
tasks.append(group.create_task(produce()))
for _ in range(WORKERS):
tasks.append(group.create_task(consume()))
except* ValueError as errors:
status = "failed"
print(f"{mode}: {sorted(str(e) for e in errors.exceptions)}")
except TimeoutError:
status = "batch-timeout"

print(f"{mode}: status={status}, workers_closed={closed}, "
f"tasks_done={all(task.done() for task in tasks)}")
if status == "ok":
ordered = sorted(results.items())
total = sum(value for value in results.values()
if isinstance(value, int))
print(f"results={ordered}, total={total}")

async def main():
for mode in ("normal", "item-timeout", "failure", "batch-timeout"):
await run_batch(mode)

asyncio.run(main())

The run on Python 3.13.15 produced this output, with results sorted by input number:

normal: status=ok, workers_closed=2, tasks_done=True
results=[(1, 1), (2, 4), (3, 9), (4, 16), (5, 25), (6, 36)], total=91
item-timeout: status=ok, workers_closed=2, tasks_done=True
results=[(1, 1), (2, 4), (3, 'timeout'), (4, 16), (5, 25), (6, 36)], total=82
failure: ['invalid item 1']
failure: status=failed, workers_closed=2, tasks_done=True
batch-timeout: status=batch-timeout, workers_closed=2, tasks_done=True

The normal total is 1 + 4 + 9 + 16 + 25 + 36 = 91; after item 3 times out, it is 91 - 9 = 82. Here status=ok means the batch was handled according to its policy. The item timeout remains explicit in the results; this status does not mean every operation succeeded.

run_batch prints the status and results and returns None. After catching TimeoutError, the worker checks the timeout context’s expired() method to establish whether its own deadline expired. If the operation itself raises the same exception before the deadline expires, the worker re-raises it; the task group cancels the other work and propagates the failure to the caller.

The producer waits for queue.join() before sending one stop marker per worker. Every successful get has one corresponding task_done, including stop markers, timeouts, and exceptions. A get that has not succeeded does not decrement the count. This accounting coordinates draining; it does not establish business success.

If a worker fails, the task group cancels the producer waiting in join or put, along with the other workers. The failure path does not wait for the queue to drain: unread items are discarded with this local queue. except* ValueError handles ValueError and its subclasses as batch failures in this example; unmatched exceptions propagate to the caller in an exception group, and external cancellation also propagates. When the batch deadline expires, the same group finishes before printing the timeout status.

Each workers_closed=2 and tasks_done=True shows that both workers passed through their exit cleanup and all three child tasks ended. Aggregation happens after the group ends. Workers write distinct numbered entries without an intervening await, so this example needs no additional lock. A shared counter with a wait between its read and write would need the synchronization demonstrated earlier.

To connect a real service, replace operation with a cancellable async call while retaining task ownership and both deadline scopes. The example does not retry. Before adding retries, decide the overall deadline, attempt limit, and idempotency of side effects: task-group cancellation will not roll back a completed remote operation.

Explore connectionsOpen network