多进程并行化for循环:如何关联对象与对应分数?
使用Multiprocessing并行计算时如何对应结果与原对象?
问题描述
我有一个函数get_score,会按顺序处理对象列表并返回每个对象的分数,串行代码如下:
def get_score(a): # do something return score objects = [obj0, obj1, obj3] results = np.zeros(len(objects)) for i in range(len(results)): results[i] = get_score(objects[i])
我希望用Multiprocessing库并行执行该函数,但有个疑问:没有共享结果列表的情况下,怎么确定某个分数对应哪个对象?
解决方案
核心思路是让每个并行任务携带原对象的索引信息,或者利用某些并行方法的顺序保证特性,以下是几种实用实现方式:
1. 给任务添加索引,返回(索引, 分数)元组
修改任务函数,让它接收包含索引和对象的参数,返回带索引的结果,最后按索引映射到结果列表:
import multiprocessing as mp import numpy as np def get_score(a): # 原有的分数计算逻辑 return score # 包装函数,接收(索引, 对象)元组 def get_score_with_index(args): idx, obj = args return (idx, get_score(obj)) objects = [obj0, obj1, obj3] results = np.zeros(len(objects)) if __name__ == '__main__': with mp.Pool() as pool: # 把索引和对象打包成任务列表 tasks = list(enumerate(objects)) # 并行执行并获取带索引的结果 indexed_results = pool.map(get_score_with_index, tasks) # 根据索引填充结果 for idx, score in indexed_results: results[idx] = score
2. 利用imap保持结果顺序
Pool.imap方法会严格按照任务提交的顺序返回结果,无需额外传递索引,直接按顺序赋值即可:
import multiprocessing as mp import numpy as np def get_score(a): # 原有的分数计算逻辑 return score objects = [obj0, obj1, obj3] results = np.zeros(len(objects)) if __name__ == '__main__': with mp.Pool() as pool: # imap返回结果的顺序与输入objects的顺序完全一致 for idx, score in enumerate(pool.imap(get_score, objects)): results[idx] = score
注意:如果用
imap_unordered,结果会按任务完成顺序返回,这时候还是需要用带索引的方式。
3. 用apply_async的回调函数绑定索引
通过apply_async提交任务时,给回调函数绑定当前对象的索引,回调函数直接把结果写入对应位置:
import multiprocessing as mp import numpy as np def get_score(a): # 原有的分数计算逻辑 return score objects = [obj0, obj1, obj3] results = np.zeros(len(objects)) if __name__ == '__main__': with mp.Pool() as pool: for idx, obj in enumerate(objects): # 提交任务时绑定索引,回调函数负责赋值 pool.apply_async( get_score, args=(obj,), callback=lambda score, i=idx: results.__setitem__(i, score) ) pool.close() pool.join()
内容的提问来源于stack exchange,提问作者ZchGarinch
相关产品推荐
相关产品推荐

