【发布时间】: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)
如果我再次收集 RDD 和 parallelize 需要大约 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