Asyncio as_completed方法SIGINT信号异常处理问题求助
问题分析与解决方案
你遇到的问题核心在于原生signal模块和asyncio事件循环的协同冲突,以及直接sys.exit()导致的任务取消异常未被捕获。当你在as_completed准备阶段(大量任务初始化/排队时)发送SIGINT,原生信号处理直接终止进程,会让asyncio来不及优雅处理所有待处理任务的取消流程,从而抛出大量未捕获的CancelledError。
下面是具体的修改方案,我们一步步解决这个问题:
关键修改点
- 替换原生信号处理为asyncio事件循环的专属处理:asyncio提供了和自身事件循环协同的信号注册方法,避免粗暴打断循环。
- 优雅停止事件循环而非直接退出:通过
loop.stop()让事件循环正常收尾,而不是sys.exit(),这样asyncio能妥善处理所有任务的取消逻辑。 - 捕获任务取消异常并清理资源:在
main函数中捕获CancelledError,关闭进度条等资源,避免输出混乱。
修改后的完整代码
import asyncio import tqdm import datetime import sys import signal import random async def heavy_load(i): await asyncio.sleep(random.random()) return None async def main(): length = int(sys.argv[1]) inputs = list(range(length)) pbar = tqdm.tqdm( total=len(inputs), position=0, leave=True, bar_format='#PROGRESS {desc}: {percentage:.3f}%|{bar}| {n_fmt}/{total_fmt} [{elapsed}<{remaining}]' ) tasks = [heavy_load(i) for i in inputs] try: for future in asyncio.as_completed(tasks): _ = await future pbar.update(1) pbar.refresh() except asyncio.CancelledError: # 捕获任务取消异常,清理进度条 tqdm.tqdm.write('#INFO '+datetime.datetime.now().strftime('%Y-%m-%d %H:%M:%S')+' 正在终止所有任务...') pbar.close() # 主动取消所有未完成的任务,并允许异常返回(避免未处理报错) for task in tasks: if not task.done(): task.cancel() await asyncio.gather(*tasks, return_exceptions=True) raise # 重新抛出异常,让事件循环正常结束 def sigint_handler(): loop = asyncio.get_running_loop() tqdm.tqdm.write('#INFO '+datetime.datetime.now().strftime('%Y-%m-%d %H:%M:%S')+' user aborted execution!') # 停止事件循环,而非直接终止进程 loop.stop() if __name__=='__main__': loop = asyncio.new_event_loop() asyncio.set_event_loop(loop) # 用asyncio事件循环注册信号处理,替代原生signal模块 loop.add_signal_handler(signal.SIGINT, sigint_handler) try: loop.run_until_complete(main()) finally: loop.close()
代码细节解释
- 信号处理优化:改用
loop.add_signal_handler注册SIGINT处理函数,函数内部调用loop.stop()让事件循环平稳终止,避免了原生信号处理直接杀进程的粗暴行为。 - 异常捕获与清理:在
main中捕获CancelledError后,先关闭进度条避免输出混乱,再主动取消所有未完成任务并等待它们结束(return_exceptions=True会让gather把异常作为返回值,而不是抛出),最后重新抛出异常让事件循环正常收尾。 - 事件循环管理:手动创建和管理事件循环,比
asyncio.run()更灵活,方便提前注册信号处理逻辑。
额外优化建议
如果length真的很大(比如10000),一次性创建10000个任务会占用较多内存,还会导致as_completed准备阶段过长。你可以用**信号量(Semaphore)**限制并发数,让任务分批执行:
async def main(): length = int(sys.argv[1]) inputs = list(range(length)) pbar = tqdm.tqdm(...) # 限制同时运行100个任务,可根据机器性能调整 semaphore = asyncio.Semaphore(100) async def bounded_heavy_load(i): async with semaphore: return await heavy_load(i) tasks = [bounded_heavy_load(i) for i in inputs] # 后续逻辑和之前一致...
这样既减少了内存占用,也能让进度条更快开始更新,避免长时间卡在0%的情况。
内容的提问来源于stack exchange,提问作者glezo
相关产品推荐
相关产品推荐

