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