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

如何向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,但实现复杂度较高,适合对性能要求极高且需要跨平台的场景。

大致实现思路

  1. 将big_cache中的集合、字典序列化为字节数据。
  2. 创建共享内存区域,写入序列化后的数据。
  3. 子进程从共享内存读取数据,反序列化为本地结构。
  4. 使用完毕后释放共享内存。

这种方法需要处理序列化/反序列化的开销,以及共享内存的管理,适合熟悉底层内存操作的场景。


内容的提问来源于stack exchange,提问作者Blake-Zeros

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.07 07:50:17