【问题标题】:Kafka get to know when related messages are consumedKafka 了解相关消息何时被消费
【发布时间】:2020-09-11 17:52:27
【问题描述】:

在 Kafka 中有什么方法可以在消费了几条相关消息后生成一条消息? (无需在应用程序代码处手动控制...)

用例是选择一个大文件,将其拆分为几个块,为主题中的每个块发布一条消息,一旦所有这些消息被消耗,就会产生另一条消息,通知另一个主题的结果。

我们可以使用数据库或 REDIS 来控制状态,但我想知道是否有任何更高级别的方法仅利用 Kafka 生态系统。

【问题讨论】:

  • 您的消费者在“一旦所有这些消息都被消费”中的样子是怎样的?它也是 Kafka Streams 应用程序还是其他什么?
  • 最初它是一个 Spring-Boot Kotlin 应用程序消费者,但我们愿意接受选择......

标签: apache-kafka apache-kafka-streams batching


【解决方案1】:

方法可以如下:

  1. 在消费每个区块后,应用程序应生成带有状态(已消费和区块编号)的消息
  2. 第二个应用程序(Kafka Streams 一次)应该聚合结果,当处理带有所有块的消息产生最终消息时,处理该文件。

【讨论】:

  • 这对我来说确实有意义并且听起来很有希望,但是我们如何在一次 kafka 流上知道所有块都已处理(请原谅我的无知,从未真正使用过流)。您有任何文档或 sn-p 指向吗?
  • 例如:带有块状态的消息可以如下:` (key: fileUniqueName, value: chunkNumber, numberOfChunks)`。在 Kafka 流应用程序中,您可以使用 ProcessorApi (kafka.apache.org/10/documentation/streams/developer-guide/…) 并以自定义方式聚合它 - 使用状态存储您可以保持有关已处理块数的状态
【解决方案2】:

您可以使用ConsumerGroupCommand 来检查某个消费者组是否已处理完特定主题中的所有消息:

  1. $ kafka-consumer-groups --bootstrap-server broker_host:port --describe --group chunk_consumer

  1. $ kafka-run-class kafka.admin.ConsumerGroupCommand ...

每个分区的零延迟表示消息已被成功消费,并且消费者已提交偏移量。

或者,您可以选择订阅__consumer_offsets 主题并自己处理来自该主题的消息,但使用ConsumerGroupCommand 似乎是一种更直接的解决方案。

【讨论】:

  • 据我了解,消费者组将绑定到特定的应用程序,而不是为每个文件动态创建。我怀疑另一个答案,使用 kafka 流对特定用例更有意义。
  • 不确定我是否遵循 - 消费者提交的偏移量是确认消息已成功使用。因此,如果从生产者方面,您监控偏移量并确保提交所有偏移量,您就会知道所有“块”都已被消耗。一旦发生这种情况,您就可以发布确认信息或做任何其他您需要做的事情。
猜你喜欢
  • 2020-08-11
  • 2016-06-12
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2017-12-21
  • 2017-09-23
  • 1970-01-01
  • 2022-06-15
相关资源
最近更新 更多