重写multiprocessing.queues.Queue的put方法时遇到的问题
实现不重复元素的multiprocessing.Queue
问题背景
想要实现一个不会添加已存在元素的multiprocessing.Queue,但两次尝试均失败:
第一次尝试代码及错误
最初的代码尝试复用标准库Queue的实现思路,但抛出属性错误:
from multiprocessing.queues import Queue from multiprocessing import get_context class CustomQueue(Queue): def put(self, obj, block=True, timeout=None): if obj not in self: return super().put(obj, block, timeout) def __contains__(self, item): with self.mutex: return item in self.queue custom_queue = CustomQueue(ctx=get_context())
执行时出现错误:AttributeError: 'CustomQueue' object has no attribute 'mutex'
第二次尝试代码及问题
查看multiprocessing.queues.Queue源码后修改代码,但obj not in self._buffer始终返回True,无法阻止重复添加:
class CustomQueue(Queue): def put(self, obj, block=True, timeout=None): if self._closed: raise ValueError(f"Queue {self!r} is closed") if not self._sem.acquire(block, timeout): raise Full with self._notempty: if self._thread is None: self._start_thread() if obj not in self._buffer: self._buffer.append(obj) self._notempty.notify()
问题分析
- 第一次错误原因:
multiprocessing.queues.Queue的内部同步原语并非mutex,内部存储也不是self.queue,而是使用_sem、_notempty等锁,以及_buffer作为待发送元素的临时队列。 - 第二次失败原因:
self._buffer仅存储待后台线程发送到管道的元素,已经发送到管道但未被消费的元素不会出现在_buffer中,因此只检查_buffer会漏掉大量已存在的元素,导致重复判断失效。
解决方案
要实现真正的去重,需要跟踪所有已放入队列但未被消费的元素,且必须保证这个跟踪机制在多进程环境下是安全的。可以借助multiprocessing.Manager提供的进程安全集合来实现:
from multiprocessing.queues import Queue from multiprocessing import get_context, Manager class UniqueQueue(Queue): def __init__(self, *args, **kwargs): ctx = kwargs.pop('ctx', get_context()) super().__init__(*args, ctx=ctx, **kwargs) # 创建进程安全的集合,Manager的集合本身已实现跨进程安全 self._manager = Manager() self._seen = self._manager.set() def put(self, obj, block=True, timeout=None): # 先检查元素是否已存在,避免重复添加 if obj not in self._seen: self._seen.add(obj) super().put(obj, block, timeout) def get(self, block=True, timeout=None): obj = super().get(block, timeout) # 元素被消费后,从集合中移除,允许后续重新添加 if obj in self._seen: self._seen.remove(obj) return obj def __contains__(self, item): return item in self._seen
关键说明
- 使用
multiprocessing.Manager().set()创建跨进程可见的集合,确保所有进程都能共享元素的存在状态。 - 重写
put方法:添加元素前先检查集合,不存在则加入集合并调用父类put。 - 重写
get方法:取出元素后从集合中移除,保证元素被消费后可以再次被添加。 - 该实现完全兼容
multiprocessing.Queue的原有接口,同时实现了去重功能。
注意事项
- 如果队列中的元素是不可哈希类型(比如列表),无法存入
set,需要先将元素转为可哈希的形式(比如元组),或者自定义哈希逻辑。 multiprocessing.Manager会启动一个额外的进程来管理共享对象,带来少量性能开销,若对性能要求极高,可考虑使用multiprocessing.Array或multiprocessing.Value结合锁来实现自定义的共享集合。
内容的提问来源于stack exchange,提问作者J Agustin Barrachina
相关产品推荐
相关产品推荐

