Handbooks / Python asyncio / Chapter 3
Async I/O Patterns
42 pages · ~81 min✓ Reviewed
Builds on Tasks, Cancellation & Timeouts. Next up: Async HTTP Clients & Servers.
Part 1 · Async I/O Patterns
Real-World asyncio: Queues, Limits, Streams and the Bugs to Avoid
Most asyncio tutorials stop at async def and await. Real programs need more: a way to hand work from fast producers to slow consumers, a cap on how many requests are in flight, a safe way to share state across an await, and a clean way to talk over a TCP socket. This chapter covers those tools, because that is where async code stops being a toy and starts being a service.
You will meet them in order. Queues give you producer/consumer pipelines with backpressure. Semaphores cap concurrency. Lock and Event coordinate tasks that share state or wait on a signal. Streams speak TCP without callbacks. Async iterators, async context managers and async generators let you plug I/O into async for and async with. Along the way you will see which tool fits which job, and the handful of bugs behind most broken async code: forgotten await, tasks that vanish, a blocked loop, and a swallowed CancelledError.
By the end you can build a bounded worker pool that shuts down cleanly, throttle 1000 requests to 10 at a time, guard a cache against a thundering herd, write a small echo server, and stream paginated data with an async generator. You can also pick between a Queue, gather, a TaskGroup and a thread for a given job, and spot a never-awaited coroutine before it ships.
You need Python 3.11 or newer and nothing beyond the standard library. You should be comfortable with async def, await, asyncio.run() and asyncio.create_task(). If those feel shaky, the first section reviews them. Create each example as its own file and run it with python file.py, and try python -X dev to get asyncio's debug warnings while you learn.
Part 2 · Event Loop and Coroutine Fundamentals
Coroutines and the event loop
A coroutine is what you get from a function declared with async def. The surprising part is that calling such a function does not run its body. It builds a coroutine object, a paused recipe, and hands it back to you. The body only starts once something awaits that object or schedules it on the loop.
import asyncio async def fetch(): return 42 c = fetch() # nothing has run yet print(type(c).__name__) print(asyncio.run(c)) # now the body runs
Calling is not running
coroutine
42Something has to drive these recipes, and that is the event loop. It is a single-threaded scheduler. On each pass it runs every callback that is ready. When nothing is ready, it polls for I/O using the operating system's selector (epoll on Linux, kqueue on macOS and BSD, IOCP on Windows). It then wakes whichever tasks had their I/O complete.
Because there is only one thread, only one coroutine executes at any instant. Switching happens through cooperative multitasking: await hands control back to the loop until the awaitable finishes. It is the only place a coroutine can be switched out. Code between two awaits always runs to completion without interruption.
| Situation | What the loop does |
|---|---|
await aw that is still pending | Runs the next ready task |
await asyncio.sleep(0) | Yields once, then resumes on the next pass |
A stretch of code with no await | Never switches away, so everything else waits |
| Ready queue empty | Polls the selector for I/O |
await asyncio.sleep(0) is the simplest cooperative checkpoint. It does no waiting at all; it just gives the loop a chance to run someone else before you continue. In the demo below, two tasks take turns because each one yields after every step.
import asyncio async def worker(name): for i in range(2): print(name, i) await asyncio.sleep(0) # checkpoint async def main(): t1 = asyncio.create_task(worker('A')) t2 = asyncio.create_task(worker('B')) await t1 await t2 asyncio.run(main())
A 0 B 0 A 1 B 1
A long loop without await starves every other task, and time.sleep(1) blocks the whole loop. Use await asyncio.sleep(...), add an occasional await asyncio.sleep(0) in long loops, and move heavy CPU work to an executor.
Running and scheduling work
asyncio.run(main()) is the normal entry point. It does more than call your function, as the chain below shows. That cleanup is why you should use it instead of managing a loop yourself.
- 1Create a new loopfresh event loop for this call
- 2Run main() to completionuntil it returns or raises
- 3Cancel leftover tasksanything still pending
- 4Shut down async generatorsso their cleanup runs
- 5Close the loopresources released
Awaiting a coroutine directly runs it right there, one step after another. To overlap waits, you schedule it instead with asyncio.create_task(coro). This registers the coroutine to run concurrently and returns a Task, which is a subclass of Future. The Task holds the eventual result or exception, and you await it when you need that result.
The loop keeps only weak references to tasks. If you drop the Task that create_task returned, it may be garbage-collected mid-run. Store it in a set or, better, use a TaskGroup (next page).
The table below puts the main ways of running work side by side. The later pages show each one in code.
| Tool | Returns | On error |
|---|---|---|
create_task | A Task (a Future) | Kept inside the Task until you await it |
gather | List of results, in input order | First error is raised (or collected as a value) |
TaskGroup (3.11+) | Waits for every child | Cancels siblings, raises an ExceptionGroup |
timeout / wait_for | The result | Cancels the work, raises TimeoutError |
asyncio.run() inside an already running loop (such as a Jupyter cell) raises RuntimeError. In that case, just await main().
gather and TaskGroup
asyncio.gather(*aws) runs awaitables concurrently and returns their results as a list in input order, not completion order. In the example, job 2 finishes first, but its result still sits in the second slot. With return_exceptions=True, an error becomes an ordinary value in the list instead of being raised, so one failure does not hide the other results.
import asyncio async def job(n, delay): await asyncio.sleep(delay) return n async def bad(): raise ValueError('boom') async def main(): r = await asyncio.gather( job(1, 0.2), job(2, 0.1), bad(), return_exceptions=True) print(r) asyncio.run(main())
[1, 2, ValueError('boom')]
asyncio.TaskGroup (3.11+) gives you structured concurrency. Every task you create inside the async with block must finish before the block exits. If one child fails, the group cancels its siblings and raises the collected errors together as an ExceptionGroup, which you catch with except*.
| Mode | First error | Other tasks |
|---|---|---|
gather (default) | Raised to the caller | Keep running |
gather(..., return_exceptions=True) | Returned in the list | Keep running |
TaskGroup | ExceptionGroup | Cancelled |
import asyncio async def ok(): try: await asyncio.sleep(1) except asyncio.CancelledError: print('ok cancelled') raise async def bad(): await asyncio.sleep(0.1) raise ValueError('boom') async def main(): try: async with asyncio.TaskGroup() as tg: tg.create_task(ok()) tg.create_task(bad()) except* ValueError as eg: print(type(eg).__name__, eg.exceptions[0]) asyncio.run(main())
Python 3.11+
ok cancelled ExceptionGroup boom
The payoff of concurrency is time. The timeline below shows two one-second sleeps overlapping, and the code that follows measures it: the group finishes in about one second, not two.
| Time | Task A | Task B | Loop |
|---|---|---|---|
| 0.0 s | runs, hits await | runs, hits await | Both parked |
| 0 to 1 s | sleeping | sleeping | Selector waits |
| 1.0 s | wakes, done | wakes, done | Group exits |
| Total | overlaps | overlaps | about 1 s, not 2 s |
import asyncio, time async def main(): start = time.perf_counter() async with asyncio.TaskGroup() as tg: tg.create_task(asyncio.sleep(1)) tg.create_task(asyncio.sleep(1)) print(f'{time.perf_counter() - start:.0f}s total') asyncio.run(main())
1s total
Deadlines and timeouts
Slow work needs a deadline. asyncio.timeout(s) (3.11+) is an async context manager that covers a whole block. asyncio.wait_for(aw, s) puts the deadline on a single awaitable. In both cases, work that runs past the deadline is cancelled and a TimeoutError is raised, so you can catch it and fall back.
import asyncio async def slow(): await asyncio.sleep(5) async def main(): try: async with asyncio.timeout(0.1): await slow() except TimeoutError: print('block timed out') try: await asyncio.wait_for(slow(), 0.1) except TimeoutError: print('wait_for timed out') asyncio.run(main())
block timed out wait_for timed out
The flowchart condenses the page into a decision you can reuse when you have several awaitables to run.
Writing fetch() without await creates a coroutine that never runs, and Python warns that it was never awaited. Remember that a coroutine is a recipe, and await or a task is what runs it.
Create tasks to overlap waits, await them to get results, and wrap slow calls in timeout(). Only await can switch tasks, so keep every stretch between awaits short.
Part 3 · asyncio.Queue: Producer/Consumer
A FIFO for coroutines, and backpressure
An asyncio.Queue is a first-in, first-out buffer that lets one group of coroutines hand work to another. Producers put items in, consumers get them out, and items leave in the order they arrived. The queue belongs to a single event loop. It is not thread-safe, so a plain thread must never call its methods directly.
You size it with asyncio.Queue(maxsize=n). With maxsize=0, the default, the queue is unbounded and put never has to wait. With a positive maxsize, the queue holds at most that many items, and that limit is what makes the pattern useful.
When the queue is full, await q.put(item) suspends the producer until a consumer frees a slot. That pause is backpressure: a fast producer is slowed to the speed of its consumers instead of racing ahead. In the same way, await q.get() suspends a consumer while the queue is empty. If you never want to wait, put_nowait and get_nowait raise an exception instead.
| Call | When the queue is... | Result |
|---|---|---|
await q.put(x) | full | waits for a free slot |
await q.get() | empty | waits for an item |
q.put_nowait(x) | full | raises QueueFull |
q.get_nowait() | empty | raises QueueEmpty |
The program below fills a queue of size 2, shows put_nowait failing, then starts a blocked put as a task. It stays pending until a get makes room.
import asyncio async def main(): q = asyncio.Queue(maxsize=2) await q.put("a") await q.put("b") print("size", q.qsize(), "full", q.full()) try: q.put_nowait("c") except asyncio.QueueFull: print("QueueFull") waiting = asyncio.create_task(q.put("c")) await asyncio.sleep(0) print("put done?", waiting.done()) print("got", await q.get()) await waiting print("put done?", waiting.done()) q.get_nowait() q.get_nowait() try: q.get_nowait() except asyncio.QueueEmpty: print("QueueEmpty") asyncio.run(main())
A full queue blocks put; the nowait calls raise instead
size 2 full True QueueFull put done? False got a put done? True QueueEmpty
put_nowait with a QueueFull handler is the way to drop work when overloaded. get_nowait suits polling code that must never block. q.qsize() never waits, but treat it as a metric only.
Workers, task_done and join
Getting an item out of the queue does not mean it has been processed. To track real completion, the queue keeps a counter of unfinished work. Every put adds one. Every call to q.task_done() subtracts one. await q.join() waits until that counter reaches zero, meaning every item ever put has been marked done.
- 1put(x)counter +1
- 2get()counter unchanged
- 3process(x)your code
- 4task_done()counter -1
- 5join() returnscounter is 0
The standard consumer is a worker coroutine that loops forever. You start N copies as tasks. Each one gets an item, processes it, and calls task_done() in a finally block. That way the counter drops even when processing raises an exception.
The producer puts everything in, then awaits join(). Only after join() returns are the workers cancelled. Cancelling first would throw away items still sitting in the queue. Keep a reference to every worker task, since you need them to cancel.
import asyncio results = [] async def worker(q): while True: item = await q.get() try: await asyncio.sleep(0.01) results.append(item * item) finally: q.task_done() async def main(): q = asyncio.Queue(maxsize=3) workers = [asyncio.create_task(worker(q)) for _ in range(3)] for x in range(1, 7): await q.put(x) await q.join() for w in workers: w.cancel() await asyncio.gather(*workers, return_exceptions=True) print(sorted(results)) print("unfinished", q.qsize()) print("cancelled", all(w.cancelled() for w in workers)) asyncio.run(main())
Three workers, a queue of 3, six items: join first, then cancel
[1, 4, 9, 16, 25, 36] unfinished 0 cancelled True
The producer is limited to three items in flight by maxsize=3. It keeps pausing inside put while the three workers drain the queue. This is backpressure in practice.
If process(item) raises and task_done() sits after it instead of in finally, that item is never marked done. The counter never reaches zero and await q.join() hangs forever. Always write try: ... finally: q.task_done().
Calling w.cancel() before await q.join() drops whatever is still queued. Join first, then cancel.
Shutting workers down
Workers that loop on while True never end on their own, so you need a deliberate shutdown. There are three ways, and the right one depends on your Python version and on whether you want workers to finish by returning or by being cancelled.
| Way | How | Catch |
|---|---|---|
| join + cancel | await q.join(), then w.cancel() on each worker | you must keep the task references |
| Sentinel | put one None per worker; a worker returns when it sees it | worker must check for None |
q.shutdown() (3.13+) | further put and get raise QueueShutDown | needs Python 3.13 or newer |
The sentinel approach puts a marker value, usually None, on the queue. Because the queue is FIFO, every real item ahead of the marker is processed first. You need one None per worker. A single marker is consumed by only one worker, and the others wait forever. If you also use join(), call task_done() for the sentinel too.
import asyncio results = [] async def worker(q): while True: item = await q.get() try: if item is None: return results.append(item * 10) finally: q.task_done() async def main(): q = asyncio.Queue() workers = [asyncio.create_task(worker(q)) for _ in range(2)] for x in range(1, 5): await q.put(x) for _ in workers: await q.put(None) await asyncio.gather(*workers) print(sorted(results)) print("workers finished", all(w.done() for w in workers)) asyncio.run(main())
One None per worker; no cancel is needed
[10, 20, 30, 40] workers finished True
On Python 3.13 and newer, q.shutdown() removes the need for sentinels. After it is called, every put raises QueueShutDown. A get keeps handing out the items still queued, then raises QueueShutDown once the queue is drained. Passing immediate=True discards whatever is left. The worker just catches the exception and returns. This example only defines the worker, because it needs 3.13 to run.
import asyncio async def worker(q, handle): try: while True: await handle(await q.get()) except asyncio.QueueShutDown: return # producer, when finished: q.shutdown()
Python 3.13+: the producer calls q.shutdown() when it has no more items
Putting a single None on a queue shared by four workers lets one worker exit while three wait forever. Put exactly one per worker.
Variants and the cost of unbounded queues
Two subclasses change the order items come out. asyncio.PriorityQueue always returns the smallest entry first, so you put (priority, item) tuples where a lower number means more urgent. asyncio.LifoQueue is a stack: the most recently put item comes out first.
| Variant | Order out | Item shape |
|---|---|---|
Queue | first in, first out | any |
PriorityQueue | smallest first | (priority, item) |
LifoQueue | last in, first out | any |
Tuples compare element by element. If two entries have the same priority, Python compares the items themselves, which can raise TypeError for objects that cannot be ordered. The fix is to add a counter as a tie-breaker: (priority, n, item).
import asyncio async def main(): pq = asyncio.PriorityQueue() jobs = [(2, 0, "email"), (0, 1, "page oncall"), (1, 2, "report"), (0, 3, "disk alarm")] for job in jobs: await pq.put(job) while not pq.empty(): prio, n, name = await pq.get() print(name) lq = asyncio.LifoQueue() for x in (1, 2, 3): await lq.put(x) print("lifo", await lq.get()) asyncio.run(main())
Priority 0 first, ties broken by the counter; the stack returns 3
page oncall
disk alarm
report
email
lifo 3The last design choice is whether to bound the queue. An unbounded queue (maxsize=0) never makes put wait, so a producer faster than its consumers keeps adding items and memory grows without limit. A bounded queue turns that growth into backpressure, so memory stays flat and the producer simply runs at the consumers' pace.
Bounded (maxsize=n) | Unbounded (maxsize=0) | |
|---|---|---|
| Fast producer | put waits, producer slows down | queue keeps growing |
| Memory | stays flat | unlimited |
| Overload | backpressure or QueueFull | hidden until out of memory |
Set maxsize. Start N workers and keep their task references. Put task_done() in finally. Shut down with join then cancel, a None per worker, or shutdown() on 3.13+. Never share the queue with other threads.
Part 4 · Semaphore for Rate Limiting
How a semaphore counts permits
An asyncio.Semaphore(n) is a counter with a waiting line. The counter starts at n and each unit of it is a permit. A task that wants to do guarded work calls acquire(), which takes one permit by lowering the counter. When the counter is already 0, the task parks until another task gives a permit back with release(), which raises the counter and wakes one waiter.
Here is the counter for Semaphore(2) with three tasks A, B and C asking for a permit. C has to wait, and the moment A releases, the permit goes straight to C instead of being left sitting in the counter.
| Step | Counter | What happens |
|---|---|---|
| Start, Semaphore(2) | 2 | nobody holds a permit |
| A acquires | 1 | A runs |
| B acquires | 0 | B runs |
| C acquires | 0 | C waits |
| A releases | 1 then 0 | C wakes and takes the permit |
The semaphore never limits how many tasks you create, only how many get past the acquire() line at the same moment. That makes it a cap on concurrency: at most n pieces of guarded work are in flight. It says nothing about how many calls start per second. If each call finishes in a few milliseconds, ten permits can be recycled hundreds of times every second.
The safe way to hold a permit is async with sem:. It acquires on entry and releases on exit, and the release happens even when the body raises. A bare acquire() followed by release() skips the release as soon as something in between fails, and that permit is gone for good.
import asyncio async def job(sem, stats): async with sem: stats["now"] += 1 stats["peak"] = max(stats["peak"], stats["now"]) await asyncio.sleep(0.01) stats["now"] -= 1 async def main(): sem = asyncio.Semaphore(3) stats = {"now": 0, "peak": 0} await asyncio.gather(*(job(sem, stats) for _ in range(12))) print("peak in flight:", stats["peak"]) print("still running:", stats["now"]) asyncio.run(main())
Twelve tasks, three permits: the peak never passes 3
peak in flight: 3 still running: 0
Now the failure case. The body below raises, yet the permit still comes back, so the semaphore is not locked afterwards.
import asyncio async def main(): sem = asyncio.Semaphore(1) async def bad(): async with sem: raise ValueError("boom") try: await bad() except ValueError as e: print("caught", e) print("locked?", sem.locked()) asyncio.run(main())
caught boom
locked? FalseFanning out over 1000 URLs
The classic use is a crawler. You have a long list of URLs and want them all fetched, but opening a thousand sockets at once would exhaust file descriptors or get you blocked by the server. Create Semaphore(10), wrap the body of fetch(url) in async with sem, and hand every call to gather.
gather still starts 1000 tasks, and all of them begin running. But 990 of them stop at the async with line and wait their turn, so only 10 sockets are open at any instant. The sleep below stands in for the network call.
import asyncio async def fetch(sem, url, stats): async with sem: stats["now"] += 1 stats["peak"] = max(stats["peak"], stats["now"]) await asyncio.sleep(0.001) # pretend network call stats["now"] -= 1 return url.upper() async def main(): urls = [f"site-{i}" for i in range(1000)] sem = asyncio.Semaphore(10) stats = {"now": 0, "peak": 0} pages = await asyncio.gather(*(fetch(sem, u, stats) for u in urls)) print("fetched", len(pages), "pages") print("peak sockets:", stats["peak"]) asyncio.run(main())
The semaphore is created inside main()
fetched 1000 pages peak sockets: 10
url 1 to url 10
each holds a permit
each has an open socket
url 11 to url 1000
parked at async with
no socket yet
Leaks are the main hazard of a semaphore. If code releases a permit it never took, or releases one twice, the counter climbs above n and the cap quietly stops being n. BoundedSemaphore is the same primitive with one extra check: it remembers its starting value, and release() raises ValueError if that would push the counter past it. A bug that a plain Semaphore hides then shows up as a crash at the faulty line.
| Primitive | Extra release() | Result |
|---|---|---|
| Semaphore(n) | allowed | counter grows silently, cap is now n+1 |
| BoundedSemaphore(n) | ValueError | the leak is caught at once |
| Lock | RuntimeError | unlocked lock cannot be released |
import asyncio async def main(): sem = asyncio.BoundedSemaphore(2) await sem.acquire() sem.release() # matches the acquire, fine try: sem.release() # one release too many except ValueError: print("leak caught: ValueError") asyncio.run(main())
leak caught: ValueError
Semaphore(1) behaves roughly like a Lock, since only one task gets in at a time. The differences matter, though. A semaphore has no owner: any task can release a permit another task took, and nothing stops a second release. The Lock refuses a release when it is not held. So Semaphore(1) works as a mutex only until a double release turns it into a two-permit gate.
import asyncio async def main(): sem = asyncio.Semaphore(1) sem.release() # extra release, no error await sem.acquire() await sem.acquire() # second holder gets in too print("two holders on Semaphore(1)") lock = asyncio.Lock() try: lock.release() except RuntimeError: print("Lock refuses release when not held") asyncio.run(main())
two holders on Semaphore(1)
Lock refuses release when not heldReal rate limiting and where to create the semaphore
If an API allows 10 requests per second, Semaphore(10) alone does not respect that. With fast replies, ten permits get recycled many times a second and you can send 100 or more requests in that second. A semaphore bounds how many calls overlap in time, not how often calls begin.
The simplest true rate limit is to keep holding the permit after the call returns and sleep for 1/rate seconds before releasing. A permit then cannot be reused sooner than that interval. With Semaphore(1) and rate = 10, a new request can start at most every 0.1 seconds. With Semaphore(k) you get roughly k requests per interval, which suits an API that allows a few at a time.
import asyncio, time RATE = 10 # requests per second async def call(sem, n): async with sem: result = f"reply {n}" # fast request await asyncio.sleep(1 / RATE) # keep the permit return result async def main(): sem = asyncio.Semaphore(1) t0 = time.perf_counter() await asyncio.gather(*(call(sem, n) for n in range(4))) elapsed = time.perf_counter() - t0 print(f"4 calls took about {round(elapsed, 1)} s") asyncio.run(main())
The sleep happens while the permit is still held
4 calls took about 0.4 s
Sleeping while holding the permit gives an even spacing, but no bursts. When the API accepts bursts, such as 20 requests at once followed by a steady trickle, use a token bucket. The bucket holds up to capacity tokens and refills at rate tokens per second. Each request spends one token, and when the bucket is empty the request waits for a refill. Pair the bucket with a Semaphore if you also want to cap how many requests are in flight.
- 1Refilladd rate times elapsed seconds, up to capacity
- 2Enough tokens?at least 1 token available
- 3Spend onerequest goes out, bursts allowed
- 4Otherwise sleepshort await, then refill again
| Goal | Tool | What it bounds |
|---|---|---|
| At most n calls open at once | Semaphore(n) | concurrency |
| Even spacing, e.g. 10 per second | hold permit + sleep(1/rate) | requests per second |
| Bursts, then a steady rate | token bucket | average rate with a burst size |
| Both rate and open sockets | bucket + Semaphore | rate and concurrency |
Where you create the semaphore matters. In Python versions before 3.10, asyncio primitives bind to an event loop when they are created, and building one at import time, before asyncio.run() starts a loop, attaches it to the wrong loop. The result is errors such as a future attached to a different loop. The habit that works everywhere is to create the semaphore inside main() or another coroutine running on the loop, and pass it down to the functions that need it.
Writing sem = asyncio.Semaphore(10) at the top of the file works on newer Pythons but fails on older ones with a loop-mismatch error. Create it inside main() and pass it to fetch(sem, url).
Semaphore or worker pool, and mistakes to avoid
A semaphore is not the only way to cap concurrency. The other common design is a worker pool: start N long-lived tasks that loop forever, pulling items from an asyncio.Queue. The cap is simply the number of workers. Both designs limit work in flight, but they differ in how many tasks exist and how input arrives.
import asyncio async def worker(q, seen): while True: item = await q.get() try: await asyncio.sleep(0.01) seen.append(item) finally: q.task_done() async def main(): q = asyncio.Queue() seen = [] workers = [asyncio.create_task(worker(q, seen)) for _ in range(3)] for i in range(9): q.put_nowait(i) await q.join() for w in workers: w.cancel() print("processed", len(seen), "items with", len(workers), "tasks") asyncio.run(main())
Three reusable workers handle nine items
processed 9 items with 3 tasks
| Semaphore + gather | Queue worker pool | |
|---|---|---|
| Setup | a few lines, no extra wiring | queue, workers, shutdown |
| Tasks created | one per item | N long-lived tasks |
| Input | a list you already have | a stream, possibly unbounded |
| Memory | grows with the number of items | stays near N |
| Best for | gather-style fan-out | reusing workers on a stream |
Prefer the semaphore when you have a finite list and want the shortest code. Its cost is that every item becomes a task up front, so a million URLs means a million tasks sitting in memory. Move to a queue with workers when input arrives over time, or when the item count is big enough that one task per item is wasteful.
Calling await sem.acquire() and then sem.release() after the work leaks the permit when the work raises. Use async with sem: so the release always runs.
Semaphore(10) limits open requests, not requests per second. Add await asyncio.sleep(1/rate) while holding the permit, or use a token bucket.
It has no owner and does not object to a double release, which silently adds a second permit. Use asyncio.Lock when you need real mutual exclusion.
A semaphore caps concurrency, and async with sem: keeps its counter honest. Use BoundedSemaphore to catch leaks, add a sleep or a token bucket for true requests per second, create the semaphore inside main(), and switch to a queue worker pool when input is a stream.
Part 5 · Lock and Event
Lock: One Task at a Time
Tasks in asyncio take turns on a single thread, and a task only gives up its turn at an await. That is usually safe, but it creates one specific hazard. If you read shared state, await something, and then write the state back, another task can run in the gap and change the state you just read. An asyncio.Lock closes that gap. It is mutual exclusion between coroutines: while one task holds the lock, every other task that asks for it is parked until the holder lets go.
The idiom is async with lock:. It acquires the lock, runs the indented block, and releases the lock even if the block raises. Put the whole read, await and write sequence inside that block so no other task can sneak between the steps.
- 1Task A acquireslock is now held
- 2A reads, awaits, writesB asks for the lock and is parked
- 3A releasesleaving the async with block
- 4B acquiresB now sees A's finished write
The example below is a lost update. Each deposit reads the balance, yields at an await, then writes bal + n. Without a lock both deposits read 100 and both write 110, so one deposit vanishes. With the lock the second deposit starts only after the first has written, and both count.
import asyncio balance = 100 async def audit(): await asyncio.sleep(0) async def deposit_unsafe(n): global balance bal = balance await audit() balance = bal + n async def deposit_safe(lock, n): global balance async with lock: bal = balance await audit() balance = bal + n async def main(): global balance await asyncio.gather(deposit_unsafe(10), deposit_unsafe(10)) print('without lock:', balance) balance = 100 lock = asyncio.Lock() await asyncio.gather(deposit_safe(lock, 10), deposit_safe(lock, 10)) print('with lock:', balance) asyncio.run(main())
The await inside audit() is the race window
without lock: 110 with lock: 120
When you do not need a lock
A lock is only needed when an await sits between reading shared state and writing it. Code with no await in that span runs from start to finish without any other task getting a turn, because the loop can only switch tasks at an await. So counter += 1 on its own is safe in asyncio, and wrapping it in a lock only adds overhead. Ask one question: can another task run between my read and my write? If the answer is no, skip the lock.
| Code shape | Can another task interleave? | Lock needed? |
|---|---|---|
| Read and write with no await between | No | No |
| Read, await, then write using the read value | Yes, at the await | Yes |
| Check a key, await a fetch, store the result | Yes, at the await | Yes |
Not reentrant
Unlike threading.RLock, an asyncio Lock does not remember who holds it. If a task already holds the lock and calls a helper that tries to acquire the same lock, the helper waits for a release that can only come from the task that is waiting. That is a deadlock, and nothing raises an error. The example uses a timeout so you can watch it happen without hanging; cancellation unwinds the async with and frees the lock.
import asyncio async def inner(lock): async with lock: return 'inner done' async def outer(lock): async with lock: return await inner(lock) async def main(): lock = asyncio.Lock() try: await asyncio.wait_for(outer(lock), timeout=0.2) except TimeoutError: print('deadlocked: timed out waiting for own lock') print('locked:', lock.locked()) asyncio.run(main())
deadlocked: timed out waiting for own lock locked: False
Calling a function that does async with lock: from inside another async with lock: block in the same task hangs forever. Split the code into a locked public function and an unlocked inner helper that assumes the caller already holds the lock.
A lock around code with no await protects nothing, since that code can never be interleaved. It only adds noise and hides where the real race window is.
Cache Stampede and Event
Example: a cache that avoids the thundering herd
An async cache is the classic place for a check-then-act race. Several tasks ask for the same missing key at the same moment. Each sees a miss, each starts the slow fetch, and the backend receives many identical requests. That stampede is called the thundering herd. The fix is to hold a lock across the check, the fetch and the store. The first task fetches and fills the cache, and every later task finds the value waiting when its turn comes.
import asyncio cache = {} loads = 0 async def load(k): global loads loads += 1 await asyncio.sleep(0.05) return k.upper() async def get_naive(k): if k not in cache: cache[k] = await load(k) return cache[k] async def get_locked(lock, k): async with lock: if k not in cache: cache[k] = await load(k) return cache[k] async def main(): global loads await asyncio.gather(*(get_naive('a') for _ in range(5))) print('naive loads:', loads) cache.clear() loads = 0 lock = asyncio.Lock() await asyncio.gather(*(get_locked(lock, 'a') for _ in range(5))) print('locked loads:', loads) asyncio.run(main())
The await in load() is the race window
naive loads: 5 locked loads: 1
One lock for the whole cache has a cost: a slow fetch of key a also makes tasks wanting key b queue up. To avoid that, keep a dictionary that maps each key to its own lock, so only tasks after the same key wait on each other.
Event: a flag you can wait on
An asyncio.Event is a boolean flag that starts out False. await ev.wait() blocks until the flag becomes True, and returns immediately if it already is. ev.set() turns it True, ev.clear() turns it back to False, and ev.is_set() reads it. The key behavior is that set() wakes every waiter at once, so an Event works as a broadcast signal. Typical uses are a 'server ready' announcement or a 'shutdown requested' notice that many tasks all need to hear.
import asyncio async def worker(name, ready): await ready.wait() print(name, 'serving') async def main(): ready = asyncio.Event() tasks = [asyncio.create_task(worker(f'w{i}', ready)) for i in range(3)] await asyncio.sleep(0.01) print('is_set:', ready.is_set()) ready.set() await asyncio.gather(*tasks) print('is_set:', ready.is_set()) ready.clear() print('after clear:', ready.is_set()) asyncio.run(main())
is_set: False w0 serving w1 serving w2 serving is_set: True after clear: False
The flag stays set until someone calls clear(). That is useful, because a task that arrives late still passes straight through wait(). It is also the source of a common bug, covered below.
If an Event is reused for repeated signals and you never call ev.clear(), every later await ev.wait() returns instantly because the flag is still True. Clear it once the signal has been handled, or create a fresh Event for each round.
Condition, Threads and Choosing a Tool
Condition: wait until the state is right
An Event answers 'has something happened?'. Sometimes the real question is about shared data, such as 'does the buffer hold at least two items?'. An asyncio.Condition combines a lock with the ability to wait for that kind of predicate. You must hold the lock, using async with cond:, before you call wait(); otherwise you get a RuntimeError. wait() releases the lock while the task sleeps and reacquires it before returning. notify(n) wakes up to n waiters, and notify_all() wakes all of them.
await cond.wait_for(pred) is the safe way to wait. It checks the predicate, sleeps if it is false, and checks again after every wake-up, so a task never proceeds on a state that changed back in the meantime.
import asyncio async def consumer(cond, buf): async with cond: await cond.wait_for(lambda: len(buf) >= 2) print('consumer saw', buf) async def main(): cond = asyncio.Condition() buf = [] task = asyncio.create_task(consumer(cond, buf)) for x in ('a', 'b'): await asyncio.sleep(0.01) async with cond: buf.append(x) print('produced', x) cond.notify() await task asyncio.run(main())
The consumer wakes after the first item, finds the predicate false, and goes back to sleep
produced a produced b consumer saw ['a', 'b']
Calling cond.wait() outside async with cond: raises RuntimeError because the condition's lock is not held. Always take the lock first, and prefer wait_for(pred) over a bare wait() so the predicate is re-checked.
Signalling from another thread
The asyncio primitives are not thread-safe. They assume that only the event loop's own thread touches them. If a worker thread calls ev.set() directly, it manipulates loop internals from the wrong thread, and the waiting task may not wake up promptly or at all. Instead, hand the call to the loop with loop.call_soon_threadsafe(ev.set). That queues set to run on the loop's thread and also wakes the loop if it is idle in its selector. Capture the loop with asyncio.get_running_loop() in async code before you start the thread.
import asyncio import threading import time async def main(): loop = asyncio.get_running_loop() done = asyncio.Event() def work(): time.sleep(0.05) loop.call_soon_threadsafe(done.set) threading.Thread(target=work).start() await done.wait() print('thread finished') asyncio.run(main())
thread finished
Writing done.set() directly inside the thread function looks harmless and often seems to work, but it is not thread-safe. Always go through loop.call_soon_threadsafe.
Picking the right primitive
The four coordination tools differ in what they let through and what a task waits for. Semaphore appears here for comparison; it caps concurrency at n instead of one.
| Primitive | Lets through | A task waits for | Typical use |
|---|---|---|---|
| Lock | One at a time | The lock to be free | Critical section with an await |
| Semaphore | n at a time | A permit to be free | Cap concurrency |
| Event | Everyone, once set | A signal | Ready or shutdown broadcast |
| Condition | Whoever sees the predicate true | State plus a notify | Wait for data to reach a shape |
One at a time? Use a Lock. Waiting for a signal? Use an Event. Waiting for state to change? Use a Condition. And remember that if no await sits between the read and the write, you need none of them.
Part 6 · Streams: open_connection and start_server
The reader and writer pair
Streams are asyncio's high-level interface to TCP. You never touch sockets or callbacks. You await calls on two objects, and the event loop runs other tasks while the bytes are in flight.
await asyncio.open_connection(host, port) connects to a server and returns a pair, reader, writer. The pair belongs to that one connection, so a new connection gives you a new pair.
| StreamReader | StreamWriter | |
|---|---|---|
| Direction | Incoming bytes | Outgoing bytes |
| Main calls | read, readline, readuntil, readexactly | write, drain, close, wait_closed |
| Needs await | Every read call | drain() and wait_closed() only |
Writing: write, then drain
writer.write(data) is an ordinary function, not a coroutine. It only copies the bytes into an in-memory send buffer and returns at once. The loop sends them to the socket in the background.
If the peer reads slowly, that buffer keeps growing. await writer.drain() is the flow-control point (backpressure). It returns immediately while the buffer is small. Once the buffer passes a high-water mark, it pauses your task until the buffer has emptied enough. Pair every write() with await writer.drain().
The example below writes 10 MB in 50 kB chunks, calling drain() after each one. The server counts what it receives and replies with the total. write_eof() tells the server that nothing more is coming, so its read returns b''.
import asyncio async def sink(reader, writer): total = 0 while chunk := await reader.read(65536): total += len(chunk) writer.write(f"{total}\n".encode()) await writer.drain() writer.close() await writer.wait_closed() async def main(): server = await asyncio.start_server(sink, "127.0.0.1", 0) port = server.sockets[0].getsockname()[1] async with server: reader, writer = await asyncio.open_connection("127.0.0.1", port) for _ in range(200): writer.write(b"x" * 50_000) await writer.drain() writer.write_eof() print("server counted", (await reader.readline()).decode().strip()) writer.close() await writer.wait_closed() asyncio.run(main())
Port 0 asks the OS for any free port
server counted 10000000Closing
Closing takes two steps. writer.close() starts the shutdown and returns at once. await writer.wait_closed() then waits until the transport has actually finished closing. Skip the second step and the socket may still be closing when your code moves on.
- 1write + drainlast bytes handed to the socket
- 2writer.close()starts shutdown, returns at once
- 3await wait_closed()transport is really closed
Choosing a read call
A TCP connection is a stream of bytes with no message boundaries. Two write() calls on one side can arrive as one chunk, and one large write can arrive in pieces. So the reader gives you four calls, each with a different idea of when to stop reading.
| Call | Reads | At EOF |
|---|---|---|
read(n) | Up to n bytes, whatever has arrived | Returns b'' |
readline() | Through the next \n | Partial line, or b'' if nothing is left |
readuntil(sep) | Through the separator sep | Raises IncompleteReadError |
readexactly(n) | Exactly n bytes, no fewer | Raises IncompleteReadError |
Pick the call from the shape of your protocol. Line-based chat protocols want readline(). Protocols with another terminator want readuntil(). Length-prefixed binary protocols want readexactly(). read(n) suits raw byte forwarding, such as a proxy, where boundaries do not matter.
The next example uses a helper that starts a server which sends a fixed payload and hangs up. The client reads the same bytes with three different calls. After the last byte, read() returns b'', which is how EOF shows up.
import asyncio async def serve_bytes(payload): async def handle(reader, writer): writer.write(payload) await writer.drain() writer.close() await writer.wait_closed() server = await asyncio.start_server(handle, "127.0.0.1", 0) return server, server.sockets[0].getsockname()[1] async def main(): server, port = await serve_bytes(b"key=value;rest\n") async with server: reader, writer = await asyncio.open_connection("127.0.0.1", port) print(await reader.readuntil(b"=")) print(await reader.readuntil(b";")) print(await reader.readline()) print(await reader.read(100)) writer.close() await writer.wait_closed() asyncio.run(main())
b'key=' b'value;' b'rest\n' b''
Length-prefixed frames with readexactly
A common binary framing is a fixed-size length header followed by that many bytes of body. readexactly(n) fits it, because it keeps waiting until all n bytes are in. If the peer closes first, it raises IncompleteReadError. The exception carries partial, the bytes that did arrive, and expected, the count you asked for.
The server below sends one good frame, then a frame whose header promises 10 bytes but delivers only 3 before hanging up. It reuses serve_bytes from the previous example.
async def read_frame(reader): header = await reader.readexactly(4) size = int.from_bytes(header, "big") return await reader.readexactly(size) async def main(): good = (5).to_bytes(4, "big") + b"hello" short = (10).to_bytes(4, "big") + b"abc" server, port = await serve_bytes(good + short) async with server: reader, writer = await asyncio.open_connection("127.0.0.1", port) print(await read_frame(reader)) try: await read_frame(reader) except asyncio.IncompleteReadError as e: print("short frame:", e.partial, "expected", e.expected) writer.close() await writer.wait_closed() asyncio.run(main())
b'hello' short frame: b'abc' expected 10
read(4) may return 1, 2 or 3 bytes if the data arrives in pieces, and your header parsing then breaks at random. Use readexactly(4) whenever you need a specific number of bytes. Also wrap reads from a peer in asyncio.timeout(), so a dead connection cannot hang the task forever.
Servers: start_server and the echo handler
server = await asyncio.start_server(handle, host, port) binds the port and starts listening. For every client that connects, asyncio creates a new task running handle(reader, writer). Clients are therefore served concurrently, and a slow one never blocks the others. The handler gets a fresh stream pair for its own connection.
The classic first server is an echo server. The handler loops: read one line, write it back, drain. An empty result from readline() means the client hung up, so the loop ends and the handler closes its side.
import asyncio async def handle(reader, writer): while True: data = await reader.readline() if not data: break writer.write(data) await writer.drain() writer.close() await writer.wait_closed() async def main(): server = await asyncio.start_server(handle, "127.0.0.1", 0) port = server.sockets[0].getsockname()[1] async with server: reader, writer = await asyncio.open_connection("127.0.0.1", port) for word in (b"hello\n", b"streams\n"): writer.write(word) await writer.drain() print(await reader.readline()) writer.close() await writer.wait_closed() asyncio.run(main())
b'hello\n' b'streams\n'
Keeping the server alive
A real server should run until it is stopped. await server.serve_forever() does exactly that, and wrapping it in async with server: closes the listening socket when the block exits. That covers normal exit, an exception and cancellation. The demo below stops the serving after 0.2 seconds with asyncio.timeout(), only so that the program can finish. It reuses handle from the echo example.
async def main(): server = await asyncio.start_server(handle, "127.0.0.1", 0) try: async with server, asyncio.timeout(0.2): await server.serve_forever() except TimeoutError: print("stopped after 0.2s, still serving:", server.is_serving()) asyncio.run(main())
stopped after 0.2s, still serving: False
TLS and Unix sockets
The stream API stays the same for encrypted and local connections. Only the setup call changes. For TLS, pass ssl=ssl_context to either open_connection or start_server. For Unix domain sockets, swap in the path-based variants.
| Transport | Client | Server | What changes |
|---|---|---|---|
| TCP | open_connection(host, port) | start_server(handle, host, port) | Nothing extra |
| TCP + TLS | open_connection(host, 443, ssl=ctx) | start_server(handle, host, port, ssl=ctx) | Pass an ssl.SSLContext |
| Unix socket | open_unix_connection(path) | start_unix_server(handle, path) | A file path replaces host and port |
For a client, ssl.create_default_context() gives a context that verifies certificates and host names. A server context must load a certificate chain and key with load_cert_chain().
Pitfalls and a checklist
write() never waits. A loop that writes megabytes without await writer.drain() piles everything into the send buffer, so memory grows without limit while a slow peer catches up. Always follow write() with await writer.drain(), as in the 10 MB example.
write is not a coroutine, so awaiting its result raises TypeError. Call it plainly, then await drain() on the next line.
close() only starts the shutdown. Without await writer.wait_closed() the socket may still be closing when your code carries on, and errors during shutdown go unnoticed.
When the peer disconnects, readline() and read() return b'' forever. A loop without if not data: break then spins at full speed and never ends.
| Rule | Why |
|---|---|
Every write() is followed by await drain() | Backpressure keeps the send buffer bounded |
Every close() is followed by await wait_closed() | The transport is really finished |
Check for b'' on every read loop | It is the only signal that the peer has gone |
Use readexactly for fixed frames | read(n) can return fewer bytes than asked for |
Run the server under async with server | It closes cleanly on exit or cancel |
A handler runs as one task per client, and the same rules apply to the client and server sides. Write and drain, read until b'', then close and wait for the close.
Part 7 · Async Iterators and Async Context Managers
The async iterator protocol
A normal for loop calls __next__ and gets the answer straight away. Many real sources cannot do that: the next row is still on the network, or the next line has not been written yet. Async iteration fixes this by letting each step be a coroutine. An async iterator is any object with two methods. __aiter__ is a plain def that returns self. async def __anext__ returns the next item, and raises StopAsyncIteration once the source is exhausted.
| Sync | Async | |
|---|---|---|
| Loop | for | async for |
| Get the iterator | __iter__ | __aiter__ (still a plain def) |
| Get the next item | __next__ | async def __anext__ |
| End of data | StopIteration | StopAsyncIteration |
| Built-in helper | next(it, d) | await anext(it, d) |
The statement async for x in aiter_obj: is where the protocol pays off. On every step it awaits __anext__. While that await is pending, the event loop is free to run other tasks. So the loop body can receive items that arrive slowly, and the iterator can do I/O (a fetch, a socket read, a sleep) between one item and the next. Here aiter_obj is any object that implements the two methods above.
- 1aiter_obj.aiter()called once, returns the iterator
- 2await anext()I/O may happen here; other tasks run
- 3x = resultloop body runs with the item
- 4StopAsyncIterationloop ends quietly
The class below counts down, sleeping briefly before each item to stand in for real I/O. The loop body never sees the sleep, only the items it produces.
import asyncio class Countdown: def __init__(self, start): self.n = start def __aiter__(self): return self async def __anext__(self): if self.n == 0: raise StopAsyncIteration await asyncio.sleep(0.01) # I/O between items self.n -= 1 return self.n + 1 async def main(): aiter_obj = Countdown(3) async for x in aiter_obj: print('got', x) asyncio.run(main())
Countdown is reused in later examples
got 3 got 2 got 1
aiter() and anext()
Python 3.10 added the built-ins aiter() and anext(), the async twins of iter() and next(). aiter(obj) calls obj.__aiter__(), and await anext(it) pulls a single item by hand. The useful extra is the default form: await anext(it, default) returns the default instead of raising StopAsyncIteration when the iterator is finished, so you can skip a try/except.
async def main(): it = aiter(Countdown(2)) print(await anext(it)) print(await anext(it)) print(await anext(it, 'done')) # no StopAsyncIteration asyncio.run(main())
2 1 done
__aiter__ must be a plain def that returns the iterator. If you write async def __aiter__, async for receives a coroutine instead of an iterator and fails.
Comprehensions and async context managers
Comprehensions can use async features too, but only inside an async def; anywhere else they are a SyntaxError. [x async for x in src] awaits __anext__ for each item. [await f(x) for x in xs] awaits each call in turn, so the calls run one after another, not at the same time. Set and dict comprehensions and if filters work the same way. For real concurrency, hand the coroutines to asyncio.gather.
async def double(x): await asyncio.sleep(0.01) return x * 2 async def main(): a = [x async for x in Countdown(3)] b = [await double(x) for x in a] # serial c = await asyncio.gather(*(double(x) for x in a)) # concurrent print(a, b, c) asyncio.run(main())
[3, 2, 1] [6, 4, 2] [6, 4, 2]
Async context managers
A context manager guarantees cleanup, and cleanup often needs I/O: closing a connection, flushing a buffer, releasing a remote lock. An async context manager defines async def __aenter__, which does the setup and returns the resource, and async def __aexit__(self, exc_type, exc, tb), which does the cleanup. You use it with async with, which awaits both ends.
__aexit__ also decides what happens to an exception raised in the body. A truthy return value suppresses it, exactly like the sync __exit__. Returning None or False lets it propagate, and raising a new error replaces the original.
| aexit returns | Exception from the body |
|---|---|
None or False | re-raised after cleanup |
True | suppressed, code continues after the block |
| raises a new error | replaces the original |
class Conn: def __init__(self, name, suppress=False): self.name = name self.suppress = suppress async def __aenter__(self): print('open', self.name) return self async def __aexit__(self, exc_type, exc, tb): print('close', self.name, exc_type.__name__ if exc_type else None) return self.suppress async def main(): async with Conn('a') as c: print('using', c.name) async with Conn('b', suppress=True): raise ValueError('boom') print('still running') asyncio.run(main())
open a
using a
close a None
open b
close b ValueError
still runningA stray return True at the end of __aexit__ hides every error from the body. Return None unless you mean to swallow a specific exception.
Generator style and dynamic stacks
For a one-off manager, writing a class is heavy. @contextlib.asynccontextmanager turns an async generator into one: the code before yield is setup, the yielded value is the resource, and the code after it is cleanup. Put the yield inside try and the cleanup in finally so it runs even when the body raises.
When the number of managers is only known at runtime, a fixed async with a, b: is not enough. contextlib.AsyncExitStack collects managers as you go with await stack.enter_async_context(cm), then closes them all in reverse order when the stack exits.
from contextlib import asynccontextmanager, AsyncExitStack @asynccontextmanager async def connect(name): print('connect', name) try: yield name.upper() finally: print('close', name) async def main(): async with connect('db') as r: print('got', r) async with AsyncExitStack() as stack: rs = [await stack.enter_async_context(connect(n)) for n in ('a', 'b')] print(rs) asyncio.run(main())
connect db got DB close db connect a connect b ['A', 'B'] close b close a
Need state or reuse? Write a class with __aenter__ and __aexit__. One-off setup and teardown? Use @asynccontextmanager.
A paginated API client, and the usual mistakes
A paginated API is the classic use for an async iterator. The caller wants a flat stream of items, but the server hands them out a page at a time. The iterator keeps a buffer of the current page. Each call to __anext__ pops from the buffer, and only when the buffer is empty does it spend an await fetching the next page. If there is no next page either, it raises StopAsyncIteration.
PAGES = {1: (['a', 'b'], 2), 2: (['c'], None)}
async def fetch_page(n):
await asyncio.sleep(0.01)
print(f' fetch page {n}')
return PAGES[n]
class Pages:
def __init__(self):
self.buf = []
self.next_page = 1
def __aiter__(self):
return self
async def __anext__(self):
if not self.buf and self.next_page:
items, self.next_page = await fetch_page(self.next_page)
self.buf = list(items)
if not self.buf:
raise StopAsyncIteration
return self.buf.pop(0)
async def main():
async for item in Pages():
print(item)
asyncio.run(main())I/O happens only when the buffer runs dry
fetch page 1 a b fetch page 2 c
| Call | Buffer before | Next page | Action |
|---|---|---|---|
| 1 | empty | 1 | fetch page 1, return a |
| 2 | b | 2 | return b, no I/O |
| 3 | empty | 2 | fetch page 2, return c |
| 4 | empty | none | raise StopAsyncIteration |
The loop in main has no idea pages exist. That is the point of the protocol: the paging and the waiting stay inside the iterator.
Mixing up sync and async
with Conn('a'): raises TypeError saying the object does not support the context manager protocol. for item in Pages(): raises TypeError saying the object is not iterable. The fix is async with and async for, which are only allowed inside async def.
An async object takes async with and async for. Async generators, covered in the next section, write __aiter__ and __anext__ for you.
Part 8 · Async Generators
An async def with yield
An async generator is an async def function whose body contains yield. Calling it does not run the body. It returns an async generator object, and you pull values out of it with async for. In the previous section you wrote __aiter__ and __anext__ by hand. An async generator gives you that protocol for free, because Python builds both methods from the function body.
The body can await between yields, so every item can involve I/O or a pause. The classic example is async def ticks(): while True: await asyncio.sleep(1); yield time.time(), which hands out a timestamp once a second for as long as the consumer keeps asking. The runnable version below stops after three items so it finishes quickly.
import asyncio async def countdown(n): while n > 0: await asyncio.sleep(0.01) # I/O or waiting is fine here yield n n -= 1 async def main(): async for x in countdown(3): print("tick", x) asyncio.run(main())
The loop awaits the sleep, then receives each yielded value
tick 3 tick 2 tick 1
Two restrictions come from the fact that this is a generator and a coroutine at the same time. Python rejects both of them when it compiles the function, not when you run it.
| Written inside an async generator | Result | Do this instead |
|---|---|---|
yield from other() | SyntaxError | async for x in other(): yield x |
return 42 | SyntaxError | yield 42, then a bare return |
bare return | Allowed | Ends the generator, and async for stops |
await something() | Allowed | Use it freely between yields |
Trying to hand back a final result with return value. An async generator has no return value, so the compiler raises a SyntaxError. Yield the value as the last item, or keep it in an object the caller can read.
Driving it by hand: asend, athrow, aclose
async for is only one way to drive a generator. The object also has three methods you can await. Each of them resumes the generator at the paused yield, but they differ in what they inject there.
| Method | What happens at the paused yield | Returns |
|---|---|---|
await g.asend(v) | The yield expression evaluates to v | The next yielded value |
await g.athrow(exc) | exc is raised at the yield | The next yielded value, if the generator catches it and yields again |
await g.aclose() | GeneratorExit is raised at the yield | Nothing; finally blocks run |
The first call to asend must pass None, because the generator has not reached a yield yet and there is nothing to receive a value. The aclose() method matters most for cleanup. It raises GeneratorExit at the paused yield, so the try/finally around that point runs and can release files or sockets. The example below uses all three.
import asyncio async def echo(): got = None try: while True: try: got = yield got # asend(v) makes this expression v except ValueError: got = "recovered" # athrow(ValueError()) lands here finally: print("cleanup ran") # aclose() triggers this async def main(): g = echo() print(await g.asend(None)) # start it, run to the first yield print(await g.asend("hi")) print(await g.athrow(ValueError())) await g.aclose() print("closed") asyncio.run(main())
None
hi
recovered
cleanup ran
closedNotice that the inner except ValueError does not catch GeneratorExit. That is deliberate. If a generator swallowed GeneratorExit and kept yielding, aclose() would fail with a RuntimeError.
Early break and deterministic cleanup
Here is the trap. When you break out of an async for, the loop simply stops asking for items. It does not call aclose(). The generator stays paused at its last yield, and its finally block has not run. Cleanup waits until the object is garbage collected or until the event loop shuts down. Whatever the generator holds, such as a socket, a file or a database cursor, stays open for that whole time.
The fix, available since Python 3.10, is contextlib.aclosing. Wrapping the generator in async with aclosing(gen) as g: guarantees that await g.aclose() runs when the block exits, whether by break, by an exception or by normal exhaustion. The finally block then runs at a known point in your code.
- 1breakasync for stops pulling
- 2Plain loopgenerator stays paused; finally waits for GC or loop shutdown
- 3aclosing blockasync with exit awaits aclose()
- 4GeneratorExitraised at the paused yield; finally runs now
As a safety net, asyncio.run calls loop.shutdown_asyncgens() as it finishes. That closes every async generator that is still alive and was never finalized, so their finally blocks do run before the process exits. The next example keeps a reference to the plain-break generator so you can see this happen at the very end.
import asyncio from contextlib import aclosing keep = [] async def rows(name): try: for i in range(1, 10): await asyncio.sleep(0) yield i finally: print("closed:", name) async def main(): g = rows("plain") keep.append(g) async for x in g: if x == 2: break print("after plain break") async with aclosing(rows("aclosing")) as h: async for x in h: if x == 2: break print("after aclosing") print("main done") asyncio.run(main())
after plain break closed: aclosing after aclosing main done closed: plain
The aclosing generator was closed the moment its block ended, before the next line of main. The plain one was closed only at the very end, by shutdown_asyncgens, long after the loop that abandoned it.
Breaking out of async for over a generator that holds a resource and assuming it is released. Wrap any generator you might abandon early in aclosing, or make sure it is always run to exhaustion.
Pipelines, and generator versus Queue
Because an async generator is also something you can async for over, generators chain naturally. Write a source generator, then a filter that loops over the source and yields only some items, then a map that loops over the filter and transforms each one. The consumer pulls from the last stage. Each pull travels back up the chain, so no stage does work until the one after it asks.
import asyncio async def source(): for n in range(1, 5): await asyncio.sleep(0) print("source", n) yield n async def evens(src): async for x in src: if x % 2 == 0: yield x async def squares(src): async for x in src: yield x * x async def main(): async for v in squares(evens(source())): print("got", v) asyncio.run(main())
source 1 source 2 got 4 source 3 source 4 got 16
The output interleaves because the work is pull-based. Source 3 is not produced until the consumer has finished with the first result. That is why a pipeline needs no buffer and holds only one item per stage at a time.
A Queue solves a different problem. It is push-based: producers put items in whenever they are ready, and any number of workers take them out. Choose by who is on each end.
| Aspect | Async generator | asyncio.Queue |
|---|---|---|
| Flow | Pull: the consumer asks for the next item | Push: producers put items in |
| Producers | One (the generator body) | Many |
| Consumers | One async for loop | Many workers |
| Buffering | Lazy, none between items | Optional, via maxsize |
| Shutdown | aclose() or aclosing | join(), sentinels or shutdown() |
One reader walking over one lazy stream, with transformation stages in between: use async generators. Many producers or several workers sharing the load: use a Queue.
Part 9 · Choosing the Right Tool: Trade-offs
Fan-out: Queue, gather and TaskGroup
asyncio gives you several ways to run many things at once. They look interchangeable on a small demo and behave very differently under load or failure. This section compares them by the questions that decide the choice: how big is the input, what happens when one job fails, and who is waiting for which result.
Queue + workers vs gather + Semaphore
With gather plus a Semaphore, you create one coroutine or task for every item before anything runs. The semaphore only limits how many are in flight. The rest sit parked, and each one still costs memory. With a million items that is a million task objects held at once, so memory grows as O(n).
A queue with N workers inverts this. You start N long-lived tasks, and the producer feeds items in as they arrive. If the queue has a maxsize, put() blocks when it is full, so a fast producer is slowed to the speed of the consumers. That is backpressure, and it is why the queue is the right shape for streaming or unbounded input, where the full list does not exist yet.
| Queue + workers | gather + Semaphore | |
|---|---|---|
| Input | Streaming or unbounded | Known list up front |
| Tasks created | N workers | One per item, all up front |
| Memory | O(workers) | O(n) |
| Backpressure | Yes, put() waits when full | None |
| Setup | More wiring | Simplest |
A short list you already hold in memory: gather plus a semaphore is the least code. An endless feed, a file too big to load, or a producer faster than the consumers: use a bounded queue.
gather vs TaskGroup on failure
The two differ most when something goes wrong. If one awaitable passed to gather raises, the exception reaches the caller immediately, but the sibling tasks keep running in the background unless you cancel them yourself. A TaskGroup (3.11+) treats the block as one unit. When a child fails, it cancels the siblings, waits for them to finish unwinding, and raises all the errors together as an ExceptionGroup.
The example below runs the same pair of jobs both ways. slow logs how it ends, and bad fails after 50 ms.
import asyncio async def slow(log): try: await asyncio.sleep(0.2) log.append('slow finished') except asyncio.CancelledError: log.append('slow cancelled') raise async def bad(): await asyncio.sleep(0.05) raise ValueError('boom') async def with_gather(): log = [] try: await asyncio.gather(slow(log), bad()) except ValueError as e: log.append(f'gather raised {e}') await asyncio.sleep(0.3) # let the orphan finish return log async def with_group(): log = [] try: async with asyncio.TaskGroup() as tg: tg.create_task(slow(log)) tg.create_task(bad()) except* ValueError as eg: log.append(f'group raised {len(eg.exceptions)} error(s)') return log async def main(): print(await with_gather()) print(await with_group()) asyncio.run(main())
['gather raised boom', 'slow finished'] ['slow cancelled', 'group raised 1 error(s)']
With gather, the caller saw the error while slow carried on and finished anyway, doing work nobody wanted. With TaskGroup, slow was cancelled before the block exited.
| On failure | gather | TaskGroup |
|---|---|---|
| Siblings | Keep running | Cancelled |
| Errors | First one raised | All grouped in an ExceptionGroup |
| Cleanup | You cancel them yourself | Automatic |
| Keep going anyway | return_exceptions=True | Not its job; catch inside each child |
Using bare gather for a fail-fast batch and assuming the other tasks stopped. They did not, and they may still be touching files, sockets or shared state after the caller has moved on. Pick TaskGroup for fail-fast, or gather(..., return_exceptions=True) when you want every result even if some fail.
Waiting Helpers: wait and as_completed
gather is all-or-nothing: you get one list, in input order, once everything is done. Sometimes you want to react earlier. Two helpers cover that, and they answer different questions.
wait: finer control with sets
asyncio.wait(aws, return_when=...) does not return results. It returns two sets of tasks, (done, pending), and you decide what to do with each. With FIRST_COMPLETED it returns as soon as any one task finishes, which suits a race where the first answer wins and the rest should be cancelled. The other modes are FIRST_EXCEPTION and ALL_COMPLETED. Pass tasks, not bare coroutines.
as_completed: fastest first
asyncio.as_completed(aws) hands you awaitables in the order they finish, not the order you supplied them. The quickest result arrives first, so you can start processing it while slower ones are still running.
import asyncio async def fetch(name, delay): await asyncio.sleep(delay) return name async def main(): tasks = [asyncio.create_task(fetch('a', 0.3)), asyncio.create_task(fetch('b', 0.1)), asyncio.create_task(fetch('c', 0.2))] done, pending = await asyncio.wait( tasks, return_when=asyncio.FIRST_COMPLETED) print([t.result() for t in done], len(pending)) for t in pending: t.cancel() jobs = [fetch('a', 0.3), fetch('b', 0.1), fetch('c', 0.2)] for fut in asyncio.as_completed(jobs): print(await fut) asyncio.run(main())
['b'] 2 b c a
The input order was a, b, c, but b was fastest, so it came out first in both halves. The first half stopped early and cancelled the two stragglers. The second half let everything finish but delivered each result as soon as it was ready.
| Call | Returns | Order | Reach for it when |
|---|---|---|---|
gather | List of results | Input order | You need every result, lined up with the inputs |
wait | (done, pending) sets | None; you inspect the sets | Stop early, race, or cancel the rest |
as_completed | Awaitables, one per input | Completion order | Process fast results while slow ones run |
TaskGroup | Nothing; read each task's .result() | Not applicable | Fail fast and clean up automatically |
Forgetting that wait leaves the pending tasks running. If you only wanted the first result, cancel the pending set yourself, as the example does, or those tasks keep working in the background.
asyncio vs Threads, and Offloading Work
Cooperative vs preemptive
asyncio switches between tasks only at an await. Between two awaits your code runs uninterrupted, so the points where another task can interleave are visible in the source and deterministic. Threads are preemptive: the OS can pause one at any instruction, so shared data needs locks everywhere. Tasks are also tiny, roughly a few KB each, while a thread carries an OS stack measured in MB. Switching between tasks costs on the order of microseconds, and a single loop comfortably holds 10k+ idle connections, where thread-per-connection servers top out in the hundreds.
| asyncio | Threads | |
|---|---|---|
| Switch point | Only at await | Anywhere, preemptive |
| Cost each | ~KB per task | ~MB of stack |
| Switch cost | ~µs, a function-call-sized hop | OS context switch |
| Idle connections | 10k+ on one loop | Hundreds |
| CPU-bound work | No speed-up | Still bound by the GIL |
| Shared-state races | Only across an await | Anywhere |
The flip side of cooperative switching is that a task that never awaits blocks everyone. asyncio helps with waiting, not with computing. A loop of pure CPU work gets no faster under asyncio, and while it runs, every other task is frozen.
Offloading CPU-bound and blocking work
There are two escape hatches, and the right one depends on why the call is slow. For CPU-bound work, hand the function to another process with loop.run_in_executor(ProcessPoolExecutor(), fn, *args). Separate processes have separate GILs, so they run in parallel. For a blocking library call such as requests.get or a legacy driver, use await asyncio.to_thread(fn, *args) (3.9+). It runs the call in the loop's default thread pool, and the loop stays free. Threads are fine here because blocking calls release the GIL while they wait.
import asyncio, time from concurrent.futures import ProcessPoolExecutor def blocking_io(): time.sleep(0.2) # stands in for requests.get(...) return 'io done' def crunch(n): return sum(i * i for i in range(n)) async def main(): loop = asyncio.get_running_loop() with ProcessPoolExecutor() as px: results = await asyncio.gather( asyncio.to_thread(blocking_io), asyncio.to_thread(blocking_io), loop.run_in_executor(px, crunch, 100_000), ) print(results) if __name__ == '__main__': asyncio.run(main())
The __main__ guard matters: process pools re-import the module on some platforms.
['io done', 'io done', 333328333350000]
The three calls were awaited together, and the loop stayed responsive the whole time. The two blocking sleeps ran side by side in worker threads, and the number crunching ran in a separate process.
Using asyncio to speed up CPU math, which gives no gain. Calling requests.get straight from a coroutine, which stalls every other task until it returns. Wrap such calls in to_thread.
Streams vs Protocols, and Quick Picks
Streams vs Protocols
There are two layers for network code. Streams (open_connection, start_server) let you write straight-line code with await reader.readline() and writer.write(...), which is easy to read and debug. asyncio.Protocol works with callbacks instead: the loop calls data_received(data) as bytes arrive, with no coroutine or task per read. That has the lowest overhead, which matters for high-throughput servers, but your logic is split across callbacks and you track state yourself.
import asyncio class Echo(asyncio.Protocol): def connection_made(self, transport): self.transport = transport def data_received(self, data): self.transport.write(data.upper()) async def main(): loop = asyncio.get_running_loop() server = await loop.create_server(Echo, '127.0.0.1', 0) port = server.sockets[0].getsockname()[1] r, w = await asyncio.open_connection('127.0.0.1', port) w.write(b'hello\n') await w.drain() print(await r.readline()) w.close() await w.wait_closed() server.close() await server.wait_closed() asyncio.run(main())
A Protocol server, talked to with the stream API.
b'HELLO\n'| Streams | asyncio.Protocol | |
|---|---|---|
| Style | await calls | Callbacks |
| Readability | High | Lower |
| Overhead | A little more | Lowest |
| Best for | Most applications | High-throughput servers |
Start with streams. Move to a Protocol only when profiling shows per-read overhead is the bottleneck.
Quick picks
| Situation | Reach for |
|---|---|
| Streaming or unbounded input | Queue + workers |
| Known list, cap on concurrency | gather + Semaphore |
| Fail fast together | TaskGroup |
| All results, some may fail | gather(..., return_exceptions=True) |
| First result wins | wait(..., FIRST_COMPLETED) |
| Handle results as they finish | as_completed |
| Blocking library call | to_thread |
| CPU-heavy function | run_in_executor + ProcessPoolExecutor |
| Simple network app | Streams |
On 3.11+, start with TaskGroup for fan-out, a bounded Queue for streams, and to_thread for blocking calls. Add asyncio.timeout() around any of them for a deadline, then reach for the specialised tools above only when a real need shows up.
Part 10 · Common Bugs and How to Catch Them
Work That Quietly Never Happens
This section covers the bugs that make asyncio code misbehave without crashing it. The first three are the quietest: a missing await, a task that gets garbage-collected, and a background failure nobody reads. In all three the program keeps running, so you only notice when a result is missing.
The forgotten await
Calling an async def function does not run it. It returns a coroutine object, and the body runs only when something awaits it or wraps it in a task. If you write fetch() and drop the result, nothing happens, and Python prints RuntimeWarning: coroutine 'fetch' was never awaited when the object is thrown away. The same thing happens when you store the coroutine and use it as if it were the answer, for example data = fetch() followed by data['id'].
Lost tasks and silent failures
create_task() schedules work, but the event loop keeps only a weak reference to the task. If you do not hold your own reference, the garbage collector may delete the task halfway through its run, and the work just stops. The fix is to keep each task in a set and remove it when it finishes, or to use a TaskGroup, which holds its tasks for you.
A related symptom is Task exception was never retrieved. A background task raised an error, nobody awaited it, and the exception was only reported when the task was collected. Awaiting the task, or attaching a done-callback that reads task.exception(), brings the failure back to a place where you can see it.
import asyncio async def fetch(): return 42 async def job(n): await asyncio.sleep(0.01) if n == 2: raise ValueError('job 2 failed') return n async def main(): c = fetch() print(type(c).__name__) print(await c) bg = set() tasks = [] for n in (1, 2): t = asyncio.create_task(job(n)) bg.add(t) # strong reference t.add_done_callback(bg.discard) # forget it when done tasks.append(t) for t in tasks: try: print('ok', await t) except ValueError as e: print('caught', e) await asyncio.sleep(0) print(len(bg)) asyncio.run(main())
Await the coroutine, keep every task referenced, and read every task's outcome
coroutine 42 ok 1 caught job 2 failed 0
| Symptom | What happened | Fix |
|---|---|---|
| RuntimeWarning: coroutine ... was never awaited | The coroutine object was created but never run | Add await, or wrap it in create_task |
| A background job just stops partway | The task was garbage-collected with no strong reference | Keep tasks in a set, or use a TaskGroup |
| Task exception was never retrieved | A background task failed and nobody looked | Await the task or add a done-callback |
Writing asyncio.create_task(job()) as a bare statement. The returned task is the only thing keeping it alive, so assign it to a name that lives in a set or use a TaskGroup.
Blocking the Loop and Turning On Debug Mode
The event loop runs on one thread and switches tasks only at an await. A blocking call such as time.sleep(1), requests.get(url) or a long CPU loop never gives control back, so every other task freezes until it returns. Timers drift, sockets sit unread, and a server stops answering anyone.
| Blocking call | Why it freezes the loop | Async fix |
|---|---|---|
time.sleep(1) | The thread sleeps and no task can run | await asyncio.sleep(1) |
requests.get(url) | The thread waits on the socket | An async client such as aiohttp or httpx |
| Heavy CPU loop | The thread is busy computing | await asyncio.to_thread(fn), or a process pool for pure CPU work |
| Blocking library you cannot replace | It is not written for asyncio | await asyncio.to_thread(fn, arg) |
The demo below runs a ticker task that counts ticks every 20 ms while another coroutine waits 0.2 seconds. With time.sleep called directly the ticker gets no turns at all. With the same sleep moved into to_thread, the loop stays free and the ticker keeps counting.
import asyncio, time async def ticker(counter): while True: await asyncio.sleep(0.02) counter[0] += 1 async def measure(label, work): counter = [0] t = asyncio.create_task(ticker(counter)) await asyncio.sleep(0) # let the ticker start await work() t.cancel() print(label, 'ticker ran:', counter[0] >= 5) async def blocking(): time.sleep(0.2) async def offloaded(): await asyncio.to_thread(time.sleep, 0.2) async def main(): await measure('time.sleep ', blocking) await measure('to_thread ', offloaded) asyncio.run(main())
time.sleep ticker ran: False to_thread ticker ran: True
Let asyncio tell you what is wrong
Debug mode makes the loop log any callback or task step that holds the thread for longer than 100 ms, and it reports coroutines that were never awaited with the place they were created. Turn it on with asyncio.run(main(), debug=True) or with the environment variable PYTHONASYNCIODEBUG=1. The slow-callback messages point straight at the blocking call from the table above.
| Switch | How to set it | What you get |
|---|---|---|
debug=True | asyncio.run(main(), debug=True) | Slow callbacks over 100 ms and never-awaited coroutines are logged |
PYTHONASYNCIODEBUG=1 | Set it in the shell before starting Python | The same debug mode, with no code change |
python -X dev app.py | Command line flag | Development mode, which includes asyncio debug |
python -W error::RuntimeWarning | Command line flag | A never-awaited coroutine raises instead of just warning, so tests fail |
Run the test suite with debug mode on and treat every asyncio warning as a failure. A forgotten await then breaks a test the day it is written.
Cancellation and Event Loop Ownership
Swallowing CancelledError
Cancelling a task works by raising CancelledError inside it at its current await. A bare except: or an except BaseException: that does not re-raise catches that error and carries on, so the task is never cancelled. task.cancel() then appears to do nothing, and every timeout built on cancellation, including asyncio.timeout() and wait_for(), stops working. Catch it if you need to clean up, but always raise it again.
Running the wrong loop
asyncio.run() creates a new loop, so it cannot be called when a loop is already running. This is the usual surprise in Jupyter, where a loop is already running behind the notebook, and the error is RuntimeError. In a notebook, write await main() in the cell instead.
A related bug involves an asyncio primitive, such as an Event, Lock or Queue, or a client session. These objects attach themselves to the loop that first waits on them. Using the same object from a different loop, for example after a second asyncio.run(), raises RuntimeError with a message like is bound to a different event loop or attached to a different loop. Create such objects inside the coroutine that runs on the loop that will use them.
import asyncio async def stubborn(): try: await asyncio.sleep(10) except BaseException: print('swallowed') # no re-raise: cancel is lost async def polite(): try: await asyncio.sleep(10) except asyncio.CancelledError: print('cleanup') raise # always re-raise async def cancel_demo(): for fn in (stubborn, polite): t = asyncio.create_task(fn()) await asyncio.sleep(0) t.cancel() try: await t except asyncio.CancelledError: pass print(fn.__name__, 'cancelled:', t.cancelled()) ev = asyncio.Event() async def first(): w = asyncio.create_task(ev.wait()) # binds ev to this loop await asyncio.sleep(0) ev.set() await w ev.clear() async def second(): try: await ev.wait() # a different loop now except RuntimeError as e: print('RuntimeError mentions a different loop:', 'different event loop' in str(e)) async def inner(): return 'inner' async def nested(): coro = inner() try: asyncio.run(coro) except RuntimeError as e: print(e) coro.close() print(await inner()) asyncio.run(cancel_demo()) asyncio.run(first()) asyncio.run(second()) asyncio.run(nested())
Four loop-related bugs in one run
swallowed stubborn cancelled: False cleanup polite cancelled: True RuntimeError mentions a different loop: True asyncio.run() cannot be called from a running event loop inner
Writing except Exception: is fine for cancellation on Python 3.8+, but except: and except BaseException: are not. If you must catch them, put raise at the end of the handler.
Creating a Lock, Queue or session at import time and then calling asyncio.run() more than once. Build them inside main() so each run gets its own.
Sessions, Races and a Triage Guide
Unclosed sessions
An aiohttp.ClientSession owns a pool of connections. If it is never closed, Python warns Unclosed client session when it is collected, and the sockets stay open until then. Open it with async with ClientSession() as s: so it closes even on errors. Create one session per app and pass it around, not one per request, because a new session throws away the connection pool and pays a fresh handshake each time.
aiohttp is not in the standard library, so the demo uses a small stand-in class with the same shape (async with, a get call). It counts how many sessions each style opens for three requests.
import asyncio class Session: opened = 0 def __init__(self): Session.opened += 1 async def get(self, url): await asyncio.sleep(0) return 'GET ' + url async def __aenter__(self): return self async def __aexit__(self, *exc): pass # a real session closes here urls = ['/a', '/b', '/c'] async def per_request(): for u in urls: async with Session() as s: await s.get(u) async def per_app(): async with Session() as s: await asyncio.gather(*(s.get(u) for u in urls)) async def main(): await per_request() print('sessions, one per request:', Session.opened) Session.opened = 0 await per_app() print('sessions, one per app:', Session.opened) asyncio.run(main())
sessions, one per request: 3 sessions, one per app: 1
Check-then-act races
Single-threaded does not mean race-free. Code like if key not in cache: cache[key] = await fetch(key) has an await between the check and the write, and that is where other tasks get their turn. Every task that arrives before the first fetch finishes also sees a miss, so the unguarded version produces duplicate work. Wrapping the check and the write in one asyncio.Lock makes later tasks wait and then find the value already cached.
import asyncio cache = {} calls = 0 async def fetch(k): global calls calls += 1 await asyncio.sleep(0.01) return k.upper() async def get_racy(k): if k not in cache: cache[k] = await fetch(k) return cache[k] async def get_safe(k, lock): async with lock: if k not in cache: cache[k] = await fetch(k) return cache[k] async def main(): global calls await asyncio.gather(*(get_racy('a') for _ in range(5))) print('racy fetches:', calls) cache.clear() calls = 0 lock = asyncio.Lock() await asyncio.gather(*(get_safe('a', lock) for _ in range(5))) print('locked fetches:', calls) asyncio.run(main())
racy fetches: 5 locked fetches: 1
Guarding code that has no await inside with a lock. Without an await between the check and the write no other task can run, so there is no race and the lock only slows things down.
From message to fix
| You see | Likely cause | Fix |
|---|---|---|
| coroutine ... was never awaited | A missing await | Add await or create a task |
| Task exception was never retrieved | A background task failed unread | Await it or add a done-callback |
| Unclosed client session | A session was never closed | async with ClientSession() as s: |
| bound to / attached to a different loop | A primitive or session used from another loop | Create it inside the loop that uses it |
| asyncio.run() cannot be called from a running event loop | Nested run(), such as in Jupyter | Use await main() |
| Everything freezes, debug mode logs slow callbacks | A blocking call in a coroutine | asyncio.sleep, an async client or to_thread |
Run with debug mode on, read every warning, search for time.sleep and requests, re-raise CancelledError, keep a reference to every task, and prefer TaskGroup over bare create_task.
Part 11 · asyncio Patterns Cheatsheet
Run, Concurrency and Deadlines
Every asyncio program starts at one place: asyncio.run(main()). It creates a fresh event loop, runs main() to completion, cancels leftover tasks, shuts down async generators and closes the loop. Inside main() you pick one of three ways to overlap work, and each reacts differently when a task fails.
| Tool | Returns | When one task fails |
|---|---|---|
tg.create_task(c) in a TaskGroup | Task objects; the block waits for all | Siblings are cancelled and you get an ExceptionGroup |
await gather(a, b) | List of results in input order | First error is raised; the others keep running |
create_task(c) | A single Task you await later | Error is stored in the Task until someone awaits it |
A bare create_task needs care: the loop keeps only a weak reference to the task, so one nobody holds can be garbage-collected mid-run. Keep a reference in a set and let a done-callback remove it. Deadlines work the same way for blocks and single awaitables: async with asyncio.timeout(5): bounds everything inside it, and await asyncio.wait_for(aw, 5) bounds one awaitable. Both cancel the slow work and raise TimeoutError.
import asyncio async def job(name, delay): await asyncio.sleep(delay) return f"{name} done" async def main(): background = set() t = asyncio.create_task(job("bg", 0.03)) background.add(t) t.add_done_callback(background.discard) async with asyncio.TaskGroup() as tg: a = tg.create_task(job("a", 0.02)) b = tg.create_task(job("b", 0.01)) print(a.result(), b.result()) print(await asyncio.gather(job("x", 0.02), job("y", 0.01))) try: async with asyncio.timeout(0.05): await asyncio.sleep(1) except TimeoutError: print("timed out") try: await asyncio.wait_for(asyncio.sleep(1), 0.05) except TimeoutError: print("wait_for timed out") print(await t) asyncio.run(main())
a done b done ['x done', 'y done'] timed out wait_for timed out bg done
Calling asyncio.run() inside a running loop (a notebook, for instance) raises RuntimeError. There, just await main(). And calling job() without await runs nothing at all.
Queues and Concurrency Limits
For producer/consumer work, put a bounded asyncio.Queue(maxsize=n) between the two sides. When the queue is full, await q.put(x) waits, which is backpressure: a fast producer is slowed down instead of filling memory. Each worker calls await q.get() and then q.task_done() exactly once per item. The producer's await q.join() returns when every item that was put has been marked done.
import asyncio async def worker(q, results): while True: item = await q.get() try: results.append(item * 2) finally: q.task_done() async def main(): q = asyncio.Queue(maxsize=2) results = [] workers = [asyncio.create_task(worker(q, results)) for _ in range(2)] for i in range(5): await q.put(i) await q.join() for w in workers: w.cancel() print(sorted(results)) asyncio.run(main())
[0, 2, 4, 6, 8]
Python 3.13 adds a third way to stop workers besides cancelling them or sending one None per worker. The producer calls q.shutdown(); get() hands out whatever is still queued and then raises QueueShutDown, while any further put() raises it straight away. The worker catches it and returns.
async def worker(q): try: while True: item = await q.get() try: await handle(item) finally: q.task_done() except asyncio.QueueShutDown: return
3.13+. The producer calls q.shutdown() when it has put everything.
To limit concurrency without a queue, create sem = asyncio.Semaphore(n) inside main() and wrap the guarded call in async with sem:. The permit is released even if the body raises. Remember that this caps calls in flight, not requests per second; for a rate, sleep while holding the permit or use a token bucket.
async def fetch(sem, url): async with sem: return await get(url)
At most n fetches run at once, however many tasks gather starts.
Forgetting task_done() when handle() raises makes join() hang forever, so it always belongs in a finally. Also, maxsize=0 means unbounded, and one None per worker is needed, not one in total.
Lock, Event, Condition and TCP Streams
Pick the smallest coordination tool that fits. A Lock gives mutual exclusion: async with lock: lets one task at a time into a critical section, and you only need it when an await sits between a read and the write that depends on it. An Event is a broadcast signal: every task parked in await ev.wait() wakes when someone calls ev.set(), and it stays set until ev.clear(). A Condition is a lock plus a wait for a predicate: async with cond: then await cond.wait_for(lambda: ...), while the other side changes state and calls cond.notify().
| Need | Primitive | Call |
|---|---|---|
| One task at a time (mutual exclusion) | Lock | async with lock: |
| Wake everyone on a signal (broadcast) | Event | ev.set() / await ev.wait() |
| Wait until a state is true (predicate) | Condition | await cond.wait_for(pred) / cond.notify() |
| Up to n at a time | Semaphore | async with sem: |
import asyncio async def main(): lock = asyncio.Lock() ready = asyncio.Event() balance = 100 async def spend(n): nonlocal balance await ready.wait() async with lock: seen = balance await asyncio.sleep(0) balance = seen - n tasks = [asyncio.create_task(spend(10)) for _ in range(5)] await asyncio.sleep(0) ready.set() await asyncio.gather(*tasks) print(balance) asyncio.run(main())
50A task that takes the same Lock twice deadlocks, because the Lock is not reentrant. Calling cond.wait() without holding the lock raises RuntimeError. Calling ev.set() from another thread is unsafe; use loop.call_soon_threadsafe(ev.set).
For TCP, open_connection(host, port) returns a (reader, writer) pair. writer.write(data) only buffers, so follow it with await writer.drain(), and finish with writer.close() followed by await writer.wait_closed(). On the server side, start_server(handler, host, port) runs handler(reader, writer) as a new task per client, and await server.serve_forever() keeps it running. Below, port 0 lets the OS pick a free port so the example can connect to itself.
import asyncio async def handle(reader, writer): data = await reader.readline() writer.write(data.upper()) await writer.drain() writer.close() await writer.wait_closed() async def main(): server = await asyncio.start_server(handle, "127.0.0.1", 0) port = server.sockets[0].getsockname()[1] async with server: reader, writer = await asyncio.open_connection("127.0.0.1", port) writer.write(b"ping\n") await writer.drain() print(await reader.readline()) writer.close() await writer.wait_closed() asyncio.run(main())
b'PING\n'Async Iterators, Context Managers and Generators
An async iterator has a plain def __aiter__ that returns self and an async def __anext__ that returns the next item, then raises StopAsyncIteration at the end. aiter(obj) and await anext(it) are the async counterparts of iter and next, and await anext(it, default) returns the default instead of raising. An async context manager defines __aenter__ and __aexit__, or you write one async def with yield inside try/finally and decorate it with @asynccontextmanager. When the number of managers is only known at runtime, enter each into an AsyncExitStack; they close in reverse order.
An async generator is an async def containing yield, which gives you async for support without writing the protocol. If you break out early, it stays paused and its finally waits for garbage collection or loop shutdown. Wrapping it in contextlib.aclosing() calls aclose() as soon as the block exits, so cleanup is deterministic.
import asyncio from contextlib import asynccontextmanager, aclosing, AsyncExitStack class Countdown: def __init__(self, n): self.n = n def __aiter__(self): return self async def __anext__(self): if self.n == 0: raise StopAsyncIteration await asyncio.sleep(0) self.n -= 1 return self.n + 1 @asynccontextmanager async def resource(name): print("open", name) try: yield name finally: print("close", name) async def letters(): try: for ch in "abc": await asyncio.sleep(0) yield ch finally: print("letters cleaned up") async def main(): it = aiter(Countdown(2)) print(await anext(it), await anext(it), await anext(it, "end")) async with AsyncExitStack() as stack: for name in ("db", "cache"): await stack.enter_async_context(resource(name)) print("working") async with aclosing(letters()) as gen: async for ch in gen: print(ch) break asyncio.run(main())
2 1 end open db open cache working close cache close db a letters cleaned up
| Protocol | You define | Ends with | Helper |
|---|---|---|---|
| Iterator | __aiter__ returning self, async __anext__ | StopAsyncIteration | aiter() / anext() |
| Context manager | __aenter__ / __aexit__ | __aexit__ always runs | @asynccontextmanager, AsyncExitStack |
| Generator | async def + yield | aclose() | aclosing(gen) |
Writing async def __aiter__ hands async for a coroutine and fails; __aiter__ is a plain method. A return True in __aexit__ silently swallows every exception.
Blocking Work and Debugging
Anything that blocks the thread freezes every task, because only await lets the loop switch. Hand blocking I/O to a worker thread with await asyncio.to_thread(fn, *args). For CPU-bound work, threads do not help because of the GIL, so use a process pool: await loop.run_in_executor(ProcessPoolExecutor(), fn, arg), with the loop from asyncio.get_running_loop().
To debug, run with asyncio.run(main(), debug=True), or set PYTHONASYNCIODEBUG=1, or start Python with -X dev. Debug mode logs any callback that takes longer than 100 ms, which exposes loop-blocking calls, and it reports coroutines that were never awaited. Running tests with python -W error::RuntimeWarning turns the never-awaited warning into a failure.
| Message | Cause | Fix |
|---|---|---|
| coroutine ... was never awaited | A coroutine was called but not awaited or scheduled | Add await or wrap it in a task |
| Task exception was never retrieved | A background task failed and nobody looked | Await the task or add a done-callback |
| Unclosed client session | A ClientSession was never closed | Use async with ClientSession(); create one session per application and reuse it, not one per request |
| attached to a different loop | A Lock, Event or session was created under another loop | Create it inside the running loop, for example inside main() |
| Slow callback taking more than 100 ms | Blocking call on the loop | Move it to to_thread or a process pool |
A bare except: or except BaseException: that does not re-raise swallows CancelledError, which breaks task.cancel() and every timeout. Clean up if you must, then raise again.
Turn on debug mode, read every warning, search for time.sleep and requests, re-raise CancelledError, and keep a reference to every task you create.
Part 12 · Check yourself
Quiz
Work out each answer yourself first, then read the reasoning printed under the question to compare. Every question asks you to predict what happens or to spot the bug, so tracing the code matters more than remembering names.
This worker pool processes five items, and item 3 makes process raise an error that the worker does not catch. The main task calls await q.join(). What happens, and what is the fix?
join()never returns. The exception escapes beforeq.task_done()runs, so the unfinished-items counter never reaches 0.- The worker task also dies on that exception, so items 4 and 5 are never taken from the queue either.
- Fix: call
q.task_done()in afinallyblock so everyget()is matched, and catch and log the error inside the loop so the worker stays alive. - After
join()returns, cancel the workers, because they are still parked onawait q.get().
async def worker(q): while True: item = await q.get() process(item) # raises on item 3 q.task_done()
You must stay under 80 requests per second. Your code uses sem = asyncio.Semaphore(10) and ends each request with await asyncio.sleep(1 / 80) before leaving the async with sem: block. Is the 80/s limit now guaranteed?
- No. The semaphore still has 10 permits, and each permit is busy for (latency + 1/80) seconds, so the rate is about 10 / (latency + 0.0125).
- With 100 ms replies that is about 89 requests/s, already over the limit. With very fast replies it climbs toward 800/s.
- The semaphore caps requests in flight. It does not cap requests per second, and a sleep inside a 10-permit block does not change that.
- Fix: space the request starts. Use a gate
Semaphore(1)held only forawait asyncio.sleep(1 / 80), then send the request after leaving the gate. Or use a token bucket refilled at 80 tokens/s. - Keep a separate
Semaphore(10)around the request itself if you also want to cap how many are in flight.
sem = asyncio.Semaphore(10) async def fetch(url): async with sem: r = await get(url) await asyncio.sleep(1 / 80) return r
What does this program print? Then say how the output would differ if asyncio.gather were replaced by a TaskGroup.
- It prints
caughtfirst. At 0.1 sbad()raises, and the defaultgatherre-raises that error tomainright away. - It then prints
ok done.gatherdoes not cancel its siblings, sook()keeps running and finishes at 0.2 s whilemainsleeps. - With a
TaskGroup,ok()would be cancelled whenbad()fails, sook donewould never print. The error would arrive as an ExceptionGroup, caught withexcept*.
async def ok(): await asyncio.sleep(0.2) print('ok done') async def bad(): await asyncio.sleep(0.1) raise ValueError async def main(): try: await asyncio.gather(ok(), bad()) except ValueError: print('caught') await asyncio.sleep(0.3) asyncio.run(main())
An async generator opens a socket and closes it in a finally block. The consumer does async for x in gen(): if x > 5: break. Why is the socket not closed as soon as the loop exits, and how do you make it deterministic?
breakleaves the generator paused at itsyield. Nothing callsaclose(), so thefinallywaits for garbage collection or for loop shutdown.- The socket stays open until then, which can exhaust descriptors in a long-running program.
- Fix: wrap the generator in
contextlib.aclosing(gen())withasync with. Leaving the block, early or not, callsaclose(), which raises GeneratorExit at theyieldsofinallyruns immediately.
async with aclosing(gen()) as g: async for x in g: if x > 5: break
Summary
- A coroutine runs only when awaited or scheduled, and tasks switch only at
await, so one blocking call freezes everything. - Bound every Queue, call
task_done()infinally, andjoin()before cancelling workers or sending one sentinel per worker. - A Semaphore caps concurrency, not requests per second. For a rate, space the starts with a gate or a token bucket.
- Use Lock when an await sits between a read and a write, Event for a broadcast signal, and Condition to wait on state.
- Pair every
write()withawait drain()and everyclose()withawait wait_closed(), and frame data withreadexactly. - Wrap async generators in
aclosing(), keep a reference to every task, and never swallowCancelledError. - Default to TaskGroup for fan-out, Queue for streams, and
to_threador a process pool for blocking and CPU work.