【问题标题】:Serialization and Custom Spark RDD Class序列化和自定义 Spark RDD 类
【发布时间】:2015-05-08 16:49:36
【问题描述】:

我正在使用 Scala 编写自定义 Spark RDD 实现,并且正在使用 Spark shell 调试我的实现。我现在的目标是:

customRDD.count

在没有异常的情况下成功。现在这就是我得到的:

15/03/06 23:02:32 INFO TaskSchedulerImpl: Adding task set 0.0 with 1 tasks
15/03/06 23:02:32 ERROR TaskSetManager: Failed to serialize task 0, not attempting to retry it.
java.lang.reflect.InvocationTargetException
    at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method)
    at sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:57)
    at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)
    at java.lang.reflect.Method.invoke(Method.java:606)
    at org.apache.spark.serializer.SerializationDebugger$ObjectStreamClassMethods$.getObjFieldValues$extension(SerializationDebugger.scala:240)

...

Caused by: java.lang.ArrayIndexOutOfBoundsException: 1
    at java.io.ObjectStreamClass$FieldReflector.getObjFieldValues(ObjectStreamClass.java:2050)
    at java.io.ObjectStreamClass.getObjFieldValues(ObjectStreamClass.java:1252)
    ... 45 more

“未能序列化任务 0”引起了我的注意。我对我正在做的事情没有一个突出的心理画面customRDD.count,而且非常不清楚什么不能被序列化。

我的自定义 RDD 包括:

  • 自定义 RDD 类
  • 自定义分区类
  • 自定义(scala)迭代器类

我的 Spark shell 会话如下所示:

import custom.rdd.stuff
import org.apache.spark.SparkContext

val conf = sc.getConf
conf.set(custom, parameters)
sc.stop
sc2 = new SparkContext(conf)
val mapOfThings: Map[String, String] = ...
myRdd = customRDD(sc2, mapOfStuff)
myRdd.count

... (exception output) ...

我想知道的是:

  • 为了创建自定义 RDD 类,哪些内容需要“可序列化”?
  • 就 Spark 而言,“可序列化”是什么意思?这类似于 Java 的“Serializable”吗?
  • 从我的 RDD 的迭代器返回的所有数据(由compute 方法返回)是否也需要可序列化?

非常感谢您对此问题的任何澄清。

【问题讨论】:

    标签: scala hadoop serialization apache-spark rdd


    【解决方案1】:

    在 Spark 上下文中执行的代码必须存在于工作节点的同一进程边界内,任务被指示在该工作节点上执行。这意味着必须注意确保 RDD 自定义中引用的任何对象或值都是可序列化的。如果对象是不可序列化的,那么您需要确保它们的范围正确,以便每个分区都有该对象的新实例。

    基本上,您不能共享在 Spark 驱动程序上声明的对象的不可序列化实例,并期望将其状态复制到集群上的其他节点。

    这是一个无法序列化不可序列化对象的示例:

    NotSerializable notSerializable = new NotSerializable();
    JavaRDD<String> rdd = sc.textFile("/tmp/myfile");
    
    rdd.map(s -> notSerializable.doSomething(s)).collect();
    

    下面的例子可以正常工作,因为它是在 lambda 的上下文中,它可以正确地分布到多个分区,而不需要序列化不可序列化对象的实例的状态。这也适用于作为 RDD 自定义(如果有)的一部分引用的不可序列化的传递依赖项。

    rdd.forEachPartition(iter -> {
      NotSerializable notSerializable = new NotSerializable();
    
      // ...Now process iter
    });
    

    更多详情请看这里:http://databricks.gitbooks.io/databricks-spark-knowledge-base/content/troubleshooting/javaionotserializableexception.html

    【讨论】:

      【解决方案2】:

      除了肯尼的解释,我建议你打开序列化调试,看看是什么导致了问题。通常,仅通过查看代码就无法弄清楚。

      -Dsun.io.serialization.extendedDebugInfo=true
      

      【讨论】:

      • 谢谢你。我和OP有同样的问题。通常,Spark 'improveException' 例程会打印出有问题的类,但在 OP 和我的情况下它会失败。将此选项添加到 $SPARK_CONF/java-opts 为我提供了前进所需的信息。
      • 请问什么是OP?
      • "OP" = "开幕海报",在这种情况下是@llovett。
      • 其实准确来说是原帖/原海报(来源:en.wikipedia.org/wiki/Internet_forum#Post
      • 谢谢:非常有用的标志添加!
      【解决方案3】:

      问题是您在 customRdd 方法 (customRDD(sc2, mapOfStuff)) 中传递了 SparkContex(Boiler plate)。确保您的班级还序列化了哪些制作 SparkContext。

      【讨论】:

        猜你喜欢
        • 2016-05-07
        • 1970-01-01
        • 2022-09-23
        • 2017-07-22
        • 1970-01-01
        • 1970-01-01
        • 2018-04-08
        • 1970-01-01
        • 2015-08-07
        相关资源
        最近更新 更多