Python multiprocessing.Array 只读共享数据读取性能优化咨询
问题根因
- 你用
multiprocessing.Manager.Array慢是因为Manager本质是单独的守护进程,所有对Array的访问都要跨进程和Manager通信,你构造bytearray时相当于逐个字节请求Manager返回数据,1MB数据就要发起上百万次IPC调用,自然耗时极高。 - 直接用
multiprocessing.Array报错是因为原生共享内存对象默认只能通过进程创建时的继承关系传递,不能通过Pool任务参数序列化传递,所以触发了运行时错误。
解决方案
方案1:fork+全局变量共享原生共享内存(性能最高,适合macOS/Linux)
macOS/Linux下Python多进程默认用fork模式创建子进程,你可以在创建Pool之前先把序列化后的数据存入multiprocessing.Array,赋值给全局变量,子进程fork后会自动继承这个共享内存的访问权限,不需要传参,直接访问即可。
修改后的示例代码:
import marshal import os from multiprocessing import Pool, Array, get_context from time import time from timeit import timeit from pympler.asizeof import asizeof # 全局变量用于继承共享内存 shared_data_buf = None def marshal_worker(): pid = os.getpid() print(f" -> [pid {pid}] Marshal worker ready at {time()}") t = time() # 直接访问全局的共享内存,用切片直接转字节,不需要逐个元素复制 ba = bytes(shared_data_buf) print(f" -> [pid {pid}] Building bytearray took {time() - t} s") t = time() data = marshal.loads(ba) print(f" -> [pid {pid}] Marshal loads() took {time() - t} s") return len(data) def test_marshal_ipc(data): global shared_data_buf print("Running test_marshal_ipc():") n_processes = os.cpu_count() - 1 or 1 # 提前序列化数据 data_bytes = marshal.dumps(data) # 创建原生共享内存,只读场景不需要加锁 shared_data_buf = Array('B', data_bytes, lock=False) # 用fork上下文创建Pool,确保子进程继承全局变量 with get_context('fork').Pool(processes=n_processes) as pool: async_results = [pool.apply_async(marshal_worker) for _ in range(n_processes)] subprocess_results = [res.get() for res in async_results] return subprocess_results # 其余测试代码保持不变
这个方案可以完全避免重复序列化、重复传输数据的开销,所有子进程直接读取同一份共享内存,性能会比原始的标准IPC方案高10倍以上。
方案2:fork+直接共享原始数据(无序列化开销)
因为你的数据是只读的,fork模式下操作系统会用写时复制机制,只要子进程不修改数据,所有子进程会直接共享父进程的原始数据物理内存,连marshal序列化/反序列化的开销都可以完全省掉,是性能最高的方案:
import os from multiprocessing import Pool, get_context from time import time from timeit import timeit from pympler.asizeof import asizeof raw_data = None def raw_worker(): pid = os.getpid() print(f" -> [pid {pid}] Worker ready at {time()}") # 直接使用原始数据,没有任何序列化/反序列化开销 return len(raw_data) def test_raw_ipc(data): global raw_data raw_data = data n_processes = os.cpu_count() - 1 or 1 with get_context('fork').Pool(processes=n_processes) as pool: async_results = [pool.apply_async(raw_worker) for _ in range(n_processes)] return [res.get() for res in async_results] # 其余测试代码保持不变
1~2GB的数据用这个方案,子进程启动后可以立刻开始计算,完全没有等待时间。
方案3:mmap共享文件(兼容spawn启动模式,跨平台)
如果你需要用spawn启动模式(比如Windows,或者macOS强制用spawn),可以把序列化后的数据写入临时文件,子进程用只读mmap映射该文件,系统会自动让所有进程共享同一份物理缓存,不会重复加载文件内容:
import marshal import os import mmap import tempfile from multiprocessing import Pool from time import time def mmap_worker(file_path, data_len): pid = os.getpid() print(f" -> [pid {pid}] Worker ready at {time()}") with open(file_path, 'rb') as f: # 只读映射文件 with mmap.mmap(f.fileno(), length=data_len, access=mmap.ACCESS_READ) as mm: t = time() data = marshal.loads(mm) print(f" -> [pid {pid}] Load data took {time() - t} s") return len(data) def test_mmap_ipc(data): n_processes = os.cpu_count() - 1 or 1 data_bytes = marshal.dumps(data) # 创建临时文件,使用后自动删除 with tempfile.NamedTemporaryFile(delete=False) as f: f.write(data_bytes) temp_path = f.name try: with Pool(processes=n_processes) as pool: async_results = [pool.apply_async(mmap_worker, (temp_path, len(data_bytes))) for _ in range(n_processes)] return [res.get() for res in async_results] finally: os.unlink(temp_path)
内容的提问来源于stack exchange,提问作者Sirion
相关产品推荐
相关产品推荐

