如何按条件删除重复数据:保留每组item_id和country_id的最小level值
解决Spark DataFrame分组保留最小level行的问题
嗨,你遇到的这种按item_id和country_id分组,保留每组中level最小值对应行的需求,确实没法直接用dropDuplicates解决,不过Spark提供了两种很实用的方案,我给你详细说下:
方法一:窗口函数(推荐,支持保留多列)
如果你的DataFrame后续可能还有其他列需要保留,窗口函数是最稳妥的选择,它能精准定位每组里level最小的那一行:
from pyspark.sql import Window import pyspark.sql.functions as F # 定义窗口规则:按item_id和country_id分组,按level升序排序(最小的排第一) window_spec = Window.partitionBy("item_id", "country_id").orderBy(F.col("level").asc()) # 给每组的行添加行号,level最小的行号为1 # 过滤出行号=1的行,再删掉临时的行号列 result_df = df.withColumn("row_num", F.row_number().over(window_spec)) \ .filter(F.col("row_num") == 1) \ .drop("row_num") # 查看结果 result_df.show()
方法二:分组聚合(简洁,适合仅保留示例三列的场景)
如果你的需求只需要保留item_id、country_id和最小的level,直接分组聚合会更简洁:
import pyspark.sql.functions as F # 按item_id和country_id分组,聚合取每组的最小level result_df = df.groupBy("item_id", "country_id") \ .agg(F.min("level").alias("level")) \ .orderBy("item_id", "country_id") # 可选,用来和示例输出顺序对齐 # 查看结果 result_df.show()
两种方法都能得到你想要的输出,窗口函数的优势是可以保留原DataFrame的其他列(比如如果有price、name这类字段,也能一起保留下来),而分组聚合更轻量,适合只需要这三列的场景。
内容的提问来源于stack exchange,提问作者Markus
相关产品推荐
相关产品推荐

