【问题标题】:Spark: broadcasting jackson ObjectMapperSpark:广播杰克逊 ObjectMapper
【发布时间】:2015-09-10 07:38:38
【问题描述】:

我有一个 spark 应用程序,它从文件中读取行并尝试使用 jackson 对它们进行反序列化。 为了让这段代码正常工作,我需要在 Map 操作中定义 ObjectMapper(否则我会得到 NullPointerException)。

我有以下正在运行的代码:

val alertsData = sc.textFile(rawlines).map(alertStr => {
      val mapper = new ObjectMapper()
      mapper.configure(DeserializationFeature.FAIL_ON_UNKNOWN_PROPERTIES, false)
      mapper.registerModule(DefaultScalaModule)
      broadcastVar.value.readValue(alertStr, classOf[Alert])
    })

但是,如果我在地图之外定义映射器并广播它,它会失败并出现 NullPointerException。

此代码失败:

val mapper = new ObjectMapper()
    mapper.configure(DeserializationFeature.FAIL_ON_UNKNOWN_PROPERTIES, false)
    mapper.registerModule(DefaultScalaModule)
    val broadcastVar = sc.broadcast(mapper)

    val alertsData = sc.textFile(rawlines).map(alertStr => {
      broadcastVar.value.readValue(alertStr, classOf[Alert])
    })

我在这里错过了什么?

谢谢, 艾丽莎

【问题讨论】:

  • 我认为广播是面向不可变数据的,可能你的对象被广播为浅拷贝,在ObjectMapper中没有使用一些状态或依赖。

标签: scala apache-spark jackson


【解决方案1】:

事实证明你可以广播映射器。有问题的部分是mapper.registerModule(DefaultScalaModule),它需要在每个从机(执行器)上执行,而不仅仅是在驱动程序上。

所以这段代码可以工作:

val mapper = new ObjectMapper()
mapper.configure(DeserializationFeature.FAIL_ON_UNKNOWN_PROPERTIES, false)
val broadcastVar = sc.broadcast(mapper)

val alertsData = sc.textFile(rawlines).map(alertStr => {
      broadcastVar.value.registerModule(DefaultScalaModule)
      broadcastVar.value.readValue(alertStr, classOf[Alert])
})

我进一步优化了代码,每个分区只运行一次 registerModule(而不是为 RDD 中的每个元素)。

val mapper = new ObjectMapper()
mapper.configure(DeserializationFeature.FAIL_ON_UNKNOWN_PROPERTIES, false)

val broadcastVar = sc.broadcast(mapper)
val alertsRawData = sc.textFile(rawlines)

val alertsData = alertsRawData.mapPartitions({ iter: Iterator[String] => broadcastVar.value.registerModule(DefaultScalaModule)
      for (i <- iter) yield broadcastVar.value.readValue(i, classOf[Alert]) })

艾丽莎

【讨论】:

  • 有没有办法使用 spark sql UDF 在每个分区的基础上执行此操作?
  • 附加花絮:如果您尝试在 main 对象中广播 val 并且您的 main 方法已通过使用 ... extends App 继承,则广播失败。我在这里找到了解决方案:stackoverflow.com/a/31323435/1628839
【解决方案2】:

确实,objectMapper 并不适合广播。它本质上是不可序列化的,也不是值类。我建议改为广播DeserializationConfig,并将其从地图操作中的广播变量传递给 ObjectMapper 的构造函数。

【讨论】:

  • FWIW,ObjectMapper 实际上是java.io.Serializable,应该可以工作,但你是对的,这很少是你应该做的(在有限的情况下你可能想要这样做,这就是为什么它是可序列化的,但值得写一篇博文)。特别是对于 Spark(或 Storm/Trident 等),它属于我认为应该在本地实例化的一类东西。
猜你喜欢
  • 2019-04-16
  • 2021-01-09
  • 1970-01-01
  • 2021-01-23
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多