登录
推荐 文章 Go 技术 课程 下载 专题 AI
首页 >  文章 >  python教程

asyncio TaskGroup 让并发任务在首错时一起收敛

来源:17golang原创

时间:2026-10-07 05:43:49 246浏览 收藏

asyncio.TaskGroup 的核心价值不是少写几行 await,而是把一组相关任务放进同一个生命周期边界:第一个任务以非 CancelledError 异常失败时,组会取消其余未完成任务;退出 async with 前会等待它们完成清理;最后再把可传播的异常组合成 ExceptionGroup 抛给调用方。这正是“首错时一起收敛”的含义。

官方文档:https://docs.python.org/3.14/library/asyncio-task.html#task-groups

TaskGroup 自 Python 3.11 加入。Python 3.13 改进了内部取消与外部取消同时发生时的处理,并正确保留取消计数;Python 3.14 的 TaskGroup.create_task() 会把所有关键字参数继续传给事件循环的 create_task()。如果项目需要兼容 3.10 或更早版本,仍要准备显式取消与回收方案。

TaskGroup 解决的不是少写几行 await

直接调用 asyncio.create_task() 时,任务的引用、等待、取消与异常回收都由调用方负责。任何一项遗漏,都可能出现后台任务继续运行、异常无人读取或资源没有及时释放。TaskGroup 把这些责任集中到异步上下文管理器中:用 tg.create_task() 创建的任务由组持有,离开上下文时统一等待。

TaskGroup 对调用方、组内任务和结果引用的静态所有权关系
图1:TaskGroup 任务所有权结构图。组边界持有子任务,并在退出上下文前等待所有任务进入终态。

这和 asyncio.gather() 的默认失败语义不同。官方文档说明,gather(..., return_exceptions=False) 会把第一个异常立即传播给等待者,但不会因此自动取消其他 awaitable;那些任务会继续运行。TaskGroup 则提供更强的结构化并发保证:一个任务失败时,会取消其余已调度任务并等待它们收尾。

能力TaskGroupgather 默认行为
任务所有权由组持有并统一等待调用方管理传入对象与结果
首个业务异常取消同组剩余任务异常先传播,其他任务继续
多个异常组合为 ExceptionGroup通常传播第一个异常,或作为结果返回
适用场景有共同成败边界的相关任务允许各任务独立完成的结果聚合

先确认版本和 create_task 参数

最基础的支持范围是 Python 3.11 及以上。任务组必须先进入 async with 才是活跃状态;在尚未进入、已经退出或正在关闭的任务组上调用 create_task() 会失败。当前 Python 3.14 文档还明确:当组不活跃时,传入的协程会被关闭,并抛出 RuntimeError。

import asyncio
import sys


def require_task_group() -> None:
    # TaskGroup 从 Python 3.11 开始提供,启动时提前给出明确提示。
    if sys.version_info  None:
    require_task_group()
    # 只有进入 async with 后,任务组才处于可创建任务的活跃状态。
    async with asyncio.TaskGroup() as tg:
        tg.create_task(asyncio.sleep(0.1), name="warm-up")


asyncio.run(main())

name 适合给任务添加便于诊断的名称,context 可指定 contextvars.Context。Python 3.14 还允许相关关键字参数透传到事件循环;但如果库需要跨多个 Python 小版本运行,应只使用目标版本明确支持的参数,不要把“最新文档可用”误当成“所有部署环境都可用”。

最小可用写法:在上下文退出后读取结果

下面的示例并发读取两个数据源。任务对象可以保存为局部变量,但应在 async with 结束后读取 result()。走到上下文外部时,TaskGroup 已经等待所有任务完成;如果其中一个失败,控制流不会执行到读取结果的位置,而是先进入异常处理。

import asyncio


async def fetch_price(symbol: str, delay: float) -> dict[str, float | str]:
    # 用可取消的 await 模拟网络 I/O;真实请求也应配置连接和读取超时。
    await asyncio.sleep(delay)
    return {"symbol": symbol, "price": 100.0 + delay}


async def load_dashboard() -> list[dict[str, float | str]]:
    async with asyncio.TaskGroup() as tg:
        task_a = tg.create_task(fetch_price("A", 0.2), name="price-A")
        task_b = tg.create_task(fetch_price("B", 0.3), name="price-B")

    # 离开上下文意味着组内任务都已结束,可以安全读取结果。
    return [task_a.result(), task_b.result()]


async def main() -> None:
    rows = await load_dashboard()
    # 示例只展示聚合结果;生产代码可在这里进入后续业务处理。
    print(rows)


asyncio.run(main())

TaskGroup 允许组内任务在组尚未关闭时继续向同一个组添加任务,例如把 tg 传给某个协程,由它创建相关子任务。一旦最后一个任务完成且上下文退出,就不能再向组中添加任务。这个边界很重要:它让任务树的所有权仍然可追踪,而不是任意扩散到全局事件循环。

首错如何让其他任务收敛

“首错”特指第一个非 asyncio.CancelledError 异常。它出现后,TaskGroup 会取消组内其他任务,不再接受新任务,并等待已取消任务执行清理逻辑。如果 async with 的主体仍在运行,包含该上下文的父任务也会收到一次取消以唤醒退出流程,但这个内部取消不会作为普通 CancelledError 穿出该 async with。

import asyncio


async def worker(name: str, fail: bool = False) -> str:
    try:
        await asyncio.sleep(0.2)
        if fail:
            # 业务异常会触发任务组取消其他未完成任务。
            raise ValueError(f"{name} 数据无效")

        await asyncio.sleep(5)
        return f"{name} 完成"
    finally:
        # 无论正常完成、失败还是被取消,都在这里释放连接或临时资源。
        print(f"{name} 已清理")


async def run_workers() -> None:
    async with asyncio.TaskGroup() as tg:
        tg.create_task(worker("task-A", fail=True))
        tg.create_task(worker("task-B"))
        tg.create_task(worker("task-C"))


asyncio.run(run_workers())

任务收到取消后,会在下一次可取消点抛出 CancelledError。协程应使用 try/finally 可靠清理资源。如果确实捕获了 CancelledError,清理完成后通常应继续 raise。官方文档特别提醒,TaskGroup 和 asyncio.timeout() 都在内部使用取消;吞掉 CancelledError 可能破坏这些结构化并发组件的行为。

async def cancellable_worker() -> None:
    resource = await open_resource()
    try:
        await resource.process()
    except asyncio.CancelledError:
        # 可以记录取消或补充清理,但不要把取消伪装成正常完成。
        await resource.abort()
        raise
    finally:
        # finally 确保连接、文件或锁最终被释放。
        await resource.close()

用 ExceptionGroup 接住并发错误

取消其余任务后,TaskGroup 会等待所有任务完成。如果存在一个或多个非取消异常,会按情况组合为 ExceptionGroup 或 BaseExceptionGroup 再抛出。Python 的 except* 可以按异常类型拆分处理,不必手动递归遍历嵌套异常组。

TaskGroup 中业务异常、取消、清理和 ExceptionGroup 的静态关系
图2:TaskGroup 异常边界图。业务异常触发同组取消,清理完成后由 ExceptionGroup 汇总可传播的异常。
import asyncio


class RemoteDataError(Exception):
    """表示某个远端数据源返回了不可用结果。"""


async def query(source: str) -> str:
    await asyncio.sleep(0)
    # 示例让两个数据源产生不同类型的业务异常。
    if source == "inventory":
        raise RemoteDataError("库存服务返回无效数据")
    raise TimeoutError("价格服务响应超时")


async def run_queries() -> None:
    try:
        async with asyncio.TaskGroup() as tg:
            tg.create_task(query("inventory"))
            tg.create_task(query("price"))
    except* RemoteDataError as group:
        # 这里只处理匹配 RemoteDataError 的异常子组。
        for error in group.exceptions:
            print(f"数据错误:{error}")
    except* TimeoutError as group:
        # 其他类型可以用独立 except* 分支分类处理。
        for error in group.exceptions:
            print(f"超时错误:{error}")


asyncio.run(run_queries())

并发调度存在竞争:某个任务失败后,其他任务可能在收到取消前已经成功、失败或进入清理。因此不能假设每次只会收到一个异常。except* 的意义正是让调用方按类型处理实际收集到的异常集合。

KeyboardInterrupt 和 SystemExit 是特殊情况。TaskGroup 仍会取消并等待其余任务,但随后重新抛出最初的 KeyboardInterrupt 或 SystemExit,而不是把它们包装成普通异常组。应用层不要用宽泛捕获把进程退出信号悄悄吞掉。

Python 3.10 及更早版本的兼容处理

旧版本没有标准库 TaskGroup 时,可以用 create_task() 加 gather() 模拟最关键的“失败后取消并回收”语义。重点是三个动作必须同时存在:保存任务引用、异常时取消所有未完成任务、再次 await 以回收取消结果,然后重新抛出原异常。

import asyncio
from collections.abc import Awaitable
from typing import TypeVar

T = TypeVar("T")


async def gather_cancel_on_error(*aws: Awaitable[T]) -> list[T]:
    # 保存强引用,便于统一取消和回收每个任务。
    tasks = [asyncio.create_task(aw) for aw in aws]
    try:
        return await asyncio.gather(*tasks)
    except BaseException:
        # 首个异常出现后,主动取消仍未结束的同批任务。
        for task in tasks:
            task.cancel()

        # 必须再次等待,确保取消和 finally 清理已经完成。
        await asyncio.gather(*tasks, return_exceptions=True)
        raise

这段兼容代码只覆盖常见平面任务集合,无法完整复制 TaskGroup 对嵌套任务组、内部与外部同时取消、取消计数和异常组的全部语义。若结构化并发是核心能力,升级到 Python 3.11 以上通常比持续维护自定义任务组更可靠。

性能与资源边界注意事项

TaskGroup 管的是生命周期,不会自动限制并发数量,也不会替你设置网络超时。一次创建几万个任务仍可能造成连接、内存和下游压力。批量任务应配合 asyncio.Semaphore、有界队列或分批调度;外部 I/O 应配合客户端超时或 asyncio.timeout()。

import asyncio


async def bounded_call(sem: asyncio.Semaphore, item: str) -> str:
    # 信号量限制同时进入外部依赖的任务数量。
    async with sem:
        async with asyncio.timeout(3):
            # timeout 会把超时转换为可在上下文外捕获的 TimeoutError。
            return await call_remote(item)


async def process_batch(items: list[str]) -> list[str]:
    sem = asyncio.Semaphore(20)
    async with asyncio.TaskGroup() as tg:
        tasks = [
            tg.create_task(bounded_call(sem, item), name=f"item-{item}")
            for item in items
        ]

    # 任务组成功退出后,再按输入顺序读取结果。
    return [task.result() for task in tasks]

还要区分“相关任务”和“长期后台任务”。如果多个请求必须共同成功或共同取消,适合放入一个 TaskGroup;如果任务应该独立于当前调用方长期运行,就不应假装放进临时任务组。长期后台任务需要明确的应用级所有者、强引用集合、异常日志和关闭协议。

常见问题

TaskGroup 会在任何异常时取消其他任务吗?

第一个非 CancelledError 异常会触发取消。正常取消本身不按普通业务异常汇总。KeyboardInterrupt 和 SystemExit 有单独的重新抛出规则。

TaskGroup 可以替代 gather 吗?

不能机械替代。相关任务需要共同成败和统一回收时优先 TaskGroup;任务允许独立完成、需要把异常当普通结果收集时,gather(return_exceptions=True) 仍有明确用途。

为什么不应该吞掉 CancelledError?

TaskGroup 依赖取消唤醒退出和清理流程。吞掉取消会让上层误以为任务仍按正常语义运行,甚至破坏超时与嵌套任务组。捕获后完成必要清理,通常应重新抛出。

如何提前主动终止一个 TaskGroup?

标准库目前没有原生 terminate 方法。官方文档给出的模式是向组内加入一个主动抛出专用异常的任务,再用 except* 忽略该专用异常。普通业务代码优先通过取消父任务、超时或明确的停止条件表达生命周期。

简要归纳:把一组任务放进 TaskGroup 后,创建、等待、取消和异常传播都归属于同一个语法边界。正确使用它的关键不是记住 async with 形式,而是让子任务可取消、在 finally 中清理资源、不要吞掉 CancelledError,并在边界外用 except* 处理并发异常。

声明:本文转载于:17golang原创 如有侵犯,请联系study_golang@163.com删除
相关阅读
更多>
最新阅读
更多>
课程推荐
更多>