Python进程嵌套场景下跨进程队列通信问题求解
问题解决方法
你的代码主要有几个核心问题导致b1收不到数据:
- 共享队列没有正确传递:b1进程根本没拿到a2/a3/a4要写入的队列,原代码里Process_B用了未定义的全局q4,完全没关联到共享队列。
- 函数/类参数不匹配:比如SubProcesses123初始化需要name,但创建时没传;Process_A2_3_4的参数个数和启动时传的不一致。
- 缺少往共享队列写数据的逻辑:a2/a3/a4的doSomething方法里没有实际写入队列的代码。
- 用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
相关产品推荐
相关产品推荐

