【问题标题】:Read json from Kafka and write json to other Kafka topic从 Kafka 读取 json 并将 json 写入其他 Kafka 主题
【发布时间】:2018-05-07 13:15:58
【问题描述】:

我正在尝试为 Spark 流式传输准备应用程序(Spark 2.1、Kafka 0.10)

我需要从 Kafka 主题“输入”读取数据,找到正确的数据并将结果写入主题“输出”

我可以基于 KafkaUtils.createDirectStream 方法从 Kafka 读取数据。

我将 RDD 转换为 json 并准备过滤器:

val messages = KafkaUtils.createDirectStream[String, String](
  ssc,
  PreferConsistent,
  Subscribe[String, String](topics, kafkaParams)
)

val elementDstream = messages.map(v => v.value).foreachRDD { rdd =>

  val PeopleDf=spark.read.schema(schema1).json(rdd)
  import spark.implicits._
  PeopleDf.show()
  val PeopleDfFilter = PeopleDf.filter(($"value1".rlike("1"))||($"value2" === 2))
  PeopleDfFilter.show()
}

我可以使用 KafkaProducer 从 Kafka 加载数据并“按原样”写入 Kafka:

    messages.foreachRDD( rdd => {
      rdd.foreachPartition( partition => {
        val kafkaTopic = "output"
        val props = new HashMap[String, Object]()
        props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092")
        props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG,
          "org.apache.kafka.common.serialization.StringSerializer")
        props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG,
          "org.apache.kafka.common.serialization.StringSerializer")

        val producer = new KafkaProducer[String, String](props)
        partition.foreach{ record: ConsumerRecord[String, String] => {
        System.out.print("########################" + record.value())
        val messageResult = new ProducerRecord[String, String](kafkaTopic, record.value())
        producer.send(messageResult)
        }}
        producer.close()
      })

    })

但是,我无法整合这两个操作 > 在 json 中找到正确的值并将结果写入 Kafka:以 JSON 格式写入 PeopleDfFilter 以“输出”Kafka 主题。

我在 Kafka 中有很多输入消息,这就是我想使用 foreachPartition 创建 Kafka 生产者的原因。

【问题讨论】:

  • 你需要将record.value()实际解析为JSON

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


【解决方案1】:

这个过程非常简单,为什么不一直使用结构化流呢?

import org.apache.spark.sql.functions.from_json

spark
  // Read the data
  .readStream
  .format("kafka")
  .option("kafka.bootstrap.servers", inservers) 
  .option("subscribe", intopic)
  .load()
  // Transform / filter
  .select(from_json($"value".cast("string"), schema).alias("value"))
  .filter(...)  // Add the condition
  .select(to_json($"value").alias("value")
  // Write back
  .writeStream
  .format("kafka")
  .option("kafka.bootstrap.servers", outservers)
  .option("subscribe", outtopic)
  .start()

【讨论】:

  • 当我尝试在 Spark 2.1 上使用带有 writeStream 到 Kafka 的结构化流时,我得到:java.lang.UnsupportedOperationException: Data source kafka does not support streamed writing 据我所知,它适用于 Spark 2.2.x
  • 这个答案不仅告诉你如何将数据写入kafka,还告诉你如何使用spark从kafka重新读取数据。
【解决方案2】:

尝试为此使用结构化流。即使您使用的是 Spark 2.1,您也可以按如下方式实现自己的 Kafka ForeachWriter:

Kafka sink:

import java.util.Properties
import kafkashaded.org.apache.kafka.clients.producer._
import org.apache.spark.sql.ForeachWriter


 class  KafkaSink(topic:String, servers:String) extends ForeachWriter[(String, String)] {
      val kafkaProperties = new Properties()
      kafkaProperties.put("bootstrap.servers", servers)
      kafkaProperties.put("key.serializer",
        classOf[org.apache.kafka.common.serialization.StringSerializer].toString)
      kafkaProperties.put("value.serializer",
        classOf[org.apache.kafka.common.serialization.StringSerializer].toString)
      val results = new scala.collection.mutable.HashMap[String, String]
      var producer: KafkaProducer[String, String] = _

      def open(partitionId: Long,version: Long): Boolean = {
        producer = new KafkaProducer(kafkaProperties)
        true
      }

      def process(value: (String, String)): Unit = {
          producer.send(new ProducerRecord(topic, value._1 + ":" + value._2))
      }

      def close(errorOrNull: Throwable): Unit = {
        producer.close()
      }
   }

用法:

val topic = "<topic2>"
val brokers = "<server:ip>"

val writer = new KafkaSink(topic, brokers)

val query =
  streamingSelectDF
    .writeStream
    .foreach(writer)
    .outputMode("update")
    .trigger(ProcessingTime("25 seconds"))
    .start()

【讨论】:

  • 虽然此链接可能会回答问题,但最好在此处包含答案的基本部分并提供链接以供参考。如果链接页面发生更改,仅链接答案可能会失效。 - From Review
猜你喜欢
  • 1970-01-01
  • 2021-09-15
  • 2020-07-02
  • 2020-08-19
  • 2022-01-03
  • 1970-01-01
  • 2019-04-23
  • 2021-06-11
  • 2018-10-24
相关资源
最近更新 更多