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

Python多进程存函数至Manager.dict()遇PicklingError问题求助

Python多进程Manager.dict()存储函数触发Pickle错误的解决

问题现象

在实现通信API时,尝试将顶级函数存入multiprocessing.Manager.dict()时触发_pickle.PicklingError,即便使用的是顶级函数也无法解决。最小复现代码如下:

import multiprocessing as mp
from typing import Callable, Any
from uuid import uuid4
from time import sleep

def defaultInit(self): pass
def defaultLoop(self): pass

class Event():
  def __init__(self, name: str, callback: Callable[..., Any]):
    super().__init__()
    self._id = uuid4()
    self._name = name
    self._trigger = False
    self._callback = callback
    self._args = None

  def _updateEvents(self, events):
    for i, event in enumerate(events[self._name]):
      if event._id == self._id:
        tmp = events[self._name]
        tmp[i] = self
        events[self._name] = tmp
        break

  def execIfTriggered(self, owner) -> None:
    if self._trigger:
      self._trigger = False
      self._updateEvents(owner._events)
      self._callback(owner, *self._args)

  def trigger(self, owner, *args) -> None:
    self._args = args
    self._trigger = True
    self._updateEvents(owner._events)
    return self


class Target():
  def __init__(self, events):
    self._init = defaultInit
    self._loop = defaultLoop
    self._events = events
  
  def event(self, callback):
    name = callback.__name__
    return self.on(name, callback)

  def on(self, eventName, callback):
    if eventName not in self._events:
      self._events[eventName] = [Event(eventName, callback)]
    else:
      self._events[eventName] = self._events[eventName] + [Event(eventName, callback)]
    return callback

  def emit(self, event, *args):
    for ev in self._events[event]:
      ev.trigger(self, *args)

  def init(self, callback):
    self._init = callback

  def loop(self, callback):
    self._loop = callback

  def processEvents(self):
    for events in self._events.values():
      for event in events:
        event.execIfTriggered(self)

def run(target):
  target._init(target)
  while True:
    target.processEvents()
    target._loop(target)


if __name__ == '__main__':
  mgr = mp.Manager()
  t = Target(mgr.dict())
  proc = mp.Process(target=run, args=(t,))

  @t.init
  def fooInit(target):
    print('fooInit')

  @t.event
  def foo(target):
    print('fooEvent')

  proc.start()
  sleep(1)
  t.emit('foo')
  sleep(1)
  proc.join()

错误场景1:使用装饰器定义回调

触发如下错误:

Traceback (most recent call last):
  File "/home/lsuardi/smc/moscau/Acquisition_python/test.py", line 88, in <module>
    def foo(target):
  File "/home/lsuardi/smc/moscau/Acquisition_python/test.py", line 47, in event
    return self.on(name, callback)
  File "/home/lsuardi/smc/moscau/Acquisition_python/test.py", line 51, in on
    self._events[eventName] = [Event(eventName, callback)]
  File "<string>", line 2, in __setitem__
  File "/home/lsuardi/miniconda3/envs/SMC-lsuardi/lib/python3.10/multiprocessing/managers.py", line 817, in _callmethod
    conn.send((self._id, methodname, args, kwds))
  File "/home/lsuardi/miniconda3/envs/SMC-lsuardi/lib/python3.10/multiprocessing/connection.py", line 206, in send
    self._send_bytes(_ForkingPickler.dumps(obj))
  File "/home/lsuardi/miniconda3/envs/SMC-lsuardi/lib/python3.10/multiprocessing/reduction.py", line 51, in dumps
    cls(buf, protocol).dump(obj)
_pickle.PicklingError: Can't pickle <function foo at 0x7f13530775b0>: attribute lookup foo on __main__ failed

错误场景2:直接调用event方法

修改代码为:

# 替换装饰器代码
def foo(target):
  print('fooEvent')
t.event(foo)

触发如下错误:

Traceback (most recent call last):
  File "/home/lsuardi/smc/moscau/Acquisition_python/test.py", line 89, in <module>
    t.event(foo)
  File "/home/lsuardi/smc/moscau/Acquisition_python/test.py", line 47, in event
    return self.on(name, callback)
  File "/home/lsuardi/smc/moscau/Acquisition_python/test.py", line 51, in on
    self._events[eventName] = [Event(eventName, callback)]
  File "<string>", line 2, in __setitem__
  File "/home/lsuardi/miniconda3/envs/SMC-lsuardi/lib/python3.10/multiprocessing/managers.py", line 833, in _callmethod
    raise convert_to_error(kind, result)
multiprocessing.managers.RemoteError:
---------------------------------------------------------------------------
Traceback (most recent call last):
  File "/home/lsuardi/miniconda3/envs/SMC-lsuardi/lib/python3.10/multiprocessing/managers.py", line 253, in serve_client
    request = recv()
  File "/home/lsuardi/miniconda3/envs/SMC-lsuardi/lib/python3.10/multiprocessing/connection.py", line 251, in recv
    return _ForkingPickler.loads(buf.getbuffer())
AttributeError: Can't get attribute 'foo' on <module '__main__' from '/home/lsuardi/smc/moscau/Acquisition_python/test.py'>
---------------------------------------------------------------------------

原因分析

  1. 装饰器场景:使用@t.event装饰器时,函数foo在定义完成的瞬间就被传入on方法并存入Manager.dict(),此时函数还未被绑定到__main__模块的全局命名空间中,pickle尝试查找__main__.foo时自然找不到,触发PicklingError。
  2. 直接调用场景:Manager是一个独立的服务进程,它启动时不会执行主进程中if __name__ == '__main__'代码块内的内容,因此在Manager进程的__main__模块中不存在foo函数的定义。当主进程把foo序列化后发送给Manager进程时,Manager进程反序列化需要在自己的__main__中找到该函数,于是触发AttributeError。

解决办法

方法1:将回调函数移到模块顶级(不在if __name__块内)

把所有回调函数定义在if __name__ == '__main__'之外,这样Manager进程启动时会加载这些函数定义:

# 移到模块顶级
def fooInit(target):
  print('fooInit')

def foo(target):
  print('fooEvent')

if __name__ == '__main__':
  mgr = mp.Manager()
  t = Target(mgr.dict())
  proc = mp.Process(target=run, args=(t,))

  t.init(fooInit)
  t.event(foo)

  proc.start()
  sleep(1)
  t.emit('foo')
  sleep(1)
  proc.join()

方法2:用事件名映射函数,避免直接共享函数

不在Manager.dict()中存储函数,而是存储事件名,子进程通过事件名从本地函数映射表中查找回调:

# 定义全局函数映射表
FUNCTION_MAP = {}

def register_func(name, func):
    FUNCTION_MAP[name] = func

class Event():
  def __init__(self, name: str, callback_name: str):
    super().__init__()
    self._id = uuid4()
    self._name = name
    self._trigger = False
    self._callback_name = callback_name
    self._args = None

  def execIfTriggered(self, owner) -> None:
    if self._trigger:
      self._trigger = False
      self._updateEvents(owner._events)
      # 通过名称查找函数
      callback = FUNCTION_MAP[self._callback_name]
      callback(owner, *self._args)

# Target类的on方法修改为存储函数名而非函数
def on(self, eventName, callback):
    callback_name = callback.__name__
    register_func(callback_name, callback)
    if eventName not in self._events:
      self._events[eventName] = [Event(eventName, callback_name)]
    else:
      self._events[eventName] = self._events[eventName] + [Event(eventName, callback_name)]
    return callback

这种方式下,Manager.dict()只存储字符串,不存在序列化函数的问题,子进程通过本地的FUNCTION_MAP找到对应函数执行。

方法3:替换默认Pickler为cloudpickle

cloudpickle支持序列化更多类型的函数,包括在if __name__块内定义的函数。需要先安装cloudpickle,然后修改Manager的序列化方式:

import cloudpickle
from multiprocessing.managers import BaseManager

# 自定义Manager,使用cloudpickle序列化
class CloudpickleManager(BaseManager):
    pass

# 注册dict类型
CloudpickleManager.register('dict', dict)

if __name__ == '__main__':
  mgr = CloudpickleManager()
  mgr.start()
  t = Target(mgr.dict())
  # 后续代码不变

注意:这种方式需要确保子进程也能访问到函数定义,cloudpickle会序列化函数的字节码,但对于复杂场景可能存在兼容性问题。


内容的提问来源于stack exchange,提问作者Fayeure

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.17 22:11:59