【问题标题】:Scala fold for RDD [String] acting wierd [duplicate]RDD [String] 的 Scala 折叠行为怪异 [重复]
【发布时间】:2018-05-22 13:37:31
【问题描述】:

这是我的代码:

var exceptions: String = ""
  val count = failures.fold("0")((x1, x2) => {
    println(s"x1: $x1 and x2: $x2")
    if (x1 != "-433" && x2 != "-433") {
      (x1.toInt + x2.toInt).toString
    } else {
      println(s"before: $exceptions")
      exceptions = exceptions + ", " + "There is an exception in processing, check the logs of executors for actual information"
      println(s"after: $exceptions")
      if (x1 == "-433") {
        if (x2 != "-433") {
          x2
        } else "0"
      }else {
        if (x2 == "-433") {
           x1
        }else "0"
      }
    }
  })

count 是一个 RDD[String]。最奇怪的是 execptions 以“”结尾。以下是日志:

x1:0 和 x2:-433

之前:

after: , 处理中出现异常,查看executors日志获取实际信息

x1:0 和 x2:0

最终:

【问题讨论】:

  • 它是否与那些要求驱动程序变量未在并行化集合中更新的问题之一重复?请阅读有关 Spark 闭包的信息。你会得到一个更好的主意。
  • @RameshMaharjan:我不是在谈论抛出任何异常。 exceptions 变量没有得到更新。但正如@philantrovert 指出的那样,由于驱动程序变量被传递给一组执行程序

标签: scala apache-spark fold


【解决方案1】:

Spark RDD 是分布式数据结构,这意味着 RDD 的元素不在同一个节点上,并且 RDD 上的转换不会发生在同一个节点上。

您的所有转换函数都包装在它们的闭包中,然后序列化,然后发送到执行程序节点,在执行程序节点反序列化,在执行程序节点执行。包装的闭包获得使用的上下文对象的“副本”。

即使您在本地节点上运行 spark,它也会在本地节点上运行多个执行器,并且它们的行为方式相同。

因此,您的函数的每次执行都将获得自己的 exceptions 变量副本。因此,您的 exception 字符串永远不会以您期望的方式更新。

【讨论】:

    猜你喜欢
    • 2021-10-16
    • 2017-02-17
    • 2021-09-28
    • 1970-01-01
    • 2020-08-17
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-06-13
    相关资源
    最近更新 更多