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
相关产品推荐
相关产品推荐

