如何在3个及以上进程间共享队列中的同一份数据?
多进程队列数据广播解决方案
问题根因
multiprocessing.Queue是消费型消息队列,单条消息被任意进程调用get()取出后就会从队列中删除,因此多个计算进程会竞争消费队列中的数据,无法各自拿到全量的原始数据。
最优解决方案:多输入队列广播模式
给每个计算进程创建独立的专属输入队列,读进程读取到数据后,将同一份数据写入所有计算进程的输入队列即可,无需额外锁控制,实现简单且性能稳定。
修改后的代码如下:
from multiprocessing import Process, Manager def read(pathList, calc_queues): # calc_queues 为所有计算进程的输入队列列表 for path in pathList: data = readFunc(path) # 同一份数据写入所有计算进程的输入队列 for q in calc_queues: q.put(data) # 给所有队列写入结束标记 for q in calc_queues: q.put(None) return def calc0(src_q, des_q): while True: data = src_q.get() if data is None: break des_q.put(calcFunc0(data)) return def calc1(src_q, des_q): while True: data = src_q.get() if data is None: break des_q.put(calcFunc1(data)) return if __name__ == '__main__': with Manager() as m: # 为每个计算进程创建独立的输入队列 calc0_in = m.queue() calc1_in = m.queue() res0 = m.queue() res1 = m.queue() readProcess = Process(target=read, args=(readPathList, [calc0_in, calc1_in])) readProcess.start() calcProcess0 = Process(target=calc0, args=(calc0_in, res0)) calcProcess0.start() calcProcess1 = Process(target=calc1, args=(calc1_in, res1)) calcProcess1.start() readProcess.join() calcProcess0.join() calcProcess1.join()
其他适配方案
- 如果单份数据体积很大,多队列复制会占用过多内存,可以使用
multiprocessing.Array/multiprocessing.Manager.list等共享存储对象搭配RLock读写锁,读取进程写数据后通知所有计算进程读取,计算进程读完后再更新下一份数据,能减少内存复制开销,但需要自行控制读写时序,实现复杂度更高。 - 对于超大数据的多进程共享读取,也可以采用内存映射文件
mmap实现零拷贝共享,所有进程直接读取同一块内存映射区域,无需额外的进程间数据传输。
内容的提问来源于stack exchange,提问作者DennisDai
相关产品推荐
相关产品推荐

