Spark 3.5.0中pyspark.sql.functions.percentile精确分位数函数原理问询
Spark 3.5.0精确分位数函数
percentile核心实现解析 核心源码位置
Python层的percentile仅为外层封装,核心逻辑在Spark的Scala代码中,关键文件包括:
- 聚合表达式入口:
sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/expressions/aggregate/Percentile.scala - 分位数计算工具类:
sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/util/PercentileCompute.scala
实现算法原理
Spark 3.5.0的percentile采用两阶段分布式精确计算,适配百亿级数据集的核心逻辑如下:
局部分区预处理
- 每个Task对当前分区内的目标列数据排序,同时记录分区元素总数。若一次性请求多个分位数,该阶段会统一处理所有目标分位点,避免重复排序。
- 当分区数据量超出内存阈值时,自动触发外部排序(利用磁盘存储中间结果),规避内存溢出问题。
全局合并计算
- 收集所有分区的排序后数据片段及分区元素数,通过归并排序合并为全局有序数据集(归并排序的时间复杂度远低于全量乱序数据排序,大幅降低计算开销)。
- 根据全局总元素数,计算每个目标分位数对应的索引位置;若位置落在两个元素之间,按配置的插值规则(默认线性插值)生成最终结果。
- 依赖Spark分布式shuffle框架高效传输分区排序片段,避免全量原始数据的大规模网络传输。
关键特性优化
- 幂等性:基于精确排序与全局合并逻辑,输入数据固定时,重复计算结果完全一致。
- 性能适配:大分区外部排序、多分区归并排序的优化,使得百亿级数据集下的性能与近似分位数函数
percentile_approx无明显差异。
内容的提问来源于stack exchange,提问作者Jakub Kaplan
相关产品推荐
相关产品推荐

