本节目标:掌握 Task 的创建与调度、gather 与 TaskGroup 的取舍、超时与取消的正确写法,并会用 Queue 与 Semaphore 控制并发。
适用版本:Python 3.12+(实测 3.14.6)
13.2 asyncio 任务、并发与超时取消
13.1 节我们用 gather 把三个任务跑成了约 1 秒,但很多问题还没回答:任务是什么时候被调度的?一个任务炸了,其他任务会被取消吗?怎么加超时?取消到底意味着什么?本节把这些「控制」层面的细节补齐——这是写出可靠异步代码的关键。
13.2.1 create_task 立刻调度,gather 要 await 才跑
两者都能让协程并发,但调度时机不同。asyncio.create_task() 一调用就把它注册进循环,从此刻起它在后台自己跑;而 gather 收到的协程,要等到 await gather(...) 那一刻才真正开始:
import asyncio, time
async def work(name, delay):
await asyncio.sleep(delay)
print(f" {name} 完成")
return name
async def with_create_task():
t0 = time.perf_counter()
t1 = asyncio.create_task(work("A", 0.5))
t2 = asyncio.create_task(work("B", 0.5))
print(" 任务已创建,此刻尚未 await")
await t1
await t2
return time.perf_counter() - t0
async def with_gather():
t0 = time.perf_counter()
coros = [work("A", 0.5), work("B", 0.5)]
print(" 协程已构造,尚未 await gather")
await asyncio.gather(*coros)
return time.perf_counter() - t0
=== create_task ===
任务已创建,此刻尚未 await
A 完成
B 完成
耗时 0.501s
=== gather ===
协程已构造,尚未 await gather
A 完成
B 完成
耗时 0.501s
两者总耗时一样(都并发,约 0.5 秒),区别在于**「先打印还是先干活」:create_task 之后任务已经在跑,gather 之前协程只是被构造、一步没动。因此——想让任务尽快开始并可在后面单独 await / cancel,用 create_task;只是想把一批协程一起跑完拿结果**,用 gather 更简洁。
13.2.2 gather 的错误处理:return_exceptions
默认情况下,gather 里任何一个协程抛异常,gather 会立刻把该异常向外抛。如果想让「失败」也作为一种结果返回,用 return_exceptions=True:
import asyncio
async def maybe_fail(x):
await asyncio.sleep(0.1)
if x == 2:
raise ValueError(f"boom at {x}")
return x * x
async def gather_re():
return await asyncio.gather(maybe_fail(1), maybe_fail(2), maybe_fail(3),
return_exceptions=True)
默认: 抛出: boom at 2
return_exceptions=True: [1, ValueError('boom at 2'), 9]
开了 return_exceptions=True 后,结果列表里成功的是返回值、失败的是异常对象,顺序与传入一致,你可以逐个 isinstance(r, Exception) 判断。
这里有一个容易踩的坑:默认的 gather 抛异常后,并不会取消其他兄弟任务,它们会继续在后台跑完。实测中,gather(boom(), other()) 捕获到 boom 之后,other 依然把 await asyncio.sleep(0.5) 睡满并打印完成——主流程已经知道出错了,后台却还有任务在偷偷占用资源。要真正做到「一个失败、全体收工」,得换 TaskGroup。
13.2.3 TaskGroup:结构化并发与 ExceptionGroup
Python 3.11 引入了 asyncio.TaskGroup,它用 async with 划定一个任务作用域:作用域退出前,所有任务必须有结果;只要有一个任务抛异常,其余任务会被自动取消,所有异常打包成一个 ExceptionGroup 抛出,用 except* 接住:
import asyncio
async def failer():
await asyncio.sleep(0.2)
raise ValueError("task failed")
async def slow():
try:
await asyncio.sleep(5)
except asyncio.CancelledError:
print(" slow 被取消了")
raise
async def main():
try:
async with asyncio.TaskGroup() as tg:
tg.create_task(failer())
tg.create_task(slow())
except* ValueError as eg:
print("捕获 ExceptionGroup:", eg.exceptions)
slow 被取消了
捕获 ExceptionGroup: (ValueError('task failed'),)
对比 13.2.2 的 gather:这里 slow 立刻被取消,而不是继续睡满 5 秒。这正是「结构化并发」的价值——任务的生命周期被一个作用域框住,不会出现没人管的孤儿任务。
| 维度 | gather | TaskGroup |
|---|---|---|
| 引入版本 | 3.4 起 | 3.11 起 |
| 调度时机 | await 时 | create_task 时立即 |
| 一个失败时 | 其他任务继续跑 | 其他任务被取消 |
| 异常形式 | 首个异常直接抛 | ExceptionGroup + except* |
| 结果收集 | 返回结果列表 | 需自行从任务或作用域取 |
新代码优先 TaskGroup:错误语义更安全,和 3.11+ 的 except* 配套。ExceptionGroup 与 except* 的完整语法可回顾 7.1 节。
13.2.4 超时:asyncio.timeout 与 wait_for
给操作加超时有两种写法。老写法是 wait_for,新写法(3.11+)是 asyncio.timeout 上下文管理器:
import asyncio
async def slow():
await asyncio.sleep(5)
return "done"
async def wf():
try:
return await asyncio.wait_for(slow(), timeout=0.3)
except TimeoutError:
return "wait_for 超时"
async def to():
try:
async with asyncio.timeout(0.3):
return await slow()
except TimeoutError:
return "timeout() 超时"
wait_for 超时
timeout() 超时
两者都抛 TimeoutError(3.11+ 里 asyncio.TimeoutError 就是内置 TimeoutError 的别名),但适用范围不同:wait_for 只能包裹单个可等待对象;而 asyncio.timeout(...) 是一个上下文,可以框住一整段包含多个 await 的代码——比如「三步请求总共不许超过 2 秒」,这是 wait_for 做不到的。超时触发时,被包裹的协程会收到取消,实现「真正中断」,而不只是「放弃等待」。
13.2.5 shield:保护关键任务不被取消
有些任务一旦开始就不能中途夭折(比如正在写文件、提交事务)。asyncio.shield() 给它套一层保护壳——外层的取消不会传递到内层任务:
import asyncio
async def critical():
print(" critical 开始(需要 0.5s)")
await asyncio.sleep(0.5)
return "critical 完成"
async def main():
task = asyncio.create_task(critical())
try:
await asyncio.wait_for(asyncio.shield(task), timeout=0.1)
except TimeoutError:
print(" 外层等待超时,但 task 未被取消")
print(" 最终结果:", await task)
critical 开始(需要 0.5s)
外层等待超时,但 task 未被取消
最终结果: critical 完成
外层 0.1 秒就超时了,但 critical 并没有被取消,稍后仍能拿到完整结果。注意 shield 返回的是 Future 而非协程,所以别把它的结果再塞进 create_task。
13.2.6 取消的语义:CancelledError 是 BaseException
调用 task.cancel() 并不会「咔嚓」一下停掉任务,而是在任务下一次到达 await 点时,往它内部抛一个 asyncio.CancelledError。关键点:从 3.8 起,CancelledError 继承自 BaseException 而不是 Exception:
import asyncio
print("是 Exception 子类吗:", issubclass(asyncio.CancelledError, Exception))
print("是 BaseException 子类吗:", issubclass(asyncio.CancelledError, BaseException))
是 Exception 子类吗: False
是 BaseException 子类吗: True
后果很直接:except Exception 抓不到取消。这是有意设计——避免你用宽泛的异常捕获把取消「吞掉」,让任务无法真正停止。如果确实要在清理时感知取消,得显式写 except asyncio.CancelledError,并且处理完要 raise 重新抛出,否则取消会被你吃掉,任务会在「假装取消」的状态下继续运行。原则:finally 里做清理,except CancelledError 里若处理了就 raise 放行。
13.2.7 asyncio.Queue:生产者—消费者
多个协程之间传递数据,用 asyncio.Queue。它是协程安全的队列(对比线程版的 queue.Queue),await queue.get() 在队列为空时会挂起等待,而不是忙等:
import asyncio
async def producer(q, n):
for i in range(n):
await q.put(i)
await q.put(None) # 哨兵,通知结束
async def consumer(q, name):
while True:
item = await q.get()
if item is None:
break
print(f" 消费者{name} 处理 {item}")
async def main():
q = asyncio.Queue()
await asyncio.gather(producer(q, 4), consumer(q, "A"))
asyncio.run(main())
消费者A 处理 0
消费者A 处理 1
消费者A 处理 2
消费者A 处理 3
asyncio.Queue(maxsize=N) 还能给队列设上限:满了以后 put 会挂起,天然实现背压(生产者跑太快时会被消费者拖住),避免内存被无限制堆积的消息撑爆。
13.2.8 Semaphore:限制并发数
gather 一上来就把所有任务同时丢出去,如果有一万个 URL,就会瞬间开出一万条连接。asyncio.Semaphore 用来限制同时进行的数量:
import asyncio, time
T0 = 0.0
async def fetch(sem, i):
async with sem: # 最多同时 2 个
print(f" [t={time.perf_counter()-T0:.2f}] 开始 {i}")
await asyncio.sleep(0.2)
print(f" [t={time.perf_counter()-T0:.2f}] 结束 {i}")
async def main():
global T0
T0 = time.perf_counter()
sem = asyncio.Semaphore(2)
await asyncio.gather(*[fetch(sem, i) for i in range(4)])
asyncio.run(main())
[t=0.00] 开始 0
[t=0.00] 开始 1
[t=0.20] 结束 0
[t=0.20] 结束 1
[t=0.20] 开始 2
[t=0.20] 开始 3
[t=0.40] 结束 2
[t=0.40] 结束 3
时间线清楚地显示:4 个任务被分成两批、每批 2 个,总耗时 0.4 秒(而非 0.2 秒)。async with sem: 在进入时获取令牌、退出时释放,令牌用完就排队。爬虫限速、数据库连接池都靠它。
13.2.9 asyncio.Runner 与 run_in_executor
asyncio.Runner(3.11+)是 asyncio.run() 的「可复用版」:它把循环的创建与关闭收在一个上下文管理器里,适合在同一进程里多次运行协程而不必反复开关循环:
import asyncio
async def main():
return "runner ok"
with asyncio.Runner() as runner:
print("Runner:", runner.run(main()))
Runner: runner ok
最后是 run_in_executor:把同步阻塞函数丢进线程池执行,从而不卡住循环。这是连接「异步世界」和「同步库」的桥梁(13.3 节会展开):
import asyncio, time
def blocking(n):
time.sleep(0.5)
return n * n
async def via_executor():
loop = asyncio.get_running_loop()
return await asyncio.gather(*[loop.run_in_executor(None, blocking, i)
for i in range(3)])
executor 结果: [0, 1, 4] 耗时 0.506s
三个各阻塞 0.5 秒的函数,在线程池里并发执行,总耗时约 0.5 秒。run_in_executor(None, fn, *args) 的第一个参数传 None 表示用默认线程池;传 ProcessPoolExecutor() 则改用进程池(适合 CPU 密集,见 12.2 节)。
小结
create_task立刻调度、gather到await才跑;只想跑完拿结果用gather,想尽快启动/单独取消用create_task。gather默认只抛首个异常且不取消兄弟任务;TaskGroup(3.11+)会取消其余任务并抛ExceptionGroup,用except*接住。- 超时用
asyncio.timeout(3.11+,可框住多个await)或wait_for(只包单个);asyncio.shield保护关键任务不被外层取消。 CancelledError继承自BaseException,except Exception抓不到它;若捕获了取消,务必raise放行。asyncio.Queue做生产者—消费者(maxsize提供背压),Semaphore限制并发数,Runner复用循环,run_in_executor把阻塞调用送进线程池。
到这里,你手里的 asyncio 工具箱已经齐全了。但知道 API 不等于写得出对的代码——下一节我们把镜头转向「生态」:同步库怎么用、哪些写法是「假异步」、什么时候根本不该用异步。
阅读导航:上一节:async/await 与事件循环 · 下一节:异步生态与常见陷阱 。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。