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

自定义multiprocessing.Queue子类在multiprocessing.Process中无法访问属性的问题排查

问题原因分析

这个错误的核心原因是multiprocessing队列在跨进程传递时的序列化机制导致自定义属性丢失。

当你把自定义的Queue实例传给multiprocessing.Process时,Python会对这个对象进行序列化(pickle),然后在子进程中反序列化重建。但原生的multiprocessing.queues.Queue的序列化逻辑只处理了它自身的核心属性,你添加的size属性并没有被纳入这个序列化流程。所以子进程中重建的Queue实例就没有size属性,调用put时自然会抛出AttributeError。

顺便说一下,你在PyCharm调试时看到主进程里的Queue有size属性,是因为那是原对象;而子进程拿到的是反序列化后的新实例,缺失了自定义属性。

解决方案

方案1:让自定义Queue支持完整序列化

你可以通过实现__getstate__和__setstate__方法,告诉pickle如何保存和恢复你的自定义属性。修改后的代码如下:

from multiprocessing import get_context
from multiprocessing.queues import Queue as BaseQueue
from multiprocessing.sharedctypes import Value

# 先实现一个可序列化的SharedCounter
class SharedCounter:
    def __init__(self, initial_value=0):
        self.count = Value('i', initial_value)
    
    def increment(self, n=1):
        with self.count.get_lock():
            self.count.value += n
    
    def decrement(self, n=1):
        with self.count.get_lock():
            self.count.value -= n
    
    def value(self):
        with self.count.get_lock():
            return self.count.value

class Queue(BaseQueue):
    """可跨平台的multiprocessing.Queue实现"""
    def __init__(self, *args, **kwargs):
        super().__init__(*args, **kwargs, ctx=get_context())
        self.size = SharedCounter(0)
    
    def put(self, *args, **kwargs):
        self.size.increment(1)
        super().put(*args, **kwargs)
    
    def get(self, *args, **kwargs):
        item = super().get(*args, **kwargs)
        self.size.decrement(1)
        return item
    
    # 自定义序列化逻辑
    def __getstate__(self):
        # 保存当前对象的所有属性
        state = self.__dict__.copy()
        return state
    
    def __setstate__(self, state):
        # 恢复所有属性
        self.__dict__.update(state)
        # 确保父类核心状态正确初始化
        super().__init__(ctx=get_context())
        self.__dict__.update(state)

这样修改后,当Queue被序列化时,size属性会被保存,子进程反序列化时就能恢复这个属性,调用put就不会报错了。

方案2:使用原生Queue+独立共享计数器

如果你不想修改Queue子类,也可以把计数器单独作为参数传给子进程,在analyse函数里手动更新:

from multiprocessing import Queue, Process, Value

def analyse(i, name, func, grid, queue, counter):
    # 业务逻辑...
    queue.put((i, name, single, minimum, current, peak))
    # 手动增加计数器
    with counter.get_lock():
        counter.value += 1

# 主进程中
queue = Queue()
counter = Value('i', 0)
for i, (name, func) in enumerate(funcs.items()):
    p = Process(target=analyse, args=(i, name, func, grid, queue, counter))
    p.start()

这种方式不需要自定义Queue,避免了序列化问题,但需要在业务代码中额外处理计数器,不如子类封装优雅。

关于你使用Manager的选择

你提到后来改用了multiprocessing.Manager的Queue,确实能解决跨进程传递的问题,但Manager的Queue是通过中间管理进程实现的,所有操作都要经过这个中间进程,所以性能会比原生Queue差一些。如果对性能有要求,方案1是更好的选择。

内容的提问来源于stack exchange,提问作者Nguyen Thai Binh

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.01 00:53:15