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

