Asyncio结合ProcessPoolExecutor提前中断时如何优雅关闭全部任务和进程
问题根源
- 子进程默认会继承主进程的信号处理逻辑,按下Ctrl+C时主进程和所有子进程同时收到KeyboardInterrupt,主进程优先开始销毁运行资源,还在初始化阶段的子进程无法完成模块导入,抛出
init_import_size致命错误 - 你没有显式管理
ProcessPoolExecutor的生命周期,退出时进程池内的排队任务、运行中任务没有被正确清理,导致进程池管理线程尝试给已经被asyncio取消的Future设置异常,抛出InvalidStateError asyncio.gather默认收到异常后不会自动取消其他未完成的任务,大量悬停任务会持有进程池引用,延缓资源释放
修复方案
- 给进程池增加子进程初始化逻辑,让子进程忽略SIGINT信号,统一由主进程控制退出流程
- 显式管理任务生命周期,收到中断信号时先取消所有未完成的async任务
- 显式调用进程池的
shutdown方法,取消所有排队中的任务,等待运行中任务完成后再退出
完整修复代码
import asyncio import signal from concurrent.futures import ProcessPoolExecutor class TestClass: def __init__(self) -> None: self.value1 = 1 self.value2 = 2 # 子进程初始化函数:忽略SIGINT信号,由主进程统一控制退出 def init_worker(): signal.signal(signal.SIGINT, signal.SIG_IGN) async def task(executor_processes, i): print(f"[TASK {i}] Initializing Abck class") # 自动获取当前线程的事件循环,无需手动传递 new_test = await asyncio.get_running_loop().run_in_executor(executor_processes, TestClass) # 此处可添加其他TestClass的同步/异步逻辑 print(f"[TASK {i}] Finished") async def main(): # 初始化进程池时指定子进程初始化函数 executor_processes = ProcessPoolExecutor(max_workers=5, initializer=init_worker) tasks = [] try: for i in range(1, 100): tasks.append(asyncio.create_task(task(executor_processes, i))) await asyncio.gather(*tasks) finally: # 先取消所有未完成的async任务 for t in tasks: if not t.done(): t.cancel() try: await t except asyncio.CancelledError: pass # 关闭进程池:cancel_futures取消所有排队未运行的任务,wait等待运行中的任务完成 executor_processes.shutdown(wait=True, cancel_futures=True) print("所有资源已释放") if __name__ == '__main__': try: asyncio.run(main()) except KeyboardInterrupt: print("ctrl + c 触发退出") finally: print('Program finished')
注意事项
ProcessPoolExecutor.shutdown的cancel_futures参数是Python3.9新增特性,如果需要兼容Python3.8,可以删除该参数,或者手动调用executor._processes.terminate()强制终止所有子进程(会中断运行中的任务)- 默认配置下按下Ctrl+C后,会等待当前正在运行的任务执行完成后再退出,如果不需要等待可将
shutdown的wait参数改为False,会直接终止所有子进程立即退出 - 如果你的长时间运行任务支持中断,也可以在子进程的任务逻辑里增加退出信号检测,配合主进程的终止逻辑实现更平滑的退出
内容的提问来源于stack exchange,提问作者hen
相关产品推荐
相关产品推荐

