【问题标题】:java.io.NotSerializableException with Spark Streaming Checkpoint enabled启用了 Spark 流检查点的 java.io.NotSerializableException
【发布时间】:2018-04-11 01:11:51
【问题描述】:

我已在我的 spark 流应用程序中启用检查点,但在作为依赖项下载的类上遇到此错误。

没有检查点,应用程序运行良好。

错误:

com.fasterxml.jackson.module.paranamer.shaded.CachingParanamer
Serialization stack:
    - object not serializable (class: com.fasterxml.jackson.module.paranamer.shaded.CachingParanamer, value: com.fasterxml.jackson.module.paranamer.shaded.CachingParanamer@46c7c593)
    - field (class: com.fasterxml.jackson.module.paranamer.ParanamerAnnotationIntrospector, name: _paranamer, type: interface com.fasterxml.jackson.module.paranamer.shaded.Paranamer)
    - object (class com.fasterxml.jackson.module.paranamer.ParanamerAnnotationIntrospector, com.fasterxml.jackson.module.paranamer.ParanamerAnnotationIntrospector@39d62e47)
    - field (class: com.fasterxml.jackson.databind.introspect.AnnotationIntrospectorPair, name: _secondary, type: class com.fasterxml.jackson.databind.AnnotationIntrospector)
    - object (class com.fasterxml.jackson.databind.introspect.AnnotationIntrospectorPair, com.fasterxml.jackson.databind.introspect.AnnotationIntrospectorPair@7a925ac4)
    - field (class: com.fasterxml.jackson.databind.introspect.AnnotationIntrospectorPair, name: _primary, type: class com.fasterxml.jackson.databind.AnnotationIntrospector)
    - object (class com.fasterxml.jackson.databind.introspect.AnnotationIntrospectorPair, com.fasterxml.jackson.databind.introspect.AnnotationIntrospectorPair@203b98cf)
    - field (class: com.fasterxml.jackson.databind.cfg.BaseSettings, name: _annotationIntrospector, type: class com.fasterxml.jackson.databind.AnnotationIntrospector)
    - object (class com.fasterxml.jackson.databind.cfg.BaseSettings, com.fasterxml.jackson.databind.cfg.BaseSettings@78c34153)
    - field (class: com.fasterxml.jackson.databind.cfg.MapperConfig, name: _base, type: class com.fasterxml.jackson.databind.cfg.BaseSettings)
    - object (class com.fasterxml.jackson.databind.DeserializationConfig, com.fasterxml.jackson.databind.DeserializationConfig@2df0a4c3)
    - field (class: com.fasterxml.jackson.databind.ObjectMapper, name: _deserializationConfig, type: class com.fasterxml.jackson.databind.DeserializationConfig)
    - object (class com.fasterxml.jackson.databind.ObjectMapper, com.fasterxml.jackson.databind.ObjectMapper@2db07651)

我不确定如何将此类扩展为可序列化的 maven 依赖项。我在 pom.xml 中使用杰克逊核心的 v2.6.0。如果我尝试使用较新版本的 Jackson 核心,则会收到不兼容的 Jackson 版本异常。

代码

liveRecordStream
      .foreachRDD(newRDD => {
        if (!newRDD.isEmpty()) {
          val cacheRDD = newRDD.cache()
          val updTempTables = tempTableView(t2s, stgDFMap, cacheRDD)
          val rdd = updatestgDFMap(stgDFMap, cacheRDD)
          persistStgTable(stgDFMap)
          dfMap
            .filter(entry => updTempTables.contains(entry._2))
            .map(spark.sql)
            .foreach( df => writeToES(writer, df))

          cacheRDD.unpersist()
        }
      }

只有在 foreachRDD 内部发生方法调用时才会出现此问题,例如 tempTableView 在这种情况下。

临时表视图

def tempTableView(t2s: Map[String, StructType], stgDFMap: Map[String, DataFrame], cacheRDD: RDD[cacheRDD]): Set[String] = {
    stgDFMap.keys.filter { table =>
      val tRDD = cacheRDD
        .filter(r => r.Name == table)
        .map(r => r.values)
         val tDF = spark.createDataFrame(tRDD, tableNameToSchema(table))
      if (!tRDD.isEmpty()) {
        val tName = s"temp_$table"
        tDF.createOrReplaceTempView(tName)
      }
      !tRDD.isEmpty()
    }.toSet
  }

感谢任何帮助。不知道如何调试和解决问题。

【问题讨论】:

  • 您可以尝试将 "tempTableView(t2s, stgDFMap, cacheRDD)" 放入像 map 或 filter 这样的转换中吗?因为转换是在工作人员中执行的,不需要通过序列化从驱动程序转移到工作人员。

标签: scala apache-spark spark-streaming


【解决方案1】:

从您共享的代码 sn-p 中,我看不到在哪里调用了 jackson 库。但是,NotSerializableException 通常发生在您尝试通过线路发送未实现Serializable 接口的对象时。

Spark 是分布式处理引擎,这意味着它是这样工作的:跨节点有一个驱动程序和多个执行程序。只有需要计算的部分代码由driver 发送到executors(通过线路)。 Spark 转换以这种方式发生,即跨多个节点,如果您尝试将未实现 serializable 接口的类的实例传递给此类代码块(跨节点执行的块),它将抛出 NotSerializableException .

例如:

def main(args: Array[String]): Unit = {
   val gson: Gson = new Gson()

   val sparkConf = new SparkConf().setMaster("local[2]")
   val spark = SparkSession.builder().config(sparkConf).getOrCreate()
   val rdd = spark.sparkContext.parallelize(Seq("0","1"))

   val something = rdd.map(str => {
     gson.toJson(str)
   })

   something.foreach(println)
   spark.close()
}

此代码块将抛出NotSerializableException,因为我们将Gson 的实例发送到分布式函数。 map 是一个 Spark 转换操作,因此它将在执行程序上执行。以下将起作用:

def main(args: Array[String]): Unit = {

   val sparkConf = new SparkConf().setMaster("local[2]")
   val spark = SparkSession.builder().config(sparkConf).getOrCreate()
   val rdd = spark.sparkContext.parallelize(Seq("0","1"))

   val something = rdd.map(str => {
     val gson: Gson = new Gson()
     gson.toJson(str)
   })

   something.foreach(println)
   spark.close()
}

上述工作的原因是,我们在转换中实例化Gson,因此它将在执行程序中实例化,这意味着它不会从驱动程序通过线路发送,因此不需要序列化。

【讨论】:

  • 感谢您对其功能的清晰解释。为什么只有在启用检查点时才会出现这种行为?
  • 很好的解释。这是我发现的最清晰和合成的。
【解决方案2】:

问题在于尝试序列化的 jackson objectMapperobjectMapper 不应该被序列化。通过添加 @transient val objMapper = new ObjectMapper...

解决了这个问题

【讨论】:

    猜你喜欢
    • 2023-02-09
    • 2016-04-05
    • 1970-01-01
    • 1970-01-01
    • 2018-06-22
    • 1970-01-01
    • 2021-01-07
    • 2023-03-10
    • 1970-01-01
    相关资源
    最近更新 更多