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

Spark DataFrame首列重复值标记:新增首值为True的布尔列

解决方案:Spark DataFrame新增首次出现标记列

问题背景

现有Spark DataFrame df,数据及初始化代码如下:

初始化代码

num_rows = 5
num_cols = 3

data = [
    [10, 20, 30],
    [10, 50, 60],
    [70, 80, 90],
    [20, 30, 10],
    [20, 10, 20]
]

columns = [f"Column_{i+1}" for i in range(num_cols)]

df = spark.createDataFrame(data, columns)

原DataFrame展示

|Column_1|Column_2|Column_3|
+--------+--------+--------+
|      10|      20|      30|
|      10|      50|      60|
|      70|      80|      90|
|      20|      30|      10|
|      20|      10|      20|
+--------+--------+--------+

需求

新增Column_4,规则:Column_1的值首次出现标记为TRUE,后续重复出现的标记为FALSE,最终效果如下:

|Column_1|Column_2|Column_3|Column_4|
+--------+--------+--------+--------+
|      10|      20|      30|    TRUE|
|      10|      50|      60|   FALSE|
|      70|      80|      90|    TRUE|
|      20|      30|      10|    TRUE|
|      20|      10|      20|   FALSE|
+--------+--------+--------+--------+

实现步骤

1. 导入依赖

from pyspark.sql import Window
from pyspark.sql.functions import row_number, when, monotonically_increasing_id

2. 保留原始顺序并新增标记列

Spark默认不保证数据顺序,需先添加递增行ID锁定原始顺序,再通过窗口函数标记首次出现的行:

# 添加行ID以保留原始数据顺序
df_with_rowid = df.withColumn("row_id", monotonically_increasing_id())

# 定义窗口:按Column_1分组,按行ID升序排序(对应原始顺序)
window_spec = Window.partitionBy("Column_1").orderBy("row_id")

# 生成Column_4:分组内第一行标记为TRUE,其余为FALSE
result_df = df_with_rowid.withColumn(
    "Column_4",
    when(row_number().over(window_spec) == 1, "TRUE").otherwise("FALSE")
).drop("row_id")  # 删除临时行ID列

3. 验证结果

执行result_df.show()即可得到符合需求的DataFrame。

原理说明

  • monotonically_increasing_id()生成全局唯一递增ID,确保数据顺序与原始输入一致;
  • 窗口函数row_number().over(window_spec)对每个Column_1分组内的行按顺序编号,首次出现的行编号为1,据此标记TRUE,后续行编号大于1则标记FALSE。

内容的提问来源于stack exchange,提问作者Mitch

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 21:13:22