如何向concurrent.futures.ProcessPoolExecutor共享全局超大缓存字典?
问题背景
我有一个包含90万个字典对象(position)的列表,每个对象需要引用一个大型字典big_cache来计算值。big_cache包含至少50个键,对应的值是集合或字典,每个集合/字典至少有8.5万个元素,作为查询缓存使用。现在用concurrent.futures.ProcessPoolExecutor处理这些position时,每个子进程都会复制一份big_cache,导致内存占用过高,想知道如何在不复制的前提下让所有进程共享这个缓存。
原代码示例(已修正笔误)
big_cache = { 'A_LOOKUP': {i for i in range(85000)}, # 模拟8.5万个元素的集合 'B_LOOKUP': {i for i in range(10000)}, # 模拟1万个元素的集合 'C_MAP': {f'K{i}': i for i in range(2500)} # 模拟2500个键值对的字典 } positions = [ {'pos_id': 1, 'pos_type': 'a', 'CCY': 'USD', 'KEY': 'K100'}, {'pos_id': 2, 'pos_type': 'b', 'CCY': 'EUR', 'KEY': 'K200'}, # ... 共90万个类似字典 {'pos_id': 900000, 'pos_type': 'a', 'CCY': 'GBP', 'KEY': 'K2000'} ] def process_position(position): if position['pos_type'] in big_cache['A_LOOKUP']: position['VALUE'] = big_cache['C_MAP'].get(position['KEY']) if position['CCY'] in big_cache['B_LOOKUP']: position['VALUE'] = big_cache['C_MAP'].get(position['KEY']) return position def execute(positions_list): res = [] with concurrent.futures.ProcessPoolExecutor() as executor: futures = executor.map(process_position, positions_list) for fut in futures: res.append(fut) return res
解决方案
方法1:利用Unix系统的写时复制(COW)特性(推荐)
在Unix/Linux/macOS系统中,ProcessPoolExecutor默认使用fork方式创建子进程。fork会让子进程共享父进程的内存页,只有当子进程修改内存内容时才会复制对应的页。因此只要保证big_cache是只读的,所有子进程就会共享同一份内存,不会产生复制。
改进代码
无需修改核心逻辑,只需确保big_cache在进程池创建前完成初始化,且子进程不修改它:
import concurrent.futures # 全局初始化big_cache,确保在进程池创建前完成 big_cache = { 'A_LOOKUP': {i for i in range(85000)}, 'B_LOOKUP': {i for i in range(10000)}, 'C_MAP': {f'K{i}': i for i in range(2500)} } positions = [ {'pos_id': 1, 'pos_type': 'a', 'CCY': 'USD', 'KEY': 'K100'}, # ... 90万个元素 ] def process_position(position): if position['pos_type'] in big_cache['A_LOOKUP']: position['VALUE'] = big_cache['C_MAP'].get(position['KEY']) if position['CCY'] in big_cache['B_LOOKUP']: position['VALUE'] = big_cache['C_MAP'].get(position['KEY']) return position def execute(positions_list): res = [] # 进程池在big_cache初始化后创建,子进程共享父进程的big_cache内存 with concurrent.futures.ProcessPoolExecutor() as executor: for result in executor.map(process_position, positions_list): res.append(result) return res
注意:如果子进程修改big_cache,会触发写时复制,导致内存被复制,失去共享效果,因此必须确保缓存是只读的。
方法2:使用multiprocessing.Manager(兼容Windows)
Windows系统不支持fork,默认用spawn方式创建进程,父进程的全局变量不会自动传递给子进程。可以用multiprocessing.Manager创建共享的字典和集合,让所有进程通过代理访问。
改进代码
import concurrent.futures from multiprocessing import Manager def init_worker(cache): # 将共享缓存设置为全局变量,供子进程使用 global big_cache big_cache = cache def process_position(position): if position['pos_type'] in big_cache['A_LOOKUP']: position['VALUE'] = big_cache['C_MAP'].get(position['KEY']) if position['CCY'] in big_cache['B_LOOKUP']: position['VALUE'] = big_cache['C_MAP'].get(position['KEY']) return position def execute(positions_list): # 创建Manager,生成共享的缓存结构 with Manager() as manager: # 复制原缓存到共享对象中 big_cache = manager.dict() big_cache['A_LOOKUP'] = manager.set({i for i in range(85000)}) big_cache['B_LOOKUP'] = manager.set({i for i in range(10000)}) big_cache['C_MAP'] = manager.dict({f'K{i}': i for i in range(2500)}) res = [] # 用initializer和initargs将共享缓存传递给子进程 with concurrent.futures.ProcessPoolExecutor(initializer=init_worker, initargs=(big_cache,)) as executor: for result in executor.map(process_position, positions_list): res.append(result) return res
注意:Manager的共享对象通过进程间通信(IPC)访问,性能比COW方式低,适合必须兼容Windows的场景。如果缓存查询频繁,这种方法可能会成为性能瓶颈。
方法3:使用共享内存(Python 3.8+)
对于超大型缓存,可以将数据序列化后存入multiprocessing.shared_memory,每个进程从共享内存中读取并反序列化。这种方式性能接近COW,但实现复杂度较高,适合对性能要求极高且需要跨平台的场景。
大致实现思路
- 将
big_cache中的集合、字典序列化为字节数据。 - 创建共享内存区域,写入序列化后的数据。
- 子进程从共享内存读取数据,反序列化为本地结构。
- 使用完毕后释放共享内存。
这种方法需要处理序列化/反序列化的开销,以及共享内存的管理,适合熟悉底层内存操作的场景。
内容的提问来源于stack exchange,提问作者Blake-Zeros

