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

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. 该推测是否正确?
  2. 有何规避方案?这是遗留代码,无法轻易取消多进程。

解答

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.14 13:55:20