You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何获取Kafka主题消息总数及未处理消息数?(基于cppkafka/librdkafka)

获取Kafka主题消息总数与未处理消息数的方法

一、获取Kafka主题的消息总数

Kafka没有直接提供获取主题总消息数的API,因为消息总数是所有分区的(最新偏移量 - 起始偏移量)之和。你可以通过以下两种方式计算:

  1. 使用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,将所有分区的结果相加就是主题总消息数。

  2. 通过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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.19 20:57:29