如何对无法载入内存的大Parquet文件进行去重?
处理超大Parquet文件的分组去重需求(按a、b保留d/e最大行)
问题背景
现有无法载入内存的Parquet文件,含a、b、c、d、e列(均为int32类型),需按a、b列分组去重:每组保留d值最大的行;若d值相同,则保留e值最大的行。使用Arrow数据集操作时遇到以下限制:
distinct(a, b, .keep_all = TRUE)操作不支持;- 分组后
filter(n()>1)表达式不支持,需先将数据载入内存; - 分区分批写入的代码报错或未达预期效果。
Python解决方案(PyArrow)
通过分组聚合+关联筛选实现,全程无需加载全量数据,内存友好:
import pyarrow as pa import pyarrow.parquet as pq import pyarrow.dataset as ds source_path = "~/Downloads/test.parquet" dest_path = "~/Downloads/result_py.parquet" # 1. 计算每个(a,b)分组的最大d、e值 dataset = ds.dataset(source_path, format="parquet") aggregated = dataset.group_by(["a", "b"]).aggregate([ ("d", "max"), ("e", "max") ]) aggregated_table = aggregated.to_table() # 2. 关联原数据集,筛选出符合条件的行 join_expr = ( (ds.field("a") == ds.field("agg.a")) & (ds.field("b") == ds.field("agg.b")) & (ds.field("d") == ds.field("agg.d_max")) & (ds.field("e") == ds.field("agg.e_max")) ) result_dataset = dataset.join( aggregated_table, keys_left=["a", "b"], keys_right=["a", "b"], filter=join_expr, right_alias="agg" ).select(dataset.schema.names) # 保留原表所有列 # 写入结果 pq.write_to_dataset(result_dataset, dest_path, format="parquet")
方案说明
- 分组聚合仅生成每个(a,b)对应的最大值,数据量远小于原文件;
- 关联筛选操作在Dataset层面执行,无需加载全量数据;
- 完美规避了
distinct和filter(n()>1)的操作限制。
R解决方案(Arrow + dplyr)
提供两种实现方式,按需选择:
方式1:分组聚合+关联筛选(兼容性强)
library(arrow) library(dplyr) source_path <- '~/Downloads/test.parquet' dest_dir <- '~/Downloads/result_r' # 1. 计算每个(a,b)分组的最大d、e值(结果量小可安全载入内存) aggregated <- open_dataset(source_path) %>% group_by(a, b) %>% summarize(max_d = max(d), max_e = max(e), .groups = "drop") %>% collect() # 2. 关联原数据集并筛选目标行 result <- open_dataset(source_path) %>% inner_join(aggregated, by = c("a", "b")) %>% filter(d == max_d, e == max_e) %>% select(all_of(colnames(open_dataset(source_path)))) # 保留原列 # 写入结果 write_dataset(result, dest_dir, format = "parquet")
方式2:窗口函数(简洁高效,需Arrow 11.0+)
若你的Arrow版本支持窗口函数,可直接用排名筛选:
library(arrow) library(dplyr) source_path <- '~/Downloads/test.parquet' dest_dir <- '~/Downloads/result_r_window' result <- open_dataset(source_path) %>% group_by(a, b) %>% mutate( # 按d降序、e降序排名,取排名第1的行 rank = dense_rank(desc(d), desc(e)) ) %>% filter(rank == 1) %>% select(-rank) %>% ungroup() write_dataset(result, dest_dir, format = "parquet")
方案说明
- 方式1兼容性强,适用于所有Arrow版本;
- 方式2代码更简洁,无需额外关联操作,依赖较新版本Arrow的窗口函数支持。
内容的提问来源于stack exchange,提问作者falsePockets
相关产品推荐
相关产品推荐

