《Python高级编程》6.3 结构化并发与调试

对比 asyncio.TaskGroup 与 anyio 4.15.1 的结构化并发,讲清异常组如何打包与传播、asyncio.timeout 的 reschedule 语义,以及 PYTHONASYNCIODEBUG 慢回调警告等调试手段,附实测输出。

本节目标:用实测对比 asyncio.TaskGroup 与 anyio 4.15.1 的结构化并发语义,掌握异常组传播、asyncio.timeout 的 reschedule,以及慢回调等调试开关。
适用版本:Python 3.12+(实测 3.14.6;TaskGroup / asyncio.timeout 为 3.11+,anyio 4.15.1)

6.3 结构化并发与调试

高级异步一篇 用 Trio 的 nursery 讲结构化并发,并推荐用 AnyIO 做兼容层——但那篇是「怎么选、怎么用」。本节把落点放到机制与实测:TaskGroup 内部到底怎么把异常打包、except* 收到的异常组长什么样、anyio 在同一台机器上能跑哪几个后端,以及慢回调警告是怎么触发的。

6.3.1 TaskGroup:把并发装进 with 块

asyncio.TaskGroup(3.11+)把「一组并发任务」变成一个上下文管理器:with 块退出前,组内所有任务必须结束;任何一个失败,其余全部被取消,异常打包成 ExceptionGroup 抛出。实测:

import asyncio

async def worker(name, delay, fail=False):
    try:
        await asyncio.sleep(delay)
        if fail:
            raise ValueError(f"{name} 失败")
        print(f"{name} 正常完成")
    except asyncio.CancelledError:
        print(f"{name} 被取消")
        raise

async def main():
    try:
        async with asyncio.TaskGroup() as tg:
            tg.create_task(worker("A", 0.3))
            tg.create_task(worker("B", 0.05, fail=True))
            tg.create_task(worker("C", 0.3))
    except* ValueError as eg:
        print("外层 except* 捕获:", eg.exceptions)

asyncio.run(main())

实测输出(被取消任务的打印顺序不保证):

C 被取消
A 被取消
外层 except* 捕获: (ValueError('B 失败'),)

B 在 50ms 时抛错,A、C 本来要跑 300ms,立刻被取消。注意 except* 收到的是一个异常组,eg.exceptions 是元组而非单个异常——这就是结构化并发的核心承诺:子任务的异常不会被吞掉,也不会以「谁先抛谁赢」的方式丢掉其它异常。

6.3.2 anyio 4.15.1:跨后端的同一套语义

anyio.create_task_group() 提供与 TaskGroup 几乎一致的结构化语义,但可切换后端。同一段逻辑改写:

import anyio

async def worker(name, delay, fail=False):
    try:
        await anyio.sleep(delay)
        if fail:
            raise ValueError(f"{name} 失败")
        print(f"{name} 正常完成")
    except anyio.get_cancelled_exc_class():
        print(f"{name} 被取消")
        raise

async def main():
    try:
        async with anyio.create_task_group() as tg:
            tg.start_soon(worker, "A", 0.3)
            tg.start_soon(worker, "B", 0.05, True)
            tg.start_soon(worker, "C", 0.3)
    except* ValueError as eg:
        print("anyio(asyncio 后端) except*:", eg.exceptions)

anyio.run(main, backend="asyncio")

实测输出:

A 被取消
C 被取消
anyio(asyncio 后端) except*: (ValueError('B 失败'),)

行为与 TaskGroup 完全一致。anyio.get_cancelled_exc_class() 返回的是当前后端的取消异常类型——asyncio 后端下就是 asyncio.CancelledError,trio 后端下是 trio.Cancelled。这是 anyio 跨后端抽象的关键:业务代码不写死 asyncio.CancelledError。

但「跨后端」在本机有个现实边界。实测:

import anyio
print("可用后端:", anyio.get_available_backends())
print("全部后端:", anyio.get_all_backends())

async def main(): pass
try:
    anyio.run(main, backend="trio")
except Exception as e:
    print("请求 trio 后端 ->", type(e).__name__, ":", e)

实测输出:

可用后端: ('asyncio',)
全部后端: ('asyncio', 'trio')
请求 trio 后端 -> LookupError : Backend 'trio' is not available. Install it with: pip install anyio[trio]

get_all_backends() 报的是 anyio 理论上支持的后端,get_available_backends() 报的是当前环境真正可用的。本机只装了 anyio、没装 trio,所以 backend="trio" 直接抛 LookupError。结论:anyio 的跨后端是「代码可移植」,不是「装上 anyio 就自动有两套运行时」。

6.3.3 异常组如何打包与传播

except* 与 ExceptionGroup 是 PEP 654(3.11)的产物,它是结构化并发的运输工具。几个容易踩的语义点:

语义行为
except* 收到的 eg.exceptions永远是一个元组,即使只有一个异常
多个子任务同时失败全部被收集进同一个 ExceptionGroup,不丢异常
except* 可以匹配多条一次 except* 只处理匹配的子异常,其余留给后续 except*
未匹配的子异常打包成新的 ExceptionGroup 继续向上抛

多个任务同时失败时的打包效果:

import asyncio

async def boom(name):
    await asyncio.sleep(0.01)
    raise ValueError(name)

async def main():
    try:
        async with asyncio.TaskGroup() as tg:
            tg.create_task(boom("x"))
            tg.create_task(boom("y"))
    except* ValueError as eg:
        print("捕获到", len(eg.exceptions), "个:", eg.exceptions)

asyncio.run(main())

实测输出:

捕获到 2 个: (ValueError('x'), ValueError('y'))

两个子任务的异常都进了同一个组。这正是 asyncio.gather() 做不到的——gather 在第一个异常处就返回,除非显式传 return_exceptions=True,否则其余异常会「附着」到某个任务上被吞掉。结构化并发用异常组一次性把「所有出错的子任务」交代清楚。

6.3.4 asyncio.timeout 与 reschedule

asyncio.timeout()(3.11+)是 wait_for 的现代替代。它返回一个上下文管理器,块内超时会真正取消内层任务,并在退出块时抛 TimeoutError。实测:

import asyncio

async def main():
    loop = asyncio.get_running_loop()
    t0 = loop.time()
    try:
        async with asyncio.timeout(10) as cm:
            cm.reschedule(loop.time() + 0.05)   # 把 10s 改成 50ms
            await asyncio.sleep(5)
    except TimeoutError:
        print(f"reschedule 后 {loop.time()-t0:.3f}s 超时 -> TimeoutError")
        print("cm.expired() =", cm.expired())

asyncio.run(main())

实测输出:

reschedule 后 0.051s 超时 -> TimeoutError
cm.expired() = True

两个关键点:

  1. TimeoutError 在 async with 退出时才抛,不在 await asyncio.sleep(5) 那一行。所以 try/except TimeoutError 必须包住整个 async with,写在块内部是捕获不到的。
  2. cm.reschedule(when) 用绝对时间(loop.time() 基准)重设 deadline;cm.expired() 可事后查询是否因超时退出。这两个 API 让「一次请求内动态调整超时预算」成为可能,是 wait_for 没有的能力。

超时抛出后若被吞掉,循环可以继续跑——超时只取消那一段代码,不影响外层。

6.3.5 调试:慢回调警告与 debug 开关

事件循环卡住最常见的原因是某个回调里做了同步阻塞。开启 debug 后,asyncio 会对超过阈值的回调打警告。实测一个阻塞 0.25s 的回调:

import asyncio, time

async def main():
    loop = asyncio.get_running_loop()
    loop.call_soon(lambda: time.sleep(0.25))   # 阻塞事件循环 250ms
    await asyncio.sleep(0.01)

asyncio.run(main(), debug=True)

实测输出:

Executing <Handle main.<locals>.<lambda>() at ... created at ...> took 0.251 seconds

(took 后的数字即回调实际耗时,每次略有浮动;路径已省略。)

阈值来自 loop.slow_callback_duration,默认 0.1 秒。debug 开关的三种打开方式与实测默认值:

方式效果
asyncio.run(main(), debug=True)只对本次 run 的循环生效
loop.set_debug(True)运行时手动开启
环境变量 PYTHONASYNCIODEBUG=1进程内所有新建循环默认开启

实测三种方式下的 loop.get_debug():

默认(无环境变量):        get_debug() = False
PYTHONASYNCIODEBUG=1:      get_debug() = True
slow_callback_duration =   0.1

debug 模式除了慢回调警告,还会额外做三件事:记录每个 Handle 的创建溯源(警告里的 created at 就是它)、检查 call_soon 是否跨线程调用、给任务加更详细的 repr。代价是每个回调多一层包装,生产环境不要常开。

还有一条最容易被忽视的警告——协程创建了却从没被 await:

import asyncio, gc

async def never():
    return 1

never()          # 只创建协程对象,不 await
gc.collect()

实测输出:

RuntimeWarning: coroutine 'never' was never awaited
RuntimeWarning: Enable tracemalloc to get the object allocation traceback

它由协程对象的 __del__ 触发,所以不一定立刻出现,往往在 GC 时才打印。加上 tracemalloc 就能定位到创建位置。这正是 await 漏写时最容易忽略的坑。

小结

  • TaskGroup 与 anyio.create_task_group() 提供同一套结构化语义:with 块内任一任务失败,其余全部取消。
  • anyio 的跨后端是「代码可移植」;本机 get_available_backends() 只有 ('asyncio',),backend="trio" 会抛 LookupError。
  • 异常组是结构化并发的运输工具:except* 收到的 eg.exceptions 永远是元组,多个子任务的失败会全部打包,不丢异常。
  • asyncio.timeout 的 TimeoutError 在退出 async with 时才抛;reschedule() 用绝对时间重设 deadline。
  • 慢回调警告阈值默认 0.1s,由 slow_callback_duration 控制;PYTHONASYNCIODEBUG=1 可全局开启 debug。
  • 「coroutine was never awaited」由 __del__ 触发,出现时机不确定,配合 tracemalloc 才能定位。

至此,异步运行时的三块拼图——循环调度、任务与取消、结构化组织——都已讲清。下一章我们离开运行时,进入导入系统:import 语句背后,finder 与 loader 是怎么把模块对象造出来的。

阅读导航:上一节:任务、Future 与取消语义 · 下一节:导入协议与 finder / loader 。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「python」更多文章

  1. 《Python高级编程》目录
  2. 《Python高级编程》11.3 PEP 流程与版本迁移策略
  3. 《Python高级编程》11.2 嵌入式与自由线程运行时