如何从PySpark DataFrame中移除所有存在重复的记录
PySpark移除所有含重复ID的记录实现方案
需求说明:原始DataFrame中部分ID存在重复(如ID=1出现两次),需要删除所有存在重复的ID对应的行,仅保留ID唯一的记录。
原始数据示例
| ID | start_dt |
|---|---|
| 1 | 2020-02-09 |
| 1 | 2021-02-15 |
| 2 | 2022-05-04 |
| 3 | 2023-05-15 |
目标数据示例
| ID | start_dt |
|---|---|
| 2 | 2022-05-04 |
| 3 | 2023-05-15 |
实现方法
方法1:使用Left Anti Join(推荐,高效简洁)
先找出所有重复的ID,再用左反连接过滤掉这些ID的所有记录:
from pyspark.sql import SparkSession # 初始化SparkSession spark = SparkSession.builder.appName("RemoveDuplicateIDs").getOrCreate() # 创建示例DataFrame data = [(1, "2020-02-09"), (1, "2021-02-15"), (2, "2022-05-04"), (3, "2023-05-15")] df = spark.createDataFrame(data, ["ID", "start_dt"]) # 筛选出存在重复的ID duplicate_ids = df.groupBy("ID").count().filter("count > 1").select("ID") # 左反连接:保留原始表中不在重复ID列表里的记录 result_df = df.join(duplicate_ids, on="ID", how="left_anti") # 查看结果 result_df.show()
方法2:窗口函数统计次数后过滤
通过窗口函数计算每个ID的出现次数,再筛选出次数为1的记录:
from pyspark.sql import SparkSession from pyspark.sql.window import Window from pyspark.sql.functions import count spark = SparkSession.builder.appName("RemoveDuplicateIDs").getOrCreate() data = [(1, "2020-02-09"), (1, "2021-02-15"), (2, "2022-05-04"), (3, "2023-05-15")] df = spark.createDataFrame(data, ["ID", "start_dt"]) # 定义按ID分组的窗口 window_spec = Window.partitionBy("ID") # 添加每个ID的出现次数字段 df_with_count = df.withColumn("id_occurrences", count("*").over(window_spec)) # 过滤出仅出现一次的记录,并删除临时字段 result_df = df_with_count.filter(df_with_count.id_occurrences == 1).drop("id_occurrences") result_df.show()
方法3:分组统计后关联过滤
先统计每个ID的出现次数,再通过内连接保留唯一ID的记录:
from pyspark.sql import SparkSession spark = SparkSession.builder.appName("RemoveDuplicateIDs").getOrCreate() data = [(1, "2020-02-09"), (1, "2021-02-15"), (2, "2022-05-04"), (3, "2023-05-15")] df = spark.createDataFrame(data, ["ID", "start_dt"]) # 统计每个ID的出现次数,筛选出仅出现一次的ID unique_ids = df.groupBy("ID").count().filter("count == 1") # 内连接原始表,保留唯一ID的记录 result_df = df.join(unique_ids, on="ID", how="inner").drop("count") result_df.show()
内容的提问来源于stack exchange,提问作者ddd
相关产品推荐
相关产品推荐

