【问题标题】:ClassTag causes Spark to serialize objectClassTag 导致 Spark 序列化对象
【发布时间】:2015-08-07 03:15:48
【问题描述】:

以下代码因org.apache.spark.SparkException: Task not serializable 异常而失败:

import org.apache.spark.rdd.RDD
import scala.reflect.ClassTag

class Foo[T](rdd: RDD[T])(implicit kt: ClassTag[T]) {
  def die() {
    rdd.map(_ => Array[T]()).count()
  }
}

val x = sc.parallelize(Array(1, 2, 3, 4, 5))
val foo = new Foo(x)
foo.die()

因为Foo 不可序列化。为什么将函数文字传递给map 导致Foo 在仅引用隐式参数ClassTag 时被序列化?我该如何解决?当Foo 工作于Int 而不是T 时,此方法有效。我的实际代码正在尝试toArray,但这是同样的问题。谢谢!

编辑:我正在使用 spark-shell 运行它。这是序列化堆栈:

Serialization stack:
    - object not serializable (class: $iwC$$iwC$Foo, value: $iwC$$iwC$Foo@4ce1292a)
    - field (class: $iwC$$iwC$Foo$$anonfun$die$1, name: $outer, type: class $iwC$$iwC$Foo)
    - object (class $iwC$$iwC$Foo$$anonfun$die$1, <function1>)
    at org.apache.spark.serializer.SerializationDebugger$.improveException(SerializationDebugger.scala:40)
    at org.apache.spark.serializer.JavaSerializationStream.writeObject(JavaSerializer.scala:47)
    at org.apache.spark.serializer.JavaSerializerInstance.serialize(JavaSerializer.scala:81)
    at org.apache.spark.util.ClosureCleaner$.ensureSerializable(ClosureCleaner.scala:312)
    ... 56 more

【问题讨论】:

  • 能否请您在提交时将-Dsun.io.serialization.extendedDebugInfo=true添加到应用中,并显示结果?

标签: scala serialization apache-spark


【解决方案1】:

这并不特定于 implicitClassTag。以下作品:

import org.apache.spark.rdd.RDD
import scala.reflect.ClassTag

class Foo[T](rdd: RDD[T])(implicit kt: ClassTag[T]) {
  def die() {
    val localKt = kt
    rdd.map(_ => Array[T]()(localKt)).count()
  }
}

val x = sc.parallelize(Array(1, 2, 3, 4, 5))
val foo = new Foo(x)
foo.die()

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2013-08-06
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多