大内存外数据集按function_name去重trace遇OOM,求高效解法
高效处理超大数据集的分组去重方案
针对你遇到的内存不足问题,核心原因是直接调用unique(["function_name", "trace"])需要在内存中维护所有字段组合的去重集合,而长字符串trace会大幅拉高内存占用。以下是两种高效且避免OOM的解决方法:
方法一:用分组聚合替代原生unique
Polars的group_by+agg组合可以分批处理分组数据,避免一次性加载全量数据到内存,同时实现和unique等价的去重效果(保留每个function_name+trace组合的第一条记录):
import polars as pl from pathlib import Path # 调整Polars内存配置,根据你的机器内存设置合适的内存池大小 pl.Config.set_memory_pool_size(8 * 1024 ** 3) # 示例:设置为8GB pl.Config.set_tbl_rows(0) # 关闭预览数据的内存占用 # 懒加载数据集 df = pl.scan_parquet(Path("...", "dataset.parquet")) # 按function_name和trace分组,保留每个组合的第一条完整记录 df = df.group_by(["function_name", "trace"]).agg(pl.all().first()) # 流式写入结果 df.sink_parquet(Path("...", "output.parquet"))
方法二:预哈希trace字段减少内存开销
如果trace是极长字符串,直接存储和比较会消耗大量内存,可以先对trace计算哈希值(固定长度整数),用哈希值辅助去重,大幅降低内存占用:
import polars as pl from pathlib import Path pl.Config.set_memory_pool_size(8 * 1024 ** 3) df = pl.scan_parquet(Path("...", "dataset.parquet")) # 计算trace的哈希值(固定种子保证重复trace得到相同哈希) df = df.with_columns( trace_hash=pl.hash(pl.col("trace"), seed=42) ).group_by(["function_name", "trace_hash", "trace"]).agg( pl.all().exclude("trace_hash").first() ).drop("trace_hash") # 移除临时哈希字段 df.sink_parquet(Path("...", "output.parquet"))
关键说明
- 分组聚合的方式利用Polars的懒执行引擎,分批处理每个分组,无需一次性加载全量数据;
- 哈希优化通过将长字符串转化为固定长度整数,把内存占用降低几个数量级,同时固定种子的哈希能保证重复
trace的去重准确性(哈希碰撞概率极低,可忽略); - 若无需保留其他字段,可将
agg(pl.all().first())改为agg(pl.col("function_name").first(), pl.col("trace").first()),进一步减少内存开销。
内容的提问来源于stack exchange,提问作者Henry
相关产品推荐
相关产品推荐

