【发布时间】:2020-02-20 16:16:03
【问题描述】:
我们想知道给定 Kafka 主题尚未被 Kafka Consumer 消费和确认的消息数。
有没有什么方法可以从 Java 中的给定主题(主题有 10 个分区)中获取未使用的 Kafka 消息的计数??
【问题讨论】:
标签: java apache-kafka kafka-consumer-api spring-kafka
我们想知道给定 Kafka 主题尚未被 Kafka Consumer 消费和确认的消息数。
有没有什么方法可以从 Java 中的给定主题(主题有 10 个分区)中获取未使用的 Kafka 消息的计数??
【问题讨论】:
标签: java apache-kafka kafka-consumer-api spring-kafka
还有另一种方法
Kafka Consumer 拥有 API 来获取主题的每个分区的端点
List partitions = new ArrayList<>();
for (PartitionInfo p : parts) {
partitions.add(new TopicPartition(topic, p.partition()));
}
Map<TopicPartition, Long> offsets = consumer.endOffsets(partitions);
对于每个主题分区,您可以获得最新的提交偏移量。您可以使用这两个数字轻松获得未消耗的滞后
lag=(offset-latest commit 结束)
for (TopicPartition tp : offsets.keySet()) {
OffsetAndMetadata commitOffset = consumer.committed(new TopicPartition(tp.topic(), tp.partition()));
Long lag = commitOffset == null ? offsets.get(tp) : offsets.get(tp) - commitOffset.offset();
}
【讨论】:
kafka 消费者有一个每个分区的“记录滞后”指标(在 kip-92 中引入)。 您可以对主题的所有分区的这些指标求和,并获得未使用消息的数字
【讨论】:
MessageListenerContainer.metrics() 访问指标,它会返回一个包含每个消费者指标的映射,并以consumer.id 为键。