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

如何在macOS上实现multiprocessing.queues.Queue.qsize()?现有方案失效

自定义Multiprocessing Queue丢失size属性的问题解决

问题原因

在macOS等默认使用Spawn启动方式的环境中,multiprocessing传递自定义Queue子类实例时,Queue的序列化依赖内置的__reduce__机制。子进程重建对象时,不会执行子类的__init__方法,而是直接恢复底层的队列共享内存结构,导致手动添加的size属性无法被传递到子进程,从而触发AttributeError: 'MyQueue' object has no attribute 'size'。

解决方案

需要重写__reduce__方法将自定义属性纳入序列化流程,同时通过__setstate__方法在子进程中恢复该属性。

修正后的完整代码

import multiprocessing
import os
import time
from multiprocessing import get_context
from multiprocessing.queues import Queue


class SharedCounter(object):
    def __init__(self, n=0):
        self.count = multiprocessing.Value('i', n)

    def increment(self, n=1):
        with self.count.get_lock():
            self.count.value += n

    @property
    def value(self):
        return self.count.value


class MyQueue(Queue):
    def __init__(self, *args, **kwargs):
        super(MyQueue, self).__init__(*args, ctx=get_context(), **kwargs)
        self.size = SharedCounter(0)

    def put(self, *args, **kwargs):
        self.size.increment(1)
        super(MyQueue, self).put(*args, **kwargs)

    def get(self, *args, **kwargs):
        self.size.increment(-1)  # 现在取消注释不会报错
        return super(MyQueue, self).get(*args, **kwargs)

    def qsize(self):
        return self.size.value

    def empty(self):
        return not self.qsize()

    def clear(self):
        while not self.empty():
            self.get()

    def __reduce__(self):
        # 序列化时返回重建函数、参数,以及需要保存的状态
        return (
            MyQueue,
            (),
            {
                '_reader': self._reader,
                '_writer': self._writer,
                '_rlock': self._rlock,
                '_wlock': self._wlock,
                '_sem': self._sem,
                'size': self.size
            }
        )

    def __setstate__(self, state):
        # 反序列化时恢复属性
        self._reader = state['_reader']
        self._writer = state['_writer']
        self._rlock = state['_rlock']
        self._wlock = state['_wlock']
        self._sem = state['_sem']
        self.size = state['size']


def worker(queue):
    while True:
        item = queue.get()
        if item is None:
            break
        print(f'[{os.getpid()}]: got {item}')
        time.sleep(1)


if __name__ == '__main__':
    num_processes = 4
    q = MyQueue()
    pool = multiprocessing.Pool(num_processes, worker, (q,))

    for i in range(10):
        q.put("hello")
        q.put("world")

    for i in range(num_processes):
        q.put(None)

    q.close()
    q.join_thread()
    pool.close()
    pool.join()

关键修改说明

  1. __reduce__方法:
    返回的元组包含三个部分:用于重建对象的类、构造参数(空元组,因为核心属性在__setstate__中恢复)、需要保存的状态字典(包含Queue的核心内部属性和自定义的size)。
  2. __setstate__方法:
    从状态字典中依次恢复Queue的内部属性(_reader、_writer等)和自定义的size属性,确保子进程中的MyQueue实例拥有完整的属性结构。

验证效果

取消注释get方法中的self.size.increment(-1)后,运行代码不再触发AttributeError,子进程可以正常访问size属性,队列的计数功能正常工作。

内容的提问来源于stack exchange,提问作者watch-this

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.09 11:55:24