【问题标题】:Efficient calculation of entropy in sparkspark中熵的高效计算
【发布时间】:2014-07-19 07:34:31
【问题描述】:

给定一个 RDD(数据)和一个索引字段列表来计算熵。执行以下流程时,在 2MB(16k 行)的源上计算单个熵值大约需要 5 秒。

def entropy(data: RDD[Array[String]], colIdx: Array[Int], count: Long): Double = { 
  println(data.toDebugString)
    data.map(r => colIdx.map(idx => r(idx)).mkString(",") -> 1)
        .reduceByKey(_ + _)
        .map(v => {
        val p = v._2.toDouble / count
        -p * scala.math.log(p) / scala.math.log(2)
      })
        .reduce((v1, v2) => v1 + v2)
}

debugString的输出如下:

(entropy,MappedRDD[93] at map at Q.scala:31 (8 partitions)
  UnionRDD[72] at $plus$plus at S.scala:136 (8 partitions)
    MappedRDD[60] at map at S.scala:151 (4 partitions)
      FilteredRDD[59] at filter at S.scala:150 (4 partitions)
        MappedRDD[40] at map at S.scala:124 (4 partitions)
          MapPartitionsRDD[39] at mapPartitionsWithIndex at L.scala:356 (4 partitions)
            FilteredRDD[27] at filter at S.scala:104 (4 partitions)
              MappedRDD[8] at map at X.scala:21 (4 partitions)
                MappedRDD[6] at map at R.scala:39 (4 partitions)
                  FlatMappedRDD[5] at objectFile at F.scala:51 (4 partitions)
                    HadoopRDD[4] at objectFile at F.scala:51 (4 partitions)
    MappedRDD[68] at map at S.scala:151 (4 partitions)
      FilteredRDD[67] at filter at S.scala:150 (4 partitions)
        MappedRDD[52] at map at S.scala:124 (4 partitions)
          MapPartitionsRDD[51] at mapPartitionsWithIndex at L.scala:356 (4 partitions)
            FilteredRDD[28] at filter at S.scala:105 (4 partitions)
              MappedRDD[8] at map at X.scala:21 (4 partitions)
                MappedRDD[6] at map at R.scala:39 (4 partitions)
                  FlatMappedRDD[5] at objectFile at F.scala:51 (4 partitions)
                    HadoopRDD[4] at objectFile at F.scala:51 (4 partitions),colIdex,13,count,3922)

如果我再次收集 RDDparallelize 需要大约 150 毫秒来计算(对于一个简单的 2MB 文件来说这似乎仍然很高) - 并且在处理时显然会产生挑战具有多个 GB 数据。正确使用 Spark 和 Scala 我缺少什么?

我最初的实现(表现更差):

data.map(r => colIdx
  .map(idx => r(idx)).mkString(","))
  .groupBy(r => r)
  .map(g => g._2.size)
  .map(v => v.toDouble / count)
  .map(v => -v * scala.math.log(v) / scala.math.log(2))
  .reduce((v1, v2) => v1 + v2)

【问题讨论】:

  • 5 秒是 JVM 中“瞬时”的时间。可能只是序列化、反序列化和 GC 的几次传递,而不是您的代码的错误。尝试更多数据!您必须在几分钟内对其进行测量,然后才能判断它是否慢。
  • 作为数据点:我已经看到 sc.parallelize 阶段在本地 Spark 设置的
  • 明确地说,我在相同的数据上运行了多次(具有不同的 colIdx 值)。如果性能仅在第一次运行时受到影响,我会很好,但缓慢的运行时间一直贯穿整个流程(在我当前的示例中,这个特定的例程被调用了 794 次)
  • 我看到了对性能的线性影响。我尝试了一个 2GB 的数据集(1M 行),每次调用都运行了几十分钟。测试机有16核,32GB内存。

标签: performance scala apache-spark entropy information-theory


【解决方案1】:

首先,您的代码中似乎存在错误,您需要处理0 中的p,因此-p * math.log(p) / math.log(2) 应该是if (p == 0.0) 0.0 else -p * math.log(p) / math.log(2)

其次,您可以使用底数 e,您实际上不需要底数为 2。

无论如何,您的代码运行缓慢的原因可能是分区太少。每个 CPU 至少应该有 2-4 个分区,实际上我经常使用更多。你有多少个 CPU?

现在可能花费最长时间的不是熵计算,因为它非常简单 - 而是在 String 键上完成的 reduceByKey。是否可以使用其他数据类型? colIdx 实际上是什么? r 究竟是什么?

最后一个观察结果是你用这个colIdx.map(r.apply) 多次索引每条记录......你知道如果r 不是ArrayIndexedSeq 类型,这将非常慢......如果它是@ 987654331@ 这将是 O(index) 因为你必须遍历列表才能得到你想要的索引。

【讨论】:

  • 由于代码计算列值的唯一实例,因此 p 在最后一步不可能为 0。
  • 我将分区更改为 32(在 16 核机器上测试),但是我还没有看到它使用超过 200% 的 CPU。有了这个,我看到计算时间下降到~2.7s。尽管我看到处理时间的剧烈波动。我添加了 colId (Array[Int]) 和 data (RDD[Array[String]) 的定义。这使得 r 成为一个 Array[String]。
  • 这是计算指定列的唯一值组合。目前这些值是字符串。没有应用哈希(可能违反唯一检查)我不知道如何改进?
  • 我刚刚发现我们用于控制分析步骤的框架导致线程争用,这就是导致序列化处理的原因。这显然不是火花问题。然而,即使存在争用,您的洞察力也帮助减少了 30% 的运行时间。非常感谢!现在来弄清楚为什么我们的 Akka 演员行为不端。
  • 如果可以建立索引,即Map[Int, String],那么您应该能够在List[Int] 上执行reduceByKey - 我希望这会加快速度。祝您的争用问题好运,我想您可以尝试暂停所述框架以确认您的怀疑。
猜你喜欢
  • 2021-01-26
  • 2019-10-09
  • 2018-05-28
  • 2015-01-31
  • 1970-01-01
  • 1970-01-01
  • 2014-01-02
  • 1970-01-01
相关资源
最近更新 更多