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

在Scala(Spark DataFrame)的窗口上应用approxQuantile

在Spark分区窗口上计算指定分位数的解决方案

好问题!确实,Spark自带的percent_rank()只能返回每行在窗口分区内的相对排名占比,没法直接给出我们需要的特定分位数对应的具体数值。针对这个需求,我们可以根据你的Spark版本选择不同的实现方式:

方法一:Spark 3.0+ 原生支持(推荐)

Spark 3.0及以上版本已经支持将percentile_approx作为窗口函数使用,这是最简洁高效的方案,不需要额外自定义逻辑。它可以直接在分区窗口上计算指定分位数,还支持同时计算多个分位数。

举个实际例子:假设我们有一个包含category(分区列)和value(需要计算分位数的列)的DataFrame,要计算每个分类下value的90分位数:

import org.apache.spark.sql.expressions.Window
import org.apache.spark.sql.functions.percentile_approx

// 定义分区窗口
val windowSpec = Window.partitionBy("category")

// 计算每个分区的90分位数
val resultDF = df.withColumn(
  "90th_percentile",
  percentile_approx($"value", 0.9, 10000).over(windowSpec)
)

参数说明:

  • 第二个参数0.9就是我们要计算的分位数(0到1之间的数值)
  • 第三个参数10000是精度控制值,数值越大结果越准确,对应的计算开销也会略高,日常场景下10000是个不错的平衡值

如果需要同时计算多个分位数(比如四分位数),可以传入数组参数,返回结果会是数组类型:

val resultDF = df.withColumn(
  "quartiles",
  percentile_approx($"value", Array(0.25, 0.5, 0.75), 10000).over(windowSpec)
)

方法二:自定义UDF + 窗口收集列表(兼容低版本Spark)

如果你的Spark版本低于3.0,没法直接用窗口版的percentile_approx,可以通过「窗口收集分区内所有值 + 自定义分位数UDF」的方式实现。不过要注意:如果分区内的数据量极大,这种方法可能会带来内存压力,因为需要把整个分区的数据加载到内存中处理。

示例代码如下:

import org.apache.spark.sql.expressions.Window
import org.apache.spark.sql.functions.{collect_list, udf, lit}

// 自定义分位数计算UDF,这里实现的是精确分位数逻辑
val getQuantile = udf((valueList: Seq[Double], quantile: Double) => {
  if (valueList == null || valueList.isEmpty) null
  else {
    val sortedList = valueList.sorted
    val index = (quantile * (sortedList.length - 1)).toInt
    sortedList(index)
  }
})

// 定义分区窗口
val windowSpec = Window.partitionBy("category")

// 先收集每个分区的所有value,再计算分位数
val resultDF = df
  .withColumn("all_values", collect_list($"value").over(windowSpec))
  .withColumn("90th_percentile", getQuantile($"all_values", lit(0.9)))
  .drop("all_values") // 移除中间临时列

如果分区数据量很大,建议在UDF中实现近似分位数算法(比如T-Digest),来降低内存占用和计算开销。

再补充一句:你提到的percent_rank()之所以满足不了需求,是因为它的逻辑是「当前行在窗口中的排名占比」,比如某行的percent_rank为0.9,意思是该行比分区内90%的行数值大;而我们要的分位数是「分区内有90%的行数值小于等于这个值」,这是完全不同的两个概念,所以必须用专门的分位数计算函数来实现。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.21 03:54:40