自定义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

