【问题标题】:Using Class Methods in Spark RDD Operations Returns Task not serializable Exception在 Spark RDD 操作中使用类方法返回任务不可序列化异常
【发布时间】:2016-11-11 03:40:13
【问题描述】:

假设我在 Spark Scala 中有以下课程:

class SparkComputation(i: Int, j: Int) {
  def something(x: Int, y: Int) = (x + y) * i

  def processRDD(data: RDD[Int]) = {
    val j = this.j
    val something = this.something _
    data.map(something(_, j))
  }
}

当我运行以下代码时,我得到了Task not serializable Exception

val s = new SparkComputation(2, 5)
val data = sc.parallelize(0 to 100)
val res = s.processRDD(data).collect

我假设发生异常是因为 Spark 正在尝试序列化 SparkComputation 实例。为了防止这种情况发生,我将我在 RDD 操作中使用的类成员存储在局部变量中(jsomething)。但是,由于该方法,Spark 仍然尝试序列化SparkComputation 对象。无论如何将类方法传递给map 而不强制Spark 序列化整个SparkComputation 类?我知道以下代码可以正常工作:

def processRDD(data: RDD[Int]) = {
    val j = this.j
    val i = this.i
    data.map(x => (x + j) * i)
  }

因此,存储值的类成员不会导致问题。问题出在功能上。 我也尝试过以下方法,但没有成功:

class SparkComputation(i: Int, j: Int) {
  def processRDD(data: RDD[Int]) = {
    val j = this.j
    val i = this.i
    def something(x: Int, y: Int) = (x + y) * i
    data.map(something(_, j))
  }
}

【问题讨论】:

  • 这是由于类声明而发生的著名错误...将其更改为对象或扩展可序列化然后它应该可以工作。
  • @RamPrasadG 没有办法像我对其他成员(在此示例中为ij)那样将方法存储在本地值中?

标签: scala serialization apache-spark


【解决方案1】:

使类可序列化:

class SparkComputation(i: Int, j: Int) extends Serializable {
  def something(x: Int, y: Int) = (x + y) * i

  def processRDD(data: RDD[Int]) = {
    val j = this.j
    val something = this.something _
    data.map(something(_, j))
  }
}

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-11-11
    • 1970-01-01
    • 1970-01-01
    • 2018-04-08
    • 2015-09-15
    • 2015-08-10
    相关资源
    最近更新 更多