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

Python多进程中如何实现类似C++指针的共享数据结构访问?

Python多进程中高效共享大型映射表的方案

问题场景

需要用多进程拆分计算任务,核心需求有两个:

  1. 子进程需高效访问只读的大型映射表(含数万浮点数的numpy向量),避免数据复制带来的内存开销;
  2. 子进程需将计算结果写入共享的映射表,且不同进程不会写入同一键,无需同步锁机制。

当前使用普通全局字典时,子进程会复制整个字典,更新仅作用于副本;使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.09 21:01:40