PySpark中同组不同ID的指定列重复识别(排除同ID内重复)
解决PySpark DataFrame跨id重复记录识别问题
需求核心:在同一group分组内,识别**不同id_number**对应的b、c列重复记录,同时忽略同一id_number内的b、c重复。
实现步骤
- 去重同一id内的重复:先过滤掉同一
group+id_number下的b、c重复记录,只保留每个id对应的唯一b、c组合。 - 统计跨id的重复次数:按
group+b+c分区,统计该组合对应的不同id_number数量。 - 标记重复状态:如果统计的id数量≥2,说明该
b、c组合在当前group下的不同id中出现过,标记为1;否则标记为0。 - 关联回原始数据:将标记结果关联到原始DataFrame,得到最终的
dups列。
完整代码
from pyspark.sql import Window import pyspark.sql.functions as F # 创建示例DataFrame import pandas as pd pddf = pd.DataFrame({ "group": ["a", "a", "a", "a", "b", "b", "b", 'c', 'c', 'd', 'd', 'd'], "id_number": ["12", "12", "13", "13", "16", "16", "17", '20', '21', '22', '22', '23'], "b": [1000, 2000, 1000, 3000, 1100, 1300, 1100, 1000, 1100, 2000, 2000, 2100], "c": ['F', 'D', 'F', 'A', 'B','C','B', 'B', 'B', 'A', 'A', 'B'] }) df = spark.createDataFrame(pddf) # 1. 去重同一id内的b、c重复 unique_df = df.dropDuplicates(["group", "id_number", "b", "c"]) # 2. 统计每个group下b、c对应的不同id数量 window_spec = Window.partitionBy("group", "b", "c") count_df = unique_df.withColumn("id_count", F.count("id_number").over(window_spec)) # 3. 标记跨id重复状态 mark_df = count_df.withColumn("dups", F.when(F.col("id_count") >= 2, 1).otherwise(0)) # 4. 关联回原始DataFrame result_df = df.join( mark_df.select("group", "id_number", "b", "c", "dups"), on=["group", "id_number", "b", "c"], how="left" ) # 查看结果 result_df.show()
输出结果
+-----+---------+----+---+-----+ |group|id_number| b| c| dups| +-----+---------+----+---+-----+ | a| 12|1000| F| 1| | a| 12|2000| D| 0| | a| 13|1000| F| 1| | a| 13|3000| A| 0| | b| 16|1100| B| 1| | b| 16|1300| C| 0| | b| 17|1100| B| 1| | c| 20|1000| B| 0| | c| 21|1100| B| 0| | d| 22|2000| A| 0| | d| 22|2000| A| 0| | d| 23|2100| B| 0| +-----+---------+----+---+-----+
内容的提问来源于stack exchange,提问作者user10969675
相关产品推荐
相关产品推荐

