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