最根本的坑:事件循环里混入了阻塞调用
asyncio 的并发建立在「单线程 + 协作式调度」上。事件循环同一时刻只跑一个任务,只有当某个任务执行到 `await`(真正的异步 I/O)时才交出控制权。如果你把一个阻塞调用直接写进协程里,整个循环就被卡住了——所有其他协程一起等,天真的「并发」退化成串行。
典型阻塞源有三类:同步的数据库驱动(`psycopg2`、`pymysql`)、同步的 HTTP 客户端(`requests`)、以及任何 `time.sleep`、文件读写、CPU 密集计算。把它们直接放进 `async def` 里,事件循环会在那一行停住,表现是所有请求同时变慢但 CPU 很闲。
隔离方案有两种。轻量的用 `asyncio.to_thread(fn, *args)`(Python 3.9+),它自动使用默认线程池,代码最简洁。`to_thread` 是新增的便捷方法,默认执行器即 `loop.run_in_executor(None, ...)`。
需要控制线程池容量或使用自定义执行器时,用 `run_in_executor` 配合显式的 `ThreadPoolExecutor`。注意线程池的 `max_workers` 是**全局共享**的:如果并发 1000 个任务同时调用 `to_thread`,而默认执行器只有 32 个槽位,其余任务会排队,表现为「异步接口却迟迟不返回」。这时应当显式创建一个容量合适的执行器,或用信号量限制并发度。
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 与 TaskGroup:异常传播的差异
`asyncio.gather` 和 `asyncio.TaskGroup`(Python 3.11+)都能并发跑多个协程,但**失败语义完全不同**,这是最容易踩的坑之一。`gather` 默认在第一个异常时立刻返回,但**不会取消其他仍在运行的任务**——它们继续在后台跑,你只是丢掉了对它们的控制,容易产生「幽灵任务」。
要修掉这个行为,`gather` 必须传 `return_exceptions=True`,这样所有任务都会跑完,异常被收集成结果列表返回,然后由你自己决定如何处理。代价是失败场景下你仍然要等最慢的那个任务完成。
`TaskGroup` 的语义更适合结构化并发:任何一个子任务失败,任务组会**自动取消其余所有任务**,并在退出时抛出一个 `ExceptionGroup`(多个异常)或单个异常。这意味着不会出现泄漏的后台任务,也不会出现「一半完成一半被遗忘」的中间状态。
选择建议:新代码优先用 `TaskGroup`,它把「并发任务的生命周期不超过当前作用域」这条不变量强制落实了。只有在需要兼容旧版本 Python,或需要逐个检查结果而非整体失败时,才用 `gather`。如果用 `gather` 又不需要整体失败语义,记得显式传 `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,其余任务已被取消事件循环的创建与「忘记 await」
最基础的坑是根本没有启动事件循环。写了 `async def main()` 却直接调用 `main()`,得到的是一个 coroutine 对象,什么都不会发生,也没有报错——除非你开启了「协程从未被 await」的警告。正确做法是 `asyncio.run(main())`,它负责创建、运行并关闭循环。
第二个高频错误是**忘记 await**。`asyncio.gather(...)` 不加 await 会得到一个 Future 而非结果,很多 API 此时连请求都不发出;`async with session.get(url)` 忘记 await 则连 `__aenter__` 都不会被调用。这类问题的特征是「代码跑完了但什么都没发生」,且默认只发一条 RuntimeWarning。
开启警告是性价比最高的一道防线:在 Python 3.12 及以后,未被 await 的协程会产生显式警告;更早版本则可开启 `-W error::RuntimeWarning` 或使用 `PYTHONASYNCIODEBUG=1` 环境变量,后者还会输出慢回调(超过 100ms)的诊断信息。
开发模式下还有第三件值得做的事:检查有没有把 async 版本的库换成了同步版本,例如安装了 `requests` 而不是 `httpx`、`psycopg` 而不是 `psycopg2-async`。这类「装错包」的问题静态检查很难发现,但一上线就表现为并发完全无效。
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.py超时控制与并发限流
异步代码最危险的地方在于**失败会静默传播**。一个永不返回的 await 会让整个请求永远挂起。没有超时控制的异步服务,本质上是把「快速失败」退化成了「缓慢泄漏」。
单次调用的超时用 `asyncio.wait_for`。注意它在 Python 3.11 前的行为差异:超时时会取消内部任务并抛 `TimeoutError`,但如果被取消的代码里有 `except Exception` 或不响应取消的清理逻辑,实际返回可能明显晚于设定值。3.11+ 改为直接用 `asyncio.timeout()` 上下文管理器,语义更清晰。
整体超时用 `asyncio.timeout()`(3.11+)包住一组操作,超时会转成 `TimeoutError`,且能正确处理内部的取消传播。对于需要在超时后仍然拿到部分结果的场景,`asyncio.wait` 更合适——它只等待、不取消。
限流则用 `asyncio.Semaphore`。语义是「同时最多 N 个」,排队等待的任务仍会全部完成。这一点与很多限流中间件不同,值得在设计时想清楚:Semaphore 限制的是并发度而不是速率,如果还需要速率限制(例如每秒最多 10 次请求),要另外实现令牌桶。
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]连接池:容量与生命周期的正确姿势
连接池的参数最容易配错。httpx 的 `Limits(max_connections=N)` 控制**总连接数**,`max_keepalive_connections` 控制空闲保留数;aiohttp 的 `TCPConnector(limit=N)` 同理。默认值往往很小,几十个并发请求就会触发排队。
池容量的正确估算方法是:并发请求数 ÷ 平均响应时间 得到的是所需吞吐,再结合上游服务能承受的并发数取小值。上限最好与上游的应用级限制对齐——把 500 个并发请求全部砸给一个只允许 50 并发的上游,只会加剧超时与重试风暴。
生命周期上要注意**池与应用进程的对应关系**。连接池应该绑定到应用生命周期,而不是请求生命周期:在 Web 框架里通常在应用启动时创建、全局复用、在关闭时统一释放。如果每个请求新建一个池,握手开销会吃掉大量延迟,且 TIME_WAIT 连接数暴涨。
最后一个常见问题是**连接泄漏**。如果对每个请求都调用了 `client.get()` 但没有关闭响应体,连接会一直占用直到池耗尽。使用 `async with client.stream(...)` 或确保响应被完整读取,是防止泄漏最有效的编码习惯。连接池耗尽的典型表现是「前几千个请求正常,之后全部超时」。
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)