当前位置:首页 > 文章列表 > 文章 > python教程 > asyncio TaskGroup 让并发任务在首错时一起收敛

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

来源:17golang原创 2026-10-07 05:43:49 0浏览 收藏

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删除
Goroutine 数量持续上涨却没有报错,如何定位泄漏入口Goroutine 数量持续上涨却没有报错,如何定位泄漏入口
上一篇
Goroutine 数量持续上涨却没有报错,如何定位泄漏入口
把无界并发改造成带容量限制的工作池
下一篇
把无界并发改造成带容量限制的工作池
查看更多
最新文章
查看更多
课程推荐
  • 前端进阶之JavaScript设计模式
    前端进阶之JavaScript设计模式
    设计模式是开发人员在软件开发过程中面临一般问题时的解决方案,代表了最佳的实践。本课程的主打内容包括JS常见设计模式以及具体应用场景,打造一站式知识长龙服务,适合有JS基础的同学学习。
    543次学习
  • GO语言核心编程课程
    GO语言核心编程课程
    本课程采用真实案例,全面具体可落地,从理论到实践,一步一步将GO核心编程技术、编程思想、底层实现融会贯通,使学习者贴近时代脉搏,做IT互联网时代的弄潮儿。
    516次学习
  • 简单聊聊mysql8与网络通信
    简单聊聊mysql8与网络通信
    如有问题加微信:Le-studyg;在课程中,我们将首先介绍MySQL8的新特性,包括性能优化、安全增强、新数据类型等,帮助学生快速熟悉MySQL8的最新功能。接着,我们将深入解析MySQL的网络通信机制,包括协议、连接管理、数据传输等,让
    500次学习
  • JavaScript正则表达式基础与实战
    JavaScript正则表达式基础与实战
    在任何一门编程语言中,正则表达式,都是一项重要的知识,它提供了高效的字符串匹配与捕获机制,可以极大的简化程序设计。
    487次学习
  • 从零制作响应式网站—Grid布局
    从零制作响应式网站—Grid布局
    本系列教程将展示从零制作一个假想的网络科技公司官网,分为导航,轮播,关于我们,成功案例,服务流程,团队介绍,数据部分,公司动态,底部信息等内容区块。网站整体采用CSSGrid布局,支持响应式,有流畅过渡和展现动画。
    485次学习
查看更多
AI推荐
  • PubMedQA数据集详解:生物医学问答基准、功能与应用指南
    PubMedQA
    深入了解PubMedQA生物医学问答数据集,涵盖其核心功能、使用方法及在临床决策、药物研发等场景的应用,助力提升NLP模型性能。
    360次使用
  • H2O EvalGPT:开源LLM大模型评估与排行榜工具
    H2O EvalGPT
    H2O EvalGPT是H2O.ai推出的开源LLM评估平台,提供详细的大模型性能排行榜、行业特定基准测试及A/B测试功能,助您快速选择最适合项目的高性能大语言模型。
    417次使用
  • LMArena是什么?伯克利AI模型评估平台使用指南与功能解析
    LMArena
    LMArena是加州大学伯克利分校推出的AI模型匿名评测平台。通过盲测投票机制,用户可对比不同大模型回答并生成实时排行榜,助力开发者优化模型及用户选择最佳AI工具。
    430次使用
  • 斯坦福HELM:大语言模型Holistic Evaluation整体评估框架详解
    HELM
    深入了解斯坦福推出的HELM(Holistic Evaluation of Language Models)大模型评测体系。本文解析其核心功能、安装配置步骤及应用场景,涵盖准确性、公平性、鲁棒性等多维度指标,助力开发者全面优化语言模型性能。
    381次使用
  • MMBench详解:多模态大模型基准测试、功能特点与使用指南
    MMBench
    MMBench是由上海人工智能实验室等机构联合推出的多模态基准测试平台,提供细粒度能力评估、大规模数据集及VLMEvalKit工具。本文详细介绍其核心功能、安装使用方法及应用场景,助力开发者全面评估多模态模型性能。
    208次使用
微信登录更方便
  • 密码登录
  • 注册账号
登录即同意 用户协议 和 隐私政策
返回登录
  • 重置密码