Python asyncio TaskGroup, Semaphore and Queue Patterns
Structured asyncio in Python: TaskGroup, a Semaphore to bound concurrency, Queue workers, async for and async with, and running blocking code off the loop.
- Course: Python study plan
- Module: Concurrency — threads, processes and asyncio
- Kind: Lesson
- Reading time: 14 min
- Runtime: CPython 3.11
What is asyncio.TaskGroup in Python?
asyncio.TaskGroup, new in Python 3.11, is an async context manager that owns the tasks created with tg.create_task inside it. The async with block waits for all of them; if one raises, the others are cancelled and the errors are raised together as an ExceptionGroup, caught with except*. That is structured concurrency: no task outlives its block and no exception is lost.
Lesson
The basics give you coroutines and gather; real programs need the shapes built on them: a group of tasks that fails together, a bounded number of concurrent requests, a producer feeding consumers through a queue, an async for over a stream, an async with around a connection, and a way to call the synchronous world without freezing the loop. This lesson covers TaskGroup (new in 3.11), Semaphore for rate limiting, asyncio.Queue workers, as_completed and its ordering caveat, async iterators, generators and context managers, to_thread and run_in_executor, and the guidance on when asyncio is the wrong tool.
TaskGroup
async def main():
async with asyncio.TaskGroup() as tg: # 3.11+
t1 = tg.create_task(fetch("a", 0.2))
t2 = tg.create_task(fetch("b", 0.1))
print(t1.result(), t2.result()) # both done when the block exits
A TaskGroup owns the tasks created inside it: the async with block waits for all of them, and if one raises, the others are cancelled and the exceptions are raised together as an ExceptionGroup (except* catches by type). This is structured concurrency — no task outlives its block, no exception is lost — and it is the replacement for gather when tasks should fail together. gather(..., return_exceptions=True) is the other choice: run everything, collect exceptions as values.
Bounding concurrency
sem = asyncio.Semaphore(5) # at most five in flight
async def fetch_limited(url):
async with sem:
return await fetch(url)
results = await asyncio.gather(*(fetch_limited(u) for u in urls))
A thousand coroutines are free to create; a thousand simultaneous connections are not free for the server. Semaphore(n) is the standard limiter, and gather still returns the results in order. asyncio.Lock, Event and Condition mirror their threading namesakes; they exist for coordination between coroutines and are needed less often, because no switch can occur between two plain statements.
Producer and consumers with a queue
async def producer(q, items):
for item in items:
await q.put(item)
await q.put(None)
async def consumer(q, out):
while (item := await q.get()) is not None:
out.append(item * 2)
q.task_done()
async def main():
q = asyncio.Queue(maxsize=10) # maxsize: back-pressure on the producer
out = []
await asyncio.gather(producer(q, range(5)), consumer(q, out))
print(out) # [0, 2, 4, 6, 8]
asyncio.Queue is the coroutine version of queue.Queue: await q.put, await q.get, task_done, await q.join(). With several consumers, each needs a sentinel (or the parent cancels the workers after q.join()), and each consumer's results should go into its own slot and be merged in order afterwards.
as_completed
for coro in asyncio.as_completed([fetch("a", 0.3), fetch("b", 0.1)]):
result = await coro # "b ready" first
For progress reporting; not for ordered answers. Sort before printing, or use gather.
Async iteration and context managers
async def ticker(n): # an async generator
for i in range(n):
await asyncio.sleep(0)
yield i
async def main():
async for i in ticker(3): # await between items
print(i)
squares = [i * i async for i in ticker(3)] # async comprehension
async with aiohttp.ClientSession() as session: # __aenter__/__aexit__ may await
...
async for drives an object with __aiter__/__anext__ (an async generator is the easy way to write one); async with drives __aenter__/__aexit__. Both exist because the protocol methods need to await — a database cursor fetching rows, a connection being opened and closed. contextlib.asynccontextmanager writes an async context manager from an async generator with one yield, exactly like its sync twin.
Bridging to blocking and CPU-bound code
data = await asyncio.to_thread(read_big_file, path) # blocking I/O → a worker thread
loop = asyncio.get_running_loop()
total = await loop.run_in_executor(process_pool, cpu_heavy, n) # CPU → a process pool
to_thread (3.9+) is the everyday bridge for a library with no async version; run_in_executor accepts a specific executor, including a process pool for CPU work. Both return awaitables, so the loop keeps serving other coroutines while the work runs elsewhere. The reverse direction — calling into a running loop from a thread — is asyncio.run_coroutine_threadsafe(coro, loop) or loop.call_soon_threadsafe.
Choosing asyncio
Reach for asyncio when the program is mostly waiting on many things at once and the libraries involved are async-native — web servers and clients, websockets, chat, scrapers, anything with thousands of connections. Prefer threads when the count is small and the libraries are blocking (a handful of requests calls, a database driver with no async version): ThreadPoolExecutor.map is three lines and needs no rewrite. Prefer processes for CPU. Mixing is normal: an async web handler that hands a CPU stage to a process pool with run_in_executor is the standard shape.
Pitfalls
- A
TaskGrouptask's exception silently cancelling siblings that were expected to finish — usegather(return_exceptions=True)when independence matters. - Unbounded fan-out against a rate-limited service.
- Consumers that never get a sentinel and wait forever.
- Printing from
as_completedand expecting input order. forinstead ofasync foron an async generator (TypeError: 'async_generator' object is not iterable).- A blocking library call inside a coroutine because "it's fast".
Key takeaways
TaskGroupgives structured concurrency: the block waits for its tasks, one failure cancels the rest, errors arrive as anExceptionGroup.Semaphore(n)bounds concurrency whilegatherkeeps result order;Lock/Event/Conditionmirror threading.asyncio.Queuewith sentinels connects producers and consumers;maxsizegives back-pressure.async for/async withdrive awaitable protocols; async generators andasynccontextmanagerwrite them.to_threadandrun_in_executorkeep the loop alive while blocking or CPU work runs elsewhere; asyncio is for many async-native waits, threads for a few blocking ones, processes for CPU.
Common questions
How do you limit concurrency in asyncio?
Create sem = asyncio.Semaphore(n) and wrap each request in async with sem:, so at most n run at once while gather still returns results in input order. A thousand coroutines are cheap to create, but a thousand simultaneous connections are not cheap for the server.
Should I use a TaskGroup or asyncio.gather?
Use a TaskGroup when the tasks should fail together, since one error cancels the rest. Use gather(..., return_exceptions=True) when the tasks are independent and you want every result, with exceptions collected as values.
How do you call blocking code from asyncio?
await asyncio.to_thread(fn, *args) runs a blocking function on a worker thread and awaits its result, keeping the event loop free. For CPU-bound work, await loop.run_in_executor(process_pool, fn, arg) sends it to a process pool instead.
Why do I get "'async_generator' object is not iterable"?
An async generator must be consumed with async for inside a coroutine, not with a plain for. The same applies to comprehensions over it, which must be async comprehensions such as [x async for x in ticker(3)].
When should I not use asyncio?
When the work is CPU-bound, use processes. When there are only a few waits and the libraries are blocking, such as a handful of requests calls, ThreadPoolExecutor.map is simpler and needs no rewrite. asyncio pays off with many concurrent waits and async-native libraries.
Exercises
Bounded fan-out
Read limit on the first line and job ids on the second. Create asyncio.Semaphore(limit) and a coroutine run(job_id) that, inside async with sem:, increments a shared in_flight counter, records the peak, awaits asyncio.sleep(0) twice, decrements the counter and returns job_id * job_id. Gather every job and print results <space-separated squares in input order> and peak <highest in_flight seen>.
Input: limit, then the ids. Output: two lines.
2
1 2 3 4 5
prints
results 1 4 9 16 25
peak 2Async producer and consumers
Read C on the first line and words on the second. With an asyncio.Queue(maxsize=2), a producer puts every word then C None sentinels; consumer i takes items until a sentinel and appends word.upper() to its own list slots[i]. Run the producer and the consumers under an asyncio.TaskGroup, then merge the slots, sort, and print consumers <C>, processed <n> and sorted <space-separated words>.
Input: C, then the words. Output: three lines.
2
gil loop task
prints
consumers 2
processed 3
sorted GIL LOOP TASKIn this module: Concurrency — threads, processes and asyncio
- The GIL and the three models — threads, processes, asyncio
- Threads — Thread, Lock, Event, Queue and the race you must see once
- concurrent.futures — executors, futures, map and as_completed
- multiprocessing — Process, Pool, pickling, queues and shared state
- asyncio basics — coroutines, await, tasks and gather
- asyncio patterns — TaskGroup, queues, semaphores, async iteration and bridging (this lesson)
- Checkpoint — Concurrency
← asyncio basics — coroutines, await, tasks and gather · Checkpoint — Concurrency →