JEPA4Japan · tutorials

Async I/O with asyncio

1,338 words 7 min read #Python

Coordinate coroutines, tasks, timeouts, and cancellation without blocking the event loop.

Course progress Course outline 24 of 24 lessons available

One event loop passes one baton

Calling an async def function creates a coroutine object. Its body has not run yet. Await it directly, or wrap it in a Task so the event loop can schedule it.

  1. Coroutine object work that can be started
  2. Task scheduled coroutine
  3. Event loop chooses a ready task
  4. await may hand over the baton
Asyncio usually coordinates many waits on one event-loop thread.
import asyncio


async def hello():
    print("body runs")
    return "done"


async def main():
    coroutine = hello()
    print(type(coroutine).__name__)
    print(await coroutine)


asyncio.run(main())
coroutine
body runs
done

asyncio.run() creates one entry-point loop, runs main(), and closes it. Do not nest it inside an existing loop.

Await hands over control only when work suspends

await means “I need this result.” Another task gets a turn only if the awaited operation really suspends. Awaiting a coroutine that immediately returns may not yield control at all.

  1. Task A runs ordinary code keeps the baton
  2. Real wait I/O, event, queue, timer
  3. Task B runs the loop passes the baton
Cooperation happens at a true suspension point, not at every written await.

Blocking calls and long CPU loops freeze the event-loop thread; adding async alone changes neither cost nor scheduling.

Tasks need owners

asyncio.create_task() schedules work now. Keep a strong reference and decide who awaits its result, observes failure, and cancels it during shutdown.

asyncio.gather() schedules awaitables and returns results in argument order, not completion order:

import asyncio


async def label(name, delay):
    await asyncio.sleep(delay)
    return name


async def main():
    results = await asyncio.gather(
        label("slow", 0.02),
        label("fast", 0.01),
    )
    print(results)


asyncio.run(main())
['slow', 'fast']

On an ordinary child failure, default gather() propagates the first observed exception; it does not automatically cancel the other submitted children. return_exceptions=True turns failures into result values, which is appropriate only for an explicit per-item failure contract.

Python 3.11+ provides TaskGroup for structured concurrency:

async def collect_structured(sources):
    tasks = []
    async with asyncio.TaskGroup() as group:
        for source in sources:
            tasks.append(group.create_task(fetch(source)))
    return [task.result() for task in tasks]

Every child belongs to the async with scope. If one child fails with an exception other than CancelledError, the group cancels the remaining children, waits for their cleanup, then raises an ExceptionGroup as the scope exits. That is a fail-together boundary; do not publish partial success accidentally.

Deadlines cancel an owned scope

Python 3.11+ asyncio.timeout() places one deadline around one or several awaits:

async def load_with_deadline():
    try:
        async with asyncio.timeout(0.5):
            return await fetch_everything()
    except TimeoutError:
        return "deadline exceeded"

The timeout cancels the current task inside the context and converts that cancellation to TimeoutError when leaving it, so catch TimeoutError outside the context. It limits the owner’s wait; it cannot prove a remote operation rolled back.

Cancellation is cooperative. task.cancel() arranges for CancelledError at a later suspension point. Release resources in finally, or catch only for cancellation-specific cleanup and then re-raise:

async def worker(resource):
    try:
        await resource.use()
    except asyncio.CancelledError:
        await resource.close()
        raise
  1. Deadline owner stops waiting
  2. Cancel request delivered at suspension
  3. Cleanup release owned resources
  4. Propagate re-raise cancellation
Cleanup is part of cancellation; swallowing it breaks structured owners.

CancelledError directly inherits from BaseException. Structured tools use it internally, so suppress it only when a component deliberately owns a complete alternative policy.

Move blocking I/O, not heavy CPU work

asyncio.to_thread(blocking_function, ...) adapts a legacy blocking I/O call so the event-loop thread stays responsive. Bound simultaneous calls and use the blocking API’s native timeout. Cancelling the awaiting task does not force arbitrary code already running in the worker thread to stop.

Do not send substantial pure-Python CPU work to to_thread() just because it is synchronous. Under the standard CPython 3.12 GIL, it still does not normally run Python bytecode in parallel and can hurt loop latency. Keep tiny calculations inline; consider a bounded process pool or separate worker for measured, independent CPU jobs.

Tiny project: Offline Metadata Collector

Create metadata_collector.py. This Python 3.9+ standard-library project uses local delays, not a network. Explicit task references, gather(), sorting, and wait_for() give deterministic results and a testable deadline. On Python 3.11+, a timeout scope can replace wait_for() when several awaits share one budget.

import asyncio
from dataclasses import dataclass


@dataclass(frozen=True)
class Source:
    slug: str
    title: str
    words: int
    delay: float


async def fetch(source):
    await asyncio.sleep(source.delay)
    return (
        source.slug,
        source.title.strip(),
        max(1, (source.words + 199) // 200),
    )


async def collect(sources):
    sources = list(sources)
    slugs = [source.slug for source in sources]
    if len(slugs) != len(set(slugs)):
        raise ValueError("slugs must be unique")

    tasks = [asyncio.create_task(fetch(source)) for source in sources]
    results = await asyncio.gather(*tasks)
    return sorted(results)


async def main():
    sources = [
        Source("models", " Open models ", 620, 0.02),
        Source("agents", "Reliable agents", 340, 0.01),
        Source("evals", "Practical evals", 199, 0.00),
    ]
    items = await collect(reversed(sources))

    print("Article metadata")
    for slug, title, minutes in items:
        print(f"- {slug}: {title} ({minutes} min)")

    timed_out = False
    try:
        await asyncio.wait_for(
            fetch(Source("slow", "Slow", 10, 10)),
            timeout=0.01,
        )
    except asyncio.TimeoutError:
        timed_out = True

    print(f"Sorted: {[slug for slug, _, _ in items]}")
    print(f"Timeout observed: {timed_out}")


if __name__ == "__main__":
    asyncio.run(main())

Run python3 metadata_collector.py:

Article metadata
- agents: Reliable agents (2 min)
- evals: Practical evals (1 min)
- models: Open models (4 min)
Sorted: ['agents', 'evals', 'models']
Timeout observed: True

Input is reversed and delays differ, yet published order is sorted. Boundary reminder: task creation is cheap, downstream capacity is not. A real service needs a semaphore or other explicit concurrency limit, its own timeout, and a documented retry policy.

Three tiny missions and the Chapter 24 checklist

  1. Watch the baton. Compare awaiting three fetches one by one with gather() while preserving the same result order.
  2. Prove cleanup. Cancel a task after an Event confirms startup; assert its finally block ran and awaiting it still raises CancelledError.
  3. Fail together. On Python 3.11+, make one TaskGroup child fail and verify its waiting sibling cleans up before the group raises.

You are ready for Chapter 24 when:

  • you distinguish a coroutine object, Task, and event loop;
  • you know await yields only when the awaitable really suspends;
  • every created task has an owner that awaits and observes it;
  • you know gather() preserves argument order but does not fail together;
  • you know a TaskGroup owns children and cancels siblings on failure;
  • you place deadlines around the intended scope;
  • cleanup runs before CancelledError continues upward;
  • to_thread() adapts bounded blocking I/O, not CPU parallelism;
  • CPU-heavy work never blocks the event-loop thread;
  • you can run the collector and reproduce its stable output.

Next, you will measure a complete program, improve only proven bottlenecks, package it, and ship it with a release checklist.