如何从多进程读取的队列中向主程序返回数据?
Python多进程队列读取:从独立Reader进程高效传递数据到主程序
问题描述
我用Python多进程做数据计算,Worker进程把结果存入队列,原本主程序用这段代码读队列:
result=[] while not Q.empty(): result.append(Q.get())
但队列数据量极大时,我想在Worker运行期间启动1-2个Reader进程读取队列,避免主进程阻塞。找到的Reader进程代码能持续读队列直到Worker通知结束,但没法把数据传回主程序。目前只找到用multiprocess.Manager()创建共享列表的方法,但速度远不如主进程直接读队列快。想问有没有其他高效方案,能从独立Reader进程把队列数据传给主程序?
附原始示例代码:
import multiprocessing as mp from datetime import datetime def worker(numbers, start, end, qu): """A worker function to calculate squares of numbers.""" res=[] for i in range(start, end): res.append(numbers[i] * numbers[i]) qu.put(res) def reader(q, outputlist): """Read from the queue; this spawns as a separate Process""" #returnedlist=[] while True: msg = q.get() # Read from the queue and do nothing if msg == "DONE": break [outputlist.append(x) for x in msg] # comment out if using 1st method (reading queue in the main program) return def start_reader_procs(q, num_of_reader_procs, L): """Start the reader processes and return all in a list to the caller""" all_reader_procs = list() for ii in range(0, num_of_reader_procs): reader_p = mp.Process(target=reader, args=((q),L,)) reader_p.daemon = True reader_p.start() # Launch reader_p() as another proc all_reader_procs.append(reader_p) return all_reader_procs def main(core_count): numbers = range(50000) # A larger range for a more evident effect of multiprocessing segment = len(numbers) // core_count processes = [] #Q = mp.Queue() m = mp.Manager() Q = m.Queue() # comment out if using 1st method (reading queue in the main program) #---------------------------------- num_of_reader_procs=2 L = m.list() all_reader_procs = start_reader_procs(Q, num_of_reader_procs, L) #----------------------------------------- for i in range(core_count): start = i * segment if i == core_count - 1: end = len(numbers) # Ensure the last segment goes up to the end else: end = start + segment # Creating a process for each segment p = mp.Process(target=worker, args=(numbers, start, end, Q)) processes.append(p) p.start() for p in processes: p.join() print("All worker processes terminated") # comment out if using 1st method (reading queue in the main program) #---------------------------------- ### Tell all readers to stop... for ii in range(0, num_of_reader_procs): Q.put("DONE") for idx, a_reader_proc in enumerate(all_reader_procs): print(" Waiting for reader_p.join() index %s" % idx) a_reader_proc.join() # Wait for a_reader_proc() to finish print(" reader_p() idx:%s is done" % idx) result=list(L) #---------------------------------- # result=[] # while not Q.empty(): # result.append(Q.get()) # result = [x for L in result for x in L] # flatten the list of lists return result if __name__ == '__main__': for core_count in [1, 2, 4]: starttime = datetime.now() print(f"Using {core_count} core(s):") result = main(core_count) print(f"First 10 squares: {list(result)[:10]}") # Display the first 10 results as a sample endtime = datetime.now() print ("Total computation time : {:.1f} sec".format((endtime-starttime).total_seconds())) print()
高效解决方案
核心思路
放弃Manager共享列表,改用原生进程间通信机制(普通队列/管道)实现Reader到主进程的数据传递,避免中间代理进程的开销和锁竞争问题。
方案1:二级普通队列传递(推荐)
让Reader进程读取Worker的结果队列后,将数据写入另一个由主进程创建的mp.Queue(),主进程异步读取这个二级队列,同时等待Worker结束。
优化后代码
import multiprocessing as mp from datetime import datetime import time def worker(numbers, start, end, qu): """Worker计算平方后将结果列表存入队列""" res = [] for i in range(start, end): res.append(numbers[i] * numbers[i]) qu.put(res) def worker_monitor(num_workers, done_signal_q, worker_q): """监控所有Worker结束,通知Reader处理剩余数据""" for _ in range(num_workers): done_signal_q.get() # 发送Worker全部结束的信号 done_signal_q.put("ALL_WORKERS_DONE") def reader(input_q, output_q, done_signal_q): """Reader从Worker队列读数据,写入主进程输出队列,直到收到结束信号""" while True: # 优先检查结束信号,避免遗漏 if not done_signal_q.empty(): signal = done_signal_q.get() if signal == "ALL_WORKERS_DONE": # 处理队列剩余数据 while not input_q.empty(): msg = input_q.get() output_q.put(msg) # 通知主进程本Reader已完成 output_q.put("READER_DONE") break # 非阻塞读取Worker队列数据 try: msg = input_q.get_nowait() output_q.put(msg) except mp.queues.Empty: time.sleep(0.01) # 空队列时短暂休眠,避免CPU空转 def main(core_count): numbers = range(50000) segment = len(numbers) // core_count processes = [] # Worker用的普通队列(比Manager.Queue效率高) worker_q = mp.Queue() # Reader传给主进程的输出队列 output_q = mp.Queue() # 用于传递Worker结束信号的队列 done_signal_q = mp.Queue() # 启动Reader进程 num_reader_procs = 2 reader_procs = [] for _ in range(num_reader_procs): rp = mp.Process(target=reader, args=(worker_q, output_q, done_signal_q)) rp.daemon = True rp.start() reader_procs.append(rp) # 启动Worker进程 for i in range(core_count): start = i * segment end = len(numbers) if i == core_count -1 else start + segment p = mp.Process(target=worker, args=(numbers, start, end, worker_q)) processes.append(p) p.start() # 给监控进程发Worker启动标记,用于计数 done_signal_q.put(f"WORKER_{i}_STARTED") # 启动Worker监控进程 monitor_proc = mp.Process(target=worker_monitor, args=(core_count, done_signal_q, worker_q)) monitor_proc.start() # 主进程异步读取输出队列,收集结果 result = [] reader_done_count = 0 while reader_done_count < num_reader_procs: try: item = output_q.get_nowait() if item == "READER_DONE": reader_done_count +=1 else: result.extend(item) except mp.queues.Empty: time.sleep(0.01) # 等待所有进程结束 for p in processes: p.join() monitor_proc.join() for rp in reader_procs: rp.join() print("All processes terminated") return result if __name__ == '__main__': for core_count in [1,2,4]: starttime = datetime.now() print(f"Using {core_count} core(s):") result = main(core_count) print(f"First 10 squares: {result[:10]}") endtime = datetime.now() print(f"Total computation time : { (endtime-starttime).total_seconds():.1f} sec") print()
方案2:管道传递数据
用mp.Pipe()创建双向管道,主进程持有一端,Reader进程持有另一端。Reader读取队列数据后直接写入管道,主进程从管道接收数据,效率和队列接近。
关键代码片段
def reader(input_q, pipe_end): while True: msg = input_q.get() if msg == "DONE": pipe_end.send("READER_DONE") break pipe_end.send(msg) # 主进程中创建管道 parent_pipe, child_pipe = mp.Pipe() # 启动Reader时传入child_pipe rp = mp.Process(target=reader, args=(worker_q, child_pipe)) rp.start() # 主进程读取管道 while True: data = parent_pipe.recv() if data == "READER_DONE": break result.extend(data)
方案优势说明
- 效率提升:使用原生
mp.Queue()或mp.Pipe(),避免Manager共享对象的中间代理开销和锁竞争,传输速度接近主进程直接读队列。 - 实时性:主进程异步读取数据,无需等待所有Worker结束再处理,能更早开始后续数据处理逻辑。
- 可靠性:通过信号机制确保Reader处理完队列所有剩余数据后再退出,不会丢失结果。
内容的提问来源于stack exchange,提问作者samR
相关产品推荐
相关产品推荐

