【问题标题】:Tracking topic size and consumer lag with kafka 0.8.2.0使用 kafka 0.8.2.0 跟踪主题大小和消费者滞后
【发布时间】:2020-05-14 05:08:49
【问题描述】:

自 kafka 0.8.2.0 以来,跟踪消费者滞后和主题大小似乎变得非常困难

您如何在 kafka 中跟踪偏移量(主题大小)和滞后?当您的生产者插入一条消息时,您是否在某处增加一个计数器并在您的消费者确认一条消息时增加另一个计数器?

我正在使用airbnb's kafka-statsd-metrics2 - 但出于某种原因,关于主题大小的所有指标总是0,这可能是他们的错误报告,但你是如何做到的?

我们的消费者和生产者是使用kafka-python 用 python 编写的,他们声明他们不支持 ConsumerCoordinator 偏移 API,所以我整理了一个解决方案来查询 zookeeper 并将这些指标发送到 statsd 实例(看起来很尴尬) ,但我仍然缺少主题大小指标。

我们正在使用 collectd 收集系统指标,我没有使用 JMX 的经验,并且在 collectd 中配置它似乎很复杂,我尝试了几次,所以我找到了一些不这样做的方法。

如果您有任何意见,我很乐意听到,即使是:“这属于 x stackexchange-site”

【问题讨论】:

    标签: python apache-kafka


    【解决方案1】:

    如果我理解正确,您可以使用FetchResponse 中的HighwaterMarkOffset。这样,您将知道分区末尾的偏移量是多少,并且能够将其与您当前确认的偏移量或此FetchResponse 中最后一条消息的偏移量进行比较。

    详情here

    【讨论】:

    • 当你第一次来到这里时,我将你的答案标记为我的答案,它也是正确的。我最终创建了一个程序来跟踪我们的消费者从 Zookeeper 的偏移量,因为那是 kafka-python 仍在使用的。这个解决方案的好处是可以很容易地将结果通过管道传输到我们的石墨设置中:) 非常感谢您的帮助!
    【解决方案2】:

    您是否尝试过使用https://github.com/quantifind/KafkaOffsetMonitor 来监控消费者延迟。它适用于 0.8.2.0

    【讨论】:

    • 这很奇怪。我从 0.8.2.0 之前就已经安装了它,但是在我升级后它从来没有工作过,但我只是想怎么回事并再次尝试,现在它可以很好地工作而无需任何更新......
    • KafkaOffsetMonitor 似乎不适用于简单的消费者,因为它试图从似乎只由高级消费者设置的 Zookeeper 检索分区所有者。我认为... ymmv
    【解决方案3】:

    这里是代码 sn-p,请确保在活动控制器中运行它。 BOOTSTRAP_SERVERS 是活动控制器 IP。

    client = KafkaAdminClient(bootstrap_servers=BOOTSTRAP_SERVERS, request_timeout_ms=300)
          list_groups_request  = client.list_consumer_groups()
    
          for group in list_groups_request:
            if group[1] == 'consumer':
              list_mebers_in_groups = client.describe_consumer_groups([group[0]])
              (error_code, group_id, state, protocol_type, protocol, members) = list_mebers_in_groups[0]
    
              if len(members) !=0:
                for member in members:
                  (member_id, client_id, client_host, member_metadata, member_assignment) = member
                  member_topics_assignment = []
                  for (topic, partitions) in MemberAssignment.decode(member_assignment).assignment:
                    member_topics_assignment.append(topic)
    
                  for topic in member_topics_assignment:
                    consumer = KafkaConsumer(
                              bootstrap_servers=BOOTSTRAP_SERVERS,
                              group_id=group[0],
                              enable_auto_commit=False
                              )
                    consumer.topics()
    
                    for p in consumer.partitions_for_topic(topic):
                      tp = TopicPartition(topic, p)
                      consumer.assign([tp])
                      committed = consumer.committed(tp)
                      consumer.seek_to_end(tp)
                      last_offset = consumer.position(tp)
                      if last_offset != None and committed != None:
                        lag = last_offset - committed
                        print "group: {} topic:{} partition: {} lag: {}".format(group[0], topic, p, lag)
    

    【讨论】:

      猜你喜欢
      • 2020-08-08
      • 1970-01-01
      • 1970-01-01
      • 2023-04-11
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2020-03-25
      相关资源
      最近更新 更多