You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.26 15:23:30