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

Multiprocessing Pool子进程异常时如何终止全脚本所有进程

问题原因

os._exit(1)仅作用于当前抛出异常的子进程,主进程无法感知子进程的异常状态,因此会继续等待其他子进程执行完成,无法实现全局终止。

实现方案(改动最小,兼容现有逻辑)

通过跨进程事件mp.Event做全局终止信号,主进程轮询信号状态,一旦触发就终止所有进程。

1. 修改 helper.py

import functools
import traceback
import os
import multiprocessing as mp

# 子进程全局变量,存储终止事件
_terminate_event = None

def init_terminate_event(event):
    global _terminate_event
    _terminate_event = event

def trace_unhandled_exceptions(func):
    @functools.wraps(func)
    def wrapped_func(*args, **kwargs):
        # 已触发终止信号的子进程直接跳过执行
        if _terminate_event is not None and _terminate_event.is_set():
            return
        try:
            return func(*args, **kwargs)
        except:
            print('Exception in '+func.__name__)
            traceback.print_exc()
            # 触发全局终止信号
            if _terminate_event is not None:
                _terminate_event.set()
            os._exit(1)
    return wrapped_func

2. 修改 scraper.py

import multiprocessing as mp
import time
from helper import trace_unhandled_exceptions, init_terminate_event

start_block = 100
end_block = 50000

@trace_unhandled_exceptions
def main(block_num):
    block = blah_blah(block_num)
    return block

if __name__ == "__main__":
    cpus = min(8, mp.cpu_count()-1 or 1)
    # 初始化跨进程终止事件
    terminate_event = mp.Event()
    # 进程池初始化时注入事件到所有子进程
    pool = mp.Pool(cpus, initializer=init_terminate_event, initargs=(terminate_event,))
    result = pool.map_async(main, range(start_block - 20, end_block), chunksize=cpus)
    pool.close()
    # 轮询任务状态和终止信号
    while not result.ready():
        if terminate_event.is_set():
            # 终止所有子进程
            pool.terminate()
            break
        time.sleep(0.5)
    pool.join()
    # 异常场景下主进程也返回错误状态码
    if terminate_event.is_set():
        exit(1)

简化方案(不需要自定义异常打印时可用)

如果不需要保留现有装饰器的异常打印逻辑,可直接让子进程异常向上抛出,主进程捕获后终止:

# scraper.py 主进程部分修改
if __name__ == "__main__":
    cpus = min(8, mp.cpu_count()-1 or 1)
    pool = mp.Pool(cpus)
    result = pool.map_async(main, range(start_block - 20, end_block), chunksize=cpus)
    pool.close()
    try:
        # 等待任务执行完成,有异常会直接抛出
        result.get()
    except Exception as e:
        print(f"捕获到子进程异常:{e}")
        pool.terminate()
        exit(1)
    pool.join()

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.03 00:36:05