【问题标题】:Track of consumed messages in Spark structured streaming跟踪 Spark 结构化流中的消费消息
【发布时间】: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 的消息时记录了什么样的进度?

【问题讨论】:

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


【解决方案1】:

您可以在writeStream 中将检查点位置设置为一个选项:

[...]
.writeStream
  .format("kafka")
  .option("kafka.bootstrap.servers", "host1:port1,host2:port2")
  .option("topic", "topic1")
  .option("checkpointLocation", "/path/to/dir")
  .trigger(Trigger.Continuous("1 second"))
  .start();

从 Kafka 读取时跟踪进度意味着跟踪 TopicPartition 中消耗的偏移量。设置检查点位置将使您的应用程序能够将该信息作为 JSON 对象存储在给定路径中,例如

{
  "topic1":{
    "0":11, 
    "1":101
  }
}

这意味着应用程序已经消耗了主题topic1的分区0中的偏移量10和分区1中的偏移量100。检查点是“提前”写入的(使用 write-ahead-logs),因此应用程序将继续从 Kafka 读取消息,在预期或意外(失败)重新启动之前中断。

Trigger.Continuous 从 Spark 2.3 版开始可用。并且现在标记为 experimental。与微批处理方法相比,它会在 Kafka 中的每条消息到达主题后立即获取它,而无需尝试将其与其他消息进行批处理。这可以改善延迟,但很可能会降低您的整体吞吐量。

参数(例如1 seconds)确定检查点的频率。

当使用这种触发模式时,至少要有与主题分区一样多的可用内核是很重要的。否则,申请将不会有任何进展。你可以阅读更多关于它here:

“例如,如果您正在读取具有 10 个分区的 Kafka 主题,那么集群必须至少有 10 个内核才能进行查询。”

【讨论】:

  • 所以你的意思是它会在失败后重新启动后从 Kafka 获取所有(已经使用 + 新的)消息?
  • 它将继续使用失败前中断的消息。有关已消费消息的信息会提前写入,并且仅在处理成功时才会提交。请记住,结构化流中的“提交”并不意味着提交回 Kafka,而是提交到检查点位置。
猜你喜欢
  • 2020-05-27
  • 2018-06-01
  • 1970-01-01
  • 2020-10-27
  • 1970-01-01
  • 2019-01-29
  • 2019-09-24
  • 2019-07-22
  • 2012-11-17
相关资源
最近更新 更多