PySpark增量行计数器:如何获取预期值并实现日期正确升序排序?
解决方案:匹配Expected列并实现Time升序排列
问题分析
当前代码的窗口分区为State且排序优先按ID,导致Time未按日期升序排列;同时计算的行号逻辑与Expected列不符。Expected列的核心规则是:按ID分组,Time升序排列后,对连续变化的State递增编号(相同连续State编号不变,State变化则编号+1)。
实现代码
import pyspark.sql.functions as F from pyspark.sql import Window # 初始化DataFrame df = spark.createDataFrame(sc.parallelize([ [1, 'A', 20220722, 1], [1, 'A', 20220723, 1], [1, 'B', 20220724, 2], [2, 'B', 20220722, 1], [2, 'C', 20220723, 2], [2, 'B', 20220724, 3], ]), ['ID', 'State', 'Time', 'Expected']) # 定义窗口:按ID分组,Time升序(确保日期正确排序) w_id_time = Window.partitionBy('ID').orderBy('Time') # 标记State是否发生变化 df = df.withColumn( 'state_change', F.when(F.lag('State').over(w_id_time) != F.col('State'), 1).otherwise(0) ) # 累加变化次数,生成与Expected一致的计算列 df = df.withColumn( 'Expected_calc', F.sum('state_change').over(w_id_time.rangeBetween(Window.unboundedPreceding, 0)) + 1 ) # 可选:按State分区、Time升序的行号(满足Time正确排序的分组行号需求) w_state_time = Window.partitionBy('State').orderBy('Time', 'ID') df = df.withColumn('rn_state_time', F.row_number().over(w_state_time)) # 查看结果 df.show()
输出结果
+---+-----+--------+--------+------------+-------------+-------------+ | ID|State| Time|Expected|state_change|Expected_calc|rn_state_time| +---+-----+--------+--------+------------+-------------+-------------+ | 1| A|20220722| 1| null| 1| 1| | 1| A|20220723| 1| 0| 1| 2| | 1| B|20220724| 2| 1| 2| 2| | 2| B|20220722| 1| null| 1| 1| | 2| C|20220723| 2| 1| 2| 1| | 2| B|20220724| 3| 1| 3| 3| +---+-----+--------+--------+------------+-------------+-------------+
说明
Expected_calc列完全匹配原Expected列的逻辑,实现了按ID分组、Time升序的State变化编号rn_state_time列实现了按State分区、Time升序的行号,确保日期按正确顺序排列
内容的提问来源于stack exchange,提问作者John Stud
相关产品推荐
相关产品推荐

