使用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())
核心问题分析
queue.task_done()未被执行:当任务抛出异常时,await t直接崩溃,后续的queue.task_done()不会运行,导致队列的未完成任务计数始终不为0,queue.join()会无限等待。- Worker异常无法向上传递:
asyncio.create_task()创建的协程任务是独立运行的,其异常不会自动传播到创建它的func()中,除非主动await该任务或获取其结果。 - 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ł
相关产品推荐
相关产品推荐

