适用场景
一个聚合接口同时读取用户资料与订单摘要,只有两项都成功才返回页面。某个依赖提前失败后,其他查询的结果已经没有用途,却仍占用连接和并发名额。本篇用纯标准库示例讨论这种“共同成功、共同结束”的请求边界。
示例需要 Python 3.11 或以上版本,本文在 Python 3.14 环境执行验证。演示使用内存事件替代数据库和 HTTP,不需要账号、联网或第三方依赖;它是可复现的模拟场景,不代表某次真实生产事故。
现象与排查入口
常见现象是接口已经返回失败,下游查询仍继续执行;压测停止后,活动任务数也没有及时回落。先把一次请求的 trace_id、任务名称、开始时间和结束时间关联起来,区分三种情况:任务还在执行、任务正在清理资源、任务已经结束但连接池指标存在采样延迟。
排查时沿调用链检查:
- 并发工作是否由裸
create_task()启动,启动者是否保存引用并等待其结束? - 是否误以为
gather()抛出异常后,其余工作一定已经停止? - 是否捕获了
CancelledError后继续执行,或者在协程内调用无法让出事件循环的阻塞操作? - 连接、响应流、锁是否通过上下文管理器或
finally释放?
可以临时观察任务名称和状态,但不要持续打印完整调用栈、请求体或局部变量;其中可能包含敏感数据。任务名称也应使用固定业务名称,避免写入用户身份信息。
明确并发边界
默认 gather() 会向等待者传播第一个异常,但不会因此自动取消其余任务。TaskGroup 中一个子任务发生非取消异常时,会取消组内其他任务,并等待它们结束,再以异常组传播失败。asyncio.timeout() 通过取消实现超时;协程应在清理后继续传播取消异常。以上语义可查阅 Python 官方协程与任务文档。
这不等于数据库事务。已经提交的写入或远端已接收的请求不能靠取消撤销。因此下面只聚合读取操作;涉及写入时,仍需业务自己的事务、幂等或补偿设计。
实现:让请求拥有全部子任务
将以下内容保存为 taskgroup_demo.py:
"""演示聚合读取的任务边界、超时和失败传播。"""
import asyncio
import math
from collections.abc import Awaitable, Callable
Reader = Callable[[], Awaitable[str]]
async def build_page(
read_profile: Reader,
read_orders: Reader,
*,
timeout_seconds: float = 2.0,
) -> dict[str, str]:
"""只有全部读取成功才返回结果,并等待请求内子任务结束。"""
if (
isinstance(timeout_seconds, bool)
or not isinstance(timeout_seconds, (int, float))
or not math.isfinite(timeout_seconds)
or timeout_seconds <= 0
):
raise ValueError("超时配置必须是有限正数")
async with asyncio.timeout(timeout_seconds):
async with asyncio.TaskGroup() as group:
profile_task = group.create_task(read_profile(), name="read_profile")
orders_task = group.create_task(read_orders(), name="read_orders")
return {"profile": profile_task.result(), "orders": orders_task.result()}
Reader 是可替换的外部调用边界:生产中注入真正的异步读取函数,测试中注入内存协程。timeout_seconds 是整个聚合操作的预算,不是分别给每个依赖的一份预算。只有成功离开两个上下文,代码才会读取结果。
生产读取函数应在自己的边界管理资源,例如使用数据库驱动提供的异步连接上下文。客户端连接超时、连接池等待超时和响应读取超时也应单独配置,避免仅靠最外层预算掩盖具体瓶颈。
验证:不仅检查报错,还检查子任务结束
将以下内容保存为同一目录下的 test_taskgroup_demo.py。测试使用事件控制任务顺序,不依赖真实网络速度:
"""验证聚合读取的成功、失败、超时与外部取消边界。"""
import asyncio
import unittest
from taskgroup_demo import build_page
class PageTests(unittest.IsolatedAsyncioTestCase):
"""每个测试使用独立事件循环检查任务生命周期。"""
async def test_success(self) -> None:
"""全部依赖成功时返回完整结果。"""
async def read_value() -> str:
return "可用"
result = await build_page(read_value, read_value)
self.assertEqual(result, {"profile": "可用", "orders": "可用"})
async def test_failure_waits_for_cleanup(self) -> None:
"""一个依赖失败后,另一个依赖已完成退出清理。"""
started = asyncio.Event()
cleaned = asyncio.Event()
async def read_slow() -> str:
try:
started.set()
await asyncio.Event().wait()
return "不会返回"
finally:
cleaned.set()
async def read_failed() -> str:
await started.wait()
raise LookupError("订单摘要不存在")
with self.assertRaises(ExceptionGroup) as caught:
await build_page(read_slow, read_failed)
self.assertTrue(cleaned.is_set())
self.assertEqual(len(caught.exception.exceptions), 1)
self.assertIsInstance(caught.exception.exceptions[0], LookupError)
async def test_timeout_leaves_no_running_readers(self) -> None:
"""预算到期后,已启动的读取任务都不再运行。"""
running: set[str] = set()
async def read_waiting() -> str:
task = asyncio.current_task()
assert task is not None
name = task.get_name()
running.add(name)
try:
await asyncio.Event().wait()
return "不会返回"
finally:
running.remove(name)
with self.assertRaises(TimeoutError):
await build_page(read_waiting, read_waiting, timeout_seconds=0.01)
self.assertEqual(running, set())
async def test_external_cancellation_propagates(self) -> None:
"""调用者取消请求时,不把取消伪装成成功结果。"""
started = asyncio.Event()
cleaned = asyncio.Event()
async def read_waiting() -> str:
try:
started.set()
await asyncio.Event().wait()
return "不会返回"
finally:
cleaned.set()
async def read_ready() -> str:
return "可用"
request_task = asyncio.create_task(build_page(read_waiting, read_ready))
await started.wait()
request_task.cancel()
with self.assertRaises(asyncio.CancelledError):
await request_task
self.assertTrue(cleaned.is_set())
async def test_invalid_timeout(self) -> None:
"""无效预算在启动依赖前被拒绝。"""
async def read_unused() -> str:
self.fail("无效配置不应调用依赖")
return "不会返回"
for value in (0.0, -1.0, float("nan"), float("inf"), True):
with self.subTest(value=value):
with self.assertRaises(ValueError):
await build_page(read_unused, read_unused, timeout_seconds=value)
if __name__ == "__main__":
unittest.main()
在这两个文件所在目录执行:
python -m unittest -v test_taskgroup_demo
预期看到五项测试通过。失败路径检查的是“异常已传播且清理已执行”,而不只是收到一个异常。超时测试不断言精确耗时,避免把机器负载差异误报为业务缺陷。IsolatedAsyncioTestCase 的使用方法见 Python 官方 unittest 文档。
示例中的 cleaned 事件只代表清理分支执行过;它不能证明某个数据库驱动真的归还了连接。接入真实依赖后,应在隔离环境补充“取消后连接重新可用、响应流已关闭”等集成验证。
上线方案与注意事项
先选择一个两项读取必须共同成功的聚合接口进行灰度迁移,记录迁移前后的失败率、请求结束后的活跃子任务数量及连接池等待时间。故障注入至少覆盖一个依赖立即失败、一个依赖不返回、调用者主动取消三类情况。
异常组需要在接口边界按已知类型转换成业务响应;不要把所有异常都变成空列表或 HTTP 200。若使用 except* 处理指定异常,只处理真正可恢复的类型,剩余错误继续传播;日志只记录脱敏后的异常类型、任务名称与关联标识。
还需要注意:
- 超时是一种协作式取消。阻塞代码或不响应取消的依赖会拖延结束,不能把两秒预算宣传为绝对两秒返回。
- 清理逻辑也可能失败或再次被取消。真实资源应使用驱动推荐的管理方式,并对清理失败单独告警。
- 可选推荐模块失败后仍允许页面返回时,应在该模块边界实现明确降级,不要让它和必要依赖共用“任何一个失败就整体失败”的策略。
- 一个请求产生上千个子任务时,任务组本身不限制并发;还需设置并发上限或采用有界工作队列。
- 请求生命周期之外的长期工作应交给有持久化、重试和关闭机制的任务系统,不要偷偷从组内再启动无人管理的任务。
总结
处理并发失败时,验收标准应包括:结果是否正确、错误是否可识别、子任务是否结束、资源是否可再次使用。通过明确聚合边界、总预算和清理责任,再用失败与取消测试验证,才能避免接口早已结束、下游工作仍不断堆积的问题。
Discussion
评论