【问题标题】:How to execute function over each element in RDD and join outputs?如何对 RDD 中的每个元素执行函数并连接输出?
【发布时间】:2017-12-07 08:48:07
【问题描述】:

我是 Spark 的新手,我使用的是 Spark 1.6.0。

我有一个 RDD,它是:RDD[Array[Array[String], Long]]

我想通过一个以Array[Array[String], Long] 作为输入并返回ListBuffer[Array[Int]] 作为输出的函数来运行RDD 中的每个元素。

RDD 中每个元素的计算可以并行完成,它们不相互依赖。但是,一旦 RDD 的所有元素都通过函数运行,我想将所有 ListBuffer[Array[Int]] 输出连接到一个 ListBuffer[Array[Int]] 中(此处的顺序也不相关,但它们都应该在相同的数据结构)。

最好的方法是什么?我可以 foreach RDD 并通过函数运行它们,但是我不确定如何处理输出,然后在驱动程序中执行此合并。

这似乎可以通过累加器来实现。前面提到的函数不仅仅是一行代码,它大概是 20+ 行。所以如果我们有这个功能:

def func(data: Array[Array[String], Long]): ListBuffer[Array[Int]] {
    // create ListBuffer
    // iterate over data
        // do some operations on an element in data
        // add some entry to the ListBuffer
    // add the ListBuffer to the Accumulator or return ListBuffer?
}

我怎么能把这一切都包起来呢?我可以这样做吗:

// create Accumulator
// RDD.foreach() // call the func and pass Accumulator as argument?

或者:

val accum = // a ListBuffer[Array[Int]] accumulator
RDD.foreach(x => accum.add(func(x)))

【问题讨论】:

  • 您可以添加示例输入和预期输出吗?另外,你的 RDD 的大小是多少?它可以安全地收集到司机身上吗?你试过使用累加器吗?
  • 很遗憾,我不能提供太多细节。输入本质上是一个 csv 行,而 Long 是附加到它的索引。输出是一个索引矩阵。 RDD 的大小相当大,一个 RDD 可以安全地存储在驱动程序上,但是会有多个这样的程序并行运行(因为 Spark 会监听一个输入源)。我现在要查找 Accumulators。
  • 累加器似乎可以解决单独计算部分输出并将输出添加到累加器的问题。
  • 查看我的编辑,我添加了一些可能使其更清晰的内容

标签: scala apache-spark rdd


【解决方案1】:

TL;DR如果您需要驱动程序上的结果,请使用map,后跟collect

map 将函数应用于 RDD 中的每个元素。

ma​​p[U](f: (T) ⇒ U)(implicit arg0: ClassTag[U]): RDD[U] 通过对 this 的所有元素应用一个函数返回一个新的 RDD RDD。

在您的情况下,TArray[Array[String], Long]UListBuffer[Array[Int]]。如果你有一个函数f 可以转换T 类型的元素,那么map 就是你的朋友。

【讨论】:

  • @osk Pleasure 是我的。如果这对你有用,你能接受它作为答案吗?谢谢。
  • 是的,我正要这样做,只是在测试它。现在感觉就像开始掌握 Spark 的窍门。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2016-10-19
  • 1970-01-01
  • 2014-02-15
  • 2021-03-19
  • 1970-01-01
相关资源
最近更新 更多