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

如何从多进程读取的队列中向主程序返回数据?

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)

方案优势说明

  1. 效率提升:使用原生mp.Queue()或mp.Pipe(),避免Manager共享对象的中间代理开销和锁竞争,传输速度接近主进程直接读队列。
  2. 实时性:主进程异步读取数据,无需等待所有Worker结束再处理,能更早开始后续数据处理逻辑。
  3. 可靠性:通过信号机制确保Reader处理完队列所有剩余数据后再退出,不会丢失结果。

内容的提问来源于stack exchange,提问作者samR

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.16 11:02:06