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

如何用Python asyncio高效处理嵌套异步操作及任务取消与错误处理

Python Asyncio 异步任务错误处理与任务取消实现

核心问题分析

原代码采用串行执行异步任务,一个任务完成后才启动下一个,因此出错时没有待处理任务需要取消;同时错误仅在main_task内部捕获,未传递到主事件循环,也没有实现任务失败后的自动重启逻辑。

解决方案实现

1. 并发任务重构

将串行执行的任务改为并发启动,使用asyncio.create_task创建任务并存储到列表,确保多个任务同时运行,出错时存在未完成任务可取消。

2. 错误触发时的任务取消

当捕获到异常时,遍历所有已创建的任务:

  • 通过task.done()判断任务是否已完成
  • 对未完成任务调用task.cancel()触发取消
  • 等待任务取消完成(await task),避免事件循环抛出未处理的CancelledError警告

3. 错误向上传递到主事件循环

在main_task的异常处理块中重新抛出异常,让asyncio.run将错误传递到主循环的try-except块,实现错误的层级反馈。

4. 任务自动重复执行

主循环添加while True循环,任务完成或失败后等待指定时间,自动重启任务,实现服务的恢复能力。

完整实现代码

import asyncio

async def fetch_data_from_service(service_name):
    await asyncio.sleep(1)  # 模拟IO操作
    if service_name == "ServiceB":
        raise Exception(f"Error fetching data from {service_name}")
    return f"Data from {service_name}"

async def process_data(data):
    await asyncio.sleep(1)  # 模拟数据处理
    if data == "Data from ServiceC":
        raise Exception("Error processing data from ServiceC")
    return f"Processed {data}"

async def main_task():
    # 并发创建所有数据获取任务
    fetch_tasks = [
        asyncio.create_task(fetch_data_from_service("ServiceA")),
        asyncio.create_task(fetch_data_from_service("ServiceB")),
        asyncio.create_task(fetch_data_from_service("ServiceC")),
    ]
    process_tasks = []
    
    try:
        # 等待所有数据获取完成,任一任务出错立即终止
        dataA, dataB, dataC = await asyncio.gather(*fetch_tasks)
        
        # 并发创建数据处理任务
        process_tasks = [
            asyncio.create_task(process_data(dataA)),
            asyncio.create_task(process_data(dataB)),
            asyncio.create_task(process_data(dataC)),
        ]
        processedA, processedB, processedC = await asyncio.gather(*process_tasks)
        
        print(processedA, processedB, processedC)
    except Exception as e:
        print(f"Exception caught in main_task: {e}")
        # 取消所有未完成的任务
        for task in fetch_tasks + process_tasks:
            if not task.done():
                task.cancel()
                # 等待任务取消完成,消除未处理异常警告
                try:
                    await task
                except asyncio.CancelledError:
                    pass
        # 重新抛出异常,传递到主事件循环
        raise

# 主事件循环,支持任务自动重启
if __name__ == "__main__":
    while True:
        try:
            asyncio.run(main_task())
            print("Task completed successfully. Restarting in 5 seconds...")
        except Exception as e:
            print(f"Event loop caught exception: {e}. Restarting in 5 seconds...")
        # 等待后重启任务
        asyncio.run(asyncio.sleep(5))

关键细节说明

  • asyncio.gather(*tasks):等待所有任务完成,若任一任务抛出异常,gather会立即抛出该异常,同时取消其他未完成任务(默认行为),结合手动取消逻辑,确保所有任务都被妥善终止。
  • 任务取消后的await task:必须等待取消操作完成,否则事件循环会检测到未处理的CancelledError并发出警告。
  • 自动重启机制:通过while True循环和asyncio.sleep实现任务失败后的延迟重启,避免频繁重试。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.21 10:45:03