Python 基础体系 · 第 52/112 篇。示例统一以 Python 3.14 为语言基线;第三方库使用与其兼容的现代稳定版本,版本敏感行为会单独说明。

Python 异步任务:create_task、TaskGroup、异常组和结构化并发

asyncio 中,协程函数只是可执行逻辑的定义;调用协程函数得到的是协程对象,并不会自动让它运行。协程对象只有在被 await、交给 asyncio.run(),或包装成 Task 后,才会被事件循环调度。Task 则是事件循环中负责驱动协程执行、保存结果并传播异常的对象。(docs.python.org)

这一区分决定了异步程序的并发边界:

await fetch("a")
await fetch("b")

这里两个操作是顺序等待的,只有 fetch("a") 完成后才会开始等待 fetch("b")

而下面的代码会先创建两个任务:

task_a = asyncio.create_task(fetch("a"))
task_b = asyncio.create_task(fetch("b"))

result_a = await task_a
result_b = await task_b

fetch("a")fetch("b") 都已经被安排到事件循环中运行,两个任务可以在各自等待 I/O 时交替执行。事件循环采用协作式调度:同一时刻运行一个任务;当当前任务在 await 处等待 Future、I/O 或定时器时,事件循环才有机会运行其他任务。(docs.python.org)

本文围绕四个概念展开:

  • asyncio.create_task():显式创建一个独立任务;
  • asyncio.TaskGroup:把相关任务绑定到一个生命周期范围内;
  • ExceptionGroupBaseExceptionGroupexcept*:表示并处理多个并发异常;
  • 结构化并发:让任务的创建、等待、取消和失败传播具有可推导的边界。

一、从协程对象到 Task

1. 协程对象不等于正在运行

import asyncio


async def work():
    print("work started")
    await asyncio.sleep(0.1)
    print("work finished")


async def main():
    work()  # 只创建协程对象,不会执行函数体


asyncio.run(main())

这段代码不会输出 work started。因为 work() 的调用结果是协程对象,而不是已经调度的任务;如果协程对象最终没有被等待,Python 通常还会报告类似下面的警告:

RuntimeWarning: coroutine 'work' was never awaited

正确的顺序等待方式是:

async def main():
    await work()

这里 await 做了两件事:

  1. 驱动 work() 开始执行;
  2. 当前协程暂停,直到 work() 返回结果或抛出异常。

如果希望多个操作并发进行,则需要创建任务。


2. Task 是协程的可管理执行实体

task = asyncio.create_task(work())

可以把这条语句理解为:

  1. 创建一个 Task
  2. work() 绑定到该任务;
  3. 把任务提交给当前正在运行的事件循环;
  4. 返回任务对象;
  5. 后续可以通过 await task 等待结果,或者通过 task.cancel() 请求取消。

asyncio.create_task() 必须在运行中的事件循环上下文中调用;如果当前线程没有运行中的事件循环,会抛出 RuntimeError。(docs.python.org)

一个完整示例:

import asyncio
import time


async def say_after(delay: float, text: str) -> str:
    await asyncio.sleep(delay)
    print(text)
    return text


async def main():
    started = time.perf_counter()

    task_a = asyncio.create_task(say_after(1, "A"))
    task_b = asyncio.create_task(say_after(2, "B"))

    result_a = await task_a
    result_b = await task_b

    elapsed = time.perf_counter() - started
    print(result_a, result_b)
    print(f"elapsed: {elapsed:.1f}s")


asyncio.run(main())

典型输出:

A
B
A B
elapsed: 2.0s

如果改成顺序等待:

async def main():
    await say_after(1, "A")
    await say_after(2, "B")

总等待时间约为 1 + 2 = 3 秒;并发版本的总等待时间接近:

Tconcurrentmax(TA,TB)T_{\text{concurrent}} \approx \max(T_A, T_B)

其中:

  • TAT_A 是任务 A 的总执行时间;
  • TBT_B 是任务 B 的总执行时间;
  • max 表示两项工作同时开始时,整体完成时间由较慢者决定。

这不是线程并行。若任务执行的是纯 Python 的长时间 CPU 计算,并且中间没有 await,它会持续占用事件循环线程,其他任务无法获得调度机会。


二、Task 的生命周期与结果传播

一个任务至少可以从以下几个状态理解:

创建
  │
  ▼
pending ──────► running ──────► done
  │                │              │
  │                │              ├── 正常返回结果
  │                │              ├── 抛出异常
  │                │              └── 被取消
  │                │
  └──── cancel() ──┘

任务完成的条件不是只有“成功返回”。协程正常返回、抛出异常,或者最终确认被取消,都会使任务进入完成状态。task.done() 可以检查任务是否完成;task.result() 返回结果或重新抛出任务中的异常;如果任务被取消,result() 会抛出 CancelledError。(docs.python.org)

import asyncio


async def fail():
    await asyncio.sleep(0)
    raise ValueError("bad input")


async def main():
    task = asyncio.create_task(fail())

    print(task.done())  # 通常为 False

    try:
        await task
    except ValueError as exc:
        print("caught:", exc)

    print(task.done())       # True
    print(task.cancelled())  # False
    print(type(task.exception()).__name__)  # ValueError


asyncio.run(main())

这里 await task 是异常传播点。任务内部的 ValueError 不会在 create_task() 调用处立即抛出,而是在等待任务结果时重新抛出。

因此下面这种代码存在明显问题:

async def main():
    asyncio.create_task(fail())
    print("main finished")

调用者既没有保存任务引用,也没有等待任务结果。失败可能只在任务被回收时通过日志报告:

Task exception was never retrieved

事件循环对任务只保留弱引用。若应用需要“发起后仍可靠运行”的后台任务,必须保存强引用;但即使保存到集合中,如果从不读取任务异常,仍然可能遗漏失败。官方文档建议对于需要可靠等待和异常传播的相关任务使用 TaskGroup。(docs.python.org)


三、create_task() 的能力与边界

1. create_task() 只负责创建,不负责收尾

create_task() 的职责是把一个协程变成任务并进行调度。它不会自动:

  • 等待任务结束;
  • 收集任务结果;
  • 取消兄弟任务;
  • 在父协程结束时保证子任务已经完成;
  • 把多个任务的异常聚合起来。

例如:

import asyncio


async def worker(name: str, delay: float):
    try:
        await asyncio.sleep(delay)
        print(name, "done")
    finally:
        print(name, "cleanup")


async def main():
    asyncio.create_task(worker("slow", 10))
    await asyncio.sleep(0.1)
    print("main returns")


asyncio.run(main())

main() 返回后,asyncio.run() 会处理事件循环的关闭过程,未完成任务通常会收到取消请求;但这不是 create_task() 本身提供的父子生命周期关系。若代码运行在一个长期存活的事件循环中,任务可能继续存在,甚至成为泄漏的后台工作。

2. 保存任务引用的“裸任务”模式

如果确实需要独立后台任务,可以使用集合维护强引用:

background_tasks: set[asyncio.Task[object]] = set()


def start_background(coro) -> asyncio.Task[object]:
    task = asyncio.create_task(coro)
    background_tasks.add(task)
    task.add_done_callback(background_tasks.discard)
    return task

但这个模式还需要处理异常:

def report_task_result(task: asyncio.Task[object]) -> None:
    try:
        task.result()
    except asyncio.CancelledError:
        pass
    except Exception as exc:
        print(f"background task failed: {exc!r}")


def start_background(coro) -> asyncio.Task[object]:
    task = asyncio.create_task(coro)
    background_tasks.add(task)
    task.add_done_callback(background_tasks.discard)
    task.add_done_callback(report_task_result)
    return task

这类代码适合真正独立于当前请求或当前操作的服务,例如:

  • 独立指标上报;
  • 生命周期长于单次请求的消费循环;
  • 明确设计为脱离调用者完成的后台服务。

如果任务和当前操作有共同目标,裸 create_task() 往往缺少足够的生命周期约束,TaskGroup 更合适。


四、TaskGroup:把任务绑定到一个范围

TaskGroup 是一个异步上下文管理器。通过 tg.create_task() 创建的任务属于该组;离开 async with 时,任务组会等待组内任务完成。TaskGroup 在 Python 3.11 引入。(docs.python.org)

import asyncio


async def say_after(delay: float, text: str) -> str:
    await asyncio.sleep(delay)
    print(text)
    return text


async def main():
    async with asyncio.TaskGroup() as tg:
        task_a = tg.create_task(say_after(1, "A"))
        task_b = tg.create_task(say_after(2, "B"))

    # 离开 async with 后,两个任务都已经完成
    print(task_a.result(), task_b.result())


asyncio.run(main())

这里的等待是隐式的:

async with asyncio.TaskGroup() as tg:
    ...
# 到这里,tg 中的所有任务都已结束

但“隐式等待”不等于“忽略结果”。tg.create_task() 仍然返回 Task,因此可以在上下文退出后读取任务结果。

1. TaskGroup 的核心故障规则

设一个任务组中有任务:

G={T1,T2,,Tn}G = \{T_1, T_2, \ldots, T_n\}

当某个任务 TiT_i 首次抛出非取消异常时,TaskGroup 会:

  1. 取消其他尚未完成的任务;
  2. 禁止继续向该组添加新任务;
  3. 如果 async with 的主体仍在执行,则取消直接包含该 async with 的父任务;
  4. 等待所有任务完成清理;
  5. 将任务异常组合成异常组并抛出。

CancelledError 不会触发这个“失败兄弟任务”规则;KeyboardInterruptSystemExit 则有特殊传播规则。(docs.python.org)

可以用下面的程序观察故障路径:

import asyncio


async def worker(name: str, delay: float, error: Exception | None = None):
    try:
        print(name, "started")
        await asyncio.sleep(delay)

        if error is not None:
            raise error

        print(name, "finished")
        return name
    except asyncio.CancelledError:
        print(name, "cancelled")
        raise
    finally:
        print(name, "cleanup")


async def main():
    async with asyncio.TaskGroup() as tg:
        tg.create_task(worker("ok", 2))
        tg.create_task(worker("failed", 0.2, ValueError("boom")))
        tg.create_task(worker("slow", 10))


try:
    asyncio.run(main())
except* ValueError as eg:
    print("handled:", eg)

典型输出顺序类似:

ok started
failed started
slow started
failed cleanup
slow cancelled
slow cleanup
ok cancelled
ok cleanup
handled: unhandled errors in a TaskGroup (1 sub-exception)

具体打印顺序可能因任务到达调度点的时机而不同,但因果关系应保持:

failed 抛出 ValueError
        │
        ▼
TaskGroup 取消 ok 和 slow
        │
        ▼
等待它们执行 finally
        │
        ▼
TaskGroup 抛出 ExceptionGroup

任务取消是一个请求,不是立即强制终止。调用 task.cancel() 会使任务在下一次可取消的执行点收到 CancelledError;协程应在 finally 中释放资源,并且通常应在清理结束后重新抛出 CancelledError。吞掉取消异常会破坏 TaskGroupasyncio.timeout() 等结构化并发组件的内部控制流程。(docs.python.org)

正确的清理模式是:

async def use_resource():
    resource = await acquire()
    try:
        await operate(resource)
    finally:
        await release(resource)

不推荐这样写:

async def bad_worker():
    try:
        await operate()
    except asyncio.CancelledError:
        print("cancelled")
        return  # 错误:取消被吞掉

如果业务确实要抑制取消,不仅要明确其语义,还通常需要配合 uncancel() 清除任务的取消状态;一般应用代码不应随意调用它。(docs.python.org)


五、create_task()TaskGroup 的差异

两者都能并发运行协程,但管理模型不同。

维度 asyncio.create_task() TaskGroup.create_task()
创建任务 可以 可以
自动等待所有任务 不可以 离开上下文时自动等待
保存任务引用 调用者负责 任务组负责
一个任务失败时取消兄弟任务 不会自动做
多个异常聚合 调用者负责 自动形成异常组
生命周期边界 由调用者自行约定 async with 词法范围约定
适合场景 独立后台任务、特殊调度 一组共同完成某个操作的任务

TaskGroup 并不意味着所有任务都必须成功。它表达的是更强的故障语义:

这些任务属于同一个操作;一个任务失败时,其他任务通常已经失去继续执行的意义。

例如并行读取三个服务:

async def load_dashboard():
    async with asyncio.TaskGroup() as tg:
        user_task = tg.create_task(load_user())
        order_task = tg.create_task(load_orders())
        quota_task = tg.create_task(load_quota())

    return {
        "user": user_task.result(),
        "orders": order_task.result(),
        "quota": quota_task.result(),
    }

如果 load_user() 失败,其他两个任务会被取消。对于一个必须同时拥有三类数据的页面,这是合理的。

但如果需求是“尽可能收集所有独立结果”,就不能机械地使用 TaskGroup 的失败即取消语义。此时可以让任务自己把失败转换成值:

async def safe_call(coro):
    try:
        return await coro
    except Exception as exc:
        return exc


async def load_all():
    async with asyncio.TaskGroup() as tg:
        tasks = [
            tg.create_task(safe_call(load_user())),
            tg.create_task(safe_call(load_orders())),
            tg.create_task(safe_call(load_quota())),
        ]

    return [task.result() for task in tasks]

这里任务本身没有把异常抛出到 TaskGroup;异常被转换成普通返回值,所以任务组会等待全部任务完成。


六、为什么 gather() 不能等同于 TaskGroup

asyncio.gather() 也能并发等待多个 awaitable,但默认情况下,第一个异常会立即传播给等待者,其他 awaitable 不会因此自动取消,而会继续运行。只有 gather() 自己被取消时,尚未完成的 awaitable 才会被取消。(docs.python.org)

async def main():
    try:
        await asyncio.gather(
            fail_fast(),
            long_running(),
        )
    except ValueError:
        print("caught")

fail_fast() 抛出异常后,gather() 可能已经把异常交给 main(),但 long_running() 仍在后台继续执行。若此时调用者认为“整个操作已经结束”,就可能出现:

  • 后台任务继续修改共享状态;
  • 任务持有的连接没有及时释放;
  • 后续日志与当前请求失去关联;
  • 进程关闭时出现未完成任务。

TaskGroup 的语义更接近:

创建一组相关任务
        │
        ▼
等待全部任务完成
        │
        ├── 全部成功:正常离开
        │
        └── 任一失败:取消兄弟任务
                          │
                          ▼
                    等待清理完成
                          │
                          ▼
                    抛出异常组

因此,gather() 更像是一个“并发结果收集器”,而 TaskGroup 是一个“带故障传播和生命周期管理的任务作用域”。

gather(return_exceptions=True) 的边界

results = await asyncio.gather(
    call_a(),
    call_b(),
    return_exceptions=True,
)

此时异常会出现在结果列表中,而不是立即抛出:

[
    "result-a",
    ValueError("call_b failed"),
]

这种方式适合每项工作彼此独立、调用者确实希望收到成功值和失败值的场景。但它不会自动建立父子生命周期,也不会因为一个子任务失败而取消其他子任务。


七、异常组:为什么需要 ExceptionGroup

同步代码中,一个 try 语句通常面对一个当前异常:

try:
    operation()
except ValueError:
    ...

异步任务并发运行时,多个任务可能在同一个等待周期内分别失败:

task_a ── ValueError
task_b ── OSError
task_c ── TypeError

如果只保留其中一个异常,其他故障信息就会丢失。ExceptionGroup 用一个异常对象包装多个互不相关的异常;BaseExceptionGroup 可以包装更广泛的异常类型。ExceptionGroup 继承自 Exception,只能包含 Exception 子类;BaseExceptionGroup 继承自 BaseException,可以包含包括 KeyboardInterrupt 在内的更底层异常。(docs.python.org)

errors = ExceptionGroup(
    "batch failed",
    [
        ValueError("bad value"),
        OSError("network failed"),
    ],
)

raise errors

打印时会看到异常组树状结构,而不是普通异常的一条线性 traceback。

1. ExceptionGroup 的嵌套结构

异常组本身可以包含另一个异常组:

raise ExceptionGroup(
    "outer",
    [
        ValueError("v1"),
        ExceptionGroup(
            "inner",
            [
                OSError("o1"),
                TypeError("t1"),
            ],
        ),
    ],
)

这意味着异常处理不能简单地把 eg.exceptions 当成平面列表。异常组保留了嵌套结构;subgroup()split() 会在保留结构的基础上提取匹配部分。(docs.python.org)


八、exceptexcept* 的区别

普通 except 匹配的是整个异常对象

try:
    raise ExceptionGroup(
        "errors",
        [ValueError("v"), TypeError("t")],
    )
except Exception as exc:
    print(type(exc).__name__)

输出:

ExceptionGroup

它不会自动把组内的 ValueErrorTypeError 拆开。

except* 则匹配异常组内部的成员:

try:
    raise ExceptionGroup(
        "errors",
        [
            ValueError("bad value"),
            TypeError("bad type"),
            OSError("network error"),
        ],
    )
except* ValueError as eg:
    print("values:", eg.exceptions)
except* OSError as eg:
    print("os errors:", eg.exceptions)

输出类似:

values: (ValueError('bad value'),)
os errors: (OSError('network error'),)

剩余的 TypeError 没有被处理,最终会继续传播。except* 的匹配过程可以形式化为:

E=MT(E)RT(E)E = M_T(E) \cup R_T(E)

其中:

  • EE 是原始异常组;
  • MT(E)M_T(E) 是匹配异常类型 TT 的子组;
  • RT(E)R_T(E) 是不匹配的剩余子组。

每个 except* T 处理当前剩余组中的匹配部分,未匹配部分继续交给后续处理器;所有处理器执行结束后,未处理异常与处理器中重新抛出的异常会重新合并并传播。(docs.python.org)

except* 不是普通 except 的“多次执行版”

下面两个特性非常重要。

第一,try 中不能同时使用 exceptexcept*

# SyntaxError
try:
    ...
except ValueError:
    ...
except* OSError:
    ...

第二,多个 except* 子句不是“匹配到一个就停止”。它们分别从异常组中提取自己的匹配子组,因此同一个异常组可以同时触发多个处理器。(docs.python.org)

except* 中也不能使用 returnbreakcontinue,因为多个异常分支需要完成统一的合并过程。(docs.python.org)


九、处理 TaskGroup 的异常组

下面是一个完整的并发批处理示例:

import asyncio


async def fetch_item(item_id: int):
    await asyncio.sleep(0.1 * item_id)

    if item_id == 2:
        raise ValueError(f"invalid item: {item_id}")

    if item_id == 3:
        raise ConnectionError(f"connection failed: {item_id}")

    return item_id * 10


async def batch():
    async with asyncio.TaskGroup() as tg:
        tasks = [
            tg.create_task(fetch_item(item_id))
            for item_id in range(1, 5)
        ]

    return [task.result() for task in tasks]


async def main():
    try:
        values = await batch()
    except* ValueError as eg:
        for exc in eg.exceptions:
            print("validation error:", exc)
    except* ConnectionError as eg:
        for exc in eg.exceptions:
            print("network error:", exc)


asyncio.run(main())

需要注意一个并发时序问题:TaskGroup 在第一个非取消异常出现后会取消其他任务。因此只有在多个任务在取消生效前都已经抛出异常时,异常组中才会包含多个异常。

例如:

async def fail(name: str):
    await asyncio.sleep(0)
    raise ValueError(name)

如果多个任务几乎同时到达失败点,它们可能都先抛出异常,随后任务组才开始取消兄弟任务,于是最终异常组可能包含多个 ValueError

而如果一个任务先失败,其他任务正在长时间等待:

async def fail_first():
    await asyncio.sleep(0.1)
    raise ValueError("first")


async def slow():
    await asyncio.sleep(10)
    raise TypeError("usually cancelled first")

slow() 通常会先收到取消请求,因此它的 TypeError 不一定会出现于最终异常组中。这里的“通常”来自并发时序,而不是异常组保证某种固定顺序。


十、异常处理器中应该做什么

except* 适合按异常类别执行不同的恢复、记录或转换逻辑:

async def main():
    try:
        async with asyncio.TaskGroup() as tg:
            tg.create_task(validate_payload())
            tg.create_task(write_to_database())
            tg.create_task(call_remote_service())
    except* ValueError as eg:
        record_validation_errors(eg)
    except* ConnectionError as eg:
        mark_dependency_unavailable(eg)

但处理器不应假定所有错误都已经被处理。下面的代码会让未匹配异常继续传播:

try:
    ...
except* ValueError:
    print("handled value errors")

如果异常组中还包含 RuntimeError,该 RuntimeError 所在的子组会在 except* 处理结束后重新抛出。

如果处理器需要重新抛出当前匹配子组,可以直接使用 raise

try:
    ...
except* ConnectionError as eg:
    log_network_errors(eg)
    raise

如果在处理器中抛出新的异常,新的异常会和原异常组中未处理的部分合并。这样可以把底层错误转换成领域异常,但必须明确保留因果关系:

class DependencyUnavailable(Exception):
    pass


try:
    ...
except* ConnectionError as eg:
    raise DependencyUnavailable("remote service unavailable") from eg

在生产日志中,建议记录:

  • 异常组的消息;
  • 每个叶子异常的类型与上下文;
  • 任务名称;
  • 任务对应的业务标识;
  • 是否由取消引起;
  • 清理阶段是否又产生了异常。

不要把 ExceptionGroup 直接格式化成一行字符串后丢弃其结构,否则不同任务的故障来源可能无法区分。


十一、TaskGroup 的嵌套与故障边界

结构化并发的价值在嵌套任务时更加明显:

async def load_page():
    async with asyncio.TaskGroup() as outer:
        outer.create_task(load_user_section())
        outer.create_task(load_recommendation_section())


async def load_recommendation_section():
    async with asyncio.TaskGroup() as inner:
        inner.create_task(load_recommendation_a())
        inner.create_task(load_recommendation_b())

这里形成两个作用域:

load_page
└── outer TaskGroup
    ├── load_user_section
    └── load_recommendation_section
        └── inner TaskGroup
            ├── load_recommendation_a
            └── load_recommendation_b

如果 load_recommendation_a() 失败:

  1. 内层组取消 load_recommendation_b()
  2. 内层组等待其清理;
  3. 内层组形成自己的异常组;
  4. 异常传播到外层任务;
  5. 外层组可能因此取消 load_user_section()
  6. 外层组等待自己的任务完成;
  7. 外层组再向上层传播异常。

Python 3.13 改进了嵌套任务组同时发生内部、外部取消时的处理,并正确保留任务的取消计数;任务组还会区分自身用于唤醒 __aexit__() 的内部取消与外部对父任务发出的取消请求。(docs.python.org)

这使得每一层都可以拥有局部故障处理策略:

async def load_recommendation_section():
    try:
        async with asyncio.TaskGroup() as tg:
            tg.create_task(load_recommendation_a())
            tg.create_task(load_recommendation_b())
    except* ConnectionError:
        # 推荐模块失败时返回降级结果,
        # 不让这个子系统的故障升级为整个页面失败
        return []

是否在内层处理异常,取决于该内层模块是否有独立的业务降级语义。如果没有,就应让异常继续向外传播。


十二、TaskGroup 的动态加任务行为

任务组并不要求所有任务都必须在 async with 的第一行创建。只要任务组仍然活跃,组内任务也可以继续添加任务:

async def parent(tg: asyncio.TaskGroup):
    tg.create_task(child("nested-a"))
    await asyncio.sleep(0)
    tg.create_task(child("nested-b"))


async def main():
    async with asyncio.TaskGroup() as tg:
        tg.create_task(parent(tg))

任务组只有在最后一个任务完成并退出上下文后才关闭;一旦任务组正在关闭,不能再添加新任务。若任务组不活跃,Python 3.13 起会关闭传入的协程对象,避免协程对象泄漏;Python 3.14 起,TaskGroup.create_task() 会把关键字参数传递给底层 loop.create_task()。(docs.python.org)

这种动态添加能力适合树形任务结构,例如:

  • 目录遍历;
  • 依赖项展开;
  • 递归抓取;
  • 任务生产者生成子任务。

但动态添加也会使任务数量失控。结构化并发解决的是生命周期和故障传播,不会自动限制并发度。并发上限仍需要通过 SemaphoreQueue 或固定 worker 数量实现。


十三、显式终止 TaskGroup

标准库没有提供一个直接的“正常终止任务组”方法。官方文档给出的方式是:向任务组添加一个专门抛出异常的任务,然后在外部通过 except* 忽略这个控制异常。(docs.python.org)

import asyncio


class StopGroup(Exception):
    pass


async def stop_group():
    raise StopGroup()


async def worker(name: str):
    try:
        while True:
            print(name, "working")
            await asyncio.sleep(0.5)
    finally:
        print(name, "cleanup")


async def main():
    try:
        async with asyncio.TaskGroup() as tg:
            tg.create_task(worker("A"))
            tg.create_task(worker("B"))

            await asyncio.sleep(1.2)
            tg.create_task(stop_group())

    except* StopGroup:
        print("group stopped normally")


asyncio.run(main())

执行路径是:

  1. stop_group() 抛出 StopGroup
  2. TaskGroup 按普通非取消异常处理;
  3. 取消其他 worker;
  4. 等待 worker 的 finally
  5. 形成包含 StopGroup 的异常组;
  6. except* StopGroup 将该控制异常处理掉。

这不是任务组原生的正常返回,而是利用了其故障传播机制实现“带清理的提前结束”。如果业务中频繁需要这种控制流,通常应重新设计为显式的 Event、队列关闭信号或取消协议,使“停止”成为任务逻辑的一部分,而不是伪装成异常。


十四、CancelledError 为什么不能随便捕获

asyncio.CancelledError 直接继承自 BaseException,而不是普通的 Exception。因此:

try:
    await operation()
except Exception:
    ...

不会捕获取消异常。

这通常是有意设计:取消是任务生命周期控制信号,而不是一般业务失败。任务应进行必要清理,然后继续传播取消:

async def worker():
    try:
        await do_work()
    except asyncio.CancelledError:
        await close_partial_state()
        raise

如果写成:

async def worker():
    try:
        await do_work()
    except asyncio.CancelledError:
        await close_partial_state()
        return

父任务可能误以为该 worker 正常完成。对于 TaskGroup 而言,这还可能导致其内部取消协调失效。官方文档明确指出,吞掉 CancelledError 可能使结构化并发组件行为异常。(docs.python.org)

取消和业务异常的关系可以这样区分:

情况 语义
ValueError 任务执行失败
ConnectionError 外部依赖失败
CancelledError 调用者要求任务停止
TimeoutError 某个超时上下文将取消转换成的业务可见错误

例如 asyncio.timeout() 在超时时会取消当前任务,并在上下文管理器外部把取消转换为 TimeoutErrorTimeoutError 应在 async with 外部捕获。(docs.python.org)

async def request():
    try:
        async with asyncio.timeout(2):
            await slow_operation()
    except TimeoutError:
        return "fallback"

如果把 slow_operation() 放进 TaskGroup,超时发生在哪一层,会决定取消和异常组的边界:

async def request():
    try:
        async with asyncio.timeout(2):
            async with asyncio.TaskGroup() as tg:
                tg.create_task(call_a())
                tg.create_task(call_b())
    except TimeoutError:
        return "timeout"

外层超时取消的是当前任务;任务组会参与清理其子任务。最终由超时上下文把当前任务的取消转换成 TimeoutError


十五、Python 3.14 中的 eager_start

Python 3.14 的 asyncio.create_task() 签名包含:

asyncio.create_task(
    coro,
    *,
    name=None,
    context=None,
    eager_start=None,
    **kwargs,
)

eager_start 用于指定任务是否在 create_task() 调用期间立即开始执行;如果没有显式传递,则使用事件循环任务工厂的设置。Python 3.14 的变化之一是 create_task() 接受并传递相关关键字参数;TaskGroup.create_task() 也会把所有关键字参数传给 loop.create_task()。(docs.python.org)

async def immediate():
    print("inside coroutine")
    return 42


async def main():
    task = asyncio.create_task(immediate(), eager_start=True)
    print("after create_task")
    print(await task)

启用 eager 执行后,如果协程在第一次阻塞前就返回或抛出异常,它可能在任务创建期间完成,而不是先进入事件循环后再执行。官方文档强调,这会改变任务执行顺序,因此它不仅是性能选项,也是语义变化。(docs.python.org)

不要根据下面这种输出顺序编写依赖调度细节的业务逻辑:

inside coroutine
after create_task

或者:

after create_task
inside coroutine

两者都可能取决于任务工厂、eager_start 取值和协程是否在首次 await 前完成。正常业务代码应依赖显式的同步原语和任务结果,而不是依赖任务创建后的打印顺序。


十六、诊断任务失败与挂起

任务对象提供了一些诊断方法:

task.done()
task.cancelled()
task.exception()
task.result()
task.get_name()
task.get_stack()

get_stack() 可以查看:

  • 对于尚未完成的任务:任务当前挂起的位置;
  • 对于异常结束的任务:异常 traceback 的栈帧;
  • 对于正常完成或取消的任务:通常为空。(docs.python.org)

可以给任务命名:

task = asyncio.create_task(
    fetch_user(42),
    name="fetch-user-42",
)

在诊断代码中列出当前未完成任务:

import asyncio


def dump_tasks():
    current = asyncio.current_task()

    for task in asyncio.all_tasks():
        if task is current:
            continue

        print(
            task.get_name(),
            "done=", task.done(),
            "cancelling=", task.cancelling(),
        )

        for frame in task.get_stack(limit=1):
            print("  suspended at:", frame.f_code.co_name, frame.f_lineno)

asyncio.current_task() 返回当前任务,asyncio.all_tasks() 返回指定事件循环中尚未完成的任务。(docs.python.org)

诊断时应重点区分三种状态:

  1. 任务未完成且没有栈变化:可能卡在 I/O、锁或队列;
  2. 任务已经完成但异常未读取:可能出现 Task exception was never retrieved
  3. 任务正在 cancelling:说明已经收到取消请求,但清理还没有完成。

十七、一个端到端的结构化并发示例

下面模拟一个订单详情请求,需要并发加载用户、订单和库存;如果任一依赖失败,取消其余任务;对网络错误做统一降级,对数据错误继续抛出。

import asyncio
from dataclasses import dataclass


@dataclass
class OrderPage:
    user: str
    orders: list[str]
    stock: dict[str, int]


async def load_user() -> str:
    await asyncio.sleep(0.2)
    return "alice"


async def load_orders() -> list[str]:
    await asyncio.sleep(0.4)
    return ["order-1001", "order-1002"]


async def load_stock() -> dict[str, int]:
    await asyncio.sleep(0.3)
    return {"item-1": 8, "item-2": 0}


async def load_page() -> OrderPage:
    async with asyncio.timeout(1.0):
        async with asyncio.TaskGroup() as tg:
            user_task = tg.create_task(
                load_user(),
                name="load-user",
            )
            orders_task = tg.create_task(
                load_orders(),
                name="load-orders",
            )
            stock_task = tg.create_task(
                load_stock(),
                name="load-stock",
            )

        return OrderPage(
            user=user_task.result(),
            orders=orders_task.result(),
            stock=stock_task.result(),
        )


async def main():
    try:
        page = await load_page()
    except TimeoutError:
        print("page loading timed out")
    except* ConnectionError as eg:
        print("dependency unavailable:", eg)
    else:
        print(page)


asyncio.run(main())

上面代码不能这样写:

try:
    ...
except TimeoutError:
    ...
except* ConnectionError:
    ...

因为普通 exceptexcept* 不能出现在同一个 try 语句中。应当分层处理:

async def main():
    try:
        await load_page()
    except TimeoutError:
        print("page loading timed out")
    except ExceptionGroup as eg:
        try:
            raise eg
        except* ConnectionError as network_errors:
            print("dependency unavailable:", network_errors)

不过更自然的写法是让内部函数处理异常组,再让外层处理普通超时:

async def load_page_with_network_fallback():
    try:
        async with asyncio.TaskGroup() as tg:
            user_task = tg.create_task(load_user())
            orders_task = tg.create_task(load_orders())
            stock_task = tg.create_task(load_stock())
    except* ConnectionError as eg:
        print("network errors:", eg)
        return None

    return OrderPage(
        user=user_task.result(),
        orders=orders_task.result(),
        stock=stock_task.result(),
    )


async def main():
    try:
        async with asyncio.timeout(1.0):
            page = await load_page_with_network_fallback()
    except TimeoutError:
        print("timed out")
    else:
        print(page)

这里的异常边界更清晰:

timeout scope
└── load_page_with_network_fallback
    └── TaskGroup
        ├── load_user
        ├── load_orders
        └── load_stock
  • TaskGroup 负责子任务的并发、取消、等待和异常聚合;
  • load_page_with_network_fallback() 负责网络异常的业务降级;
  • asyncio.timeout() 负责整个页面操作的时间预算;
  • 最外层负责把超时转化为接口层结果。

十八、结构化并发的准确含义

结构化并发不是“把多个任务放进一个列表”,也不只是“使用 TaskGroup 的风格建议”。它要求并发任务具有可以从代码结构中推导出的生命周期:

  1. 子任务在父作用域内创建;
  2. 父作用域不会在子任务未完成时正常离开;
  3. 子任务失败时,相关兄弟任务会收到取消;
  4. 取消后,父作用域等待所有清理结束;
  5. 子任务异常不会悄悄脱离调用链;
  6. 嵌套作用域分别处理自己的故障,再向外传播。

TaskGroup 直接实现了这些关键行为:它通过异步上下文管理器建立范围,通过退出时等待保证收尾,通过失败时取消兄弟任务,并通过异常组保留并发失败信息。(docs.python.org)

相比之下,裸 create_task() 只提供任务创建能力。它并非错误 API,而是更底层、更自由的工具。当任务确实具有独立生命周期时,使用它是合理的;当任务属于某个请求、批处理、连接或业务操作时,应优先让任务进入明确的 TaskGroup 作用域。

最终可以用下面的判断区分两者:

这个任务是否必须随当前操作一起完成或取消?
        │
        ├── 是:使用 TaskGroup
        │
        └── 否:才考虑独立 create_task,
                 同时自行保存引用、读取异常并设计关闭流程

create_task() 解决“如何启动一个并发任务”;TaskGroup 解决“这一组任务属于谁、何时结束、如何失败”;异常组解决“多个并发失败如何完整表达和分类处理”。三者合起来,才构成 Python asyncio 中可维护的任务生命周期模型。


系列导航与关联阅读

官方资料

本文依据 Python 官方文档、相关 PEP 与生态项目官方文档重新梳理;正文、示例与工程清单由 WR BLOG 编写。