【问题标题】:How to compute percentiles in Apache Spark如何在 Apache Spark 中计算百分位数
【发布时间】:2015-05-02 13:22:58
【问题描述】:

我有一个整数 rdd(即RDD[Int]),我想做的是计算以下十个百分位数:[0th, 10th, 20th, ..., 90th, 100th]。最有效的方法是什么?

【问题讨论】:

    标签: apache-spark


    【解决方案1】:

    你可以:

    1. 通过 rdd.sortBy() 对数据集进行排序
    2. 通过 rdd.count() 计算数据集的大小
    3. 带索引的压缩以方便百分位检索
    4. 通过 rdd.lookup() 检索所需的百分位数,例如对于第 10 个百分位 rdd.lookup(0.1 * size)

    要计算中位数和第 99 个百分位数: getPercentiles(rdd, new double[]{0.5, 0.99}, size, numPartitions);

    在 Java 8 中:

    public static double[] getPercentiles(JavaRDD<Double> rdd, double[] percentiles, long rddSize, int numPartitions) {
        double[] values = new double[percentiles.length];
    
        JavaRDD<Double> sorted = rdd.sortBy((Double d) -> d, true, numPartitions);
        JavaPairRDD<Long, Double> indexed = sorted.zipWithIndex().mapToPair((Tuple2<Double, Long> t) -> t.swap());
    
        for (int i = 0; i < percentiles.length; i++) {
            double percentile = percentiles[i];
            long id = (long) (rddSize * percentile);
            values[i] = indexed.lookup(id).get(0);
        }
    
        return values;
    }
    

    请注意,这需要对数据集进行排序,O(n.log(n)) 并且在大型数据集上可能会很昂贵。

    建议简单计算直方图的另一个答案不会正确计算百分位数:这是一个反例:一个由 100 个数字组成的数据集,99 个数字为 0,一个数字为 1。您最终得到所有 99 个 0在第一个 bin 中,1 在最后一个 bin 中,中间有 8 个空 bin。

    【讨论】:

    • 如果 N 百分比很小,比如 10、20%,那么我将执行以下操作:
    【解决方案2】:

    t-digest怎么样?

    https://github.com/tdunning/t-digest

    一种新的数据结构,用于准确在线累积基于等级的统计数据,例如分位数和修剪均值。 t-digest 算法还对并行非常友好,使其在 map-reduce 和并行流应用程序中非常有用。

    t-digest 构造算法使用一维 k-means 聚类的变体来生成与 Q-digest 相关的数据结构。这种 t-digest 数据结构可用于估计分位数或计算其他排名统计信息。 t-digest 优于 Q-digest 的优点是 t-digest 可以处理浮点值,而 Q-digest 仅限于整数。通过小的更改,t-digest 可以处理任何有序集合中的任何值,这些值类似于均值。尽管 t-digests 存储在磁盘上时更紧凑,但 t-digests 产生的分位数估计的准确性可能比 Q-digests 产生的精确度高几个数量级。

    总之,t-digest 的特别有趣的特点是它

    • 摘要比 Q-digest 小
    • 适用于双精度数和整数。
    • 为极端分位数提供百万分之一的准确度,为中间分位数提供通常
    • 速度很快
    • 很简单
    • 有一个测试覆盖率 > 90% 的参考实现
    • 可以很容易地与 map-reduce 一起使用,因为可以合并摘要

    使用 Spark 中的参考 Java 实现应该相当容易。

    【讨论】:

    • 实际上这里有 Erik Erlandson 的 spark 实现:github.com/isarn/isarn-sketches-spark。它工作得很好。我发现的唯一问题是您无法将 TDigest 对象保存为镶木地板格式。只要您只是扔进大量数据并要求一些百分位结果,这就是您要寻找的:)
    【解决方案3】:

    我发现了这个要点

    https://gist.github.com/felixcheung/92ae74bc349ea83a9e29

    包含以下功能:

      /**
       * compute percentile from an unsorted Spark RDD
       * @param data: input data set of Long integers
       * @param tile: percentile to compute (eg. 85 percentile)
       * @return value of input data at the specified percentile
       */
      def computePercentile(data: RDD[Long], tile: Double): Double = {
        // NIST method; data to be sorted in ascending order
        val r = data.sortBy(x => x)
        val c = r.count()
        if (c == 1) r.first()
        else {
          val n = (tile / 100d) * (c + 1d)
          val k = math.floor(n).toLong
          val d = n - k
          if (k <= 0) r.first()
          else {
            val index = r.zipWithIndex().map(_.swap)
            val last = c
            if (k >= c) {
              index.lookup(last - 1).head
            } else {
              index.lookup(k - 1).head + d * (index.lookup(k).head - index.lookup(k - 1).head)
            }
          }
        }
      }
    

    【讨论】:

      【解决方案4】:

      这是我在 Spark 上的 Python 实现,用于计算包含感兴趣值的 RDD 的百分位数。

      def percentile_threshold(ardd, percentile):
          assert percentile > 0 and percentile <= 100, "percentile should be larger then 0 and smaller or equal to 100"
      
          return ardd.sortBy(lambda x: x).zipWithIndex().map(lambda x: (x[1], x[0])) \
                  .lookup(np.ceil(ardd.count() / 100 * percentile - 1))[0]
      
      # Now test it out
      import numpy as np
      randlist = range(1,10001)
      np.random.shuffle(randlist)
      ardd = sc.parallelize(randlist)
      
      print percentile_threshold(ardd,0.001)
      print percentile_threshold(ardd,1)
      print percentile_threshold(ardd,60.11)
      print percentile_threshold(ardd,99)
      print percentile_threshold(ardd,99.999)
      print percentile_threshold(ardd,100)
      
      # output:
      # 1
      # 100
      # 6011
      # 9900
      # 10000
      # 10000
      

      另外,我定义了以下函数来获取第 10 到第 100 个百分位数。

      def get_percentiles(rdd, stepsize=10):
          percentiles = []
          rddcount100 = rdd.count() / 100 
          sortedrdd = ardd.sortBy(lambda x: x).zipWithIndex().map(lambda x: (x[1], x[0]))
      
      
          for p in range(0, 101, stepsize):
              if p == 0:
                  pass
                  # I am not aware of a formal definition of 0 percentile, 
                  # you can put a place holder like this if you want
                  # percentiles.append(sortedrdd.lookup(0)[0] - 1) 
              elif p == 100:
                  percentiles.append(sortedrdd.lookup(np.ceil(rddcount100 * 100 - 1))[0])
              else:
                  pv = sortedrdd.lookup(np.ceil(rddcount100 * p) - 1)[0]
                  percentiles.append(pv)
      
          return percentiles
      
      randlist = range(1,10001)
      np.random.shuffle(randlist)
      ardd = sc.parallelize(randlist)
      get_percentiles(ardd, 10)
      
      # [1000, 2000, 3000, 4000, 5000, 6000, 7000, 8000, 9000, 10000]
      

      【讨论】:

      • 不应该在sortedrdd 定义get_percentiles 中将ardd 替换为rdd 吗?以及添加import numpy as np。物联网似乎不适用于numpy 1.11.3
      【解决方案5】:

      如果您不介意将 RDD 转换为 DataFrame 并使用 Hive UDAF,则可以使用 percentile。假设您已将 HiveContext hiveContext 加载到作用域中:

      hiveContext.sql("SELECT percentile(x, array(0.1,0.2,0.3,0.4,0.5,0.6,0.7,0.8,0.9)) FROM yourDataFrame")

      我在 this answer. 发现了这个 Hive UDAF

      【讨论】:

        【解决方案6】:

        将您的 RDD 转换为 Double 的 RDD,然后使用 .histogram(10) 操作。见DoubleRDD ScalaDoc

        【讨论】:

        • .histogram(bucketCount) 不计算百分位数,它“使用 bucketCount 的桶数计算数据的直方图 在 RDD 的最小值和最大值之间均匀间隔
        【解决方案7】:

        如果 N 百分比很小,比如 10、20%,那么我将执行以下操作:

        1. 计算数据集的大小,rdd.count(),跳过它也许你已经知道它并作为参数。

        2. 与其对整个数据集进行排序,不如从每个分区中找出top(N)。为此,我必须找出 N = rdd.count 的 N%,然后对分区进行排序并从每个分区中取 top(N)。现在你有一个小得多的数据集要排序。

        3.rdd.sortBy

        4.zipWithIndex

        5.filter(索引

        【讨论】:

          【解决方案8】:

          根据Median UDAF in Spark/Scala 此处给出的答案,我使用 UDAF 计算 spark 窗口 (spark 2.1) 上的百分位数:

          首先是用于其他聚合的抽象通用 UDAF

          import org.apache.spark.sql.Row
          import org.apache.spark.sql.expressions.{MutableAggregationBuffer, UserDefinedAggregateFunction}
          import org.apache.spark.sql.types._
          
          import scala.collection.mutable
          import scala.collection.mutable.ArrayBuffer
          
          
          abstract class GenericUDAF extends UserDefinedAggregateFunction {
          
            def inputSchema: StructType =
              StructType(StructField("value", DoubleType) :: Nil)
          
            def bufferSchema: StructType = StructType(
              StructField("window_list", ArrayType(DoubleType, false)) :: Nil
            )
          
            def deterministic: Boolean = true
          
            def initialize(buffer: MutableAggregationBuffer): Unit = {
              buffer(0) = new ArrayBuffer[Double]()
            }
          
            def update(buffer: MutableAggregationBuffer,input: org.apache.spark.sql.Row): Unit = {
              var bufferVal = buffer.getAs[mutable.WrappedArray[Double]](0).toBuffer
              bufferVal+=input.getAs[Double](0)
              buffer(0) = bufferVal
            }
          
            def merge(buffer1: MutableAggregationBuffer, buffer2: org.apache.spark.sql.Row): Unit = {
              buffer1(0) = buffer1.getAs[ArrayBuffer[Double]](0) ++ buffer2.getAs[ArrayBuffer[Double]](0)
            }
          
            def dataType: DataType
            def evaluate(buffer: Row): Any
          
          }
          

          然后是为十分位数定制的百分位 UDAF:

          import org.apache.spark.sql.Row
          import org.apache.spark.sql.expressions.{MutableAggregationBuffer, UserDefinedAggregateFunction}
          import org.apache.spark.sql.types._
          
          import scala.collection.mutable
          import scala.collection.mutable.ArrayBuffer
          
          
          class DecilesUDAF extends GenericUDAF {
          
            override def dataType: DataType = ArrayType(DoubleType, false)
          
            override def evaluate(buffer: Row): Any = {
              val sortedWindow = buffer.getAs[mutable.WrappedArray[Double]](0).sorted.toBuffer
              val windowSize = sortedWindow.size
              if (windowSize == 0) return null
              if (windowSize == 1) return (0 to 10).map(_ => sortedWindow.head).toArray
          
              (0 to 10).map(i => sortedWindow(Math.min(windowSize-1, i*windowSize/10))).toArray
          
            }
          }
          

          然后,UDAF 被实例化并在一个分区和有序的窗口上调用:

          val deciles = new DecilesUDAF()
          df.withColumn("mt_deciles", deciles(col("mt")).over(myWindow))
          

          然后您可以使用 getItem 将结果数组拆分为多列:

          def splitToColumns(size: Int, splitCol:String)(df: DataFrame) = {
            (0 to size).foldLeft(df) {
              case (df_arg, i) => df_arg.withColumn("mt_decile_"+i, col(splitCol).getItem(i))
            }
          }
          
          df.transform(splitToColumns(10, "mt_deciles" ))
          

          UDAF 比原生 spark 函数慢,但只要每个分组包或每个窗口相对较小并且适合单个执行器,应该没问题。主要优点是使用火花并行。 不费吹灰之力,这段代码就可以扩展到 n 分位数。

          我用这个函数测试了代码:

          def testDecilesUDAF = {
              val window = W.partitionBy("user")
              val deciles = new DecilesUDAF()
          
              val schema = StructType(StructField("mt", DoubleType) :: StructField("user", StringType) :: Nil)
          
              val rows1 = (1 to 20).map(i => Row(i.toDouble, "a"))
              val rows2 = (21 to 40).map(i => Row(i.toDouble, "b"))
          
              val df = spark.createDataFrame(spark.sparkContext.makeRDD[Row](rows1++rows2), schema)
          
              df.withColumn("deciles", deciles(col("mt")).over(window))
                .transform(splitToColumns(10, "deciles" ))
                .drop("deciles")
                .show(100, truncate=false)
            }
          

          前 3 行输出:

          +----+----+-----------+-----------+-----------+-----------+-----------+-----------+-----------+-----------+-----------+-----------+------------+
          |mt  |user|mt_decile_0|mt_decile_1|mt_decile_2|mt_decile_3|mt_decile_4|mt_decile_5|mt_decile_6|mt_decile_7|mt_decile_8|mt_decile_9|mt_decile_10|
          +----+----+-----------+-----------+-----------+-----------+-----------+-----------+-----------+-----------+-----------+-----------+------------+
          |21.0|b   |21.0       |23.0       |25.0       |27.0       |29.0       |31.0       |33.0       |35.0       |37.0       |39.0       |40.0        |
          |22.0|b   |21.0       |23.0       |25.0       |27.0       |29.0       |31.0       |33.0       |35.0       |37.0       |39.0       |40.0        |
          |23.0|b   |21.0       |23.0       |25.0       |27.0       |29.0       |31.0       |33.0       |35.0       |37.0       |39.0       |40.0        |
          

          【讨论】:

            【解决方案9】:

            另一种替代方法是在 double 的 RDD 上使用 top 和 last。例如,val percentile_99th_value=scores.top((count/100).toInt).last

            此方法更适合单个百分位数。

            【讨论】:

              【解决方案10】:

              这是我的简单方法:

              val percentiles = Array(0.1, 0.2, 0.3, 0.4, 0.5, 0.6, 0.7, 0.8, 0.9, 1)
              val accuracy = 1000000
              df.stat.approxQuantile("score", percentiles, 1.0/accuracy)
              

              输出:

              scala> df.stat.approxQuantile("score", percentiles, 1.0/accuracy)
              res88: Array[Double] = Array(0.011044141836464405, 0.02022990956902504, 0.0317261666059494, 0.04638145491480827, 0.06498630344867706, 0.0892181545495987, 0.12161539494991302, 0.16825592517852783, 0.24740923941135406, 0.9188197255134583)
              

              accuracy:accuracy 参数(默认值:10000)是一个正数字面量,它以内存为代价控制近似精度。精度值越高,精度越高,1.0/accuracy 是近似值的相对误差。

              【讨论】:

                猜你喜欢
                • 1970-01-01
                • 2013-08-31
                • 2016-04-12
                • 2016-07-12
                • 2015-03-09
                • 2011-12-29
                • 2013-06-20
                • 2017-10-12
                • 2017-10-19
                相关资源
                最近更新 更多