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

Spark实现带重置的行号(RowNumber with Reset)异常问题排查

解决Spark中连续相同State的重置行号问题

你的问题出在rank()的分区逻辑上——只按id分区的话,所有同ID的记录会被放在同一个排序池里,后面的P会和前面的P序列累计计算排名,自然没法按预期重置。要实现连续相同State时递增、State变化时重置的行号,得先给连续相同的State生成分组标识,再基于分组做行号。

修正后的代码如下:

from pyspark.sql import functions as F
from pyspark.sql.window import Window

# 初始化窗口(按id分区、time排序)
w = Window.partitionBy('id').orderBy('time')

# create df
df = spark.createDataFrame(sc.parallelize([
    [1, 'P', 20220722, 1],
    [1, 'P', 20220723, 2],
    [1, 'P', 20220724, 3],
    [1, 'P', 20220725, 4],
    [1, 'D', 20220726, 1],
    [1, 'O', 20220727, 1],
    [1, 'D', 20220728, 1],
    [1, 'P', 20220729, 2],
    [1, 'P', 20220730, 3],
    [1, 'P', 20220731, 4],   
]),
                           ['ID', 'State', 'Time', 'Expected'])

# 1. 生成lagState
df = df.withColumn('lagState', F.lag('State').over(w))

# 2. 生成分组标识:当State和前一行不同时标记为1,否则0,累加后得到分组ID
df = df.withColumn('group_id', 
                   F.sum(F.when(F.col('State') != F.col('lagState'), 1).otherwise(0))
                   .over(w.rangeBetween(Window.unboundedPreceding, 0))
                  )
# 处理第一行的null情况(lagState为null时,group_id设为1)
df = df.withColumn('group_id', F.coalesce('group_id', F.lit(1)))

# 3. 基于id和group_id分区,生成行号
df = df.withColumn('rank', F.row_number().over(Window.partitionBy('id', 'group_id').orderBy('time')))

# view
df.show()

运行后rank列会和Expected列完全匹配:

  • 连续相同的State会按顺序递增行号
  • 每次State变化时,行号重置为1

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.18 14:15:41