如何优化Polars的unique()方法以降低内存占用?解决OOM问题
unique()/drop_nulls()内存溢出问题分析与解决方案 一、为什么自定义逐列过滤null能规避OOM?
Polars原生的drop_nulls(subset)是一次性检查所有指定列的null值,执行时需要同时加载所有目标列到内存并生成联合布尔掩码,叠加列本身的内存开销后,容易瞬间触发内存阈值。
而你的自定义方法是逐列递进过滤:
def iterative_drop_nulls(expr: pl.Expr, subset: list[str]) -> pl.LazyFrame: for col in subset: expr = expr.filter(~pl.col(col).is_null()) return expr
每一步只针对单列过滤,过滤后数据集逐步缩小,后续步骤处理的是已精简的数据,内存占用呈递减趋势,不会出现一次性加载多列+生成联合掩码的峰值开销,因此能避开OOM。
二、低内存unique()实现方案
针对unique()的内存问题,可尝试以下几种方案:
1. 分阶段分区去重(流式处理优化)
利用Polars的分区扫描能力,先对每个分区单独去重,再对全局结果二次去重,避免一次性加载全量数据生成哈希表:
df_unique = ( pl.scan_parquet("data.parquet") .map_batches(lambda batch: batch.unique(subset=["col1", "col2"])) .unique(subset=["col1", "col2"]) .collect(streaming=True) )
map_batches会在每个数据分区上单独执行去重,先缩小各分区的数据量,最后全局去重仅处理分区去重后的结果,内存压力大幅降低。
2. 哈希分桶手动拆分处理
如果上述方法仍有压力,可手动按目标列的哈希值拆分数据集,分批次去重后合并:
def hash_bucket_unique(df_lazy: pl.LazyFrame, subset: list[str], num_buckets: int = 10) -> pl.DataFrame: # 添加临时哈希桶列 df_with_bucket = df_lazy.with_columns( (pl.hash(*subset) % num_buckets).alias("_bucket") ) # 按桶分批处理 parts = [] for bucket in range(num_buckets): batch_unique = ( df_with_bucket.filter(pl.col("_bucket") == bucket) .drop("_bucket") .unique(subset=subset) .collect() ) parts.append(batch_unique) # 合并后二次去重(避免跨桶重复) return pl.concat(parts).unique(subset=subset) # 调用示例 df_unique = hash_bucket_unique(pl.scan_parquet("data.parquet"), ["col1", "col2"])
通过哈希分桶把数据拆分成多个小批次,每个批次单独去重,最后全局再做一次去重确保无遗漏,内存占用被分摊到各个小批次。
3. 升级Polars版本
你当前使用的是0.18.4版本,Polars后续版本(如0.20+)对unique()的内存效率做了大量优化,包括流式处理下的哈希表内存控制、分区去重的自动适配,升级后原生unique().collect(streaming=True)可能直接解决问题。
三、关于Pandas中转方案的补充
你提到普通Pandasdrop_duplicates可行,但PyArrow扩展数组方式会OOM,原因是PyArrow扩展数组转换时保留了Polars的内存布局,而普通Pandas转换会将数据转为原生numpy数组,其哈希去重的内存管理逻辑刚好避开了当前Polars版本的内存峰值问题,但这种中转方式性能较低,仅适合临时救急。
内容的提问来源于stack exchange,提问作者stressed

