asyncio.TaskGroup 实战:告别 gather,让并发任务的异常与取消终于不乱

2026-09-14 0 663

写异步代码时间长了,迟早会撞上这么一幕:用 asyncio.gather 并发发出去五个请求,其中一个因为对方服务挂了抛了异常,另外四个已经正常返回。结果你以为处理得很干净,实际上那个异常被搜集到返回的列表里,然后你在循环里一 await,整个函数就此中断,剩下三个已经发出去的请求变成了无人认领的悬空任务,事件循环退出的时候还在后台默默跑。日志里什么都看不到,但下游偶尔会抱怨”你们在半夜发了一堆请求”。

这不是 gather 的 bug,它的设计本来就是这样:所有任务一起发出去,异常像商品一样堆在结果列表里,等你自己去取。问题在于,绝大多数时候我们并不想要这种”自己管理生命周期”的自由,我们想要的是:一起发出去,一起结束,出错就一起收场。

Python 3.11 引入的 asyncio.TaskGroup 就是为这个需求设计的。它属于 PEP 654 提出的结构化并发里比较早落地的一块,到今天已经足够稳定,值得拿出来认真用。

先看一段典型的老代码

需求很简单:并发地从三个下游拉数据,任意一个失败,整个操作算失败。

async def fetch_all(ids):
    tasks = [asyncio.create_task(fetch_one(i)) for i in ids]
    try:
        results = await asyncio.gather(*tasks)
        return results
    except Exception as e:
        # 想把剩下还在跑的都取消掉
        for t in tasks:
            t.cancel()
        raise

写到这里其实还没完,因为 gather 的取消语义很微妙。它只会在”外部取消 gather 本身”的时候,把取消传递给子任务;如果只是某个子任务自己抛了异常,gather 会等其它任务跑完之后再抛出,剩下的任务并不会因为你 t.cancel() 就立刻停下——因为它已经完成了。

而且还有一个更隐蔽的问题:asyncio.create_task 创建的任务必须有人 await,否则当它抛异常时可以产生 “Task exception was never retrieved” 的警告,有时候你根本没在日志里看到这个警告,因为可能被其它输出的刷屏盖掉了。

TaskGroup:给并发任务划一个明确的边界

同样的事情,用 TaskGroup 重写是这样的:

async def fetch_all(ids):
    results = []
    async with asyncio.TaskGroup() as tg:
        tasks = [tg.create_task(fetch_one(i)) for i in ids]
    return [t.result() for t in tasks]

这段代码短是短,但真正重要的是它背后那套语义。

退出 async with 块的时候,TaskGroup 会等里面所有任务跑完。 这是它的核心承诺——代码块结束时,你可以确定地说”里面已经没有任何活着的任务了”。这是结构化并发最基本的那条约束。

任何一个子任务抛异常,TaskGroup 会取消其余尚未完成的任务,然后抛出一个 ExceptionGroup 你不用自己写取消逻辑,也不用担心漏掉哪个任务还在后台。

没有”未取回的异常”这种状态。 所有子任务的异常都会被聚合进 ExceptionGroup,由 async with 负责往外抛。你的 except 分支一定看得见。

取消是可以传播的。 如果外层这个函数自己被取消(比如超时到点),取消会正确穿透进去,TaskGroup 里的每个任务都会收到 CancelledError,整个结构能完整收场。这一块以前用 gather 手写特别容易漏。

ExceptionGroup:一个异常,多个原因

上面一直说 ExceptionGroup,得把它单独讲清楚,不然用起来会一头雾水。

用 TaskGroup 的时候,如果两个任务同时失败了,你实际上会收到一个”装着两个异常的容器”。它的类型是 ExceptionGroup,可以理解成一个特殊的异常,里面包着一组异常。

async def boom(n, delay):
    await asyncio.sleep(delay)
    raise RuntimeError(f"任务 {n} 失败")

async def main():
    async with asyncio.TaskGroup() as tg:
        tg.create_task(boom(1, 0.1))
        tg.create_task(boom(2, 0.2))

try:
    asyncio.run(main())
except Exception as e:
    print(type(e))             # <class 'ExceptionGroup'>
    print(len(e.exceptions))   # 2

处理这种异常,有两种思路。

第一种是传统 try/except,直接捕获 ExceptionGroup,然后自己去遍历 e.exceptions。适合你确实需要拿到所有异常、逐一分析的时候。

第二种是 Python 3.11 同时引入的 except* 语法,它可以用多个分支按类型分别处理内部的各种异常。

try:
    asyncio.run(main())
except* TimeoutError as eg:
    print(f"超时 {len(eg.exceptions)} 个")
except* (RuntimeError, ValueError) as eg:
    for e in eg.exceptions:
        print(f"业务异常:{e}")

这段代码的意思是:把 ExceptionGroup 里的异常按类型”分流”到不同分支去处理。TimeoutError 归第一个分支管,RuntimeError 和 ValueError 归第二个。这种匹配方式有点像模式匹配,但专门针对异常组。

有两条规则一定要记住:

用了 except*,就不能在同一个 try 里再用普通的 except Python 直接抛 SyntaxError,没有商量余地。要么全用 except*,要么全用 except

except* 分支里 raise 一个异常,会把它重新打包成 ExceptionGroup 往外抛。 也就是说,即使你的分支里只处理了一个异常,只要它没被吞掉,最终传播出去的还是一个 Group。接收方如果不熟悉这个特性,往往会被搞懵。

顺带说一句,ExceptionGroup 是可以嵌套的。TaskGroup 的 ExceptionGroup 里如果某个异常本身也是 ExceptionGroup(比如嵌套的并发结构),你拿到的是一个树状结构。except* 会自动递归匹配,不用自己写深度遍历。

实战:并发拉取用户资料的聚合接口

做一个更完整一点的例子。场景是一个前端传入用户 ID 列表,后端需要并发去拉每个用户的详细资料,然后把成功的返回,失败的记录日志但不影响整体响应。

这个需求如果直接用 TaskGroup,会遇到一个矛盾:TaskGroup 语义是”一个失败全部取消”,而我们要的是”能拿多少是多少”。所以不能天真地直接用它。

正确的做法是让每个子任务自己处理异常,把”是否成功”作为一个数据结果传到外面。

import asyncio
from dataclasses import dataclass

@dataclass
class FetchResult:
    user_id: str
    data: dict | None
    error: Exception | None
    @property
    def ok(self) -> bool:
        return self.error is None

async def safe_fetch(client, user_id: str) -> FetchResult:
    try:
        data = await client.get_user(user_id)
        return FetchResult(user_id, data, None)
    except Exception as e:
        return FetchResult(user_id, None, e)

async def fetch_users(client, user_ids: list[str]) -> list[FetchResult]:
    async with asyncio.TaskGroup() as tg:
        tasks = [tg.create_task(safe_fetch(client, uid)) for uid in user_ids]
    return [t.result() for t in tasks]

这里的关键在于 safe_fetch 把所有异常都吞掉了,所以 TaskGroup 永远不会因为业务失败而提前取消其它任务。等待所有任务结束后,t.result() 一定能拿到一个 FetchResult 对象,不会抛异常。

外层调用拿到结果之后,可以根据 result.ok 分开处理:

results = await fetch_users(client, user_ids)
ok = [r.data for r in results if r.ok]
failed = [r.user_id for r in results if not r.ok]

if failed:
    logger.warning("以下用户拉取失败:%s", failed)

return {"users": ok, "failed_ids": failed}

这套写法的好处是:异常处理的位置非常清楚。业务层面的”失败”变成了一种数据,用返回值传递;编程层面的意外(比如你代码里有 bug)依然会被 TaskGroup 捕获,不会被静默吞掉。

用这种模式的时候,有个细节值得留意:safe_fetch 里不要捕 CancelledError。 Python 3.8 之后 CancelledError 继承自 BaseException 而不是 Exception,所以 except Exception 不会误捕获取消信号。但如果你不小心写成 except BaseException,那 TaskGroup 的取消传播就会被你拦住,导致整个结构卡住不退出。这类 bug 排查起来特别浪费时间。

超时与取消:一个容易搞混的地方

TaskGroup 本身不带超时参数,需要配合 asyncio.timeout 来用(这是 Python 3.11 引入的另一个好东西,代替了旧版的 wait_for)。

async def fetch_with_timeout(client, user_ids, seconds):
    try:
        async with asyncio.timeout(seconds):
            async with asyncio.TaskGroup() as tg:
                tasks = [tg.create_task(safe_fetch(client, uid)) for uid in user_ids]
        return [t.result() for t in tasks]
    except TimeoutError:
        logger.warning("整体超时 %s 秒", seconds)
        return []

这里的层次很讲究:timeout 在外,TaskGroup 在内。超时到点的时候,TaskGroup 会收到取消信号,它再把取消传给里面所有还在跑的任务。等到所有任务都退出了,async with 才让外层 timeout 的 TimeoutError 抛出去。整个收场过程是有序的。

如果反过来写,把 TaskGroup 放外面,timeout 放里面,那超时触发之后,TaskGroup 会看到某个子任务抛了 TimeoutError,于是取消其它任务并抛 ExceptionGroup——这时候你拿到的是一个 ExceptionGroup 包着 TimeoutError,处理起来就更绕。

另外 asyncio.timeout 有一个特点:它抛的是标准的 TimeoutError 而不是 asyncio.TimeoutError。Python 3.10 之前这俩是不同类,3.11 之后统一了,如果你升级上来遇到旧代码里捕 asyncio.TimeoutError 的地方,注意一下。

TaskGroup 做不到的那些事

写了这么多好处,也得说说什么时候它不合适,免得你把它当成万能工具到处用。

不能只等一部分任务。 TaskGroup 的语义是”一块进一块出”,代码块结束就必须等所有任务。如果你需要”哪个任务先完成就先处理哪个”,还得回到 asyncio.as_completed 或者 asyncio.wait

不能只取第一个成功结果。 类似”从多个镜像源并发下载,谁先成功用谁”的需求,TaskGroup 处理起来不顺手。你可以在子任务里做一次事件通信,但写出来会比较别扭。这种场景用 asyncio.wait(FIRST_COMPLETED) 循环更自然。

不适合长期后台任务。 TaskGroup 的生命周期是跟着 async with 走的。如果你要的是”启动一个后台任务,函数返回之后它还在跑”,那不是在 TaskGroup 的语义里。

动态往组里加任务有限制。 任务只能在 async with 块内创建,块结束了就不能再加。如果你要写”边处理边根据需要启动新的子任务”,得用嵌套的 TaskGroup,或者换别的机制。

和 asyncio.gather 的取舍

那是不是就说 gather 没用,应该全部换成 TaskGroup?也不是。

gather 有一个 TaskGroup 没有的特性:它会保留每个任务的结果,即使某个任务失败了,其它任务的返回值也可以从 return_exceptions=True 拿到。这在”知道一部分会失败但仍然想要全部结果”的场景下更直接。

另外 gather 的写法更轻,如果只是两三个已知的任务并发跑一下,用 TaskGroup 加一层缩进反而显得重。它的重量级在于更严格的取消语义和更清楚的错误边界——当你真的需要这些的时候,它是正确选择;如果只是在写一个玩具脚本,gather 就够了。

我自己的判断标准是这样的:

  • 需要”要么全成要么全败”的原子性,用 TaskGroup
  • 需要区分”哪些成功哪些失败”并从返回列表里分别处理,用 gather 加 return_exceptions=True
  • 任务数多余 5 个、或者任务生命周期跨越多个函数,用 TaskGroup
  • 只是一次性发两个请求然后 await 结果,用 gather 更省事

在异步代码里,”取消”一直是最容易出错的地方,TaskGroup 至少把这块做正确了。多花几行代码换来一个”不存在泄漏任务”的保证,通常划算。

最后说一句

Python 这些年一直在慢慢把”并发”这件事往更可控的方向推。结构化并发这个概念从语言层面落地到实际可用,TaskGroup 是最接地气的那一块。它不炫技,也不需要你重新学一套心智模型,只是让异步代码的边界更清晰了一点。

如果你手上还有一堆用 gather 写的、每次看到都觉得心里没底的并发代码,不妨挑一个改成 TaskGroup 试试。写完第一遍可能觉得只是缩进变了,但多看几个版本之后,你会发现自己的异步代码开始变得可预测了——出错的时候行为是确定的,取消的时候是干净的,这才是真正省时间的地方。

asyncio.TaskGroup 实战:告别 gather,让并发任务的异常与取消终于不乱
收藏 (0) 打赏

感谢您的支持,我会继续努力的!

打开微信/支付宝扫一扫,即可进行扫码打赏哦,分享从这里开始,精彩与您同在
点赞 (0)

版权声明:
本站资源有的来自互联网收集整理,本站纯免费分享提供学习使用,如果侵犯了您的合法权益,请发送邮件1506151422@qq.com联系,将会及时下架删除。
本站资源仅供研究、学习交流之用,免费开源项目不代表完全可商用,若商业用途请先咨询开发企业能否商用,否则产生的一切后果将由下载用户自行承担。
原创板块未经允许不得转载,否则将追究法律责任。

淘吗网 python asyncio.TaskGroup 实战:告别 gather,让并发任务的异常与取消终于不乱 https://www.taomawang.com/server/python/2756.html

下一篇:

已经没有下一篇了!

常见问题

相关文章

猜你喜欢
发表评论
暂无评论
官方客服团队

为您解决烦忧 - 24小时在线 专业服务