Spark窗口聚合函数添加排序后行为不符合预期的原因咨询
为什么Spark窗口添加orderBy后min值会变化?
核心原因是:当窗口定义中加入orderBy时,Spark会自动设置默认的窗口帧(Frame)范围为RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW——也就是仅包含从分区起始行到当前行的数据,而非整个分区的所有数据。
两种场景的帧范围差异
无orderBy的窗口
当窗口仅用partitionBy("id")定义时,默认帧范围是RANGE BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING,即覆盖整个分区的所有行。此时min()和max()会计算全分区的极值,所以每条记录的结果都一致,符合你的预期。添加orderBy后的窗口
加入orderBy(F.col("val").desc())后,Spark自动切换默认帧范围为「分区起始到当前行」。此时:max(val)结果看似不变,是因为分区按val降序排列后,第一行就是整个分区的最大值,后续所有行的帧都包含这个最大值,所以计算结果始终是333;min(val)则是计算当前帧内的最小值,随着行的推进,帧包含的数据越来越多,最小值也随之更新:- 第一行(val=333):帧只有[333],min=333
- 第二行(val=222):帧包含[333,222],min=222
- 第三行(val=111):帧包含[333,222,111],min=111
解决方法:手动指定帧范围
如果想要在保留orderBy的同时,仍计算全分区的极值,只需手动指定帧范围覆盖整个分区即可:
window = Window.partitionBy("id")\ .orderBy(F.col("val").desc())\ .rangeBetween(Window.unboundedPreceding, Window.unboundedFollowing)
用这个窗口定义运行代码,min_val和max_val就会回到全分区的极值,和无orderBy时的结果一致。
内容的提问来源于stack exchange,提问作者M_S
相关产品推荐
相关产品推荐

