如何查询Kafka消费者组最后提交Offset的时间?
查询Kafka消费者组Offset提交时间的方法
1. 使用kafka-consumer-groups工具(推荐)
从Kafka 2.0版本开始,kafka-consumer-groups工具新增了--verbose参数,添加后就能显示Offset的最后提交时间。修改你的命令为:
bin/windows/kafka-consumer-groups.bat --bootstrap-server localhost:9092 --group consumer-group --describe --topic com.gdn.test --verbose
执行后会额外输出LAST-COMMIT-TIME字段,这就是你需要的最后提交Offset的时间。
2. 查看__consumer_offsets内部主题(进阶方式)
消费者的Offset提交信息默认存储在__consumer_offsets内部主题中,你可以用kafka-console-consumer工具解析该主题内容:
bin/windows/kafka-console-consumer.bat --bootstrap-server localhost:9092 --topic __consumer_offsets --from-beginning --formatter "kafka.coordinator.group.GroupMetadataManager\$OffsetsMessageFormatter" --property print.key=true --property key.deserializer=org.apache.kafka.common.serialization.StringDeserializer --property value.deserializer=org.apache.kafka.common.serialization.StringDeserializer
输出结果会包含各消费者组、分区对应的Offset提交时间,你需要自行筛选目标消费者组和主题的相关信息。
3. 通过AdminClient API编程查询(自定义场景)
如果需要将查询逻辑集成到代码中,可以使用Kafka的AdminClient API获取详细信息,以下是Java示例代码:
import org.apache.kafka.clients.admin.AdminClient; import org.apache.kafka.clients.admin.AdminClientConfig; import org.apache.kafka.clients.admin.ListConsumerGroupOffsetsResult; import org.apache.kafka.common.TopicPartition; import java.util.Properties; import java.util.concurrent.ExecutionException; public class OffsetCommitTimeQuery { public static void main(String[] args) throws ExecutionException, InterruptedException { Properties props = new Properties(); props.put(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092"); try (AdminClient adminClient = AdminClient.create(props)) { ListConsumerGroupOffsetsResult offsetsResult = adminClient.listConsumerGroupOffsets("consumer-group"); offsetsResult.partitionsToOffsetAndMetadata().get().forEach((tp, offsetMeta) -> { if (tp.topic().equals("com.gdn.test")) { System.out.println("分区: " + tp.partition() + ", 最后提交时间: " + offsetMeta.commitTimestamp()); } }); } } }
这段代码会输出指定消费者组、主题下每个分区的Offset提交时间。
内容的提问来源于stack exchange,提问作者Steven Kristian
相关产品推荐
相关产品推荐

