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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.28 02:24:20