Python多进程队列引用匹配问题及MultiQueue移除方法失效解决方案咨询
Python多进程队列引用匹配问题及MultiQueue移除方法失效解决方案咨询
我正在写一个Python软件,其中单个生产者进程生成的数据需要转发给多个消费者进程。这些消费者进程并非同时启动和结束——主进程会根据外部因素(比如连接服务器、GUI上的用户操作)来启动不同的消费者。
我想要实现的效果是:当创建新消费者时,同时给MultiQueue对象添加一个新的队列,这样生产者就能把数据也转发给这个新消费者;当消费者结束时(比如GUI事件或网络事件触发),要把对应的队列从MultiQueue中移除。
下面是我写的代码框架,我对多进程还不是很有信心,可能整个思路都有问题?不过除了remove_queue()方法之外,其他部分都能正常工作。我该怎么解决这个移除方法的问题呢?
import multiprocessing import threading from time import sleep class MultiQueue(): def __init__(self, queues = []): manager = multiprocessing.get_context("spawn").Manager() self.queues = manager.list() for queue in queues: self.add_queue(queue) def add_queue(self, queue: multiprocessing.Queue): """ Add a queue """ if not isinstance(queue, multiprocessing.managers.BaseProxy): raise ValueError("Input must be a multiprocessing.managers.Queue object") self.queues.append(queue) def remove_queue(self, queue): """ Remove a queue by reference """ if queue in self.queues: self.queues.remove(queue) else: raise ValueError("Queue not found in MultiQueue") def put(self, item, block=True, timeout=None): """ Put an item into all queues """ for q in self.queues: try: q.put(item, block=False) except multiprocessing.queues.Full: pass class Producer(multiprocessing.Process): def __init__(self, queues: MultiQueue): multiprocessing.Process.__init__(self) self.queues = queues def run(self): i = 0 while True: self.queues.put(i) i+=1 sleep(1) def watcher(n, q): while True: i = q.get() print(f"[{n}]->{i}") if __name__ == "__main__": manager = multiprocessing.get_context("spawn").Manager() q1 = manager.Queue() multiq = MultiQueue([q1]) producer = Producer(multiq) producer.start() threading.Thread(target=watcher, args=(1, q1)).start() #consumer for q1 sleep(3) q2 = manager.Queue() multiq.add_queue(q2) threading.Thread(target=watcher, args=(2, q2)).start() #consumer for q2 sleep(3) multiq.remove_queue(q2) #the code runs correctly until this point
我知道队列没有实现__equal__方法,但有没有办法判断两个队列实际上指向同一个位置呢?
备注:内容来源于stack exchange,提问作者Merlo98765
相关产品推荐
相关产品推荐

