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:把相关任务绑定到一个生命周期范围内;ExceptionGroup、BaseExceptionGroup与except*:表示并处理多个并发异常;- 结构化并发:让任务的创建、等待、取消和失败传播具有可推导的边界。
一、从协程对象到 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 做了两件事:
- 驱动
work()开始执行; - 当前协程暂停,直到
work()返回结果或抛出异常。
如果希望多个操作并发进行,则需要创建任务。
2. Task 是协程的可管理执行实体
task = asyncio.create_task(work())
可以把这条语句理解为:
- 创建一个
Task; - 将
work()绑定到该任务; - 把任务提交给当前正在运行的事件循环;
- 返回任务对象;
- 后续可以通过
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 秒;并发版本的总等待时间接近:
其中:
- 是任务 A 的总执行时间;
- 是任务 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 的核心故障规则
设一个任务组中有任务:
当某个任务 首次抛出非取消异常时,TaskGroup 会:
- 取消其他尚未完成的任务;
- 禁止继续向该组添加新任务;
- 如果
async with的主体仍在执行,则取消直接包含该async with的父任务; - 等待所有任务完成清理;
- 将任务异常组合成异常组并抛出。
CancelledError 不会触发这个“失败兄弟任务”规则;KeyboardInterrupt 和 SystemExit 则有特殊传播规则。(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。吞掉取消异常会破坏 TaskGroup、asyncio.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)
八、except 与 except* 的区别
普通 except 匹配的是整个异常对象:
try:
raise ExceptionGroup(
"errors",
[ValueError("v"), TypeError("t")],
)
except Exception as exc:
print(type(exc).__name__)
输出:
ExceptionGroup
它不会自动把组内的 ValueError 和 TypeError 拆开。
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* 的匹配过程可以形式化为:
其中:
- 是原始异常组;
- 是匹配异常类型 的子组;
- 是不匹配的剩余子组。
每个 except* T 处理当前剩余组中的匹配部分,未匹配部分继续交给后续处理器;所有处理器执行结束后,未处理异常与处理器中重新抛出的异常会重新合并并传播。(docs.python.org)
except* 不是普通 except 的“多次执行版”
下面两个特性非常重要。
第一,try 中不能同时使用 except 和 except*:
# SyntaxError
try:
...
except ValueError:
...
except* OSError:
...
第二,多个 except* 子句不是“匹配到一个就停止”。它们分别从异常组中提取自己的匹配子组,因此同一个异常组可以同时触发多个处理器。(docs.python.org)
except* 中也不能使用 return、break 或 continue,因为多个异常分支需要完成统一的合并过程。(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() 失败:
- 内层组取消
load_recommendation_b(); - 内层组等待其清理;
- 内层组形成自己的异常组;
- 异常传播到外层任务;
- 外层组可能因此取消
load_user_section(); - 外层组等待自己的任务完成;
- 外层组再向上层传播异常。
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)
这种动态添加能力适合树形任务结构,例如:
- 目录遍历;
- 依赖项展开;
- 递归抓取;
- 任务生产者生成子任务。
但动态添加也会使任务数量失控。结构化并发解决的是生命周期和故障传播,不会自动限制并发度。并发上限仍需要通过 Semaphore、Queue 或固定 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())
执行路径是:
stop_group()抛出StopGroup;TaskGroup按普通非取消异常处理;- 取消其他 worker;
- 等待 worker 的
finally; - 形成包含
StopGroup的异常组; 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() 在超时时会取消当前任务,并在上下文管理器外部把取消转换为 TimeoutError;TimeoutError 应在 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)
诊断时应重点区分三种状态:
- 任务未完成且没有栈变化:可能卡在 I/O、锁或队列;
- 任务已经完成但异常未读取:可能出现
Task exception was never retrieved; - 任务正在 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:
...
因为普通 except 和 except* 不能出现在同一个 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 的风格建议”。它要求并发任务具有可以从代码结构中推导出的生命周期:
- 子任务在父作用域内创建;
- 父作用域不会在子任务未完成时正常离开;
- 子任务失败时,相关兄弟任务会收到取消;
- 取消后,父作用域等待所有清理结束;
- 子任务异常不会悄悄脱离调用链;
- 嵌套作用域分别处理自己的故障,再向外传播。
TaskGroup 直接实现了这些关键行为:它通过异步上下文管理器建立范围,通过退出时等待保证收尾,通过失败时取消兄弟任务,并通过异常组保留并发失败信息。(docs.python.org)
相比之下,裸 create_task() 只提供任务创建能力。它并非错误 API,而是更底层、更自由的工具。当任务确实具有独立生命周期时,使用它是合理的;当任务属于某个请求、批处理、连接或业务操作时,应优先让任务进入明确的 TaskGroup 作用域。
最终可以用下面的判断区分两者:
这个任务是否必须随当前操作一起完成或取消?
│
├── 是:使用 TaskGroup
│
└── 否:才考虑独立 create_task,
同时自行保存引用、读取异常并设计关闭流程
create_task() 解决“如何启动一个并发任务”;TaskGroup 解决“这一组任务属于谁、何时结束、如何失败”;异常组解决“多个并发失败如何完整表达和分类处理”。三者合起来,才构成 Python asyncio 中可维护的任务生命周期模型。
系列导航与关联阅读
- 系列入口:Python 完整学习路线:从语言模型、并发到 Web、数据、AI 与生产交付
- 上一篇:Python asyncio 完整基础:事件循环、协程、Future 与调度
- 下一篇:Python asyncio 网络流:TCP、Reader、Writer、TLS 与背压
- 延伸:Python asyncio 取消与同步:Timeout、Lock、Queue 和清理
官方资料
本文依据 Python 官方文档、相关 PEP 与生态项目官方文档重新梳理;正文、示例与工程清单由 WR BLOG 编写。

评论
0 条讨论