如何在multiprocessing Manager中实现复杂对象跨进程共享
问题根因
你遇到的问题本质上是两个原因叠加导致的:
multiprocessing.Manager提供的共享容器(比如dict)只会监测容器本身的键值变更,不会自动同步值内部的自定义类属性变化- 普通方式创建的自定义类实例存到共享容器时,仅会做一次序列化拷贝,后续对实例内部属性的修改都是对本地副本的修改,不会同步到共享容器和其他进程
可行解决方案
方案1:将自定义类注册到Manager,使用托管代理对象(最适合需要实时跨进程同步的场景)
把你用到的A、B类全部注册到自定义的Manager中,通过Manager创建的实例是全局代理对象,所有属性修改都会实时同步到所有进程。
示例改造代码:
import multiprocessing as mp from multiprocessing.managers import BaseManager # 你的自定义类 class B: def __init__(self): self.val = 0 def add(self, n): self.val += n class A: # 注意嵌套的B实例也要用托管对象,所以通过参数传入 def __init__(self, b_inst): self.b = b_inst # 自定义Manager,注册所有需要共享的类 class MyManager(BaseManager): pass MyManager.register('B', B) MyManager.register('A', A) # 进程初始化函数,你的业务逻辑可以写在这里 def init_worker(shared_dict): global d d = shared_dict def worker_task(a_id): a = d.get(a_id) a.b.add(1) # 修改内部属性会直接同步 if __name__ == '__main__': # 启动自定义Manager with MyManager() as manager: # 通过Manager创建托管实例 b = manager.B() a = manager.A(b) # 存到共享dict shared_d = manager.dict() shared_d['a1'] = a # 启动进程池 with mp.Pool(initializer=init_worker, initargs=(shared_d,)) as pool: pool.map(worker_task, ['a1']*10) # 验证结果:输出10,说明内部属性已同步 print(shared_d['a1'].b.val)
方案2:拆分自定义类为基础共享类型(适合业务逻辑简单的场景)
不用自定义类,把所有需要共享的属性全部用Manager自带的可共享结构(dict、list、Value、Array等)嵌套实现,直接操作共享结构的键值即可自动同步,不需要额外注册类。
示例:
if __name__ == '__main__': d = mp.Manager().dict() # 用嵌套dict代替A、B类 d['a1'] = mp.Manager().dict({ 'b': mp.Manager().dict({'val': 0}) }) # 修改内部属性直接同步 d['a1']['b']['val'] += 1
方案3:进程结束后汇总结果(适合不需要实时同步的场景,性能最优)
如果你的业务不需要进程运行时实时共享数据,只是需要收集所有进程的执行结果,可以完全不用共享对象:每个进程自己创建并修改对象,执行完成后把对象的状态返回给主进程,主进程统一合并结果即可,没有代理开销,运行效率更高。
注意事项
- 所有托管对象的方法调用、属性修改都会走进程间通信,性能比普通对象低,高并发场景要做好压测
- 托管对象的方法如果涉及多进程同时写,需要自行加
Manager.Lock()保证线程安全 - 嵌套的自定义类必须全部注册为托管类,不能在类内部直接普通实例化嵌套类,否则嵌套实例还是本地副本,不会同步
内容的提问来源于stack exchange,提问作者python_interest
相关产品推荐
相关产品推荐

