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

PySpark如何基于历史行计算每行累计Q1、中位数、Q3百分位数?

结论:你提供的代码不能实现需求,存在多处错误

错误点说明:

  • 语法逻辑错误:你先使用groupBy('SEQ_ID').agg()会将数据聚合为每个SEQ_ID仅对应1行的结果,完全丢失了原始数据的行结构,后续调用窗口函数over(w)会直接报错,无法得到每行对应的滚动百分位数结果
  • 字段名不匹配:百分位数计算使用的字段是Pad_Wear,但你实际存储数值的字段是RESULT,运行会报字段不存在的错误
  • 函数兼容性问题:原生percentile函数在部分低版本Spark中不支持作为窗口函数调用,滚动计算更推荐使用percentile_approx,性能更高兼容性更好

正确实现代码:

你定义的窗口规则是正确的,只需要调整调用逻辑即可:

from pyspark.sql import functions as f
from pyspark.sql.window import Window

# 窗口定义符合需求:按SEQ_ID分组、时间升序、窗口范围是分组内当前行及所有之前的行
w = Window.partitionBy('SEQ_ID')\
          .orderBy(f.col('TIME_STAMP').asc())\
          .rangeBetween(Window.unboundedPreceding, 0)

# 直接在select中调用窗口函数,保留所有原始行,每行对应计算滚动百分位数
df_result = df.select(
    'SEQ_ID',
    'TIME_STAMP',
    'RESULT',
    f.expr('percentile_approx(RESULT, 0.25)').over(w).alias('Q1'),
    f.expr('percentile_approx(RESULT, 0.5)').over(w).alias('Median'),
    f.expr('percentile_approx(RESULT, 0.75)').over(w).alias('Q3')
)

补充说明:

如果你需要完全精确的百分位数,可以将代码中的percentile_approx替换为percentile,仅需要确认你使用的Spark版本在3.0及以上即可,该版本之后percentile已支持窗口函数调用。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.29 18:54:04