使用ProcessPoolExecutor调用subprocess阻塞问题及线程池方案咨询
问题根因
- 核心语法错误:
executor.submit(client.Start(agent, buffer))的传参方式错误。submit方法要求传入可调用对象+对应参数,你当前的写法会在主线程同步执行client.Start(agent, buffer),根本没有将任务提交到进程池运行,这是程序拉起客户端后直接卡住的直接原因。正确写法为executor.submit(client.Start, agent, buffer),将函数对象和参数分开传入。 - 上下文管理器的隐式等待逻辑:你将
ProcessPoolExecutor放在with代码块中,每轮episode执行完退出with块时,会自动等待所有已提交的任务执行完成才会进入下一轮循环,这是你感知到代码会等待任务完成再往下走的原因,和你将函数封装为类、给类变量赋值的操作无关。 - 子进程调用本身阻塞:
subprocess_function中使用的p.communicate()是阻塞API,会等到你启动的批处理进程完全退出、返回所有标准输出/错误后才会结束函数执行,该任务天然会持续到批处理运行结束。
关于是否改用ThreadPoolExecutor
不需要盲目切换线程池,修正上述错误后可根据场景选择:
- 如果
CustomClient.Start是CPU密集型逻辑,继续使用ProcessPoolExecutor可以规避GIL限制,性能更优;如果是IO密集型逻辑(仅做网络通信、消息等待等),ThreadPoolExecutor内存开销更低,更适合。 - 无论是用进程池还是线程池,都必须先修正
submit的传参写法,否则都会出现同步阻塞的问题。
参考修正代码
基础修正(保留每轮等待任务完成的逻辑)
for episode in range(n_episodes): print(f"\r\n{'-' * 60}\r\n") with concurrent.futures.ProcessPoolExecutor() as executor: client = CustomClient("tcp://127.0.0.1:5556", "tcp://127.0.0.1:5555") sub_proc = executor.submit(subprocess_function) # 函数对象和参数分开传入,不直接调用函数 client_proc = executor.submit(client.Start, agent, buffer)
进阶修正(不需要每轮等待任务完成,异步提交)
# 全局初始化进程池,避免每轮重复创建销毁开销 executor = concurrent.futures.ProcessPoolExecutor() for episode in range(n_episodes): print(f"\r\n{'-' * 60}\r\n") client = CustomClient("tcp://127.0.0.1:5556", "tcp://127.0.0.1:5555") sub_proc = executor.submit(subprocess_function) client_proc = executor.submit(client.Start, agent, buffer) # 所有任务提交完成后,再手动关闭进程池并等待执行完成 executor.shutdown(wait=True)
内容的提问来源于stack exchange,提问作者B. Cratty
相关产品推荐
相关产品推荐

