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
相关产品推荐
相关产品推荐

