【问题标题】:Apache Spark and non-serializable application contextApache Spark 和不可序列化的应用程序上下文
【发布时间】:2016-05-01 19:48:50
【问题描述】:

我是 Spark 的新手。

我想使用 Spark 和 map-reduce 方法并行化我的计算。 但是,我将这个计算放入 Map 阶段的 PairFunction 实现中,需要初始化一些上下文。此上下文包括来自 3rd 方 jar 的多个单例对象,并且此对象不可序列化,因此我无法将它们分布在工作节点上,也无法在我的 PairFunction 中使用它们。

所以我的问题是:我可以使用 Apache Spark 以某种方式并行化需要不可序列化上下文的作业吗?还有其他解决方案吗?也许我可以告诉 Spark 在每个工作节点上初始化所需的上下文?

【问题讨论】:

  • 你的问题对我来说有点模棱两可。我将尝试根据我对它的理解来回答。 Spark 有两个主要的执行环境: 代码将以正常(非分布式)方式运行的驱动程序。您可以在此处初始化上下文并打开 spark 上下文。分布式代码将在worker上执行。
  • 我的问题是关于应该在工人身上执行的分布式代码。问题是这段代码必须使用不可序列化的第三方对象。所以我不能在主服务器上实例化它们一次,然后通过网络传递给工作人员。我想知道是否有任何解决方法。
  • 如果您的代码将被发送给工作人员,则应该对其进行序列化。没有变通办法。如果您在工作人员中不需要这些对象,您可以将它们声明为瞬态。

标签: java serialization apache-spark mapreduce


【解决方案1】:

您可以尝试使用 mapPartitionforeachPartition 在执行程序中初始化您的 3rd 方 jar。

rdd.foreachPartition { iter =>
  //initialize here
  val object = new XXX()
  iter.foreach { p =>
    //then you can use object here
  }
}

【讨论】:

  • 谢谢。您能否解释一下,这些 rdd 方法到底在做什么?我打开了 spark javadocs,没有太多细节:“foreachPartition - 将函数 f 应用于此 RDD 的每个分区。”就是这样。
  • foreachPartition 为每个分区执行一个函数。通过迭代器参数提供对分区中包含的数据项的访问。 Spark 将尝试在驱动程序(master)上初始化一个变量,然后序列化对象以将其发送给工作人员,如果对象不可序列化,它将失败。
猜你喜欢
  • 2017-01-30
  • 1970-01-01
  • 1970-01-01
  • 2015-06-17
  • 2023-03-12
  • 2017-10-22
  • 2015-09-21
  • 1970-01-01
  • 2015-08-31
相关资源
最近更新 更多