ProcessPoolExecutor中EXTRA_QUEUED_CALLS导致SIGINT无法捕获的问题
解决ProcessPoolExecutor捕获SIGINT时额外启动进程的问题
使用concurrent.futures.ProcessPoolExecutor实现多进程任务时,触发SIGINT(Ctrl+C)后调用executor.shutdown(wait=True, cancel_futures=True),本应终止所有未启动的进程,但由于底层EXTRA_QUEUED_CALLS = 1的默认设置,队列中会提前预加载一个任务,导致shutdown后仍会启动一个额外进程,且该进程无法捕获SIGINT信号。
问题复现代码
import signal import time from concurrent.futures import ProcessPoolExecutor def thread_worker(tenant): print("Working on tenant: ", tenant) def signal_handler(sig: signal, frame: any) -> None: print("Interrupt received in sub process") signal.signal( signal.SIGINT, lambda signum, frame: signal_handler(signum, frame), ) time.sleep(2) def multi_process(): with ProcessPoolExecutor(2) as process_executor: def signal_handler(process_executor, sig: signal, frame: any) -> None: print("Interrupt received") process_executor.shutdown(wait=True, cancel_futures=True) print("Killed all processes") signal.signal( signal.SIGINT, lambda signum, frame: signal_handler(process_executor, signum, frame), ) process_futures = [ process_executor.submit(thread_worker, tenant) for tenant in range(5) ] if __name__ == "__main__": multi_process()
执行输出
Working on tenant: 0 Working on tenant: 1 ^CInterrupt received in sub process Interrupt received in sub process Interrupt received Working on tenant: 2 Killed all processes
解决方案
方案1:自定义ProcessPoolExecutor关闭预加载机制
ProcessPoolExecutor内部的EXTRA_QUEUED_CALLS参数默认值为1,用于预加载一个任务到队列以避免worker进程空闲。我们可以自定义子类将该值改为0,彻底关闭预加载:
import signal import time from concurrent.futures import ProcessPoolExecutor class NoPreloadProcessPoolExecutor(ProcessPoolExecutor): def __init__(self, max_workers=None): super().__init__(max_workers) # 修改预加载任务数为0,禁止提前加载任务到队列 self._executor._extra_queued_calls = 0 def thread_worker(tenant): print("Working on tenant: ", tenant) def signal_handler(sig, frame): print("Interrupt received in sub process") exit(0) # 收到信号后直接退出,避免继续执行任务 signal.signal(signal.SIGINT, signal_handler) time.sleep(2) def multi_process(): with NoPreloadProcessPoolExecutor(2) as process_executor: process_futures = [ process_executor.submit(thread_worker, tenant) for tenant in range(5) ] def signal_handler(sig, frame): print("Interrupt received") # 先取消所有未完成的任务 for future in process_futures: if not future.done(): future.cancel() process_executor.shutdown(wait=True, cancel_futures=True) print("Killed all processes") signal.signal(signal.SIGINT, signal_handler) if __name__ == "__main__": multi_process()
方案2:主进程捕获信号时主动终止所有子进程
如果不想修改executor内部参数,可以在信号处理函数中直接终止所有子进程,确保即使有预加载任务启动,也会被立即终止:
import signal import time import os import psutil from concurrent.futures import ProcessPoolExecutor def thread_worker(tenant): print("Working on tenant: ", tenant) def signal_handler(sig, frame): print("Interrupt received in sub process") exit(0) signal.signal(signal.SIGINT, signal_handler) time.sleep(2) def multi_process(): with ProcessPoolExecutor(2) as process_executor: process_futures = [ process_executor.submit(thread_worker, tenant) for tenant in range(5) ] def signal_handler(sig, frame): print("Interrupt received") # 遍历并终止主进程的所有子进程 parent_process = psutil.Process(os.getpid()) for child in parent_process.children(recursive=True): child.send_signal(signal.SIGINT) # 取消未启动任务并关闭executor process_executor.shutdown(wait=True, cancel_futures=True) print("Killed all processes") signal.signal(signal.SIGINT, signal_handler) if __name__ == "__main__": multi_process()
原理说明
EXTRA_QUEUED_CALLS = 1是ProcessPoolExecutor的内部优化设置,会提前将一个任务放入队列,worker进程完成当前任务后会立即执行该预加载任务。即使调用shutdown(cancel_futures=True),已进入队列的预加载任务仍会被启动,导致出现额外进程。- 方案1通过关闭预加载机制,从根源避免shutdown后启动新任务;方案2通过主动终止所有子进程,确保任何已启动的额外进程都会被SIGINT信号终止。
内容的提问来源于stack exchange,提问作者gd vigneshwar
相关产品推荐
相关产品推荐

