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

如何使用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)

运行结果

输出完全符合预期:

idvehicletimestampeventrank
1v12023-04-10 08:19:29event_11
2v12023-04-10 08:41:01event_21
3v12023-04-10 09:41:01event_11
4v12023-04-10 10:41:01event_21
5v12023-04-10 11:41:01event_22
6v12023-04-10 12:41:01event_11
7v12023-04-10 13:41:01event_12
8v12023-04-10 14:41:01event_21

关键逻辑解释

  • lag("event"):获取当前行的上一行event,用来判断是否为连续重复的event
  • event_flag:标记每次event切换的位置,为后续分组做准备
  • sum("event_flag"):通过累加生成分组ID,将连续相同的event归为一组
  • row_number():在每个分组内生成递增的行号,即所需的连续重复排名

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.13 03:34:57