Skip to main content

Chapter 46: Async Runtime Architecture

异步运行时要回答的问题是:一个 async def 调用产生的 coroutine object,怎样被 event loop 包装成 task,怎样在 await 处暂停,怎样在 future 完成、callback 入队、queue 有空间、lock 释放、取消请求到达时恢复执行。本章用一个生产者、消费者、超时取消和关闭信号组成的短程序贯穿说明。读完之后,读者应能追踪一个 asyncio 程序的调度路径,判断卡顿、泄漏、取消失效和顺序异常来自哪个运行时对象。

asyncio 是 Python 标准库中基于 async / await 语法编写并发代码的库,官方文档把它定位为适合 I/O 密集和高层网络代码的基础设施,并提供 coroutine、task、queue、同步原语、event loop、future 等 API。这里的主语是 CPython 标准库中的 asyncio,不是所有 Python 实现或所有异步框架的通用模型。第三方 loop、框架封装和不同 Python 小版本会改变局部实现细节,但本章建立的对象关系可以迁移到大多数 asyncio 问题排查中。

本章的核心结论是:asyncio 的并发来自协作式调度。event loop 持有就绪 callback 队列、定时器、I/O 事件和 task / future 状态;task 驱动 coroutine 向前执行;await 把当前执行权交回 loop;future、queue、lock 或 timer 在条件满足时安排 task 继续运行。这个模型能解释为什么 CPU 循环会卡住整个 loop,为什么取消请求要等到下一个可注入点,为什么忘记保存 task 引用会造成后台任务不可控,为什么 queue 的 maxsize 是背压机制的一部分。

下面的贯穿材料模拟一个爬取管线。生产者把 URL 放入有限队列,消费者从队列取任务并受 semaphore 限流,主任务在超时时取消所有 worker。示例只保留运行时关系,省略真实网络 I/O。

import asyncio

async def fetch(url: str) -> str:
await asyncio.sleep(0.05)
return f"body:{url}"

async def producer(queue: asyncio.Queue[str], urls: list[str], worker_count: int) -> None:
for url in urls:
await queue.put(url)
for _ in range(worker_count):
await queue.put("STOP")

async def worker(name: str, queue: asyncio.Queue[str], gate: asyncio.Semaphore) -> None:
while True:
url = await queue.get()
try:
if url == "STOP":
return
async with gate:
body = await fetch(url)
print(name, len(body))
finally:
queue.task_done()

async def main() -> None:
queue: asyncio.Queue[str] = asyncio.Queue(maxsize=2)
gate = asyncio.Semaphore(3)
urls = [f"https://example.com/{i}" for i in range(8)]
worker_count = 2

async with asyncio.TaskGroup() as tg:
tg.create_task(producer(queue, urls, worker_count))
for i in range(worker_count):
tg.create_task(worker(f"worker-{i}", queue, gate))
await asyncio.wait_for(queue.join(), timeout=1.0)

asyncio.run(main())

这个程序中,每个 async def 调用都先创建 coroutine object。TaskGroup.create_task() 把 coroutine 包装成 task,并把 task 交给当前运行中的 event loop。queue.put()queue.get()asyncio.sleep()Semaphore.acquire()queue.join()wait_for() 都可能让当前 task 暂停。暂停以后,线程并没有同时执行另一个 Python 字节码流;event loop 从自己的就绪队列中取出下一个可以推进的 callback 或 task step。这个事实贯穿后面五节。

46.1 Event loop, task, and future

event loopasyncio 程序的调度核心。它运行异步 task 和 callback,处理网络 I/O、子进程、定时器和底层事件。应用代码通常通过 asyncio.run() 启动顶层 coroutine;库和框架作者才更频繁地直接使用 loop 对象。官方文档对 event loop 的定位可以参考 Python 3.14 event loop 文档

在贯穿程序中,asyncio.run(main()) 做了三件事:创建或取得一个 event loop,创建顶层 task 来运行 main(),运行 loop 直到这个顶层 task 完成,然后执行收尾逻辑。main() 内部调用 TaskGroup.create_task() 时,新的 producer task 和 worker task 被登记到同一个 loop。它们共享一个线程内的调度器,也共享 queuegate 这些 Python 对象。

task 是 coroutine 的调度包装器。coroutine object 只是可暂停执行体;task 才记录“这个 coroutine 已经交给 loop 管理”。官方文档说明,asyncio.create_task() 会把 coroutine 包装为 Task 并安排其执行;TaskGroup 是 Python 3.11 加入的结构化并发接口,负责把一组相关 task 的创建、等待、异常传播和取消收束在一个 async with 作用域内。coroutine、task、future 的官方术语边界可以参考 Coroutines and Tasks 文档

future 是“未来某个时刻会得到结果”的低层 awaitable。它保存完成状态、结果、异常和完成回调。task 本身也属于 future-like 对象,因为等待 task 等价于等待它包装的 coroutine 最终返回或抛出异常。普通应用代码很少直接创建 future,但 loop.run_in_executor()、协议层、transport、底层 I/O 和第三方库经常用 future 把 callback 世界接回 await 世界。

可以把三者关系压缩为一条路径:coroutine 表达可暂停的执行体,task 把 coroutine 交给 loop 调度,future 表达 task 正在等待的外部条件或内部结果。贯穿程序里的 await queue.get() 会让 worker task 暂停,并把它挂到 queue 的等待者集合;当 producer 放入 item 时,queue 会安排等待的 worker 恢复。这里恢复 worker 的动作并非 worker 自己主动轮询,而是运行时对象在条件变化后把 callback 放回 event loop 的就绪路径。

下面的图只描述本章使用的核心路径,不展开 selector、transport、系统调用和平台差异。

图中的关键状态迁移是 Task step → await target → register waiter → ready callbacks → Task step。task 每次只推进到下一个暂停点、返回点或异常点。future pending 时,task 记录自己正在等待的对象;future done 时,loop 安排 task 的下一步。await 表达式的结果来自 future 的 result 或 exception,所以 await fetch(url) 继续执行时拿到的是 fetch() coroutine 的返回值,await queue.get() 继续执行时拿到的是 queue 中的 item。

这个对象模型还解释了 task 引用问题。官方文档提醒,event loop 对 task 只保持弱引用;对“fire-and-forget”任务,调用方应保存强引用并在完成后移除。工程上,这意味着后台任务需要有所有权容器、生命周期边界和错误收集策略。把 create_task() 的返回值丢弃会让异常日志、取消时机和任务归属变得不稳定。

46.2 Cooperative scheduling

协作式调度的判断点是“当前 task 何时交回控制权”。在 asyncio 中,运行中的 coroutine 通过 await 暂停;没有 await 的长 CPU 循环会连续占用当前线程,event loop 得不到机会处理其他 task、timer、I/O ready callback 或取消请求。这也是 asyncio 适合 I/O 密集任务的直接原因:I/O 等待可以转换为 future pending,event loop 在等待期间推进其他 task。

贯穿程序里,producer()await queue.put(url) 处可能暂停。Queue(maxsize=2) 让队列最多保存两个尚未处理的 item;当队列已满,put() 会等待消费者取走 item。Python 3.14 文档说明,asyncio.Queuemaxsize 大于 0 时,await put() 会在队列达到容量上限时阻塞,直到有 item 被移除;put_nowait() 则在没有空位时抛出 QueueFull。这就是背压:下游处理速度通过 queue 状态反向限制上游生产速度。

worker 的执行路径也体现协作式调度。await queue.get() 等待输入,async with gate 内部等待 semaphore,await fetch(url) 等待 I/O 模拟。每个等待点都把 worker task 从“正在运行”转为“等待某个条件”。条件满足后,worker 回到 ready 队列,等待 loop 选择下一次 task step。这里的公平性只针对具体原语的约定,例如 asyncio.Lock.acquire() 在多个 coroutine 等待同一把锁时按开始等待的顺序唤醒;queue、semaphore、callback 和 timer 的组合还会受到就绪队列、回调批次和任务自身运行时间影响。

下面的短片段展示一个常见卡顿来源。bad_worker() 中没有等待点,其他 task 要等到循环结束才有机会运行;good_worker() 周期性调用 await asyncio.sleep(0),把执行权交还给 loop。

import asyncio

async def bad_worker() -> None:
total = 0
for i in range(20_000_000):
total += i
print(total)

async def good_worker() -> None:
total = 0
for i in range(20_000_000):
total += i
if i % 100_000 == 0:
await asyncio.sleep(0)
print(total)

这段代码的结论很直接:async def 本身并不带来抢占式调度,await 才是让出点。asyncio.sleep(0) 可以用于显式让出执行权,但它也会增加调度开销。真正的 CPU 密集任务通常应切到线程池、进程池、原生扩展或批处理设计中;asyncio 的优势是把大量 I/O 等待统一纳入一个事件循环,从而减少大量 OS 线程各自阻塞的结构成本。

协作式调度还改变了阅读代码的顺序。同步代码通常沿调用栈连续执行;异步代码在每个 await 处可能暂停,暂停期间共享对象可能被其他 task 修改。贯穿程序里的 queuegate 和 stdout 都是共享状态。读异步 bug 时,需要把每个 await 当作状态可能变化的边界:await 之前检查到的条件,在 await 之后需要重新确认,尤其是连接状态、缓存命中、任务取消状态和资源是否仍然打开。

46.3 Cancellation and async state

取消请求在 asyncio 中是状态注入。调用 task.cancel() 会请求 task 在下一次合适机会抛出 asyncio.CancelledError;已经运行的 Python 代码不会在任意字节码位置被强制中断。官方文档说明,task 被取消时,CancelledError 会在 task 的下一个机会中抛出,coroutine 应使用 try / finally 执行清理,并在显式捕获后通常继续传播该异常。

贯穿程序中,await asyncio.wait_for(queue.join(), timeout=1.0) 给整体管线设置了时间边界。超时发生时,wait_for() 会取消正在等待的 awaitable;在 TaskGroup 作用域内,如果某个 task 以普通异常失败,相关 task 会被取消并等待收束。结构化并发的价值在于所有子任务归属于同一个作用域,异常、取消和等待路径能在 async with 退出时闭合。

取消沿 await 链传播。假设 worker 正在执行 await fetch(url),外部取消 worker task 时,CancelledError 会注入到 fetch() 当前暂停处。如果 fetch() 内部正在 await asyncio.sleep(0.05),sleep 对应的 timer future 会收到取消或在恢复时让异常继续向上。worker 的 finally 仍会执行,所以 queue.task_done() 能减少 unfinished task 计数。这个 finallyqueue.join() 很关键:取出 item 后若没有调用 task_done(),join 会一直等待 unfinished task 归零。

下面的写法展示取消清理的稳定边界。代码把资源释放放进 finally,并在捕获取消后重新抛出,让外层 TaskGrouptimeout() 或调用方仍能观察到取消状态。

import asyncio

async def use_resource() -> None:
resource = "opened"
try:
await asyncio.sleep(10)
except asyncio.CancelledError:
print("record cancellation")
raise
finally:
print("close", resource)

这里要区分三种状态。第一,task 处于 pending,表示 coroutine 还没有结束。第二,取消请求已经登记,但异常尚未注入到 coroutine 内部。第三,coroutine 已经处理取消并使 task 进入 done 状态,done 状态可能带着 CancelledError、普通异常或正常结果。排查“取消没有生效”时,需要分别确认 task 当前等待的 awaitable、代码是否长期运行在无 await 的 CPU 段、是否捕获了 CancelledError 后吞掉了状态、是否存在 shield 或 timeout 组合改变了传播范围。

CancelledError 直接继承自 BaseException,这让宽泛的 except Exception 捕获不到它。这个设计使取消更像控制流信号,降低普通业务异常处理误吞取消的概率。工程代码中如果确实要把取消转为正常结果,应把这个转换写在明确边界上,并同时处理 task 的取消状态;否则上层结构化并发组件可能收到与真实状态不一致的信号。

46.4 Async queue, lock, and callback scheduling

异步同步原语解决的是同一 event loop 内 task 之间的顺序、容量和互斥问题。asyncio 的 queue、lock、event、condition、semaphore 设计为 task 级协作原语;官方同步原语文档明确说明它们不用于 OS 线程同步,超时通常通过 asyncio.wait_for() 包装。线程之间共享状态仍然需要 threading、跨线程调度入口或更明确的消息边界。

asyncio.Queue 的运行时角色是容量控制加等待者唤醒。贯穿程序中的 Queue(maxsize=2) 保存未处理 URL,put() 在队列满时让 producer task 等待,get() 在队列空时让 worker task 等待,task_done()join() 共同维护“已放入但尚未完成处理”的计数。Python 3.13 起 Queue.shutdown() 增加了显式关闭模式;它会使后续 put() 抛出 QueueShutDown,并解除已经阻塞的 putter。这个版本边界对长期服务的优雅停机有直接影响。

asyncio.Lockasyncio.Semaphore 解决共享资源进入顺序。lock 保证同一时刻只有一个 task 进入临界区;semaphore 通过内部计数限制并发度。贯穿程序用 Semaphore(3) 限制最多三个 fetch 同时进入,这个限制比 queue 的 maxsize 更靠近外部资源边界:queue 控制积压,semaphore 控制并发访问外部服务、数据库连接或文件句柄的数量。

callback scheduling 是 event loop 的底层入口之一。loop.call_soon(callback, *args) 把 callback 安排到下一轮 loop 迭代,并按登记顺序调用一次;loop.call_later()loop.call_at() 基于 monotonic clock 安排定时 callback;loop.call_soon_threadsafe() 用于其他线程向 loop 投递 callback。官方 event loop 文档同时指出,普通 call_soon() 不是线程安全的,跨线程投递应使用 call_soon_threadsafe()

这些对象会共同影响公平性与背压。queue 的容量决定上游等待点;lock 的等待队列决定互斥入口;semaphore 的计数决定并发窗口;callback 的入队顺序决定某些恢复动作的先后;每个 task 自身在恢复后运行多久决定其他 task 的等待时间。一个队列设置了 maxsize,如果消费者在拿到 item 后做长 CPU 计算,生产者仍会感到背压,但整个 loop 也会卡顿。背压解决的是容量增长问题,CPU 卡顿需要缩短无 await 段、切分工作或移出 event loop。

排查顺序问题时,可以用同一组维度比较这些原语:谁拥有等待者集合,谁负责唤醒等待者,唤醒后是立即执行还是进入 ready 队列,取消等待时是否清理等待者,资源关闭时 pending task 得到什么异常。这个比较能把“顺序怪异”落到具体状态,并跳出“异步调度不可预测”这种空泛判断。

46.5 Async runtime architecture checklist

异步问题的排查入口应回到 event loop 的调度事实:同一线程内同一时刻只有当前 task 或 callback 在运行;其他 task 要等当前代码到达让出点。遇到“服务卡住”,先确认 event loop 是否被 CPU 段、同步阻塞 I/O、长 callback 或过大的临界区占用。证据通常来自日志时间戳、debug mode 的 slow callback 报告、任务栈、请求延迟分布和 CPU 使用率。

遇到“任务泄漏”,先确认 task 的所有权。每个 create_task() 都应能回答三个问题:谁保存引用,谁等待结果或异常,谁在关闭时取消它。TaskGroup 给相关任务提供作用域边界;长期后台任务也需要集合保存、完成回调清理、异常收集和关闭时取消。悬空 task 常见于丢弃返回值、在 callback 中创建任务后缺少登记、关闭流程只关闭连接而没有取消等待者。

遇到“取消失效”,沿 await 链找注入点。检查当前 task 是否停在可取消的 awaitable 上,是否运行在长 CPU 段,是否使用 try / finally 完成清理,是否捕获 CancelledError 后继续传播,是否被 shield()wait_for()timeout()TaskGroup 改变了传播范围。取消相关问题的关键证据是 task 状态、当前等待对象、异常处理代码和作用域退出路径。

遇到“顺序不符合预期”,把 queue、lock、semaphore、callback 和 timer 分开看。queue 决定生产与消费的容量边界;lock 决定互斥入口;semaphore 决定并发窗口;callback 决定下一轮调度入口;timer 决定最早恢复时间;task 自身运行时间决定其他 ready 项的延迟。顺序异常通常来自多个边界叠加,单看某一行 await 很难定位。

遇到“吞吐低”,先判断瓶颈属于 I/O 等待、外部服务限速、连接池容量、queue 容量、semaphore 并发窗口、CPU 计算还是 callback 过重。asyncio 提高的是等待期间的调度利用率;它不会自动缩短单个 I/O 的远端响应时间,也不会让纯 Python CPU 循环并行执行。吞吐优化需要把等待、容量和执行成本分别量化。

下面给出一个可复用检查顺序:

  • 定位当前 loop:确认代码运行在哪个 event loop 和线程中,跨线程投递使用 call_soon_threadsafe()
  • 列出 task 所有权:为每个 task 找到创建点、保存引用的位置、等待结果的位置和关闭时取消的位置。
  • 标出 await 边界:把每个 await 当作共享状态可能变化点,检查恢复后是否重新验证状态。
  • 识别等待对象:区分 task 正在等待 queue、lock、semaphore、timer、I/O future 还是另一个 task。
  • 检查取消传播:确认 CancelledError 的注入点、清理路径和外层作用域是否能观察到取消。
  • 检查背压链路:确认 Queue.maxsize、连接池、semaphore、外部限流和消费者处理速度是否匹配。
  • 检查无让出代码段:定位同步阻塞 I/O、长 CPU 循环、长 callback 和过大的临界区。
  • 确认版本边界:涉及 TaskGroupQueue.shutdown()eager_start、loop policy 等行为时标注 Python 小版本。

本章使用的官方资料只支撑 API 边界和版本事实:asyncio 总览见 Python 3.14 asyncio 文档,coroutine / task / future / cancellation 见 Coroutines and Tasks,event loop 和 callback 调度见 Event loop,queue 行为见 Queues,同步原语见 Synchronization Primitives。正文中的架构判断来自这些 API 事实与 CPython asyncio 运行时对象关系的整理。

最小自检任务

阅读下面代码,判断它可能出现哪两个运行时问题,并说明应从 event loop、task、queue、取消和背压中的哪几个对象入手排查。

import asyncio

async def cpu_step(n: int) -> int:
total = 0
for i in range(n):
total += i * i
return total

async def producer(queue: asyncio.Queue[int]) -> None:
for i in range(100_000):
await queue.put(i)

async def consumer(queue: asyncio.Queue[int]) -> None:
while True:
item = await queue.get()
try:
await cpu_step(200_000)
print(item)
finally:
queue.task_done()

async def main() -> None:
queue: asyncio.Queue[int] = asyncio.Queue()
asyncio.create_task(consumer(queue))
await producer(queue)
await queue.join()

asyncio.run(main())

答案要点

第一个问题是 event loop 卡顿。cpu_step() 虽然写成 async def,内部循环没有 await,consumer task 一旦进入这个函数,就会长时间连续执行 Python 代码。其他 task、timer、callback 和取消请求都要等它返回后才有机会推进。排查时先看无让出代码段、CPU 使用率和 slow callback 证据,再决定把 CPU 工作切小、移到 executor 或改成真正的异步 I/O 等待。

第二个问题是背压失控和任务所有权不清。Queue() 没有设置 maxsize,producer 可以持续放入大量 item,内存压力会随未处理 item 增长。asyncio.create_task(consumer(queue)) 的返回值被丢弃,consumer 的异常、关闭取消和生命周期归属都缺少显式所有者。稳定写法应为 queue 设置容量,把 consumer task 放入 TaskGroup 或保存引用,在关闭时取消并等待收束,同时保留 task_done()join() 的配对关系。

本章知识点总结

  • 运行时核心asyncio 的核心对象关系是 event loop 调度 task,task 驱动 coroutine,future 表示等待结果。
  • 启动边界asyncio.run() 负责启动顶层 coroutine、运行 event loop,并在顶层 task 完成后执行收尾。
  • Task 语义:coroutine object 只有被包装成 task 或被 await 后才进入可执行调度路径。
  • Future 语义:future 保存 pending、result、exception 和 done callback,是 callback 世界接入 await 的低层对象。
  • 协作调度await 是当前 task 交回 event loop 的主要让出点,无让出 CPU 段会阻塞整个 loop。
  • 共享边界:每个 await 之后,共享对象状态都可能已经被其他 task 修改,需要重新验证关键条件。
  • 取消传播:取消请求在下一次合适机会向 task 注入 CancelledError,清理逻辑应放入 finally 并保持取消状态可被上层观察。
  • 结构化并发TaskGroup 把相关 task 的创建、等待、异常传播和取消收束在同一个 async with 作用域中。
  • Queue 背压Queue.maxsize 把下游处理速度反馈给上游生产者,是控制积压和内存增长的容量边界。
  • 同步原语:lock、event、condition、semaphore 属于 task 级同步原语,用于同一 event loop 内的互斥、通知和并发窗口控制。
  • Callback 调度call_soon()call_later()call_at()call_soon_threadsafe() 把普通 callback 接入 event loop 的就绪或定时路径。
  • 排查顺序:异步问题应先定位 loop 与 task 所有权,再检查 await 边界、等待对象、取消传播、背压链路和无让出代码段。