PySpark实现按ID去重并保留test为Y的记录
解决PySpark中删除重复ID且test为N的记录问题
嘿,我来帮你搞定这个需求!你的核心目标是:对于重复的ID,只保留test为Y的那条记录;对于唯一的ID,如果test是N就保留它。单纯用groupBy("id").count()只能拿到计数,没法直接筛选符合要求的行,咱们可以用两种简单的方法实现:
方法一:使用窗口函数(推荐)
窗口函数能让我们按ID分组后,给每条记录分配优先级,然后筛选出优先级最高的行。这里我们给test为Y的记录排第1位,N排第2位,然后取每个ID组的第一条:
from pyspark.sql import Window import pyspark.sql.functions as F # 定义窗口:按id分组,按test字段排序(Y优先于N) window_spec = Window.partitionBy("id").orderBy(F.when(F.col("test") == "Y", 1).otherwise(2)) # 添加行号列,然后筛选行号为1的记录 new_df = df.withColumn("row_num", F.row_number().over(window_spec)) \ .filter(F.col("row_num") == 1) \ .drop("row_num") # 查看结果 new_df.show()
运行后你会得到期望的输出:
+---+----+ | id|test| +---+----+ | 1| Y| | 2| Y| | 3| N| +---+----+
方法二:先筛选有Y的ID集合,再过滤数据
另一种思路是先找出所有存在test=Y的ID,然后过滤数据:要么属于这些ID且test=Y,要么不属于这些ID(也就是该ID没有Y记录,保留N):
# 找出所有有test=Y的ID has_y_ids = df.filter(F.col("test") == "Y").select("id").distinct().rdd.flatMap(lambda x: x).collect() # 过滤数据:要么是有Y的ID且test=Y,要么是没有Y的ID(直接保留) new_df = df.filter((F.col("test") == "Y") | (~F.col("id").isin(has_y_ids))) new_df.show()
这种方法也能得到同样的结果,适合数据量不大的场景(因为collect()会把数据拉到Driver端)。
为什么你的groupBy方法不行?
groupBy("id").count()只能统计每个ID的出现次数,但没法关联回原数据的test字段,所以没法直接筛选出符合要求的行。上面两种方法都是先关联ID和test的关系,再做筛选,就能精准满足你的需求啦~
内容的提问来源于stack exchange,提问作者User12345
相关产品推荐
相关产品推荐

