如何获取Kafka主题消息总数及未处理消息数?(基于cppkafka/librdkafka)
获取Kafka主题消息总数与未处理消息数的方法
一、获取Kafka主题的消息总数
Kafka没有直接提供获取主题总消息数的API,因为消息总数是所有分区的(最新偏移量 - 起始偏移量)之和。你可以通过以下两种方式计算:
使用Kafka自带脚本
运行kafka-topics.sh(Windows下为.bat),添加--describe参数查看每个分区的LogStartOffset和CurrentOffset,手动求和:kafka-topics.sh --describe --topic your_topic_name --bootstrap-server your_kafka_broker:9092单个分区的消息数为
CurrentOffset - LogStartOffset,将所有分区的结果相加就是主题总消息数。通过librdkafka代码计算
调用rd_kafka_query_watermark_offsets获取每个分区的最小(起始)和最大(最新)偏移量,遍历所有分区累加得到总数:rd_kafka_t *rk; // 省略rd_kafka_t实例初始化代码 rd_kafka_topic_partition_list_t *parts = rd_kafka_topic_partition_list_new(0); rd_kafka_topic_partition_list_add(parts, "your_topic", RD_KAFKA_PARTITION_UA); rd_kafka_metadata_t *metadata; if (rd_kafka_metadata(rk, 0, NULL, &metadata, 5000) == RD_KAFKA_RESP_ERR_NO_ERROR) { int64_t total_msg_count = 0; for (int i = 0; i < metadata->topic_cnt; i++) { rd_kafka_topic_metadata_t *topic_meta = metadata->topics[i]; if (strcmp(topic_meta->topic, "your_topic") == 0) { for (int j = 0; j < topic_meta->partition_cnt; j++) { rd_kafka_partition_metadata_t *part_meta = topic_meta->partitions[j]; int32_t partition = part_meta->id; int64_t low, high; if (rd_kafka_query_watermark_offsets(rk, "your_topic", partition, &low, &high, 5000) == RD_KAFKA_RESP_ERR_NO_ERROR) { total_msg_count += (high - low); } } } } rd_kafka_metadata_destroy(metadata); printf("主题总消息数:%" PRId64 "\n", total_msg_count); } rd_kafka_topic_partition_list_destroy(parts);
二、获取未处理消息数(消费组滞后量)
未处理消息数本质是每个分区的(最新偏移量 - 消费组已提交偏移量)之和,确实和偏移量直接相关。以下是librdkafka和cppkafka的实现方式:
1. 使用librdkafka实现
核心步骤:获取主题分区列表 → 获取消费组已提交偏移量 → 获取分区最新偏移量 → 计算差值累加
rd_kafka_t *rk; const char *group_id = "your_consumer_group"; const char *topic = "your_topic"; // 初始化消费者实例(需指定group.id配置) rd_kafka_topic_partition_list_t *parts = rd_kafka_topic_partition_list_new(0); rd_kafka_topic_partition_list_add(parts, topic, RD_KAFKA_PARTITION_UA); int64_t total_lag = 0; if (rd_kafka_committed(rk, parts, 5000) == RD_KAFKA_RESP_ERR_NO_ERROR) { for (int i = 0; i < parts->cnt; i++) { rd_kafka_topic_partition_t *part = &parts->elems[i]; int64_t high; if (rd_kafka_query_watermark_offsets(rk, topic, part->partition, NULL, &high, 5000) == RD_KAFKA_RESP_ERR_NO_ERROR) { if (part->offset != RD_KAFKA_OFFSET_INVALID) { int64_t lag = high - part->offset; total_lag += lag > 0 ? lag : 0; } } } printf("未处理消息总数:%" PRId64 "\n", total_lag); } rd_kafka_topic_partition_list_destroy(parts);
2. 使用cppkafka实现
cppkafka是librdkafka的C++封装,API更简洁:
#include <cppkafka/cppkafka.h> using namespace cppkafka; int main() { Configuration config = { {"bootstrap.servers", "your_kafka_broker:9092"}, {"group.id", "your_consumer_group"} }; Consumer consumer(config); std::string topic = "your_topic"; // 获取主题所有分区 TopicPartitionList partitions = consumer.get_partitions(topic); int64_t total_lag = 0; // 获取消费组已提交偏移量 TopicPartitionList committed_offsets = consumer.get_commit_offsets(partitions); for (const auto& partition : committed_offsets) { // 获取分区最新偏移量(高水位) int64_t high_watermark = consumer.get_high_watermark(partition); if (partition.get_offset() != Offset::INVALID) { int64_t lag = high_watermark - partition.get_offset(); total_lag += lag > 0 ? lag : 0; } } std::cout << "未处理消息总数:" << total_lag << std::endl; return 0; }
注意:如果消费组从未提交过偏移量,获取到的偏移量会是Offset::INVALID,此时可根据业务逻辑选择从起始或最新位置开始消费。
内容的提问来源于stack exchange,提问作者kreuzerkrieg
相关产品推荐
相关产品推荐

