【发布时间】:2021-05-15 23:53:04
【问题描述】:
我想设置配置,让我的应用程序跟踪来自 kafka 的消费消息。这样每当它失败时,它就可以从最后一次提交或消耗的偏移量开始选择。
readStream
.format("kafka")
.option("kafka.bootstrap.servers", "host1:port1,host2:port2")
.option("subscribe", "topic1")
.load()
.selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)")
.writeStream
.format("kafka")
.option("kafka.bootstrap.servers", "host1:port1,host2:port2")
.option("topic", "topic1")
.trigger(Trigger.Continuous("1 second")) // only change in query
.start();
我在网上看到可以设置checkpointlocation 属性,spark 可以使用该属性来跟踪偏移量。
想知道我可以在哪里设置这个属性?我可以在option 中设置上面的代码吗?请问我该如何正确设置它。
其次,我无法理解trigger(Trigger.Continuous("1 second")) 属性。 Docs 说continuous processing engine will record the progress of the query every second,它在阅读来自 kafka 的消息时记录了什么样的进度?
【问题讨论】:
-
是的,checkpointLocation 已添加到选项中,但它需要 HDFS 兼容的文件系统,如 HDFS 或 S3。 spark.apache.org/docs/latest/…
标签: java apache-spark apache-kafka spark-structured-streaming