【问题标题】:How to use foreachPartition in Spark 2.2 to avoid Task Serialization error如何在 Spark 2.2 中使用 foreachPartition 来避免任务序列化错误
【发布时间】:2018-05-29 13:57:18
【问题描述】:

我有以下使用结构化流式处理 (Spark 2.2) 的工作代码,以便从 Kafka (0.10) 读取数据。 在ForeachWriter 中使用kafkaProducer 时,我无法解决的唯一问题与Task serialization problem 有关。 在为 Spark 1.6 开发的此代码的旧版本中,我使用 foreachPartition 并为每个分区定义 kafkaProducer 以避免任务序列化问题。 如何在 Spark 2.2 中做到这一点?

val df: Dataset[String] = spark.readStream
      .format("kafka")
      .option("kafka.bootstrap.servers", "localhost:9092")
      .option("subscribe", "test") 
      .option("startingOffsets", "latest")
      .option("failOnDataLoss", "true")
      .load()
      .selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)").as[(String, String)] 
      .map(_._2)

var mySet = spark.sparkContext.broadcast(Map(
  "metadataBrokerList"->metadataBrokerList,
  "outputKafkaTopic"->outputKafkaTopic,
  "batchSize"->batchSize,
  "lingerMS"->lingerMS))

val kafkaProducer = Utils.createProducer(mySet.value("metadataBrokerList"),
                                mySet.value("batchSize"),
                                mySet.value("lingerMS"))

val writer = new ForeachWriter[String] {

    override def process(row: String): Unit = {
         // val result = ...
         val record = new ProducerRecord[String, String](mySet.value("outputKafkaTopic"), "1", result);
        kafkaProducer.send(record)
    }

    override def close(errorOrNull: Throwable): Unit = {}

    override def open(partitionId: Long, version: Long): Boolean = {
      true
    }
}

val query = df
        .writeStream
        .foreach(writer)
        .start

query.awaitTermination()

spark.stop()

【问题讨论】:

    标签: scala apache-spark apache-kafka spark-dataframe spark-streaming


    【解决方案1】:

    编写 ForeachWriter 的实现,然后使用它。 (避免使用不可序列化对象的匿名类 - 在您的情况下为 ProducerRecord)
    示例:val writer = new YourForeachWriter[String]
    这里还有一篇关于 Spark 序列化问题的有用文章:https://www.cakesolutions.net/teamblogs/demystifying-spark-serialisation-error

    【讨论】:

    • 应该 YourForeachWriter 扩展 ForeachWriter 并且也可以序列化吗?
    • 不,如果父类具有可序列化的行为,那么默认情况下子类也将具有可序列化的行为。您有例外,因为您在匿名类中使用了 ProducerRecord(不可序列化)并且 Spark 驱动程序无法对其进行序列化并传输到执行程序。如果您将使用自己的类的实例,则不会出现此异常。
    • 抱歉,我不明白如何遵循您的建议。我应该创建一个case class YourForeachWriter(s:String) extends ForeachWriter 然后将process 的内容移动到这个案例类吗?如果是这样,我如何将sparklogger 和其他变量传递给这个案例类?
    • 好的,但是我应该如何实现YourForeachWriter?它应该包含什么?
    • 我理解问题的原因,我理解你的方法。我怀念的是YourForechWriterClass的内容。我真的不明白我应该在那里实施什么。
    猜你喜欢
    • 2019-02-13
    • 1970-01-01
    • 2020-07-25
    • 2017-09-21
    • 2017-08-24
    • 1970-01-01
    • 1970-01-01
    • 2020-02-04
    • 1970-01-01
    相关资源
    最近更新 更多