You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.06 00:20:24