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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.23 09:22:27