大数据集按单列分组聚合多列求和内核重启,求Vaex或Pandas解决方案
问题原因说明
- 第一版代码参数逻辑错误:
groupby方法传入的参数是分组依据的列,你传入了所有45列的列名,相当于按所有列去重分组,计算量指数级上升直接撑爆内存。 - 第二版代码冗余度过高:循环对单列做分组聚合再join,重复生成全量分组中间表,产生大量不必要的内存占用,超过阈值触发内核崩溃。
可行实现方案
Vaex 实现(内存友好,优先推荐)
Vaex默认采用懒计算机制,不会一次性加载全量数据到内存,一次性指定所有聚合规则即可,无需循环和join:
# 指定分组键为MSISDN group_key = "MSISDN" # 筛选出所有需要求和的非分组列 agg_columns = [col for col in data_df.column_names if col != group_key] # 构造所有列的求和聚合规则 agg_rules = {col: vaex.agg.sum(col) for col in agg_columns} # 一次性完成分组聚合 df_result = data_df.groupby(group_key, agg=agg_rules) # 若需要落地结果可以执行导出,否则直接使用df_result即可 # df_result.export_parquet("sum_result.parquet")
180万行、45列的规模用上述代码,32GB内存完全可以流畅运行。
Pandas 实现
如果需要用Pandas完成,优先用原生的分组聚合逻辑,内存紧张时可以采用分块处理:
普通场景(内存足够)
import pandas as pd group_key = "MSISDN" # 直接按指定列分组,对所有数值列求和 df_result = data_df.groupby(group_key, as_index=False).sum(numeric_only=True)
内存紧张场景(分块处理)
import pandas as pd group_key = "MSISDN" chunk_size = 100000 # 每批处理10万行,可根据实际内存调整 final_result = pd.DataFrame() # 按块读取文件逐批处理 for chunk in pd.read_csv("你的数据文件路径", chunksize=chunk_size): # 先对当前块做预聚合 chunk_agg = chunk.groupby(group_key, as_index=False).sum(numeric_only=True) # 合并到全局结果后再次聚合,消除分块导致的分组键重复 final_result = pd.concat([final_result, chunk_agg], ignore_index=True) \ .groupby(group_key, as_index=False).sum(numeric_only=True)
内容的提问来源于stack exchange,提问作者aziz shaw
相关产品推荐
相关产品推荐

