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'> ---------------------------------------------------------------------------
原因分析
- 装饰器场景:使用
@t.event装饰器时,函数foo在定义完成的瞬间就被传入on方法并存入Manager.dict(),此时函数还未被绑定到__main__模块的全局命名空间中,pickle尝试查找__main__.foo时自然找不到,触发PicklingError。 - 直接调用场景:
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
相关产品推荐
相关产品推荐

