适用场景

一个聚合接口同时读取用户资料与订单摘要,只有两项都成功才返回页面。某个依赖提前失败后,其他查询的结果已经没有用途,却仍占用连接和并发名额。本篇用纯标准库示例讨论这种“共同成功、共同结束”的请求边界。

示例需要 Python 3.11 或以上版本,本文在 Python 3.14 环境执行验证。演示使用内存事件替代数据库和 HTTP,不需要账号、联网或第三方依赖;它是可复现的模拟场景,不代表某次真实生产事故。

现象与排查入口

常见现象是接口已经返回失败,下游查询仍继续执行;压测停止后,活动任务数也没有及时回落。先把一次请求的 trace_id、任务名称、开始时间和结束时间关联起来,区分三种情况:任务还在执行、任务正在清理资源、任务已经结束但连接池指标存在采样延迟。

排查时沿调用链检查:

  1. 并发工作是否由裸 create_task() 启动,启动者是否保存引用并等待其结束?
  2. 是否误以为 gather() 抛出异常后,其余工作一定已经停止?
  3. 是否捕获了 CancelledError 后继续执行,或者在协程内调用无法让出事件循环的阻塞操作?
  4. 连接、响应流、锁是否通过上下文管理器或 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* 处理指定异常,只处理真正可恢复的类型,剩余错误继续传播;日志只记录脱敏后的异常类型、任务名称与关联标识。

还需要注意:

  • 超时是一种协作式取消。阻塞代码或不响应取消的依赖会拖延结束,不能把两秒预算宣传为绝对两秒返回。
  • 清理逻辑也可能失败或再次被取消。真实资源应使用驱动推荐的管理方式,并对清理失败单独告警。
  • 可选推荐模块失败后仍允许页面返回时,应在该模块边界实现明确降级,不要让它和必要依赖共用“任何一个失败就整体失败”的策略。
  • 一个请求产生上千个子任务时,任务组本身不限制并发;还需设置并发上限或采用有界工作队列。
  • 请求生命周期之外的长期工作应交给有持久化、重试和关闭机制的任务系统,不要偷偷从组内再启动无人管理的任务。

总结

处理并发失败时,验收标准应包括:结果是否正确、错误是否可识别、子任务是否结束、资源是否可再次使用。通过明确聚合边界、总预算和清理责任,再用失败与取消测试验证,才能避免接口早已结束、下游工作仍不断堆积的问题。