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

Python处理超大规模数据集groupby操作的内存溢出问题求助

嘿,针对你处理超大规模数据集(2亿条记录)时遇到的内存溢出问题,咱们来一步步拆解优化方案——核心思路是避免一次性加载全量数据到内存,同时简化冗余的聚合逻辑:

先分析原脚本的内存瓶颈

你的原脚本在小数据上能跑通,但面对2亿条记录时,几个关键步骤会直接撑爆内存:

  • 用read_csv一次性加载全量数据,瞬间占用几GB甚至几十GB内存
  • 多次groupby.transform生成额外的大列,进一步消耗内存
  • 中间的排序、去重步骤会产生数据副本,加剧内存压力

优化方案1:Pandas分块读取+逐步聚合

这是最贴近你原有代码的改进,利用Pandas的chunksize参数分批次处理数据,先完成底层聚合,再筛选目标记录:

import pandas as pd

# 第一步:分块读取,先聚合每个(CLASS_ID, COURSE_ID)的总费用和任意一个ID
chunk_aggregates = []
chunksize = 1_000_000  # 可根据内存调整,比如100万/块,内存小就调小

for chunk in pd.read_csv('Inpt.txt', dtype={'CLASS_ID': str}, chunksize=chunksize):
    # 对当前块做局部聚合:求和费用,保留第一个出现的ID(后续取最大COURSE_ID时不影响)
    chunk_grouped = chunk.groupby(['CLASS_ID', 'COURSE_ID'], as_index=False).agg(
        COURSE_FEE=('COURSE_FEE', 'sum'),
        ID=('ID', 'first')
    )
    chunk_aggregates.append(chunk_grouped)

# 合并所有块的聚合结果,再次全局聚合(同一个(CLASS_ID, COURSE_ID)可能跨块)
full_agg = pd.concat(chunk_aggregates).groupby(['CLASS_ID', 'COURSE_ID'], as_index=False).agg(
    COURSE_FEE=('COURSE_FEE', 'sum'),
    ID=('ID', 'first')
)

# 第二步:筛选每个CLASS_ID对应最大COURSE_ID的记录
full_agg['max_course_id'] = full_agg.groupby('CLASS_ID')['COURSE_ID'].transform('max')
result = full_agg[full_agg['COURSE_ID'] == full_agg['max_course_id']].drop('max_course_id', axis=1)

# 调整列顺序匹配预期输出
result = result[['ID', 'CLASS_ID', 'COURSE_ID', 'COURSE_FEE']]

# 输出结果
result.to_csv('Op.txt', index=False)

优势:分块加载避免一次性读全量,两次聚合(块内+块间)保证数据准确性,后续筛选逻辑极简,内存占用仅为原脚本的几十分之一。


优化方案2:用Dask处理超大规模数据

如果Pandas分块还是不够,Dask是专门为大数据设计的工具,语法和Pandas几乎一致,但能自动并行处理、拆分数据集,内存占用可控:

import dask.dataframe as dd

# 用Dask读取CSV,自动拆分成分区
ddf = dd.read_csv('Inpt.txt', dtype={'CLASS_ID': str})

# 第一步:按(CLASS_ID, COURSE_ID)聚合费用和ID
agg_ddf = ddf.groupby(['CLASS_ID', 'COURSE_ID']).agg(
    COURSE_FEE=('COURSE_FEE', 'sum'),
    ID=('ID', 'first')
).reset_index()

# 第二步:找到每个CLASS_ID的最大COURSE_ID
max_course = agg_ddf.groupby('CLASS_ID')['COURSE_ID'].max().reset_index(name='max_course_id')

# 合并筛选出符合条件的记录
result_ddf = agg_ddf.merge(max_course, on='CLASS_ID')
result_ddf = result_ddf[result_ddf['COURSE_ID'] == result_ddf['max_course_id']].drop('max_course_id', axis=1)

# 调整列顺序
result_ddf = result_ddf[['ID', 'CLASS_ID', 'COURSE_ID', 'COURSE_FEE']]

# 计算结果并保存
result_ddf.compute().to_csv('Op.txt', index=False)

优势:无需手动管理块大小,自动利用多核CPU并行处理,适合TB级别的数据集,性能比Pandas分块提升数倍。


优化方案3:用Awk流式处理(最快最省内存)

如果允许使用命令行工具,Awk是最佳选择——它逐行处理数据,内存仅存储聚合后的键值对,对于2亿条记录,只要不同的(CLASS_ID, COURSE_ID)组合数量不是特别夸张(比如千万级以内),可以轻松处理:

  1. 保存以下脚本为process.awk:
BEGIN {
    FS = ","
    OFS = ","
    print "ID,CLASS_ID,COURSE_ID,COURSE_FEE"
}

NR > 1 {
    # 记录每个(CLASS_ID, COURSE_ID)的总费用和第一个出现的ID
    key = $2 "," $4
    sum_fee[key] += $3
    if (!id[key]) id[key] = $1
    
    # 记录每个CLASS_ID的最大COURSE_ID
    if ($4 > max_course[$2]) {
        max_course[$2] = $4
    }
}

END {
    # 筛选并输出符合条件的记录
    for (key in sum_fee) {
        split(key, parts, ",")
        class_id = parts[1]
        course_id = parts[2]
        if (course_id == max_course[class_id]) {
            print id[key], class_id, course_id, sum_fee[key]
        }
    }
}
  1. 运行命令:
awk -f process.awk Inpt.txt > Op.txt

优势:内存占用极低(仅存聚合结果),处理速度是Python方案的5-10倍,适合纯数据聚合场景。


方案选择建议

  • 如果你想继续用Python生态,优先选方案1,调整chunksize到适合你内存的大小;
  • 如果数据集特别大(比如超过10GB)或需要并行加速,选方案2的Dask;
  • 如果追求极致性能和低内存,选方案3的Awk。

内容的提问来源于stack exchange,提问作者Marek

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.29 07:53:22