【问题标题】:applying a function of every element of an RDD应用 RDD 的每个元素的函数
【发布时间】:2016-11-30 11:53:50
【问题描述】:

我正在使用 spark rdd。我必须对该rdd的每个元素应用一个函数。当我调用rdd.map(x=>function(x)) 时,代码没有给出想要的输出,但是当我调用rdd.collect().foreach(x=>function(x)) 时,代码工作正常。但是collect() 的问题在于它将数据带入内存,这使得大容量数据变得困难。如何在 rdd 的每个元素上调用这个函数?

【问题讨论】:

  • “代码没有给出想要的输出”是什么意思?它会抛出错误吗?它会给出错误的结果吗?什么?
  • 你在 rdd.map(x=>function(x)) 之后有没有调用任何动作?
  • colRDD5.map(x=>println(x)) 没有给出任何结果,而 colRDD5.collect().foreach(x=>println(x)) 打印每个元素
  • @RaviRanjan 那是因为转换是延迟执行的,你必须做一些动作来执行转换 - 即计数,收集

标签: function apache-spark rdd


【解决方案1】:

这是因为 RDD 是不可变的并且是惰性执行的。

当您执行rdd.map(x=>function(x)) 时,您将创建具有应用转换的新 RDD。应用并不意味着执行 - RDD 是转换和操作的谱系,当您键入 rdd.map 时,您正在创建新的 RDD,并在 RDD 图中增加了一个步骤。这就是为什么如果你这样做:

val rdd = // here reading
rdd.map (...)
rdd.collect()

collect() 的结果不会是map 函数转换后的源数据。旧的 RDD 没有改变。

这是您代码中的第一个错误。

其次,在这种情况下,转换,映射,将在触发某些操作时(收集​​、减少等)执行。

请检查:

val mapped = rdd.map(x=>function(x))
// collect is an action, so above transformation map will be executed
mapped.collect().foreach (x => println(x)) // collect will trigger `map` also

转换后会打印内容。例如,如果你这样做mapped.count()map() 也将被执行。在调用动作之前不会执行任何转换,因为 RDD 是惰性的

【讨论】:

  • 但由于内存问题,我无法调用 collext()。现在,本例中的函数在对 rdd 的每个元素运行后返回一个更新的列表。当我做 .map 的事情时,我无法访问该列表,而在其他情况下它可以工作。
  • map 和其他转换仅在您将触发操作时执行,例如collectcount,reduce。 RDD 中的值不能比 map 改变,但是在调用 action 时会懒惰地执行转换
  • 惰性,不是不可变,你不觉得吗?
  • @LostInOverflow 你是什么意思?您不能更改 RDD 本身,但会创建新的 RDD。并且调用将被懒惰地完成。我错了吗?
  • 你是对的,但你描述的更多的是懒惰评估的效果。如果 RDD 是可变的,那也是一样。
【解决方案2】:

试试 rdd.mapPartitions(func)。它应用在每个分区上采用 Iterable(ex:list) 的函数。从技术上讲,您可以对每个分区中的元素进行迭代列表并对每个元素应用所需的函数。

【讨论】:

    【解决方案3】:

    我了解您的function(x) 实际上是println(x)。您没有看到任何打印的原因是因为该函数是在您的工作节点上执行的,所以它会打印到您的工作节点上的标准输出,而 not 在您的驱动程序上。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2017-07-28
      • 2016-06-24
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2014-09-24
      相关资源
      最近更新 更多