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
相关产品推荐
相关产品推荐

