The fundamental pitfall: blocking calls inside the loop
asyncio concurrency rests on a single thread with cooperative scheduling. The loop runs one task at a time and only yields control when a task hits a real `await` on asynchronous I/O. Put a blocking call inside a coroutine and the whole loop stalls — every other coroutine waits, and your optimistic concurrency degrades into serial execution.
There are three classes of blocking sources: synchronous database drivers (`psycopg2`, `pymysql`), synchronous HTTP clients (`requests`), and anything calling `time.sleep`, touching files, or doing CPU-heavy work. Writing them directly into an `async def` parks the event loop at that line, which shows up as "all requests slow, CPU mostly idle".
There are two isolation options. The lightweight one is `asyncio.to_thread(fn, *args)` (Python 3.9+), which uses the default executor and keeps the code minimal. `to_thread` is a newer convenience method whose default executor is `loop.run_in_executor(None, ...)`.
When you need to size the pool or supply a custom executor, use `run_in_executor` with an explicit `ThreadPoolExecutor`. Remember that `max_workers` is shared globally: if 1000 concurrent tasks call `to_thread` against a default executor with 32 slots, the rest queue up and it looks like the async API hangs. Either create an executor sized for the workload, or bound concurrency with a semaphore.
import asyncio
async def fetch_one(client, url):
# 正确:异步 HTTP 客户端,await 时会让出控制权
resp = await client.get(url)
return await resp.text()
async def read_legacy_csv(path):
# 同步库:隔离到线程池,避免阻塞事件循环
def _read():
with open(path, encoding="utf-8") as f:
return f.read()
return await asyncio.to_thread(_read)
# 需要自定义容量时
from concurrent.futures import ThreadPoolExecutor
executor = ThreadPoolExecutor(max_workers=64)
data = await loop.run_in_executor(executor, blocking_fn, arg1, arg2)gather versus TaskGroup: how exceptions differ
Both `asyncio.gather` and `asyncio.TaskGroup` (Python 3.11+) run coroutines concurrently, but their failure semantics differ in a way that catches people. `gather` returns on the first exception by default yet does **not** cancel the still-running tasks — they continue in the background while you lose control over them, producing "ghost tasks".
To fix that, pass `return_exceptions=True`: every task then runs to completion and exceptions are collected into the returned results list for you to handle. The cost is that in the failure path you still wait for the slowest task.
`TaskGroup` embodies structured concurrency: when any child task fails, the group cancels all remaining tasks automatically and raises an `ExceptionGroup` (or a single exception) on exit. No leaked background tasks, no half-completed states, and no silently forgotten work.
Recommendation: prefer `TaskGroup` in new code — it enforces the invariant that concurrent tasks cannot outlive their enclosing scope. Fall back to `gather` only for older Python compatibility or when you want to inspect results individually instead of failing as a group. If you use `gather` without group-failure semantics, remember `return_exceptions=True`.
# 旧写法:一个失败,其他仍在后台跑(幽灵任务)
results = await asyncio.gather(*tasks)
# 改进:全部跑完,异常作为结果返回
results = await asyncio.gather(*tasks, return_exceptions=True)
for r in results:
if isinstance(r, Exception):
log.warning("task failed: %s", r)
# 推荐写法:TaskGroup 自动取消其余任务
async with asyncio.TaskGroup() as tg:
for url in urls:
tg.create_task(fetch_one(client, url))
# 失败时抛出 ExceptionGroup,其余任务已被取消Creating the loop and the forgotten await
The most basic mistake is never starting the loop. Defining `async def main()` and calling `main()` directly produces a coroutine object that does nothing and raises nothing — unless you enable the "coroutine was never awaited" warning. The correct entry point is `asyncio.run(main())`, which creates, runs and closes the loop.
The second common error is forgetting `await`. Calling `asyncio.gather(...)` without await yields a Future rather than results, and for many APIs no request is even issued; `async with session.get(url)` without await never calls `__aenter__`. The signature is "the code ran and nothing happened", with only a RuntimeWarning as evidence.
Turning on warnings is the best-value defence available. Python 3.12+ emits an explicit warning for un-awaited coroutines; on older versions enable `-W error::RuntimeWarning` or set the `PYTHONASYNCIODEBUG=1` environment variable, which also reports slow callbacks exceeding 100ms.
A third worthwhile development-time check is whether a sync build of a library got installed instead of the async one — `requests` rather than `httpx`, `psycopg2` rather than `psycopg`. Static analysis rarely catches this, but in production it shows up immediately as concurrency that does nothing at all.
import asyncio
async def main():
# 记得每个协程都要 await
data = await asyncio.gather(fetch_a(), fetch_b())
return data
if __name__ == "__main__":
asyncio.run(main())
# 开发时开启调试:慢回调与未 await 协程告警
# PYTHONASYNCIODEBUG=1 python app.pyTimeout control and concurrency limiting
The most dangerous property of asynchronous code is that failures propagate silently. An `await` that never returns leaves the request hanging forever. An async service without timeout control has effectively converted fast failures into slow leaks.
Use `asyncio.wait_for` for per-call timeouts. Note the pre-3.11 behaviour difference: on timeout it cancels the inner task and raises `TimeoutError`, but if the cancelled code has a broad `except Exception` or a cleanup that ignores cancellation, the real return can come noticeably later than the deadline. Python 3.11+ offers the `asyncio.timeout()` context manager with clearer semantics.
For a deadline covering a group of operations, wrap them in `asyncio.timeout()` (3.11+); it converts expiry into `TimeoutError` and handles cancellation propagation correctly. If you need partial results after a deadline, `asyncio.wait` is the better tool — it waits without cancelling.
Rate and concurrency limiting use `asyncio.Semaphore`. Its meaning is "at most N at a time"; queued waiters still all complete eventually. That differs from many rate-limiting middlewares and is worth thinking through at design time: a semaphore bounds concurrency, not rate. If you also need a velocity limit such as 10 requests per second, implement a token bucket separately.
import asyncio
async def call_with_limits(client, url, sem):
# 并发度限制:同时最多 20 个在飞
async with sem:
# 单次超时(Python 3.11+ 推荐写法)
async with asyncio.timeout(5):
resp = await client.get(url)
return resp.status
async def main(urls):
sem = asyncio.Semaphore(20)
async with asyncio.TaskGroup() as tg:
tasks = [tg.create_task(call_with_limits(c, u, sem)) for u in urls]
return [t.result() for t in tasks]Connection pools: sizing and lifecycle
Pool parameters are easy to get wrong. httpx `Limits(max_connections=N)` bounds total connections and `max_keepalive_connections` bounds idle retention; aiohttp `TCPConnector(limit=N)` is the same idea. Defaults are often small, so a few dozen concurrent requests already queue.
Size the pool by dividing concurrency by average response time to get required throughput, then take the smaller of that and what the upstream can actually accept. Aligning the limit with the upstream application-level cap matters: firing 500 concurrent requests at a service that permits 50 only amplifies timeouts and retry storms.
Regarding lifecycle, the pool belongs to the application process, not to the request. Create it at startup, share it globally, and close it during shutdown. Building a pool per request makes handshake latency dominate and drives the TIME_WAIT count through the roof.
The last recurring problem is connection leakage: calling `client.get()` without closing the response body keeps the connection checked out until the pool is exhausted. Prefer `async with client.stream(...)` or fully read every response. A pool exhaustion symptom looks like "the first few thousand requests are fine, then everything times out".
import httpx
limits = httpx.Limits(
max_connections=100, # 总连接数上限,按上游承受能力设置
max_keepalive_connections=20, # 空闲连接保留数
keepalive_expiry=30.0,
)
# 绑定到应用生命周期,全局复用,关闭时释放
client = httpx.AsyncClient(limits=limits, timeout=httpx.Timeout(10.0))
# 读取响应体必须完整,否则连接不会归还池
async def ok(url):
async with client.stream("GET", url) as resp:
chunks = [c async for c in resp.aiter_bytes()]
return b"".join(chunks)