为何Spark中基于Window.unboundedFollowing的first()比last()慢很多
耗时问题根源
- 升序排序下,
unboundedPreceding -> 当前行的last()聚合被Spark识别为可增量计算的场景,会使用RunningWindowFunction执行模式:每处理一行仅需做一次值判断更新,整个分区计算复杂度为O(n),同时支持全阶段代码生成(就是你看到的节点旁的*(3)标记),执行效率极高。 - 同升序排序下,
当前行 -> unboundedFollowing的first()聚合当前Spark没有对应的增量优化实现,只能退化为普通Window执行模式:对每一行都完整扫描从当前行到分区末尾的所有数据来查找第一个非空值,单个分区的计算复杂度为O(n²),同时不支持全阶段代码生成,在单分区5万行的规模下,计算量直接上升了几个数量级,这就是耗时飙升的核心原因。
最优修复方案
只需要调整next相关窗口的排序方向,把反向查找转化为正向的增量计算即可,修改后和原逻辑完全等价,且能复用和part1一致的高效执行模式:
# 替换原win_next定义:按rank倒序排序,仍用正向无界到当前行的窗口 win_next_opt = Window.partitionBy('PORT_TYPE', 'loss_process').orderBy(F.desc('rank')).rowsBetween(Window.unboundedPreceding, 0) # part2计算逻辑替换为用优化后的窗口取last非空值,效果和原first逻辑完全一致 df_part2 = (df_part1 .withColumn('next_rank', F.last(F.col('rank'), ignorenulls=True).over(win_next_opt)) .withColumn('next_sf', F.last(F.col('scale_factor'), ignorenulls=True).over(win_next_opt)) ).cache()
修改后part2的耗时会降到和part1同量级,整体执行时间从3分钟压缩到几秒内。
补充说明
你观察到的执行节点差异完全符合上述分析:
RunningWindowFunction对应增量优化的执行模式,支持全阶段代码生成,所以有星号标记- 普通
Window节点对应非优化的逐行全窗口扫描模式,无全阶段代码生成优化,执行效率极低
内容的提问来源于stack exchange,提问作者Alain
相关产品推荐
相关产品推荐

