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

Python进程嵌套场景下跨进程队列通信问题求解

问题解决方法

你的代码主要有几个核心问题导致b1收不到数据:

  1. 共享队列没有正确传递:b1进程根本没拿到a2/a3/a4要写入的队列,原代码里Process_B用了未定义的全局q4,完全没关联到共享队列。
  2. 函数/类参数不匹配:比如SubProcesses123初始化需要name,但创建时没传;Process_A2_3_4的参数个数和启动时传的不一致。
  3. 缺少往共享队列写数据的逻辑:a2/a3/a4的doSomething方法里没有实际写入队列的代码。
  4. 用qsize()判断队列是否有数据不可靠:多进程环境下qsize()的结果可能不准确,直接用get()阻塞等待更稳妥。

修正后的完整代码

import multiprocessing
from multiprocessing import Queue, Process

def Process_B(shared_queue: multiprocessing.Queue):
    # 直接用get()阻塞等待数据,不用qsize()判断
    while True:
        var1, var2, var3 = shared_queue.get()
        print(f"Process_B收到数据: {var1}, {var2}, {var3}")

def Process_A2_3_4(q123: multiprocessing.Queue, shared_queue: multiprocessing.Queue, name: str):
    # 修正初始化参数,传入name
    supProcesses = SubProcesses123(name)
    while True:
        var1, var2, var3 = q123.get()
        supProcesses.doSomething(var1, var2, var3, shared_queue)

class SubProcesses123:
    def __init__(self, name):
        self.name = name
    
    def doSomething(self, var1, var2, var3, shared_queue):
        print(f"{self.name} 正在处理数据: {var1}, {var2}, {var3}")
        # 往共享队列写入数据,传给Process_B
        shared_queue.put((var1, var2, var3))

class Process_A:
    def __init__(self, shared_queue):
        self.a2 = Queue()
        self.a3 = Queue()
        self.a4 = Queue()
        # 把共享队列传给a2/a3/a4
        self.startProcess234(shared_queue)
        self.doSomething()

    def startProcess234(self, shared_queue):
        # 修正参数个数,传入各自的队列、共享队列、进程名
        a2 = Process(target=Process_A2_3_4, args=(self.a2, shared_queue, "a2"))
        a3 = Process(target=Process_A2_3_4, args=(self.a3, shared_queue, "a3"))
        a4 = Process(target=Process_A2_3_4, args=(self.a4, shared_queue, "a4"))
        a2.start()
        a3.start()
        a4.start()
    
    def doSomething(self):
        # 模拟给a2/a3/a4发送数据
        import time
        count = 0
        while True:
            self.a2.put((count, count+1, count+2))
            self.a3.put((count+3, count+4, count+5))
            self.a4.put((count+6, count+7, count+8))
            count +=9
            time.sleep(1)

if __name__ == '__main__':
    # 主进程创建唯一的共享队列
    shared_queue = Queue()
    # 启动Process_B,传入共享队列
    b1 = Process(target=Process_B, args=(shared_queue,))
    b1.start()
    # 启动Process_A,传入共享队列
    retrieve = Process(target=Process_A, args=(shared_queue,))
    retrieve.start()

    # 等待子进程(可选,因为都是无限循环)
    b1.join()
    retrieve.join()

关键修正点说明

  • 共享队列的传递:主进程创建shared_queue,同时传给Process_A和Process_B,Process_A再把它传给自己启动的a2/a3/a4,确保所有需要读写的进程拿到的是同一个队列实例。
  • 去掉qsize()判断:Queue.get()会自动阻塞直到有数据,避免了多进程下qsize()的竞态问题。
  • 补全参数匹配:修正了SubProcesses123初始化、Process_A2_3_4的参数个数问题,确保代码能正常运行。
  • 添加写队列逻辑:在doSomething方法里调用shared_queue.put(),把数据传给Process_B。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.22 21:48:14