You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

等待首个asyncio任务完成后清理:吞掉CancelledError是否安全?

问题分析与解决方案

原写法的安全性问题

你的写法确实存在不安全的地方:当捕获CancelledError后直接continue而不调用uncancel(),会导致执行finally块的父任务取消计数未被重置。asyncio的取消机制用计数方式跟踪取消请求,若不调用uncancel(),父任务的取消状态会残留,后续的await操作可能意外抛出CancelledError,破坏正常执行流程。

修正后的代码实现

我们需要在捕获CancelledError后调用asyncio.current_task().uncancel()重置取消计数,同时修正原代码中的语法错误(如asyncio.await应为asyncio.wait、列表推导式括号错误):

import asyncio
import logging

logger = logging.getLogger(__name__)

async def fetch_first(data_sources):
    all_tasks = [asyncio.create_task(ds.get()) for ds in data_sources]
    try:
        done, pending = await asyncio.wait(
            all_tasks,
            return_when=asyncio.FIRST_COMPLETED
        )
        # 处理已完成的任务
        for task in done:
            if ex := task.exception():
                raise ex
            response = await task
            return response
    finally:
        # 处理待完成任务:取消并等待结束
        for task in pending:
            if task.done():
                continue
            task.cancel()
            try:
                await task
            except asyncio.CancelledError:
                # 重置当前任务的取消计数,避免状态残留
                asyncio.current_task().uncancel()
            except Exception as ex:
                logger.warn("task failed", exc_info=ex)
            else:
                # 检查任务是否有未抛出的异常
                if ex := task.exception():
                    logger.warn("task failed", exc_info=ex)
        # 检查已完成任务中未处理的异常(比如多个任务同时完成的情况)
        for task in done:
            if ex := task.exception():
                logger.warn("task failed", exc_info=ex)

关于TaskGroup的替代方案

你提到的asyncio.TaskGroup并非不能满足需求,只需调整使用方式:TaskGroup会自动管理任务生命周期,我们可以在TaskGroup内创建所有任务,用asyncio.wait等待第一个任务完成后,手动取消其他任务。TaskGroup会等待所有任务(包括被取消的)完成,但不会将子任务的取消传播给父任务,除非主动取消TaskGroup本身:

import asyncio
import logging

logger = logging.getLogger(__name__)

async def fetch_first_with_taskgroup(data_sources):
    async with asyncio.TaskGroup() as tg:
        all_tasks = [tg.create_task(ds.get()) for ds in data_sources]
        done, pending = await asyncio.wait(
            all_tasks,
            return_when=asyncio.FIRST_COMPLETED
        )
        # 处理已完成任务
        for task in done:
            if ex := task.exception():
                raise ex
            response = await task
            # 取消所有待完成任务
            for task in pending:
                task.cancel()
            # TaskGroup会自动等待所有任务结束,无需手动处理
            return response
        logger.error("No tasks completed")
        return None

这种方式更简洁,TaskGroup自动处理任务收尾工作,避免了手动编写finally块的大部分逻辑,同时从根源上避免了取消状态残留问题。

内容的提问来源于stack exchange,提问作者Charles Duffy

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.12 23:29:55