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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.21 11:27:58