【发布时间】:2015-03-26 00:07:45
【问题描述】:
我有一个 RDD 集合中的数字列表。从这个列表中,我需要创建另一个 RDD 列表,其中每个元素等于它之前所有元素的总和。如何在 Spark 中构建这样的 RDD?
以下 Scala 代码说明了我试图在 Spark 中实现的目标:
object Test {
def main(args: Array[String]) {
val lst: List[Float] = List(1, 2, 3)
val result = sum(List(), 0, lst)
println(result)
}
def sum(acc: List[Float], runningSum: Float, list: List[Float]): List[Float] = {
list match {
case List() => acc.reverse
case List(x, _*) => {
val newSum = runningSum + x
sum(newSum :: acc, newSum, list.tail)
}
}
}
运行此结果:
List(1.0, 3.0, 6.0)
此示例的等效 Spark 代码是什么?
【问题讨论】:
-
你不能用 RDD 来做这件事(至少,不能获得 RDD 的优势); RDD 的重点是并行处理其中的一部分,而在您的情况下,新列表的元素取决于原始列表的每个元素。
-
明白,但我仍然需要从 RDD 中计算这个总和。在这种情况下你会怎么做?
-
理想情况下,
.collect()RDD 并在本地计算总和,然后在需要时再次sc.parallelize。如果它太大而无法放入内存,我所能想到的就是找出某种方法来知道哪些元素比其他元素“更早”(也许通过给每个元素一个索引),@ 987654325@ RDD 本身以获得所有可能元素对,filter取出第二个条目在第一个条目“之后”的那些对,然后aggregateByKey进行求和。它会起作用,但肯定会很慢。
标签: scala apache-spark