【问题标题】:spark streaming join kafka topics火花流加入kafka主题
【发布时间】:2019-10-11 04:46:59
【问题描述】:

我们有两个来自两个Kafka主题的InputDStream,但是我们必须将这两个输入的数据连接在一起。 问题是每个InputDStream都是独立处理的,因为foreachRDD之后,什么都不能返回,到join之后。

  var Message1ListBuffer = new ListBuffer[Message1]
  var Message2ListBuffer = new ListBuffer[Message2]

    inputDStream1.foreachRDD(rdd => {
      if (!rdd.partitions.isEmpty) {
        val offsetRanges = rdd.asInstanceOf[HasOffsetRanges].offsetRanges
        rdd.map({ msg =>
          val r = msg.value()
          val avro = AvroUtils.objectToAvro(r.getSchema, r)
          val messageValue = AvroInputStream.json[FMessage1](avro.getBytes("UTF-8")).singleEntity.get
          Message1ListBuffer = Message1FlatMapper.flatmap(messageValue)
          Message1ListBuffer
        })
        inputDStream1.asInstanceOf[CanCommitOffsets].commitAsync(offsetRanges)
      }
    })


    inputDStream2.foreachRDD(rdd => {
      if (!rdd.partitions.isEmpty) {
        val offsetRanges = rdd.asInstanceOf[HasOffsetRanges].offsetRanges
        rdd.map({ msg =>
          val r = msg.value()
          val avro = AvroUtils.objectToAvro(r.getSchema, r)
          val messageValue = AvroInputStream.json[FMessage2](avro.getBytes("UTF-8")).singleEntity.get
          Message2ListBuffer = Message1FlatMapper.flatmap(messageValue)
          Message2ListBuffer

        })
        inputDStream2.asInstanceOf[CanCommitOffsets].commitAsync(offsetRanges)
      }
    })

我想我可以返回 Message1ListBuffer 和 Message2ListBuffer,将它们变成数据帧并加入它们。但这不起作用,我认为这不是最好的选择

从那里返回每个 foreachRDD 的 rdd 以进行连接的方法是什么?

inputDStream1.foreachRDD(rdd => {

})


inputDStream2.foreachRDD(rdd => {

})

【问题讨论】:

  • 什么是 Spark 版本?

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


【解决方案1】:

不确定你使用的Spark版本,用Spark 2.3+,可以直接实现。

使用 Spark >= 2.3

订阅2个你想加入的主题

val ds1 = spark
  .readStream 
  .format("kafka")
  .option("kafka.bootstrap.servers", "brokerhost1:port1,brokerhost2:port2")
  .option("subscribe", "source-topic1")
  .option("startingOffsets", "earliest")
  .option("endingOffsets", "latest")
  .load

val ds2 = spark
  .readStream 
  .format("kafka")
  .option("kafka.bootstrap.servers", "brokerhost1:port1,brokerhost2:port2")
  .option("subscribe", "source-topic2")
  .option("startingOffsets", "earliest")
  .option("endingOffsets", "latest")
  .load

格式化两个流中的订阅消息

val stream1 = ds1.selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)")
  .as[(String, String)]

val stream2 = ds2.selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)")
  .as[(String, String)]

加入两个流

resultStream = stream1.join(stream2)

更多join operations here

警告:

延迟记录不会得到连接匹配。需要稍微调整缓冲区。 more information found here

【讨论】:

  • 简短而清晰的答案:-)
猜你喜欢
  • 2015-10-13
  • 2017-04-27
  • 1970-01-01
  • 2020-06-21
  • 2017-02-12
  • 2020-04-11
  • 1970-01-01
  • 2023-03-18
  • 1970-01-01
相关资源
最近更新 更多