【问题标题】:how to get count of unread messages in a Kafka topic for a certain group如何获取某个组的 Kafka 主题中未读消息的计数
【发布时间】:2021-12-12 15:57:48
【问题描述】:

我知道,Kafka 旨在像处理无限的事件流一样处理,并且获取剩余消息计数不是内置功能。但我必须以某种方式监控我的消费者进程的运行情况以及我是否为他们提供了足够的资源。 总体场景是基本的 Kafka 使用,不同服务器上的几个生产者插入一个主题,group_a 中的消费者读取消息,执行一些 AI 工作并插入另一个主题以进行进一步处理。 传入消息的速率绝不是恒定的或可预测的,因此我需要检查我的消费者是否落后(假设group_a 在我的输入中有超过 1000 条未读消息主题)。

考虑到我可以完全控制 Kafka 设置以及消费者和生产者代码,我有哪些选择?

  • 我认为在没有任何繁重处理和计数消息的情况下阅读整个主题既不干净也不高效 (?)
  • 如果只有一对生产者/消费者,我可以计算生产和消费消息的数量并计算剩余的消息,但我正在运行多服务器设置,这不太可行。是否可以将 Kafka 本身用作共享数据存储并记录所有生产和消费的消息?
  • 为了避免 XY 问题,计算未读消息数量的正确方法是了解我的消费者是否获得了“足够”的资源?

【问题讨论】:

    标签: python apache-kafka


    【解决方案1】:

    Kafka 维护每个消费者使用的每个分区的消息偏移量的元数据。分区使您可以为单个主题拥有多个消费者。

    Lag 是消费者组偏移量与主题中最新偏移量之间的增量。

    这些元数据由 Kafka 维护,您无需在应用中维护它们。

    您可以使用例如CLI tool 来检查消费者组偏移和滞后。

    在下面的示例中,foo 组已经消费了来自主题“quickstart-events”的所有消息,因此延迟为 0。在此示例中,主题中只有一个分区和一个消费者。

    bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 --describe --group foo                                                        
    
    GROUP           TOPIC             PARTITION  CURRENT-OFFSET  LOG-END-OFFSET  LAG             CONSUMER-ID                                         HOST            CLIENT-ID
    foo             quickstart-events 0          3               3               0               consumer-foo-1-0fbf8d40-732f-483d-91d3-9b6f686b5040 /127.0.0.1      consumer-foo-1
    

    同样的信息也可以通过AdminClient API 获得,请参阅 describe_consumer_groups。

    如果您需要自动对延迟做出反应,最好为延迟设置单独的监控流程,而不是尝试检测消费者流程中的延迟。这种方法的一种工具是Kafka Lag Exporter。这种方法的好处是您可以使用通用监控工具来制作警报和仪表板,但当然需要做一些工作来设置所需的基础架构。

    为每个主题拥有足够数量的分区很重要,因为最多。一个主题的并发消费者数量由分区数量决定。

    【讨论】:

    • 这正是我想要的。它甚至有一个 Kafka-python 实现 :) 谢谢。是的,我正在使用自己的仪表板在一个单独的进程中进行监控。
    猜你喜欢
    • 2019-11-25
    • 1970-01-01
    • 2018-01-19
    • 2018-07-27
    • 2015-08-01
    • 2021-08-07
    • 1970-01-01
    • 2015-04-19
    • 1970-01-01
    相关资源
    最近更新 更多