You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.01 12:15:48