【发布时间】: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