Python多进程中如何实现类似C++指针的共享数据结构访问?
Python多进程中高效共享大型映射表的方案
问题场景
需要用多进程拆分计算任务,核心需求有两个:
- 子进程需高效访问只读的大型映射表(含数万浮点数的numpy向量),避免数据复制带来的内存开销;
- 子进程需将计算结果写入共享的映射表,且不同进程不会写入同一键,无需同步锁机制。
当前使用普通全局字典时,子进程会复制整个字典,更新仅作用于副本;使用multiprocessing.Manager虽能实现共享,但性能极低,还会因大型numpy数组引发序列化错误。
最小复现代码
from pathos.multiprocessing import ProcessingPool as Pool from multiprocessing import Manager # intended to be accessible by everything everywhere #M = Manager() #SINGLE_SOURCE_MAP = M.dict() # OR SINGLE_SOURCE_MAP = dict() # the function to be used in parallel. In practice OO is a very large def func_run_in_parallel(i): # do we get a "pointer" to OO, or a copy? print(f"job {i}, At start of parallel function: SINGLE_SOURCE_MAP address {id(SINGLE_SOURCE_MAP)}") # a dictionary using something unique to job dd = {f"string-{i}": i} # update the global dict SINGLE_SOURCE_MAP.update(dd) print(f"After attempted update of single source map: SINGLE_SOURCE_MAP address {id(SINGLE_SOURCE_MAP)}") return SINGLE_SOURCE_MAP class OrchestratingObj(object): def __init__(self): self.attribute = 1.318 SINGLE_SOURCE_MAP['Ich'] = 1.618 def run_parallel_job(self): results = None with Pool(processes=2) as pool: results = pool.map(func_run_in_parallel, list(range(3))) pool.close() pool.join() print("\nParallel Job Results") for r in results: print(r) if __name__ == "__main__": OObj = OrchestratingObj() OObj.run_parallel_job() print("\nActual of single source map") for k, v in SINGLE_SOURCE_MAP.items(): print(k, v)
高效解决方案
1. 拆分读写逻辑:子进程返回结果,父进程合并(最优选择)
既然不同进程不会写入同一键,完全可以让子进程只返回自己的计算结果片段,由父进程统一合并到主映射表。这种方式完全避免进程间同步开销,性能最高。
修改后的示例代码:
from pathos.multiprocessing import ProcessingPool as Pool # 只读大型数据(Unix/Linux下fork的子进程会共享父进程内存,只读不会触发复制) LARGE_READ_ONLY_DATA = { "big_numpy_array": ... # 实际的大型numpy向量 } def func_run_in_parallel(i): # 直接访问只读数据,无复制开销 print(f"Job {i} accessing large data shape: {LARGE_READ_ONLY_DATA['big_numpy_array'].shape}") # 返回自己的结果片段 return {f"string-{i}": i} class OrchestratingObj(object): def __init__(self): self.attribute = 1.318 self.SINGLE_SOURCE_MAP = {'Ich': 1.618} def run_parallel_job(self): with Pool(processes=2) as pool: results = pool.map(func_run_in_parallel, list(range(3))) # 父进程合并所有子进程的结果 for res in results: self.SINGLE_SOURCE_MAP.update(res) print("\nMerged Shared Map:") for k, v in self.SINGLE_SOURCE_MAP.items(): print(k, v) if __name__ == "__main__": OObj = OrchestratingObj() OObj.run_parallel_job()
优势:
- 只读数据在Unix/Linux的fork模式下是写时复制(COW),子进程读取不会复制内存;
- 无进程间同步开销,性能接近单进程计算的线性扩展;
- 避免了
Manager的序列化问题。
2. 使用共享内存存储大型只读numpy数组
对于需要跨平台支持(如Windows),或者不想依赖fork的写时复制特性,可以用multiprocessing.shared_memory(Python 3.8+)创建共享内存,让所有进程直接访问同一块内存中的numpy数组。
示例代码:
from pathos.multiprocessing import ProcessingPool as Pool from multiprocessing import shared_memory import numpy as np def func_run_in_parallel(args): shm_name, arr_shape, arr_dtype, i = args # 子进程连接已有的共享内存 existing_shm = shared_memory.SharedMemory(name=shm_name) # 基于共享内存创建numpy数组视图(无数据复制) shared_arr = np.ndarray(arr_shape, dtype=arr_dtype, buffer=existing_shm.buf) # 使用数组进行计算(只读) print(f"Job {i} array sum: {shared_arr.sum()}") # 返回自己的结果片段 res = {f"string-{i}": i} existing_shm.close() return res if __name__ == "__main__": # 父进程创建大型numpy数组并放入共享内存 big_arr = np.random.rand(10000000) # 示例大型数组 shm = shared_memory.SharedMemory(create=True, size=big_arr.nbytes) shared_arr = np.ndarray(big_arr.shape, dtype=big_arr.dtype, buffer=shm.buf) shared_arr[:] = big_arr[:] # 将数据复制到共享内存 # 准备任务参数:传递共享内存名称、数组形状、 dtype 和任务ID tasks = [(shm.name, big_arr.shape, big_arr.dtype, i) for i in range(3)] # 启动进程池 with Pool(processes=2) as pool: results = pool.map(func_run_in_parallel, tasks) # 合并结果到主映射表 SINGLE_SOURCE_MAP = {'Ich': 1.618} for res in results: SINGLE_SOURCE_MAP.update(res) print("\nFinal Shared Map:") for k, v in SINGLE_SOURCE_MAP.items(): print(k, v) # 清理共享内存 shm.close() shm.unlink()
优势:
- 真正实现跨进程内存共享,无数据复制;
- 支持Windows等非fork平台;
- 适合超大型numpy数组的场景。
3. 基于multiprocessing.Array构建轻量共享结构(进阶)
如果确实需要子进程直接写入共享结构,且无竞态,可以用multiprocessing.Array结合ctypes构建自定义的键值对结构,但实现复杂度较高,仅适合性能要求极高的场景。
关键注意事项
- 尽量避免进程间写共享数据,优先用父进程合并结果的方式,这是Python多进程中最高效的模式;
- 只读大型数据在Unix/Linux下利用fork的写时复制特性,无需额外处理即可实现共享;
- 跨平台场景下,使用
shared_memory代替fork的COW特性。
内容的提问来源于stack exchange,提问作者GoneAsync
相关产品推荐
相关产品推荐

