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

Spark窗口函数不支持median?求滚动窗口中位数计算方案

PySpark 3.4.1滚动窗口中位数计算方案

PySpark 3.4.1中median()是近似聚合函数,不支持作为窗口函数配合ORDER BY或窗口帧(ROWS/RANGE)使用,这就是报错Cannot specify order by or frame for 'median'.的核心原因——哪怕去掉ORDER BY,它也只能执行全局聚合,无法实现滚动窗口计算。以下是两种可行的替代方案:

方案1:使用percentile_approx(推荐,适合大数据量)

percentile_approx是支持窗口操作的近似分位数函数,计算0.5分位数即可得到中位数,性能优异且适配大数据场景。

SQL写法

SELECT 
  id1,
  v1,
  percentile_approx(v1, 0.5) OVER (ORDER BY id1 ROWS BETWEEN 2 PRECEDING AND CURRENT ROW) AS rolling_median
FROM x

DataFrame API写法

from pyspark.sql.window import Window
from pyspark.sql.functions import percentile_approx

# 定义滚动窗口:当前行及前2行
window_spec = Window.orderBy("id1").rowsBetween(-2, 0)

# 计算滚动中位数
df.withColumn("rolling_median", percentile_approx("v1", 0.5).over(window_spec)).show()

注:如果需要更高精度,可以添加第三个采样数参数,比如percentile_approx(v1, 0.5, 1000),数值越大结果越精确,但性能会略有下降。

方案2:自定义UDF+collect_list(精确结果,适合小窗口)

如果需要完全精确的中位数,可以通过收集窗口内的所有值,再用自定义UDF计算。但这种方法会将窗口内的数据加载到内存,窗口过大时可能引发性能问题。

from pyspark.sql.window import Window
from pyspark.sql.functions import collect_list, udf
from pyspark.sql.types import DoubleType

# 定义计算中位数的UDF
def calculate_median(arr):
    sorted_arr = sorted(arr)
    n = len(sorted_arr)
    if n % 2 == 1:
        return sorted_arr[n // 2]
    else:
        return (sorted_arr[n//2 - 1] + sorted_arr[n//2]) / 2.0

median_udf = udf(calculate_median, DoubleType())

# 定义滚动窗口
window_spec = Window.orderBy("id1").rowsBetween(-2, 0)

# 收集窗口值并计算中位数
df.withColumn("window_values", collect_list("v1").over(window_spec)) \
  .withColumn("rolling_median", median_udf("window_values")) \
  .select("id1", "v1", "rolling_median") \
  .show()

内容的提问来源于stack exchange,提问作者jangorecki

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.15 07:47:46