PySpark分组排序后标记分组内最后一行的规范实现方法
PySpark 分组排序后标记分组最后一行的规范实现
原有实现的问题
你当前使用随机数列比对的写法属于典型的hack实现,存在三个明显缺陷:
- 随机数存在理论上的碰撞概率,极端场景下会出现标记错误
- 当排序键存在重复值时,随机数的不确定性会导致每次运行的标记结果不一致,无法复现
- 额外生成、删除随机列增加了不必要的计算开销
标准实现方案
核心思路是使用窗口函数的行号定位能力,完全不需要依赖随机值,逻辑透明可解释,结果稳定可复现。
方案1:正序行号匹配分组总行数(逻辑最直观)
- 可选预操作:如果分组内存在排序键重复的情况,先生成确定性的单调递增行ID作为次级排序键,保证同排序键下的行顺序和原始数据一致,结果100%稳定
- 定义分区排序窗口,按要求的分组字段分区、排序字段排序
- 计算每行在分组内的正序行号,同时计算分组总行数,行号等于总行数的行即为分组最后一行,标记为1,其余标记为0
示例代码:
from pyspark.sql import Window import pyspark.sql.functions as F # 排序键重复时开启,生成确定性全局唯一ID,无随机性 df = df.withColumn("_row_id", F.monotonically_increasing_id()) # 分组排序窗口,加_row_id次级排序保证顺序稳定 win_sort = Window.partitionBy("bvdidnumber", "dt_year").orderBy("dt_rfrnc", "_row_id") # 纯分组窗口,用于计算分组总行数 win_part = Window.partitionBy("bvdidnumber", "dt_year") df = df.withColumn( "marker", F.when( F.row_number().over(win_sort) == F.count("*").over(win_part), 1 ).otherwise(0) ).drop("_row_id") # 用完删除临时ID即可
方案2:逆序行号取首行(写法更简洁)
如果想要更简洁的写法,可以直接对排序键做逆序排序,分组内逆序后排在第一位的行就是正序的最后一行,不需要额外计算分组总行数:
from pyspark.sql import Window import pyspark.sql.functions as F df = df.withColumn("_row_id", F.monotonically_increasing_id()) # 逆序排序窗口 win_rev = Window.partitionBy("bvdidnumber", "dt_year").orderBy(F.desc("dt_rfrnc"), F.desc("_row_id")) df = df.withColumn( "marker", F.when(F.row_number().over(win_rev) == 1, 1).otherwise(0) ).drop("_row_id")
注意事项
如果你的业务场景中可以确定dt_rfrnc在每个bvdidnumber+dt_year分组内是唯一不重复的,可以省略生成_row_id的步骤,直接按dt_rfrnc排序即可,两种方案的执行效率没有本质差异。
内容的提问来源于stack exchange,提问作者safex
相关产品推荐
相关产品推荐

