Python使用multiprocessing并行运行函数并通过Queue取输出的疑问
Python multiprocessing 并行执行函数并获取结果的正确操作顺序
核心结论
你写的「启动所有进程→从Queue获取所有结果→join所有进程」是唯一无死锁风险且效率最高的执行顺序,另外两种顺序都不推荐。
各操作顺序的问题分析
- 先join所有进程再调用get:存在死锁风险。Queue的底层缓冲区有大小上限,当子进程返回的结果总大小超过缓冲区时,子进程会阻塞等待主进程从队列取走数据才能退出,此时主进程卡在join等待子进程退出,子进程卡在等主进程取数据,直接造成永久死锁。
- 每join第i个进程再获取第i个结果:首先完全浪费了并行优势,你需要等待第i个进程跑完才能拿结果,哪怕后面的进程早就执行完成;其次依然存在和上面相同的死锁风险。
- 先拿结果再join:
queue.get()本身就是阻塞调用,队列没有数据时会自动等待子进程写入结果,不需要额外等待。所有结果都取出后,子进程基本都已经执行完成,此时再调用join只需要回收进程资源即可,不会有阻塞风险,也不会浪费并行时间。
你当前代码的修正点
你现在的函数直接return结果是无法传递到主进程的,子进程的返回值不会自动进入Queue,需要手动把结果写入队列:
from numpy.random import uniform from multiprocessing import Process, Queue # 新增queue参数,手动把结果写入队列 def function(x, queue): res = uniform(0.0, x) queue.put(res) if __name__ == "__main__": queue = Queue() processes = [] x_values = [1.0, 10.0, 100.0] # 启动所有进程 for x in x_values: process = Process(target=function, args=(x, queue, )) processes.append(process) process.start() # 获取所有结果,get会自动阻塞等待队列写入 outputs = [queue.get() for _ in range(len(x_values))] # 回收进程资源,避免出现僵尸进程 for process in processes: process.join()
简化方案
如果不需要精细控制进程逻辑,直接用multiprocessing.Pool的map方法可以省去手动维护队列、进程启动和回收的步骤,代码更简洁:
from numpy.random import uniform from multiprocessing import Pool def function(x): return uniform(0.0, x) if __name__ == "__main__": x_values = [1.0, 10.0, 100.0] # 自动管理进程池,并行执行后直接返回所有结果 with Pool(processes=len(x_values)) as pool: outputs = pool.map(function, x_values)
内容的提问来源于stack exchange,提问作者Euler_Salter
相关产品推荐
相关产品推荐

