【发布时间】: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 操作中使用的类成员存储在局部变量中(j 和something)。但是,由于该方法,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 没有办法像我对其他成员(在此示例中为
i和j)那样将方法存储在本地值中?
标签: scala serialization apache-spark