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

Python多进程处理分块字典时结果列表为空的问题

多进程处理字典时结果列表为空的问题解决

问题场景

拆分大字典为与CPU核心数等量的子字典,通过多进程并行提取每个子字典中的数据三元组,期望将各进程结果存入对应列表triplets_lists,但最终输出为空列表。

原代码

主进程代码:

if __name__ == '__main__':
    num_cores = multiprocessing.cpu_count()
    keys_list = list(largeDict.keys())
    random.shuffle(keys_list)
    subdict_size = len(keys_list) // num_cores
    subdicts = []   # 存储与核心数等量的子字典
    start_idx = 0
    triplets_lists = [[] for _ in range(num_cores)]   # 存储各进程的三元组列表
    processes = []
    for i in range(num_cores):
        end_idx = subdict_size * (i + 1)
        if i == (num_cores - 1):    # 最后一个核心处理剩余所有数据
            end_idx = len(keys_list)
        subdicts.append({key: largeDict[key] for key in keys_list[start_idx:end_idx]})
        start_idx = end_idx
        process = multiprocessing.Process(target=processing_records, args=(subdicts[i], triplets_lists[i]))
        processes.append(process)
        process.start()
    for process in processes:
        process.join()
    print(triplets_lists)   

子进程处理函数:

def processing_records(smDict, triplets_list):
    for key, value in smDict.items():
        for sentence in value:
            triplets = get_triplets(sentence)  # 提取单句中的三元组列表
            for triplet in triplets:
                triplets_list.append(triplet)

问题原因

多进程的内存隔离机制:每个子进程启动时会完整复制父进程的内存空间,传入子进程的triplets_lists[i]是子进程内存中的副本。子进程对该列表的append操作仅作用于自身副本,完全不会修改父进程中的原列表,因此父进程最终打印的还是初始状态的空列表。

解决方案

方案一:使用multiprocessing.Manager创建共享列表

通过Manager创建跨进程共享的列表对象,子进程对共享列表的修改会同步到父进程:

import multiprocessing
import random

def processing_records(smDict, triplets_list):
    for key, value in smDict.items():
        for sentence in value:
            triplets = get_triplets(sentence)
            for triplet in triplets:
                triplets_list.append(triplet)

if __name__ == '__main__':
    # 示例大字典与三元组提取函数,实际替换为你的业务逻辑
    largeDict = {"key1": ["sentence1", "sentence2"], "key2": ["sentence3"]}
    def get_triplets(sentence):
        return [(sentence, "predicate", "object")]

    num_cores = multiprocessing.cpu_count()
    keys_list = list(largeDict.keys())
    random.shuffle(keys_list)
    subdict_size = len(keys_list) // num_cores
    subdicts = []
    start_idx = 0
    
    # 使用Manager创建可跨进程共享的列表
    with multiprocessing.Manager() as manager:
        triplets_lists = [manager.list() for _ in range(num_cores)]
        processes = []
        
        for i in range(num_cores):
            end_idx = subdict_size * (i + 1)
            if i == (num_cores - 1):
                end_idx = len(keys_list)
            subdicts.append({key: largeDict[key] for key in keys_list[start_idx:end_idx]})
            start_idx = end_idx
            process = multiprocessing.Process(target=processing_records, args=(subdicts[i], triplets_lists[i]))
            processes.append(process)
            process.start()
        
        for process in processes:
            process.join()
        
        # 将共享列表转换为普通列表,方便后续处理
        triplets_lists = [list(lst) for lst in triplets_lists]
        print(triplets_lists)

方案二:使用multiprocessing.Pool的map方法(推荐)

Pool会自动处理任务分发、进程管理和结果收集,无需手动维护共享对象,让子进程直接返回处理后的三元组列表,父进程汇总结果即可,代码更简洁:

import multiprocessing
import random

def processing_records(smDict):
    triplets_list = []
    for key, value in smDict.items():
        for sentence in value:
            triplets = get_triplets(sentence)
            triplets_list.extend(triplets)
    return triplets_list

if __name__ == '__main__':
    # 示例大字典与三元组提取函数,实际替换为你的业务逻辑
    largeDict = {"key1": ["sentence1", "sentence2"], "key2": ["sentence3"]}
    def get_triplets(sentence):
        return [(sentence, "predicate", "object")]

    num_cores = multiprocessing.cpu_count()
    keys_list = list(largeDict.keys())
    random.shuffle(keys_list)
    subdict_size = len(keys_list) // num_cores
    subdicts = []
    start_idx = 0
    
    # 拆分大字典为子字典
    for i in range(num_cores):
        end_idx = subdict_size * (i + 1)
        if i == (num_cores - 1):
            end_idx = len(keys_list)
        subdicts.append({key: largeDict[key] for key in keys_list[start_idx:end_idx]})
        start_idx = end_idx
    
    # 使用Pool并行处理并收集结果
    with multiprocessing.Pool(num_cores) as pool:
        triplets_lists = pool.map(processing_records, subdicts)
    
    print(triplets_lists)

方案二避免了共享对象的同步问题,代码结构更清晰,是处理这类分块并行任务的首选方式。

内容的提问来源于stack exchange,提问作者diegobc11

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.24 04:35:22