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

