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

Python多进程:子进程异常捕获、进度交互及主进程终止实现

Python多进程实现任务调度、进度反馈与异常处理

你提出的三个需求可以通过Python的multiprocessing模块完美实现,不管是直接使用Process类(适合单任务或精细控制场景)还是Pool进程池(适合批量任务调度)都能满足,以下是两种方案的具体实现:

一、使用multiprocessing.Process实现

这种方式灵活性更高,能直接控制单个子进程的生命周期和通信逻辑。

核心逻辑

  • 借助multiprocessing.Queue实现主、子进程间的进度和异常信息传递;
  • 子进程执行任务时,定期向队列发送进度数据;
  • 子进程捕获到异常时,将异常信息写入队列,主进程监听队列时触发自身终止逻辑;
  • 主进程负责启动子进程、持续监听队列,一旦检测到异常立即退出。

代码示例

import multiprocessing
import time
import sys

def task_worker(queue):
    try:
        total_steps = 5
        for step in range(1, total_steps + 1):
            # 模拟任务执行耗时
            time.sleep(1)
            # 向主进程发送进度信息
            queue.put({"type": "progress", "data": f"完成第 {step}/{total_steps} 步"})
            # 模拟异常场景:第3步触发错误
            if step == 3:
                raise RuntimeError("子进程执行出错:第3步数据校验失败")
        # 任务完成后发送结束信号
        queue.put({"type": "finish", "data": "任务执行完毕"})
    except Exception as e:
        # 将异常信息传递给主进程
        queue.put({"type": "error", "data": str(e)})

if __name__ == "__main__":
    # 创建跨进程通信队列
    progress_queue = multiprocessing.Queue()
    
    # 启动子进程
    worker_process = multiprocessing.Process(target=task_worker, args=(progress_queue,))
    worker_process.start()
    
    # 主进程监听队列
    while True:
        # 子进程已结束且队列无消息时退出循环
        if not worker_process.is_alive() and progress_queue.empty():
            break
        try:
            # 尝试从队列获取消息,超时0.5秒避免阻塞
            msg = progress_queue.get(timeout=0.5)
            if msg["type"] == "progress":
                print(f"主进程收到进度:{msg['data']}")
            elif msg["type"] == "finish":
                print(f"主进程收到通知:{msg['data']}")
                break
            elif msg["type"] == "error":
                print(f"主进程捕获子进程异常:{msg['data']}")
                print("主进程终止自身")
                worker_process.terminate()
                sys.exit(1)
        except multiprocessing.queues.Empty:
            continue
    
    # 等待子进程正常结束
    worker_process.join()

二、使用multiprocessing.Pool实现

进程池适合批量调度多个任务,通过共享队列和线程监听可以实现进度反馈和异常处理。

核心逻辑

  • 使用multiprocessing.Manager().Queue()创建可在进程池内共享的队列(普通Queue无法在进程池的子进程间共享);
  • 主进程单独启动一个线程监听队列,实时获取进度或异常信息;
  • 子进程执行任务时发送进度到队列,若抛出异常则将异常信息写入队列,主进程检测到异常后终止进程池并退出。

代码示例

import multiprocessing
import time
import sys
from threading import Thread

def task_worker(queue, task_id):
    try:
        total_steps = 4
        for step in range(1, total_steps + 1):
            time.sleep(1)
            # 发送当前任务的进度信息
            queue.put({"type": "progress", "task_id": task_id, "data": f"完成第 {step}/{total_steps} 步"})
            # 模拟异常场景:第2步触发错误
            if step == 2:
                raise ValueError(f"任务{task_id}执行出错:第2步参数非法")
        # 返回任务成功结果
        return {"task_id": task_id, "status": "success", "data": "任务完成"}
    except Exception as e:
        # 发送异常信息到队列
        queue.put({"type": "error", "task_id": task_id, "data": str(e)})
        return {"task_id": task_id, "status": "failed", "data": str(e)}

def monitor_queue(queue, pool):
    while True:
        try:
            msg = queue.get(timeout=0.5)
            if msg["type"] == "progress":
                print(f"主进程收到任务{msg['task_id']}进度:{msg['data']}")
            elif msg["type"] == "error":
                print(f"主进程捕获任务{msg['task_id']}异常:{msg['data']}")
                print("主进程终止进程池并退出")
                pool.terminate()
                sys.exit(1)
        except multiprocessing.queues.Empty:
            # 进程池已关闭且无任务时,结束监听
            if pool._state == multiprocessing.pool.CLOSE:
                break
            continue

if __name__ == "__main__":
    # 使用Manager创建跨进程池的共享队列
    with multiprocessing.Manager() as manager:
        progress_queue = manager.Queue()
        
        # 创建进程池,指定进程数
        with multiprocessing.Pool(processes=2) as pool:
            # 启动监听队列的守护线程
            monitor_thread = Thread(target=monitor_queue, args=(progress_queue, pool))
            monitor_thread.daemon = True
            monitor_thread.start()
            
            # 提交2个任务到进程池
            tasks = [pool.apply_async(task_worker, args=(progress_queue, i)) for i in range(2)]
            
            # 等待所有任务完成(若中途出现异常会被监听线程捕获并强制退出)
            for task in tasks:
                result = task.get()
                if result["status"] == "success":
                    print(f"任务{result['task_id']}执行结果:{result['data']}")
            
            # 关闭进程池,不再接受新任务
            pool.close()
            pool.join()
    
    print("所有任务执行完成")

方案对比

  • multiprocessing.Process:适合单任务场景,对子进程的控制更直接,无需额外线程监听;
  • multiprocessing.Pool:适合批量任务调度,能复用进程资源,但需要借助线程来监听进度和异常。

两种方案都能完全满足你提出的三个需求,可根据实际任务规模选择合适的实现方式。

内容的提问来源于stack exchange,提问作者Mike Wang

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.01 23:20:36