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)组合数量不是特别夸张(比如千万级以内),可以轻松处理:
- 保存以下脚本为
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] } } }
- 运行命令:
awk -f process.awk Inpt.txt > Op.txt
优势:内存占用极低(仅存聚合结果),处理速度是Python方案的5-10倍,适合纯数据聚合场景。
方案选择建议
- 如果你想继续用Python生态,优先选方案1,调整
chunksize到适合你内存的大小; - 如果数据集特别大(比如超过10GB)或需要并行加速,选方案2的Dask;
- 如果追求极致性能和低内存,选方案3的Awk。
内容的提问来源于stack exchange,提问作者Marek
相关产品推荐
相关产品推荐

