【问题标题】:How to get batch ID in Kafka output from Spark Structured Streaming如何从 Spark Structured Streaming 获取 Kafka 输出中的批处理 ID
【发布时间】:2020-01-18 18:15:12
【问题描述】:

我正在更新模式下运行我的 Spark 结构化流式处理作业,但无法确定是否可以获得每次更新的批次 ID。例如,当您以更新模式输出到控制台时,Spark 会在输出时显示每个批次号:

-------------------------------------------
Batch: 0
-------------------------------------------
...
-------------------------------------------
Batch: 1
-------------------------------------------
...

等等。 我需要将相同的信息添加到我发送给 Kafka 的每条消息中。为此,我仅限于使用 Spark 2.3,因此我无法使用 forEachBatch。

我的工作输出一组特定维度的聚合指标。每个触发器,自上次触发器以来,指标可能已更新 - 具有更新指标的维度将在下一批中输出,因为我正在更新模式下运行。当我将这些输出到 Kafka 时,我需要知道哪个批次是最新的——因此需要批次号。我认为 forEachBatch 可以得到我需要的东西,但不幸的是我无法访问 Spark 2.4。我可以使用 forEach 来完成这个吗?我仅限于使用更新模式,因为可能会出现延迟事件并更新之前已经输出的指标。

这是我用来测试的控制台模式。此输出分别显示每个批次,以及它的编号:

StreamingQuery query = logs.writeStream()
        .format("console")
        .outputMode(OutputMode.Update())
        .start();

我想做这样的事情

StreamingQuery query = agg.WriteStream()
    .format("kafka")
    .outputMode(OutputMode.Update())
    .option("kafka.bootstrap.servers", "myconnection")
    .Option("topic", "mytopic")
    .Start();

但仍保留在 mytopic 中判断消息来自哪个批次的能力。这可能吗?

【问题讨论】:

    标签: java apache-spark apache-kafka spark-structured-streaming


    【解决方案1】:

    我认为你可以使用来自ForeachWriter的版本号long version

    您可以像这样实现自己的 KafkaCustomSink。

    
    class KafkaCustomSink(val config: Config) extends ForeachWriter[String] {
      var producer: KafkaProducer[String, String] = _
      var _version: Long = _
    
      override def open(partitionId: Long, version: Long): Boolean = {
        _version = version
        val props = new Properties()
        props.put("bootstrap.servers", config(Constant.OUTPUT_BOOTSTRAP_SERVER))
        props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer")
        props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer")
        props.put("acks", "0")
        producer = new KafkaProducer[String, String](props)
        true
      }
    
      override def process(value: String): Unit = {
        //use version here
        val record = new ProducerRecord[String, String](config(Constant.OUTPUT_TOPIC), null, "version : %s, data : %s".format(_version, value))
        producer.send(record)
      }
    
      override def close(errorOrNull: Throwable): Unit = {
        producer.close()
      }
    }
    
    

    并将其分配给

          logs
              .writeStream
              .outputMode("update")
              .foreach(new KafkaCustomSink(config))
              .trigger(Trigger.ProcessingTime(config(Constant.TRIGGER_INTERVAL).toInt, TimeUnit.SECONDS))
              .option("checkpointLocation", config(Constant.CHECKPOINT_LOCATION))
    
    

    【讨论】:

      猜你喜欢
      • 2021-05-22
      • 2020-01-05
      • 2023-03-31
      • 2020-09-03
      • 2019-06-25
      • 1970-01-01
      • 2021-12-05
      • 2019-02-19
      • 2018-02-28
      相关资源
      最近更新 更多