Spark全局窗口avg()结果随orderBy列变化的原因咨询
窗口函数avg()结果随orderBy列变化的原因解析
核心结论
当使用无分区且指定了orderBy的窗口函数时,Spark默认的计算范围是从窗口起始行到当前行的累积区间(等价于rowsBetween(Window.unboundedPreceding, Window.currentRow)),而非整个数据集的全局范围。这是导致结果随orderBy变化、呈现逐行累积值的根本原因。
针对疑问的具体解答
1. 为何orderBy选择的列会影响结果?
orderBy决定了窗口内数据行的排序顺序,不同的排序会改变每一行对应的"从起始到当前行"的数据集合,累积计算的平均值自然不同:
- 当
orderBy("value1")时,窗口内的行按value1升序排列为:(1,10.0,5.0) → (5,15.0,None) → (3,20.0,None) - 当
orderBy("id")时,窗口内的行按id升序排列为:(1,10.0,5.0) → (3,20.0,None) → (5,15.0,None)
两种排序下,每行对应的前置数据集合不同,最终累积平均值也就产生差异。
2. 为何avg()结果是逐行变化的?
因为默认窗口范围是累积到当前行,而非全局:
- 第一行:仅包含自身数据,平均值为当前行的
value1值 - 第二行:包含前两行数据的平均值
- 第三行:包含前三行数据的平均值
结合两种排序的行顺序,就出现了你看到的逐行变化结果:
- 按
value1排序后,第三行对应的数据是10.0、15.0、20.0的平均值15.0,但你最后按id排序显示,所以id=5的行(原窗口第二行)显示的是12.5(前两行10.0、15.0的平均) - 按
id排序后,第三行对应的数据是10.0、20.0、15.0的平均值15.0,和全局聚合结果一致
如何用窗口函数获取全局平均值?
如果需要在窗口中得到和df.agg(avg("value1"))一致的全局平均值,需要显式指定窗口范围为整个数据集:
from pyspark.sql.window import Window from pyspark.sql.functions import avg, col # 定义覆盖整个数据集的窗口范围 w = Window.rowsBetween(Window.unboundedPreceding, Window.unboundedFollowing) df.withColumn("AVG", avg(col("value1")).over(w))\ .sort("id", ascending=True)\ .show()
输出结果会统一显示全局平均值15.0,不受orderBy的影响(即使添加orderBy,只要范围是全局,结果依然一致)。
内容的提问来源于stack exchange,提问作者anurag86
相关产品推荐
相关产品推荐

