如何在PySpark中生成去重DataFrame及带标记的合并DataFrame
PySpark 生成指定去重DataFrame及标记重复行方案
原始DataFrame结构
column_1 column_2 column_3 column_4 column_5 column_6 column_7 34432 apple banana mango pine lemon j93jk84 98389 grape orange pine kiwi cherry j93jso3 94749 apple banana mango pine lemon ke03jfr 48948 apple banana mango pine lemon 9jef3f4 . . . . . . . 90493 pear apricot papaya plum lemon 93jd30d 90843 grape orange pine kiwi cherry 03nd920
需求说明
需要生成两个目标DataFrame:
- Dataframe_1:忽略
column_1和column_7,基于其余列去重,仅保留唯一行; - Dataframe_2:包含原所有列,新增
type列标记行是unique还是duplicated,新增Tag列将同一组重复行标记为相同编号(唯一行与对应重复行Tag一致),示例结构如下:
column_1 column_2 column_3 column_4 column_5 column_6 column_7 type Tag 34432 apple banana mango pine lemon j93jk84 unique 1 98389 grape orange pine kiwi cherry j93jso3 unique 2 94749 apple banana mango pine lemon ke03jfr duplicated 1 48948 apple banana mango pine lemon 9jef3f4 duplicated 1 . . . . . . . 90493 pear apricot papaya plum lemon 93jd30d unique 3 90843 grape orange pine kiwi cherry 03nd920 duplicated 2
完整实现方案
导入依赖
from pyspark.sql import functions as F from pyspark.sql.window import Window
1. 生成Dataframe_1
直接通过指定去重列实现,忽略column_1和column_7:
# 定义用于去重的列集合 dedup_cols = ["column_2", "column_3", "column_4", "column_5", "column_6"] # 生成Dataframe_1:基于指定列去重,保留每组的第一行 df_unique = df.dropDuplicates(dedup_cols)
2. 生成Dataframe_2
通过窗口函数实现重复组标记和类型区分,替代原代码中subtract的方式(原方式无法保留所有行并关联Tag):
# 步骤1:按去重列分组,给每组内的行添加行号(用于判断是否为唯一行) window_group = Window.partitionBy(*dedup_cols).orderBy("column_1") df_with_row_num = df.withColumn("row_num", F.row_number().over(window_group)) # 步骤2:给每个重复组分配唯一的Tag编号(dense_rank生成连续的编号) window_tag = Window.orderBy(*dedup_cols) df_with_tag = df_with_row_num.withColumn("Tag", F.dense_rank().over(window_tag)) # 步骤3:添加type列,标记行的唯一性 df_combined = df_with_tag.withColumn( "type", F.when(F.col("row_num") == 1, F.lit("unique")).otherwise(F.lit("duplicated")) ).drop("row_num") # 移除临时的行号列
关键说明
- 原代码使用
subtract获取重复行的方式存在缺陷:subtract会移除所有与df_unique重复的行,但实际上我们需要保留所有原始行并标记类型,且column_1/column_7的差异会导致subtract无法正确识别同组重复行。 - 采用窗口函数的方式可以精准关联同一重复组的所有行,并生成连续的
Tag编号,同时准确标记每行的type。
内容的提问来源于stack exchange,提问作者Bella_18
相关产品推荐
相关产品推荐

