如何使用Scala计算Spark DataFrame列中整数列表的近似分位数ApproxQuantiles
Scala实现方案
实现思路
由于是针对每行的数组单独计算分位数,我们通过自定义UDF(用户自定义函数)实现,UDF接收整数数组作为输入,返回对应分位数值,再通过withColumn方法将计算结果作为新列追加到原DataFrame中。
完整代码实现
首先导入依赖:
import org.apache.spark.sql.functions.{col, udf}
然后定义分位数计算UDF,以下代码的分位数计算逻辑可完全对齐你给出的示例结果:
// 计算25分位UDF val calculateQT25 = udf((arr: Seq[Int]) => { if (arr == null || arr.isEmpty) return null val sortedArr = arr.sorted val n = sortedArr.length val pos = (0.25 * (n - 1)).floor.toInt sortedArr(pos) }) // 计算75分位UDF val calculateQT75 = udf((arr: Seq[Int]) => { if (arr == null || arr.isEmpty) return null val sortedArr = arr.sorted val n = sortedArr.length val pos = (0.75 * (n - 1)).ceil.toInt sortedArr(pos) })
调用UDF生成结果:
// 假设原输入DataFrame名为inputDf val resultDf = inputDf .withColumn("QT25", calculateQT25(col("List_Nb_total_operations"))) .withColumn("QT75", calculateQT75(col("List_Nb_total_operations"))) // 打印结果验证 resultDf.show()
说明
分位数存在多种计算标准,若后续需要调整分位数计算规则,仅修改UDF内的pos位置计算逻辑即可适配。
内容的提问来源于stack exchange,提问作者Fahd Zaghdoudi
相关产品推荐
相关产品推荐

