【问题标题】:ClassCastException while deserializing with Java's native readObject from Spark driver使用来自 Spark 驱动程序的 Java 的本机 readObject 反序列化时发生 ClassCastException
【发布时间】:2018-07-22 13:31:19
【问题描述】:

我有两个 spark 作业 A 和 B,因此 A 必须在 B 之前运行。A 的输出必须可以从以下位置读取:

  • Spark 作业 B
  • Spark 环境之外的独立 Scala 程序(不依赖于 Spark)

我目前正在将 Java 的本机序列化与 Scala 案例类一起使用。

来自 A Spark 工作:

val model = ALSFactorizerModel(...)

context.writeSerializable(resultOutputPath, model)

带序列化方式:

def writeSerializable[T <: Serializable](path: String, obj: T): Unit = {
  val writer: OutputStream = ... // Google Cloud Storage dependant
  val oos: ObjectOutputStream = new ObjectOutputStream(writer)
  oos.writeObject(obj)
  oos.close()
  writer.close()
}

来自 B Spark 作业或任何独立的非 Spark Scala 代码:

val lastFactorizerModel: ALSFactorizerModel = context
                     .readSerializable[ALSFactorizerModel](ALSFactorizer.resultOutputPath)

带反序列化方法:

def readSerializable[T <: Serializable](path: String): T = {
  val is : InputStream = ... // Google Cloud Storage dependant
  val ois = new ObjectInputStream(is)
  val model: T = ois
    .readObject()
    .asInstanceOf[T]
  ois.close()
  is.close()

  model
}

(嵌套的)案例类:

ALSFactorizerModel:

package mycompany.algo.als.common.io.model.factorizer

import mycompany.data.item.ItemStore

@SerialVersionUID(1L)
final case class ALSFactorizerModel(
  knownItems:       Array[ALSFeaturedKnownItem],
  unknownItems:     Array[ALSFeaturedUnknownItem],
  rank:             Int,
  modelTS:          Long,
  itemRepositoryTS: Long,
  stores:           Seq[ItemStore]
) {   
}

物品商店:

package mycompany.data.item

@SerialVersionUID(1L)
final case class ItemStore(
  id:     String,
  tenant: String,
  name:   String,
  index:  Int
) {
}

输出:

  • 来自独立的非 Spark Scala 程序 => 好的
  • 来自在我的开发机器上本地运行的 B Spark 作业(Spark 独立本地节点)=> 确定
  • 从在 (Dataproc) Spark 集群上运行的 B Spark 作业 => 失败,出现以下异常:

例外:

java.lang.ClassCastException: cannot assign instance of scala.collection.immutable.List$SerializationProxy to field mycompany.algo.als.common.io.model.factorizer.ALSFactorizerModel.stores of type scala.collection.Seq in instance of mycompany.algo.als.common.io.model.factorizer.ALSFactorizerModel
  at java.io.ObjectStreamClass$FieldReflector.setObjFieldValues(ObjectStreamClass.java:2133)
  at java.io.ObjectStreamClass.setObjFieldValues(ObjectStreamClass.java:1305)
  at java.io.ObjectInputStream.defaultReadFields(ObjectInputStream.java:2251)
  at java.io.ObjectInputStream.readSerialData(ObjectInputStream.java:2169)
  at java.io.ObjectInputStream.readOrdinaryObject(ObjectInputStream.java:2027)
  at java.io.ObjectInputStream.readObject0(ObjectInputStream.java:1535)
  at java.io.ObjectInputStream.readObject(ObjectInputStream.java:422)
  at mycompany.fs.gcs.SimpleGCSFileSystem.readSerializable(SimpleGCSFileSystem.scala:71)
  at mycompany.algo.als.batch.strategy.ALSClusterer$.run(ALSClusterer.scala:38)
  at mycompany.batch.SinglePredictorEbapBatch$$anonfun$3.apply(SinglePredictorEbapBatch.scala:55)
  at mycompany.batch.SinglePredictorEbapBatch$$anonfun$3.apply(SinglePredictorEbapBatch.scala:55)
  at scala.concurrent.impl.Future$PromiseCompletingRunnable.liftedTree1$1(Future.scala:24)
  at scala.concurrent.impl.Future$PromiseCompletingRunnable.run(Future.scala:24)
  at scala.concurrent.impl.ExecutionContextImpl$AdaptedForkJoinTask.exec(ExecutionContextImpl.scala:121)
  at scala.concurrent.forkjoin.ForkJoinTask.doExec(ForkJoinTask.java:260)
  at scala.concurrent.forkjoin.ForkJoinPool$WorkQueue.runTask(ForkJoinPool.java:1339)
  at scala.concurrent.forkjoin.ForkJoinPool.runWorker(ForkJoinPool.java:1979)
  at scala.concurrent.forkjoin.ForkJoinWorkerThread.run(ForkJoinWorkerThread.java:107)

我错过了什么吗?我是否应该配置 Dataproc/Spark 以支持对此代码使用 Java 序列化?

我使用--jars &lt;path to my fatjar&gt; 提交作业,之前从未遇到过其他问题。这个Jar中不包含spark依赖,作用域是Provided

Scala 版本: 2.11.8 Spark 版本: 2.0.2 SBT 版本: 0.13.13

感谢您的帮助

【问题讨论】:

    标签: java scala apache-spark serialization google-cloud-dataproc


    【解决方案1】:

    stores: Seq[ItemStore] 替换为stores: Array[ItemStore] 为我们解决了问题。

    或者,我们可以使用另一个类加载器来进行序列化/反序列化操作。

    希望这会有所帮助。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2015-08-18
      • 1970-01-01
      • 2021-09-30
      • 2010-11-13
      • 2019-01-29
      • 2011-10-05
      • 1970-01-01
      • 2011-08-20
      相关资源
      最近更新 更多