如何在Python fork进程间共享大型数据结构?
我在Python程序中创建了一个58 GB的图结构,想通过多个Linux进程处理它。我用了带fork上下文的ProcessPoolExecutor,只给每个进程传处理块对应的起始整数:
import concurrent.futures import multiprocessing as mp with concurrent.futures.ProcessPoolExecutor(mp_context=mp.get_context('fork')) as executor: futures = [] for start in range(0, len(vertices_list), BATCH_SIZE): future = executor.submit(process_batch, start) futures += [future]
process_batch函数负责计算图顶点的指标,返回一个小型结果列表。vertices_list和graph都是全局变量,graph绑定到一个实现了所需数据结构和计算逻辑的C++ Python扩展模块。
def process_batch(start): results = [] for v in vertices_list[start:start + BATCH_SIZE]: results.append((v, graph.some_metric(v))) return results
处理小样本时一切正常,但处理58GB的结构时,机器直接卡死,内存被大量换出,通过strace能看到fork进程以1/4 MB的块调用mmap来申请更多内存。查资料发现,Python中fork的“写时复制”实际上会变成“访问时复制”——因为只要访问对象就会修改其引用计数,这和我观察到的行为一致:fork进程似乎在尝试创建整个数据结构的私有副本。
我的问题是:Python有没有办法在fork进程之间共享任意数据结构?我知道multiprocessing.shared_memory及相关工具,但它们只支持内置的Array和Value类型。
利用C++扩展模块的内存控制
你的graph是C++扩展实现的,这是关键突破口。可以修改扩展代码,让底层的图数据结构存放在匿名共享内存或者内存文件映射中,并且确保Python层的对象引用不会触发引用计数修改导致的写时复制:- 在扩展初始化时,把图数据加载到共享内存区域(比如用
mmap创建共享映射),让所有进程的graph对象指向这块共享内存; - 保证Python层对
graph的访问只做只读操作,并且避免修改Python对象的引用计数(比如把graph设计成不可变对象,或者在C++层不触发Python的引用计数变更)。
- 在扩展初始化时,把图数据加载到共享内存区域(比如用
使用
forkserver而非fork上下文forkserver上下文会先启动一个单独的服务器进程,后续的工作进程都从这个服务器进程fork出来。服务器进程可以预先加载好图数据,并且在fork后不做任何修改,这样工作进程就能通过写时复制共享内存,而不会因为引用计数修改触发复制。
修改代码中的上下文:with concurrent.futures.ProcessPoolExecutor(mp_context=mp.get_context('forkserver')) as executor: # 后续逻辑不变注意:
forkserver需要在主进程启动时就设置好,并且主进程不能在forkserver启动后再加载大内存数据,否则还是会有问题。正确的做法是:先初始化forkserver,再加载图数据,最后启动进程池。避免全局变量触发的隐式引用计数修改
不要把vertices_list和graph作为全局变量,而是在fork后的子进程中直接从共享内存重新获取或者通过只读的方式访问。比如,把graph的底层数据放在共享内存,子进程通过扩展模块的接口直接访问这块内存,而不是通过Python层的全局变量引用。内存映射文件替代内存中的数据结构
如果可以把图数据预先存储在磁盘文件中,那么所有进程都可以通过mmap把文件映射到自己的地址空间,实现真正的只读共享。这种方式下,进程间共享的是磁盘文件的映射,不会因为Python的引用计数问题触发复制。你可以修改C++扩展模块,让它从内存映射文件中加载图数据,而不是从内存中创建。
内容的提问来源于stack exchange,提问作者Diomidis Spinellis

