如何用concurrent.futures实现Pandas DataFrame的apply并行处理?
使用
concurrent.futures.ProcessPoolExecutor替代df.apply实现并行处理 针对你的需求,核心思路是将df.apply的逐行计算任务拆解为独立参数列表,通过多进程批量执行后,将结果回写至DataFrame。以下是具体实现方案和注意事项:
核心实现步骤
1. 提前整理可序列化的参数列表
多进程要求任务函数和参数必须能被pickle序列化,因此避免直接传递类实例(self),提前提取所需参数:
import os from concurrent.futures import ProcessPoolExecutor # 提取类实例中需要用到的参数,避免传递整个self p1 = self.param1 p2 = self.param2 p3_abs = abs(self.param3) # 生成参数列表:每个元素对应do_stuff的一组入参 # 假设你的do_stuff需要结合DataFrame行数据(比如row.colA、row.colB),用itertuples比iterrows更高效 params = [ ((p1, p2), p3_abs, row.colA, row.colB) for row in df.itertuples(index=False) ]
2. 定义辅助执行函数(可选但推荐)
如果do_stuff的参数是元组形式,定义一个顶层辅助函数解包参数,避免多进程中lambda序列化问题:
def run_do_stuff(args): # 解包参数并调用原函数 return do_stuff(*args)
3. 并行执行并回写结果
用ProcessPoolExecutor批量执行任务,利用map保证结果顺序与DataFrame行顺序一致:
# 初始化进程池,max_workers建议设为CPU核心数(os.cpu_count()) with ProcessPoolExecutor(max_workers=os.cpu_count()) as executor: # 批量执行任务,转换为列表获取结果 results = list(executor.map(run_do_stuff, params)) # 将结果赋值给DataFrame新列 df['new_column'] = results
关键注意事项
- 优先用ProcessPool而非ThreadPool:你的任务是CPU密集型(每行计算耗时2-5分钟),ThreadPool仅适用于IO密集型场景,ProcessPool才能真正利用多核资源。
- 序列化问题规避:
- 不要直接传递类实例
self,提前提取所需属性作为独立参数; - 确保
do_stuff是模块级函数(而非类内方法),或改为静态方法,避免绑定实例导致的序列化失败。
- 不要直接传递类实例
- 效率优化:用
df.itertuples()替代df.iterrows()遍历行,前者性能更优,适合大DataFrame。 - 进程数控制:
max_workers不建议超过CPU核心数,避免进程切换开销抵消并行收益。
内容的提问来源于stack exchange,提问作者new_to_code
相关产品推荐
相关产品推荐

