Python2中多进程下使用ThreadPoolExecutor引发死锁的问题咨询
多进程与全局ThreadPoolExecutor的冲突问题
问题场景
我在多进程环境中,让每个子进程向全局ThreadPoolExecutor提交任务,代码如下:
from concurrent.futures import ThreadPoolExecutor import time executor = ThreadPoolExecutor(max_workers=200) def fake_func(): time.sleep(3) print("??") def query(): print("submitted fake func threadpool: {}".format(executor)) result = executor.submit(fake_func) result.result() print("Done fake func: {}".format(result.done())) return [] def simple_parallel(f, args_list): import multiprocessing as mp processes = [] print("start") for args in args_list: p = mp.Process(target=f, args=(args)) processes.append(p) p.start() for p in processes: p.join(timeout=10) return [None] * len(args_list)
正常运行的调用方式
if __name__ == "__main__": data = [] for result in simple_parallel(query, [()] * 10): if result: data.append(result)
导致永久挂起的调用方式
如果在启动多进程前先向全局线程池提交任务,程序会永久挂起:
if __name__ == "__main__": result = executor.submit(fake_func) result.result() print("done dry call") data = [] for result in simple_parallel(query, [()] * 10): if result: data.append(result)
我的推测
子进程会复制父进程的全局线程池对象,而线程池中的线程本身是无法被跨进程复制的。当子进程尝试复用这些“复制过来的线程对象”时,找不到对应的实际线程,进而导致任务无法执行,程序挂起。
疑问
- 该推测是否正确?
- 有何规避方案?这是遗留代码,无法轻易取消多进程。
解答
1. 推测正确性
你的推测是正确的。
在Python的multiprocessing中,子进程创建逻辑会导致这个问题:
- Unix/Linux下用fork模式:子进程会完整复制父进程的地址空间,包括全局
ThreadPoolExecutor对象和其中已创建的线程状态。但线程是操作系统级资源,无法被复制到子进程,子进程里的线程对象只是空壳,没有对应运行的线程。当子进程用这个复制的线程池提交任务时,线程池尝试复用不存在的线程,导致任务无法调度,result.result()一直阻塞,程序挂起。 - Windows下用spawn模式:虽然会重新导入模块,但如果全局线程池已被父进程使用过,模块级别的线程池状态会不一致,同样引发任务阻塞问题。
2. 规避方案
针对无法取消多进程的遗留代码,提供以下可行方案:
方案一:延迟初始化全局线程池
不在模块级别直接创建线程池,而是在需要使用时初始化,确保每个子进程拥有独立的线程池:
from concurrent.futures import ThreadPoolExecutor import time # 全局变量占位,延迟初始化 executor = None def get_executor(): global executor if executor is None: executor = ThreadPoolExecutor(max_workers=200) return executor def fake_func(): time.sleep(3) print("??") def query(): executor = get_executor() print("submitted fake func threadpool: {}".format(executor)) result = executor.submit(fake_func) result.result() print("Done fake func: {}".format(result.done())) return [] # simple_parallel函数保持不变
父进程提前调用线程池不影响子进程,子进程启动时executor为None,会重新初始化自己的线程池。
方案二:子进程启动后重置线程池
在query函数开头判断当前进程是否为子进程,若是则重新创建线程池:
from concurrent.futures import ThreadPoolExecutor import time import multiprocessing as mp executor = ThreadPoolExecutor(max_workers=200) def fake_func(): time.sleep(3) print("??") def query(): global executor # 子进程重置线程池 if mp.current_process().name != 'MainProcess': executor = ThreadPoolExecutor(max_workers=200) print("submitted fake func threadpool: {}".format(executor)) result = executor.submit(fake_func) result.result() print("Done fake func: {}".format(result.done())) return [] # simple_parallel函数保持不变
既保留父进程的线程池使用,又让子进程拥有独立的可用线程池。
方案三:进程隔离的线程池存储
用字典按进程ID存储线程池,确保每个进程拥有独立实例:
from concurrent.futures import ThreadPoolExecutor import time import multiprocessing as mp # 按进程ID存储线程池 process_executors = {} def get_executor(): pid = mp.current_process().pid if pid not in process_executors: process_executors[pid] = ThreadPoolExecutor(max_workers=200) return process_executors[pid] def fake_func(): time.sleep(3) print("??") def query(): executor = get_executor() print("submitted fake func threadpool: {}".format(executor)) result = executor.submit(fake_func) result.result() print("Done fake func: {}".format(result.done())) return [] # simple_parallel函数保持不变
完全隔离父进程与子进程的线程池状态,避免冲突。
内容的提问来源于stack exchange,提问作者toby_yo
相关产品推荐
相关产品推荐

