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
相关产品推荐
相关产品推荐

