【问题标题】:Spark Context not Serializable?Spark 上下文不可序列化?
【发布时间】:2017-01-30 18:16:50
【问题描述】:

所以,我在 Spark 中遇到了臭名昭著的 Task Not Serializable 错误。这是相关的代码块:

val labeledPoints: RDD[LabeledPoint] = events.map(event => {

    var eventsPerEntity = try {
          HBaseHelper.scan(...filter entity here...)(sc).map(newEvent =>  {

            Try(new Object(...))
          }).filter(_.isSuccess).map(_.get)
        } catch {
          case e: Exception => {
            logger.error(s"Failed to convert event ${event}." +
              s"Exception: ${e}.")
            throw e
          }
        }
    })

基本上我想要实现的是我正在访问sc,这是我在map 中的Spark Context 对象。在运行时,我收到Task Not Serializable 错误。

这是我能想到的一个潜在解决方案:

在没有sc 的情况下查询 HBase,我可以这样做,但反过来我会得到一个列表。 (如果我尝试并行化;我必须再次使用sc)。有一个列表会导致我无法使用reduceByKey,在我的另一个问题中建议here。所以我也无法成功实现这一目标,因为我不知道如果没有reduceByKey,我将如何实现this。另外我真的很想使用 RDD :)

所以我正在寻找另一种解决方案 + 询问我是否做错了什么。提前致谢!

更新

所以基本上,我的问题变成了这样:

我有一个名为eventsRDD。这是整个 HBase 表。注意:每个event 都由performerId 执行,这又是event 中的一个字段,即event.performerId

对于events 中的每个event,我需要计算event.numericColumnevents(events 的子集)的平均值的比率,这些events 由相同的@ 执行987654344@.

我在映射events 时尝试这样做。在map 内,我试图根据他们的performerId 过滤事件。

基本上,我正在尝试将每个event 转换为LabeledPoint,上面的比率将成为我的向量中的一个特征。即对于每一个事件,我都试图得到

// I am trying to calculate the average, but cannot use filter, because I am in map block.

LabeledPoint(
  event.someColumn,
  Vectors.dense(
    averageAbove,
    ...
  )
)

我将不胜感激。谢谢!

【问题讨论】:

    标签: java scala apache-spark filter rdd


    【解决方案1】:

    如果适用,一个选项是加载 整个 HBase 表(或 - 所有可能匹配 events RDD 中的事件之一的元素,如果您有任何方法可以在没有遍历 RDD) 进入 Dataframe,然后使用 join

    要将 HBase 表中的数据加载到 Dataframe 中,您可以使用 Hortonworks 的预览版Spark-HBase Connector。然后,在两个数据帧之间执行正确的连接操作应该很容易。

    【讨论】:

    • 加入后变成了val someValue = kv._2._1._1._1._1._2,但基本上就是这样!
    【解决方案2】:

    您可以将列表添加为事件的新字段 - 通过获取新的 RDD(事件+实体列表)。然后,您可以使用常规 Spark 命令“分解”列表,从而获得多个事件+列表项记录(尽管使用 DataFrames/DataSets 比使用 RDDs 更容易做到这一点)

    【讨论】:

    • 这将是很多重复的条目,因为将添加到事件的列表只是同一实体执行的另一个事件。不会吧?
    • 那么你可以做 dedup - 或者我不明白你想要做什么
    • 基本上,对于每个event,我都试图获取由执行event 的同一实体完成的其他事件,以计算一些特征。所以,我试图通过映射事件来做到这一点,在映射中,我试图访问 SparkContext 以查询 HBase 以获取同一实体的其他相关事件。
    • 您有以下三个选项之一 - 正如 Tzach 建议的那样,您阅读所有数据库并加入两个数据集。正如我建议您从 HBase 中读取而不考虑 spark 那样,将其添加到相关行并使用 Spark 机制将其分解,或者您在 Spark 之外处理来自 HBase 的所有数据(使用“普通”scala)不是一个糟糕的选择,因为它将被分发。请注意,在选项 2,3 中,您将与 HBase 有很多连接(每个任务)
    • 如果它是同一个数据集,则意味着根据您的标准将其加入自身
    【解决方案3】:

    很简单,你不能在 RDD Closure 上使用 spark 上下文,所以找另一种方法来处理这个问题。

    【讨论】:

    • 这根本不是答案。
    • 是的,但这也不是问题,因为当您知道为什么会出现任务可序列化错误时,请更改您的问题,这完全与您的方法有关。
    • 你很粗鲁,没有任何帮助。不过感谢您的努力。
    • 我一点也不粗鲁,我只是告诉你,如果你改变关于如何解决这个问题的问题,它对回答这个问题的人更有帮助,因为当你询问任务可序列化错误时,我给出了理由如果您觉得这很粗鲁,请在抱歉之间。
    • @SandeepPurohit 好的,让我们继续吧:)
    猜你喜欢
    • 2016-05-01
    • 2018-04-06
    • 2021-03-16
    • 2019-06-25
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2015-12-16
    相关资源
    最近更新 更多