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

使用Workers时asyncio任务失败异常未传播的问题排查

Asyncio队列任务异常处理问题修复

问题重现

你的代码在任务正常执行时可以正常完成,但当任务抛出异常时,程序会无限挂起,且异常无法被上层func()的try-except捕获。核心代码如下:

import asyncio

async def task():
    await asyncio.sleep(0.001)
    raise Exception('task is corrupted')
    print('task completed')

async def worker(queue):
    while not queue.empty():
        t = await queue.get()
        await t
        queue.task_done()
    asyncio.current_task().cancel()

async def func():
    try:
        queue = asyncio.Queue()
        await queue.put(task())
        workers=[asyncio.create_task(worker(queue))]
        await queue.join()
        print('tasks execution completed')
    except Exception as err:
        raise Exception('tasks execution did not complete')

asyncio.run(func())

核心问题分析

  1. queue.task_done()未被执行:当任务抛出异常时,await t直接崩溃,后续的queue.task_done()不会运行,导致队列的未完成任务计数始终不为0,queue.join()会无限等待。
  2. Worker异常无法向上传递:asyncio.create_task()创建的协程任务是独立运行的,其异常不会自动传播到创建它的func()中,除非主动await该任务或获取其结果。
  3. Worker循环逻辑错误:while not queue.empty()的判断不可靠,队列状态是动态变化的,且当队列暂时为空时worker会提前退出,无法处理后续可能加入的任务;手动cancel()自身完全多余,worker协程在循环结束后会自然终止。

修复方案

1. 确保queue.task_done()始终执行

对每个任务单独添加try-finally块,无论任务成功或失败,都标记任务完成:

async def worker(queue):
    while True:
        # 阻塞等待获取任务,直到收到哨兵值或被取消
        task_coro = await queue.get()
        try:
            await task_coro
        except Exception as err:
            # 这里可以记录日志、收集异常,或根据需求决定是否终止worker
            print(f"任务执行失败: {str(err)}")
        finally:
            # 必须确保任务完成标记被执行
            queue.task_done()

2. 监控Worker异常并处理

使用asyncio.gather()同时等待队列完成和所有worker的执行,这样worker的异常会被捕获并向上传播;同时添加哨兵值(如None)让worker可以正常退出:

async def func():
    queue = asyncio.Queue()
    # 添加任务
    await queue.put(task())
    # 添加哨兵值,告诉worker任务已全部加入,处理完后可以退出
    await queue.put(None)

    # 创建worker任务
    workers = [asyncio.create_task(worker(queue))]

    try:
        # 同时等待队列完成和所有worker结束
        await asyncio.gather(queue.join(), *workers)
        print('tasks execution completed')
    except Exception as err:
        # 发生异常时取消所有worker,避免资源泄漏
        for w in workers:
            w.cancel()
        # 等待worker完成取消流程
        await asyncio.gather(*workers, return_exceptions=True)
        raise Exception('tasks execution did not complete') from err

3. 可选:收集所有任务异常

如果需要记录所有任务的异常情况,可以通过共享列表收集:

async def worker(queue, errors):
    while True:
        task_coro = await queue.get()
        if task_coro is None:
            # 收到哨兵值,退出循环
            queue.task_done()
            break
        try:
            await task_coro
        except Exception as err:
            errors.append(err)
        finally:
            queue.task_done()

async def func():
    errors = []
    queue = asyncio.Queue()
    await queue.put(task())
    await queue.put(None)

    workers = [asyncio.create_task(worker(queue, errors))]

    await asyncio.gather(queue.join(), *workers)
    
    if errors:
        raise Exception(f"{len(errors)}个任务执行失败") from errors[0]
    print('tasks execution completed')

最终优化代码

整合所有修复点后的完整代码:

import asyncio

async def task():
    await asyncio.sleep(0.001)
    raise Exception('task is corrupted')
    print('task completed')

async def worker(queue, errors):
    while True:
        task_coro = await queue.get()
        if task_coro is None:
            queue.task_done()
            break
        try:
            await task_coro
        except Exception as err:
            errors.append(err)
        finally:
            queue.task_done()

async def func():
    errors = []
    queue = asyncio.Queue()
    
    # 添加任务
    await queue.put(task())
    # 添加哨兵值
    await queue.put(None)

    workers = [asyncio.create_task(worker(queue, errors))]

    try:
        await asyncio.gather(queue.join(), *workers)
        if errors:
            raise Exception(f"任务执行失败,共{len(errors)}个异常") from errors[0]
        print('tasks execution completed')
    except Exception as err:
        for w in workers:
            w.cancel()
        await asyncio.gather(*workers, return_exceptions=True)
        raise

asyncio.run(func())

运行这段代码,当任务抛出异常时,会被正确捕获,程序不会无限挂起,且会向上抛出包含具体错误的异常。

内容的提问来源于stack exchange,提问作者Cezary Kozioł

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.25 19:28:19