PySpark窗口orderBy对min/max聚合函数行为的影响及原因
PySpark窗口函数中max/min结果不符合预期的原理分析
我正在使用PySpark 3.3.1版本,在窗口函数中发现max函数的输出不符合预期,B-123分区内的max_val是累积值而非全局最值0.8,想了解该现象的原理以及窗口分组内发生了什么。
测试代码
import pyspark.sql.functions as F from pyspark.sql import Window import pandas as pd df = spark.createDataFrame(pd.DataFrame( { "id": ["A-123","A-123","A-123","A-123","B-123","B-123","B-123","B-123"], "val": [0.1,0.2,0.3,0.4,0.5,0.6,0.7,0.8] } )) window_group = Window.partitionBy(F.col('id')).orderBy(F.col('val')) ( df .withColumn('min_val', F.min(F.col('val')).over(window_group)) .withColumn('max_val', F.max(F.col('val')).over(window_group)) ).show()
执行输出
+-----+---+-------+-------+ | id|val|min_val|max_val| +-----+---+-------+-------+ |A-123|0.1| 0.1| 0.1| |A-123|0.2| 0.1| 0.2| |A-123|0.3| 0.1| 0.3| |A-123|0.4| 0.1| 0.4| |B-123|0.5| 0.5| 0.5| --> should be 0.8? |B-123|0.6| 0.5| 0.6| --> should be 0.8? |B-123|0.7| 0.5| 0.7| --> should be 0.8? |B-123|0.8| 0.5| 0.8| +-----+---+-------+-------+
现象原理
核心原因是当窗口函数同时指定partitionBy和orderBy时,PySpark会默认使用"范围窗口(Range Frame)",默认的窗口范围是从分区起始行到当前行(包含当前行):
- 若未指定
orderBy,窗口默认覆盖整个分区(范围为UNBOUNDED PRECEDING到UNBOUNDED FOLLOWING),此时max/min计算的是整个分组的全局最值。 - 一旦添加
orderBy,PySpark会自动将窗口范围限定为RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW,也就是只计算从分区第一行到当前行范围内的max/min,因此你看到的是累积式的最大值(每一行的max_val是从分区开头到当前行的最大值)。
两种修复方法的原理
- 添加
rowsBetween(Window.unboundedPreceding, Window.unboundedFollowing):显式强制窗口覆盖整个分区的所有行,无论是否有orderBy,都会计算分组内的全局最值。 - 移除
orderBy子句:此时窗口默认范围回到整个分区,自然得到分组内的全局max/min值。
内容的提问来源于stack exchange,提问作者Sean L
相关产品推荐
相关产品推荐

