Python multiprocessing.Pool遇异常时如何重启单个出错进程
问题场景
使用Python multiprocessing模块开发时存在进程池任务异常自动重启需求,初始演示代码如下:
import signal import asyncio import os import random import time import multiprocessing my_list = [] for i in range(0,10): n = random.randint(1,100) my_list.append(n) async def loop_item(my_item): while True: a = random.randint(1, 2) if a == 2: print(f"process id: {os.getpid()}") raise Exception('Error') print(f"process id: {os.getpid()} - {my_item}") time.sleep(0.5) def run_loop(my_item): asyncio.run(loop_item(my_item)) def throw_error(e): os.system('bash /root/my-script.sh') #that launchs "python my-script.py" os.killpg(os.getpgid(os.getpid()), signal.SIGKILL) if __name__ == '__main__': pool = multiprocessing.Pool(processes=10) for my_item in my_list: pool.apply_async(run_loop, (my_item,), error_callback=throw_error) pool.close() pool.join()
上述代码逻辑:
- 初始化生成包含10个随机数的列表
my_list - 创建容量为10的进程池,为每个列表项分配独立进程循环打印进程ID与对应数值
- 通过随机抛出
Exception模拟业务运行时的各类异常
核心需求:单个进程运行触发异常退出时,自动在新进程中重启对应loop_item(my_item)任务,无需重启整个脚本。
现存问题与已尝试方案
- 参数传递问题:最初考虑通过Redis这类外部存储实现异常场景下的
my_item参数传递,需要更轻量的无外部依赖实现方案 - 现有重启方案效率低:当前通过
throw_error回调执行shell脚本杀掉主进程、重启整个Python脚本的方式实现故障恢复,资源开销大、恢复效率低 - 僵尸进程问题:曾尝试在
throw_error回调中新建容量为1的进程池重新提交出错任务,实现代码如下:
def throw_error(e): pool2 = multiprocessing.Pool(processes=1) pool2.apply_async(run_loop, (my_item,), error_callback=throw_error) pool2.close() pool2.join()
该方案在多次异常触发后会出现进程池失控问题,累计产生成百上千的僵尸进程,无法用于生产环境。
最终需要实现生产级可靠逻辑:进程池单个任务异常退出时仅重启对应出错任务,同时规避僵尸进程问题。
生产级实现方案
不推荐基于内置Pool的error_callback实现重启逻辑:该回调执行时无法稳定绑定出错任务的对应参数,回调中重复提交任务容易触发进程池内部状态异常,反复创建新进程池更是会直接导致进程资源泄漏,不适合常驻生产场景。
推荐采用固定常驻进程+进程安全任务队列的实现模式,核心逻辑是启动固定数量的常驻工作进程,持续从队列消费任务,单个任务异常崩溃时将任务重新放回队列等待重新执行,全程复用进程资源,从根源避免僵尸进程。
完整实现代码:
import signal import asyncio import os import random import time import multiprocessing async def loop_item(my_item): while True: a = random.randint(1, 2) if a == 2: print(f"process id: {os.getpid()}") raise Exception('Error') print(f"process id: {os.getpid()} - {my_item}") time.sleep(0.5) def run_loop(my_item): # 子进程忽略终端中断信号,避免被主进程信号意外终止 signal.signal(signal.SIGINT, signal.SIG_IGN) asyncio.run(loop_item(my_item)) def task_worker(task_queue): # 常驻工作进程,循环消费队列任务 while True: my_item = task_queue.get() try: run_loop(my_item) except Exception as e: print(f"Task for item {my_item} crashed: {str(e)}, scheduled for restart") # 异常任务重新入队,等待调度重启 task_queue.put(my_item) # 加入短暂休眠避免异常循环打满系统资源 time.sleep(0.1) if __name__ == '__main__': # 初始化任务列表 my_list = [random.randint(1,100) for _ in range(10)] # 初始化进程间安全的任务队列 task_queue = multiprocessing.Queue() # 初始任务全部入队 for item in my_list: task_queue.put(item) # 启动固定数量的常驻工作进程 worker_count = 10 workers = [] for _ in range(worker_count): p = multiprocessing.Process(target=task_worker, args=(task_queue,)) # 标记为守护进程,主进程退出时自动回收子进程,避免孤儿/僵尸进程 p.daemon = True p.start() workers.append(p) # 主进程阻塞等待,捕获退出信号做优雅关闭 try: for p in workers: p.join() except KeyboardInterrupt: print("Received exit signal, shutting down...") os.killpg(os.getpgid(os.getpid()), signal.SIGTERM)
方案核心优势:
- 无外部依赖:不需要Redis等第三方存储,所有任务参数通过
multiprocessing.Queue传递,原生支持进程间安全通信 - 重启效率高:单个任务崩溃后仅需重新入队即可被空闲工作进程拉起,不需要重启主进程或整个脚本,恢复耗时在毫秒级
- 无进程泄漏风险:全程运行固定数量的工作进程,不存在反复创建销毁进程池的逻辑,守护进程标记保证主进程退出时所有子进程资源被自动回收
- 扩展性强:后续新增任务仅需向队列提交元素即可,不需要调整进程配置;如需实现异常重试次数限制、指数退避、异常告警等逻辑,仅需在
task_worker函数中增加对应判断即可,可控性远高于内置进程池的回调模式。
内容的提问来源于stack exchange,提问作者liuru
相关产品推荐
相关产品推荐

