如何使用PySpark为连续重复的行实现自定义排名?
PySpark实现连续重复行的排名
需求很明确:给每个车辆按时间排序后的连续相同event行生成排名,连续的相同event从1开始递增,event切换时重置排名。
实现步骤
核心思路是先把连续相同的event归为同一个分组,再在分组内生成行号作为排名,具体代码如下:
1. 准备测试数据
from pyspark.sql import SparkSession from pyspark.sql import functions as F from pyspark.sql.window import Window # 初始化SparkSession spark = SparkSession.builder.appName("ConsecutiveDuplicateRank").getOrCreate() # 原始数据 data = [ (1, "v1", "2023-04-10 08:19:29", "event_1"), (2, "v1", "2023-04-10 08:41:01", "event_2"), (3, "v1", "2023-04-10 09:41:01", "event_1"), (4, "v1", "2023-04-10 10:41:01", "event_2"), (5, "v1", "2023-04-10 11:41:01", "event_2"), (6, "v1", "2023-04-10 12:41:01", "event_1"), (7, "v1", "2023-04-10 13:41:01", "event_1"), (8, "v1", "2023-04-10 14:41:01", "event_2") ] # 创建DataFrame并转换timestamp为时间类型 df = spark.createDataFrame(data, ["id", "vehicle", "timestamp", "event"]) df = df.withColumn("timestamp", F.to_timestamp("timestamp"))
2. 生成连续重复行的排名
# 定义窗口:按车辆分组,按时间戳升序排序 base_window = Window.partitionBy("vehicle").orderBy("timestamp") # 获取上一行的event,第一行用null填充 df = df.withColumn("prev_event", F.lag("event").over(base_window)) # 标记event是否发生变化:当前event和上一行不同时记为1,否则0 df = df.withColumn("event_flag", F.when(F.col("event") != F.col("prev_event"), 1).otherwise(0)) # 生成分组ID:累加event_flag,连续相同的event会被分到同一个组 df = df.withColumn("group_id", F.sum("event_flag").over(base_window.rowsBetween(Window.unboundedPreceding, Window.currentRow))) # 在每个车辆+分组ID的组内,生成排名 rank_window = Window.partitionBy("vehicle", "group_id").orderBy("timestamp") result_df = df.withColumn("rank", F.row_number().over(rank_window)).drop("prev_event", "event_flag", "group_id") # 查看结果(按id排序) result_df.orderBy("id").show(truncate=False)
运行结果
输出完全符合预期:
| id | vehicle | timestamp | event | rank |
|---|---|---|---|---|
| 1 | v1 | 2023-04-10 08:19:29 | event_1 | 1 |
| 2 | v1 | 2023-04-10 08:41:01 | event_2 | 1 |
| 3 | v1 | 2023-04-10 09:41:01 | event_1 | 1 |
| 4 | v1 | 2023-04-10 10:41:01 | event_2 | 1 |
| 5 | v1 | 2023-04-10 11:41:01 | event_2 | 2 |
| 6 | v1 | 2023-04-10 12:41:01 | event_1 | 1 |
| 7 | v1 | 2023-04-10 13:41:01 | event_1 | 2 |
| 8 | v1 | 2023-04-10 14:41:01 | event_2 | 1 |
关键逻辑解释
lag("event"):获取当前行的上一行event,用来判断是否为连续重复的eventevent_flag:标记每次event切换的位置,为后续分组做准备sum("event_flag"):通过累加生成分组ID,将连续相同的event归为一组row_number():在每个分组内生成递增的行号,即所需的连续重复排名
内容的提问来源于stack exchange,提问作者TheCurrent
相关产品推荐
相关产品推荐

