Python多进程Pool如何实现每个worker仅执行一次pickle/unpickle
默认使用multiprocessing.Pool提交绑定了实例的方法(比如示例里的a.run_step)时,Python会自动把实例本身作为任务上下文的一部分,通过pickle序列化后随任务一起发送给worker进程。提交10000次任务就会触发10000次序列化、10000次反序列化,遇到结构复杂的大对象时,这部分开销会远大于任务本身的计算耗时,直接盖过多进程的并行收益。
要实现每个worker仅执行一次pickle/load,核心逻辑是把大对象的传输、加载动作放到worker进程的启动阶段,后续所有任务直接复用进程内已经加载好的对象,不再随单个任务重复传输大对象。
方案1:利用Pool初始化器预加载对象(最推荐,生产环境通用)
multiprocessing.Pool支持传入initializer和initargs参数,会在每个worker进程启动时仅执行一次初始化逻辑。你可以在这个阶段把复杂对象反序列化到worker进程的全局命名空间,后续提交任务时只传轻量的任务函数和必要的小参数,不需要传输整个大对象。
对应给出的示例代码,改造后如下:
import multiprocessing # worker进程全局变量,存储预加载的Analysis实例 _worker_analysis = None def _init_worker(analysis_instance): # 该函数每个worker启动时仅执行1次,仅在这里触发1次反序列化 global _worker_analysis _worker_analysis = analysis_instance def _run_step_in_worker(): # 任务函数直接调用进程内已加载实例的方法,无额外序列化开销 return _worker_analysis.run_step() class Analysis: def run_step(self): print('run_step') def __getstate__(self): print('I dump') return self.__dict__ def __setstate__(self, state): print('I load') self.__dict__ = state if __name__ == '__main__': a = Analysis() # 初始化Pool时传入初始化逻辑和参数,仅会为每个worker触发1次序列化/反序列化 pool = multiprocessing.Pool(4, initializer=_init_worker, initargs=(a,)) for i in range(10): # 提交任务时仅传入轻量顶层函数,不绑定大对象,不会触发重复序列化 pool.apply_async(_run_step_in_worker) pool.close() pool.join()
运行上述代码可以看到,I dump和I load总共只会打印4次(对应4个worker进程各加载1次),和任务提交次数无关,哪怕提交10000次任务也不会额外增加序列化开销。该方案兼容Windows、Linux、macOS全平台,支持所有多进程启动模式,无第三方依赖,稳定性最高。
注意:Windows系统默认用spawn模式启动多进程,必须将Pool创建、任务提交的逻辑放到if __name__ == '__main__':代码块下,避免子进程重复导入模块引发递归启动错误。
方案2:使用fork启动模式复用主进程内存(仅类Unix系统可用)
如果运行环境是Linux、旧版macOS,可以将多进程启动模式设置为fork:fork启动子进程时会直接复制主进程的内存空间,主进程中已经创建好的复杂对象会被worker进程直接继承,完全不需要走pickle序列化流程,连每个worker一次的序列化开销都可以省掉。
使用方式非常简单,在创建Pool之前设置启动模式即可:
multiprocessing.set_start_method('fork', force=True)
注意:fork模式存在固有线程安全问题,如果主进程在创建Pool之前持有多线程锁、数据库连接、打开的文件句柄等资源,子进程可能出现死锁、资源状态异常的问题,逻辑简单的场景可以快速使用,复杂生产场景优先选择方案1。另外macOS从Python 3.8版本开始已经将spawn作为默认启动模式,使用fork需要手动设置。
方案3:替换高性能序列化后端(兜底优化方案)
如果受场景限制必须随任务传输大对象,可以将Python默认的pickle序列化器替换为性能更高的实现,比如cloudpickle,可以大幅降低单次序列化/反序列化的耗时。但该方案本质上还是会随每个任务触发一次序列化,优化幅度远低于前两个方案,仅作为无法调整任务结构时的兜底选择。
- 不要在提交的任务中直接引用主进程的大对象:只要大对象作为任务参数、或者绑定方法的隐式self传入任务队列,就一定会触发重复序列化。
- spawn启动模式下不要依赖模块级全局变量传对象:spawn模式会重新导入业务模块,主进程中定义的全局变量不会自动同步到子进程,必须通过初始化器或者共享内存的方式传递。
- 初始化器加载的对象是进程级隔离的:不同worker进程中的实例互不干扰,如果任务需要修改对象状态,不要依赖跨进程的状态自动同步,需要单独通过队列、共享内存、管理器实现状态传递。
内容的提问来源于stack exchange,提问作者Eurydice

