joblib并行化Pandas DataFrame分组迭代变慢的优化求助
优化方案
首先分析并行版本变慢的核心原因:原并行代码每次向子进程传递整个70万行的DataFrame,多进程模式下每个子进程都会复制一份完整数据,这种数据复制的开销远超过并行处理节省的时间,导致整体速度更慢。同时原串行代码的手动分组、循环concat也存在大量效率浪费。以下是针对性的优化步骤:
第一步:优先优化串行版本(效果立竿见影)
用Pandas原生的groupby.apply替代手动分组与循环,Pandas的groupby是C级优化实现,效率远高于手动筛选+循环concat:
import pandas as pd df = pd.DataFrame({"colA": ["x"]*5+["y"]*5, "colB": ["a"]*2+["b"]*2+["c"]*2+["d"]*2+["e"]*2, "colC": range(10)}) def example(df): df['colD'] = df['colC'] * 2 return df def my_optimized_serial(df): # group_keys=False避免保留分组索引,输出结构与原函数一致 return df.groupby(['colA', 'colB'], group_keys=False).apply(example) # 调用 my_optimized_serial(df)
第二步:优化并行版本(若串行仍无法满足需求)
核心是避免传递整个大DataFrame,提前拆分所有分组为子DataFrame,只将子数据传递给子进程,消除跨进程数据复制的巨大开销:
from joblib import Parallel, delayed def my_optimized_parallel(df, n_jobs): # 提前拆分所有分组为子DataFrame列表 group_list = [sub_df for _, sub_df in df.groupby(['colA', 'colB'])] # 并行处理每个子分组,仅传递子DataFrame result_list = Parallel(n_jobs=n_jobs, batch_size='auto')( delayed(example)(sub_df) for sub_df in group_list ) # 一次性合并所有结果,避免循环concat的开销 return pd.concat(result_list) # 调用 my_optimized_parallel(df, 2)
额外优化细节
batch_size='auto':让joblib自动合并小任务,减少进程调度的开销,尤其适合1400个分组这种数量较多的小任务场景。- 若你的
example函数涉及复杂操作,可检查是否能进一步向量化,避免在函数内使用循环,最大化Pandas的性能。
内容的提问来源于stack exchange,提问作者Luca Monno
相关产品推荐
相关产品推荐

