Pandas groupby处理大型DataFrame速度过慢,求优化方案(无法用pandarallel)
优化大型DataFrame的groupby聚合速度
原问题代码
def group_func(group): school_open = (group['school_open'] == True) exam = (group['exam_scheduled'] == True) attendance_required = (group['att_flag'] == True) score_mask = school_open & exam attendance_mask = attendance_required & school_open score = group.loc[score_mask, 'ind_score'].mean() attendance = group.loc[attendance_mask, 'att'].mean() active_day = group[attendance_mask]['dates'].nunique() median_score = group.loc[score_mask, 'ind_score'].median() return pd.Series({'score': score, 'attendance': attendance, 'active_day': active_day, 'median_score': median_score}) student_consolidated = student_df.groupby(['student_name', pd.Grouper(key='dates', freq='M')]).apply(group_func)
验证用DataFrame生成代码
import pandas as pd import numpy as np from faker import Faker fake = Faker() date_rng = pd.date_range(start='1/1/2023', end='12/31/2023', freq='D') data = {'student_name': [fake.name() for i in range(len(date_rng)*100)], 'dates': np.tile(date_rng, 100), 'school_open': np.random.choice([True, False], size=len(date_rng)*100), 'att_flag': np.random.choice([True, False], size=len(date_rng)*100), 'exam_scheduled': np.random.choice([0, 1], size=len(date_rng)*100), 'ind_score': np.random.randint(1, 30, size=len(date_rng)*100), 'att': np.random.choice([True, False], size=len(date_rng)*100)} student_df = pd.DataFrame(data)
优化方案
1. 用向量化聚合替代apply(最有效)
apply是逐组Python循环,速度极慢。我们可以提前计算掩码列,用where过滤无关行后,直接调用pandas内置的C实现聚合函数,速度能提升一个数量级以上:
# 提前计算掩码列 student_df['score_mask'] = student_df['school_open'] & student_df['exam_scheduled'].astype(bool) student_df['attendance_mask'] = student_df['att_flag'] & student_df['school_open'] # 过滤生成聚合用列(NaN会被聚合函数自动忽略) student_df['score_col'] = student_df['ind_score'].where(student_df['score_mask']) student_df['attendance_col'] = student_df['att'].where(student_df['attendance_mask']) student_df['active_day_col'] = student_df['dates'].where(student_df['attendance_mask']) # 分组聚合 student_consolidated = student_df.groupby( ['student_name', pd.Grouper(key='dates', freq='M')] ).agg( score=('score_col', 'mean'), attendance=('attendance_col', 'mean'), active_day=('active_day_col', 'nunique'), median_score=('score_col', 'median') ).reset_index()
2. 用Numba加速自定义聚合逻辑
如果必须保留自定义逻辑,用Numba编译聚合函数,比纯Python的apply快数倍:
from numba import jit import numpy as np @jit(nopython=True) def numba_group_func(school_open, exam_scheduled, att_flag, ind_score, att, dates): score_mask = school_open & exam_scheduled attendance_mask = att_flag & school_open # 计算分数相关指标 score_vals = ind_score[score_mask] score = score_vals.mean() if len(score_vals) > 0 else np.nan median_score = np.median(score_vals) if len(score_vals) > 0 else np.nan # 计算出勤率 att_vals = att[attendance_mask] attendance = att_vals.mean() if len(att_vals) > 0 else np.nan # 计算活跃天数(Numba不直接支持nunique,手动实现) date_vals = dates[attendance_mask] active_day = len(np.unique(date_vals)) if len(date_vals) > 0 else 0 return score, attendance, active_day, median_score # 包装函数适配pandas分组 def wrapper_func(group): return numba_group_func( group['school_open'].values, group['exam_scheduled'].values.astype(bool), group['att_flag'].values, group['ind_score'].values, group['att'].values.astype(float), group['dates'].astype(np.int64).values ) # 执行分组计算 student_consolidated = student_df.groupby( ['student_name', pd.Grouper(key='dates', freq='M')] ).apply(wrapper_func, include_groups=False).apply(pd.Series, index=['score', 'attendance', 'active_day', 'median_score'])
3. 优化数据类型减少内存占用
大型DataFrame的内存压力会直接拖慢处理速度,先做类型优化:
# 布尔类型转int8(比原生bool更节省内存,不影响逻辑计算) student_df['school_open'] = student_df['school_open'].astype('int8') student_df['att_flag'] = student_df['att_flag'].astype('int8') student_df['exam_scheduled'] = student_df['exam_scheduled'].astype('int8') student_df['att'] = student_df['att'].astype('int8') # 分数列范围小,转int8 student_df['ind_score'] = student_df['ind_score'].astype('int8')
4. 用Dask进行分布式并行计算
如果数据量超过内存上限,或者想进一步提升并行效率,用Dask自动分块处理:
import dask.dataframe as dd # 转换为Dask DataFrame,分区数根据CPU核心数设置 ddf = dd.from_pandas(student_df, npartitions=4) # 定义聚合规则 agg_dict = { 'score_col': ['mean', 'median'], 'attendance_col': ['mean'], 'active_day_col': ['nunique'] } # 执行并行聚合并转为pandas DataFrame result = ddf.groupby( ['student_name', dd.Grouper(key='dates', freq='M')] ).agg(agg_dict).compute() # 重命名列并整理结构 result.columns = ['score', 'median_score', 'attendance', 'active_day'] result = result.reset_index()
内容的提问来源于stack exchange,提问作者Lata
相关产品推荐
相关产品推荐

