如何在Python多进程中获取函数返回值并构建字典
用multiprocessing并行处理并获取返回值的几种方案
针对你的场景,这里提供几种用Python multiprocessing 并行执行任务并收集返回值的实用方案,直接适配你的代码逻辑:
方案一:用Pool.map/Starmap(最简洁)
Pool 是multiprocessing中最常用的并行工具,适合批量处理同类型任务。因为你的函数需要两个参数,有两种方式传递:
方式1:用functools.partial绑定固定参数
把df绑定到函数上,让函数只接收name参数,再用map批量处理:
import pandas as pd from multiprocessing import Pool from functools import partial df = pd.DataFrame({'Matthew': [4, 9, 6], 'Mark': [2, 3, 5], 'Luke': [10, 1, 8], 'John': [20, 22, 21]}) def sum_funct(name, df): return int(df[name].sum()) names = ['Matthew', 'Mark', 'Luke', 'John'] if __name__ == '__main__': # 绑定df到sum_funct,生成只需要name的新函数 bound_sum = partial(sum_funct, df=df) # 创建进程池,默认使用CPU核心数 with Pool() as pool: # 并行执行,结果顺序和names一致 results = pool.map(bound_sum, names) # 转成目标字典 totals_dict = dict(zip(names, results)) print(totals_dict)
方式2:用Starmap直接传递多参数元组
把每个任务的参数打包成元组,用starmap自动解包参数:
import pandas as pd from multiprocessing import Pool df = pd.DataFrame({'Matthew': [4, 9, 6], 'Mark': [2, 3, 5], 'Luke': [10, 1, 8], 'John': [20, 22, 21]}) def sum_funct(name, df): return int(df[name].sum()) names = ['Matthew', 'Mark', 'Luke', 'John'] if __name__ == '__main__': # 构造每个任务的参数元组 tasks = [(name, df) for name in names] with Pool() as pool: results = pool.starmap(sum_funct, tasks) totals_dict = dict(zip(names, results)) print(totals_dict)
方案二:用Pool.apply_async(异步获取结果)
如果需要跟踪每个任务对应的name,或者想随时获取已完成的结果,用apply_async更灵活:
import pandas as pd from multiprocessing import Pool df = pd.DataFrame({'Matthew': [4, 9, 6], 'Mark': [2, 3, 5], 'Luke': [10, 1, 8], 'John': [20, 22, 21]}) def sum_funct(name, df): return int(df[name].sum()) names = ['Matthew', 'Mark', 'Luke', 'John'] if __name__ == '__main__': totals_dict = {} with Pool() as pool: # 提交所有任务,保存每个任务的future对象 futures = {name: pool.apply_async(sum_funct, args=(name, df)) for name in names} # 逐个获取结果 for name, future in futures.items(): totals_dict[name] = future.get() print(totals_dict)
方案三:用Process+Queue(底层手动控制)
如果需要更精细地管理进程生命周期,用Process配合Queue传递结果:
import pandas as pd from multiprocessing import Process, Queue df = pd.DataFrame({'Matthew': [4, 9, 6], 'Mark': [2, 3, 5], 'Luke': [10, 1, 8], 'John': [20, 22, 21]}) def sum_funct(name, df): return int(df[name].sum()) # 定义工作进程函数,把结果放入队列 def worker(name, df, queue): result = sum_funct(name, df) queue.put((name, result)) names = ['Matthew', 'Mark', 'Luke', 'John'] if __name__ == '__main__': queue = Queue() processes = [] totals_dict = {} # 创建并启动所有进程 for name in names: p = Process(target=worker, args=(name, df, queue)) processes.append(p) p.start() # 等待所有进程执行完毕 for p in processes: p.join() # 从队列取出所有结果 while not queue.empty(): name, result = queue.get() totals_dict[name] = result print(totals_dict)
关键注意事项
- 必须加
if __name__ == '__main__'::这是Windows系统多进程的强制要求,防止子进程重复执行主模块代码导致无限递归创建进程。 - 大数据框的内存问题:多进程会复制
df到每个子进程的内存空间,如果df非常大,会占用大量内存。这种情况可以考虑让每个子进程从文件读取df,或者用共享内存(比如multiprocessing.Array)优化。 - 进程池大小:
Pool()默认使用CPU核心数,如果你的任务是IO密集型(比如涉及文件读写、网络请求),可以适当调大进程数;如果是CPU密集型,保持核心数即可,避免上下文切换开销。
内容的提问来源于stack exchange,提问作者top bantz
相关产品推荐
相关产品推荐

