单线程扛住上万并发连接,靠的不是魔法,而是一个事件循环 + 一堆会"让路"的协程。这一篇把
await底下的机制拆开给你看。
你将学到
- 事件循环、协程、
Task、await的本质,以及它们和生成器的血缘关系 asyncio.run/gather/wait/create_task各自该在什么时候用- 用一段代码实测:异步并发下载 vs 串行下载的差距
- 超时、取消与
CancelledError的正确处理方式 - 用
Queue/Semaphore做限流,避免打爆下游 - 同步阻塞代码拖垮事件循环的坑,以及
run_in_executor/to_thread的救法
前置知识
需要先理解生成器与 yield。建议先读 05 - 迭代器、生成器与协程,再读 07 - threading 与 multiprocessing 实战 做对比。
一、先搞懂 await 到底做了什么
async def 定义的函数调用时不会执行,而是返回一个协程对象(coroutine);直接调用它只会得到 <coroutine object hello at 0x...>,什么都不打印(还会触发 never awaited 警告),只有被驱动(await 或交给事件循环)它才会运行。
await 的语义是"把当前协程挂起,把控制权交还给事件循环"。 事件循环手里有一大堆就绪的协程/Task,它挑一个继续跑。所以本质是协作式调度:每个协程必须自己主动在 IO 处 await 让路,谁都不让路,事件循环就卡死了。它和生成器的血缘关系:await 底层就是 yield from 演化来的,async/await 只是给这套机制加了专用语法。
import asyncio
async def main():
print("start")
await asyncio.sleep(1) # 这里让出控制权,事件循环去干别的
print("after 1s")
asyncio.run(main())
# 输出: start
# (1 秒后)输出: after 1s
asyncio.run() 是唯一入口:它创建事件循环、跑完主协程、收尾并关闭循环。一个程序通常只调用一次。
二、事件循环与 Task
事件循环(event loop)是一个单线程的死循环:不断取出"就绪"的回调/协程执行,直到它们 await 一个未完成的 IO;然后处理 IO 事件(socket 可读可写、定时器到期),把就绪的任务放回队列。Task 是"被事件循环调度的协程",创建 Task 就等于把它排进循环立刻开跑,这是并发的关键:
import asyncio, time
async def work(name, sec):
await asyncio.sleep(sec)
return name
async def main():
# 顺序 await:串行,总耗时 = 1 + 2 = 3 秒
t0 = time.perf_counter()
await work("A", 1); await work("B", 2)
print(f"串行耗时 {time.perf_counter() - t0:.2f}s") # 输出: 3.00s
# create_task:并发,总耗时 = max(1, 2) = 2 秒
t0 = time.perf_counter()
ta = asyncio.create_task(work("A", 1))
tb = asyncio.create_task(work("B", 2))
await ta; await tb
print(f"并发耗时 {time.perf_counter() - t0:.2f}s") # 输出: 2.00s
asyncio.run(main())
区别就一句话:await coro 是"现在就跑,跑完再往下",create_task(coro) 是"丢进队列马上跑,我继续往下"。
三、gather / wait:并发编排
import asyncio
async def fetch(i, sec):
await asyncio.sleep(sec)
return f"结果{i}"
async def main():
# gather:并发运行,按输入顺序返回结果
print(await asyncio.gather(fetch(1, 0.3), fetch(2, 0.1), fetch(3, 0.2)))
# 输出: ['结果1', '结果2', '结果3'](顺序固定,不随完成时间变)
# return_exceptions=True:异常作为结果元素返回,不向外抛
print(await asyncio.gather(fetch(1, 0.1), asyncio.sleep(0, result=1/0),
return_exceptions=True))
# 输出: ['结果1', ZeroDivisionError('division by zero')]
asyncio.run(main())
# wait 更底层,返回 (done, pending),用来控流程(超时、FIRST_COMPLETED 等):
# tasks = [asyncio.create_task(fetch(i, 0.1 * i)) for i in range(1, 4)]
# done, pending = await asyncio.wait(tasks, timeout=0.25)
# for t in pending: t.cancel() # wait 不会自动取消,需自己收尾
选择指南:gather 用来看结果(顺序可控、异常策略明确),wait 用来控流程(超时、FIRST_COMPLETED 等返回条件)。
四、实测:并发下载 vs 串行下载
import asyncio, time
async def download(name, latency):
await asyncio.sleep(latency) # 模拟网络延迟
return f"{name}({latency}s)"
URLS = [("a", 0.5), ("b", 0.5), ("c", 0.5), ("d", 0.5)]
async def main():
t0 = time.perf_counter()
[await download(n, s) for n, s in URLS] # 串行
print(f"串行: {time.perf_counter() - t0:.2f}s") # 输出: 串行: 2.00s
t0 = time.perf_counter()
await asyncio.gather(*(download(n, s) for n, s in URLS)) # 并发
print(f"并发: {time.perf_counter() - t0:.2f}s") # 输出: 并发: 0.51s
asyncio.run(main())
真实场景换成 aiohttp,逻辑一模一样:
import asyncio, aiohttp
async def fetch_status(session, url):
async with session.get(url) as resp: # async with:异步上下文管理器
return resp.status
async def main():
urls = ["" for _ in range(20)]
async with aiohttp.ClientSession() as session:
statuses = await asyncio.gather(*(fetch_status(session, u) for u in urls))
print(set(statuses)) # 输出: {200} (完整运行需 asyncio.run(main()))
注意:
await后面必须跟"可等待对象"(协程、Task、Future)。requests.get()是同步函数,await requests.get(...)直接报错——这是新手最常见的撞墙点。
五、超时、取消与 CancelledError
取消是 asyncio 的核心机制,不是边角料。 当超时或你手动 task.cancel(),协程内部会在下一个 await 处被抛出 CancelledError。
import asyncio
async def slow():
try:
await asyncio.sleep(10)
except asyncio.CancelledError:
print("被取消了,做清理")
raise # ⚠️ 一定要重新抛出,否则取消语义被吞掉
async def main():
# 方式 1(3.11+):timeout 上下文管理器
try:
async with asyncio.timeout(0.5):
await slow()
except TimeoutError:
print("超时了")
# 方式 2:wait_for(兼容旧版本);方式 3:手动 task.cancel()
try:
await asyncio.wait_for(slow(), timeout=0.5)
except TimeoutError:
print("超时了")
task = asyncio.create_task(slow())
await asyncio.sleep(0.1)
task.cancel()
try:
await task # 取消后 await 会抛 CancelledError
except asyncio.CancelledError:
print("确认已被取消")
asyncio.run(main())
# 输出: 超时了 / 超时了 / 被取消了,做清理 / 确认已被取消
三条铁律: ① 捕获 CancelledError 后必须 raise,否则会破坏调用方的取消链;② 清理逻辑(关连接、释放锁)放进 finally;③ 别用 except Exception 顺手吞掉它——它继承自 BaseException,正是为了不被误捕。
六、限流:Semaphore 与 Queue
gather 一万个任务会一次性建一万个 Task、同时发起一万个请求,大概率把下游打挂。用 Semaphore(信号量) 控制并发上限:
import asyncio, aiohttp
async def fetch(session, url, sem):
async with sem: # 最多 N 个协程能同时进这里
async with session.get(url) as resp:
return resp.status
async def main():
sem = asyncio.Semaphore(10) # 同时最多 10 个请求
urls = [f"" for i in range(100)]
async with aiohttp.ClientSession() as session:
results = await asyncio.gather(*(fetch(session, u, sem) for u in urls),
return_exceptions=True)
print(len(results)) # 输出: 100
asyncio.run(main())
asyncio.Queue 则是生产者-消费者的异步版,适合"边生产边消费"的流水线:
import asyncio
async def producer(q):
for i in range(5):
await q.put(i) # 队列满则挂起等待
await q.put(None) # 哨兵
async def consumer(name, q):
while (item := await q.get()) is not None: # 队列空则挂起等待
print(f"{name} 处理 {item}")
async def main():
q = asyncio.Queue(maxsize=3)
await asyncio.gather(producer(q), consumer("w1", q), consumer("w2", q))
asyncio.run(main()) # 输出: w1/w2 轮流处理 0..4(顺序不定)
七、最大的坑:同步代码阻塞事件循环
事件循环是单线程的。你在协程里调用任何同步阻塞函数,整个循环、所有其他任务都会一起卡住——并发优势瞬间归零。
import asyncio, time
def blocking_io(): # 同步阻塞,比如 requests.get、读大文件
time.sleep(2)
async def bad():
await asyncio.sleep(0.1)
blocking_io() # ❌ 直接调用,事件循环被冻住 2 秒
async def good():
await asyncio.sleep(0.1)
await asyncio.to_thread(blocking_io) # ✅ 交给线程池,事件循环不被卡
async def main():
t0 = time.perf_counter()
await asyncio.gather(bad(), bad(), bad())
print(f"阻塞版: {time.perf_counter() - t0:.2f}s") # 输出: 6.00s!完全串行
t0 = time.perf_counter()
await asyncio.gather(good(), good(), good())
print(f"to_thread 版: {time.perf_counter() - t0:.2f}s") # 输出: 2.10s
asyncio.run(main())
# 精细控制线程池可用 loop.run_in_executor:
# await asyncio.get_running_loop().run_in_executor(ThreadPoolExecutor(4), blocking_io)
一句话原则:在协程里,凡是"会阻塞"的操作,要么换成 async 库(aiohttp / aiomysql / aiofiles),要么丢进 to_thread。绝不能直接同步调用。
常见坑
# ❌ 坑 1:定义了 async 函数却忘了 await
async def get():
return 1
async def main():
get() # 只创建了协程对象,没执行,还触发 RuntimeWarning: never awaited
# ✅ 正解:await 或 create_task
async def main_ok():
await get()
# ❌ 坑 2:在协程里用阻塞库 -> 阻塞整个事件循环
import requests
async def bad():
requests.get("")
# ✅ 正解:用异步库(aiohttp);退而求其次:await asyncio.to_thread(requests.get, url)
import aiohttp
async def good(session):
async with session.get("") as r:
return r.status
# ❌ 坑 3:吞掉 CancelledError,导致任务无法被取消(它继承 BaseException,别用兜底 except 吞)
try:
await something()
except Exception:
pass
# ✅ 正解:显式重新抛出,清理逻辑放 finally
# try: await something()
# except asyncio.CancelledError: raise
# finally: cleanup()
# ❌ 坑 4:一口气 gather 上万个任务 -> 打爆下游 / 内存
await asyncio.gather(*(fetch(u) for u in huge_list))
# ✅ 正解:用 Semaphore 限流(见第六节),或分批 gather
小结
async def调用只生成协程对象;await的本质是挂起当前协程、把控制权还给事件循环(协作式调度)。asyncio.run()是唯一入口;create_task让协程并发跑,await coro则是顺序执行。gather取结果(顺序可控、异常可配),wait控流程(超时 / 返回条件)。- 取消是核心机制:捕获
CancelledError后必须重新抛出,清理逻辑写进finally。 - 用
Semaphore限流、asyncio.Queue做流水线,别一次性铺开海量任务。 - 事件循环是单线程,任何同步阻塞调用都会冻结全部并发——换 async 库或
to_thread。
延伸阅读
- Python 官方文档:
asyncio任务的并发与取消章节 - 《Fluent Python》第 19-21 章关于协程与异步的深入讲解
文章回复
0 条公开回复