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

如何追踪Kafka未消费消息并解决重启丢消息问题

问题描述

我有一个向指定Topic发送消息的Kafka生产者,代码如下:

// Produce a message to a Kafka topic
String topic = "dev.topic.proxy";
String key = "some-key";
String value = "Hello, Kafka!";

ProducerRecord<String, String> record = new ProducerRecord<>(topic, key, value);

while (true) {
    producer.send(record);
    System.out.println("Message has been sent");
}

我的消费者基于Camel框架实现,配置如下:

<routes xmlns="http://camel.apache.org/schema/spring">
    <route id="readMsg">
        <description>Route to refresh the proxy routes</description>
       <from uri="appKafkaConsumer:dev.topic.proxy" />
        <log message="Received message with headers: ${headers}" loggingLevel="INFO" />
        <delay>
        <constant>5000</constant>
    </delay>
    </route>
</routes>

该消费者的groupId会在应用启动时自动随机生成,当我停止并重启应用后,由于分配了新的groupId,导致大量消息丢失。我希望避免消息丢失,同时想了解如何追踪Kafka中的未消费消息,请问最优解决方法是什么?

解决方案

一、避免消息丢失的核心措施

  • 固定消费者Group ID:这是解决重启后消息丢失的核心。Kafka通过Group ID维护消费偏移量,只要使用固定的Group ID,重启应用后就能从上次提交的偏移量继续消费。修改Camel路由的消费者URI,添加固定groupId参数:
    <from uri="appKafkaConsumer:dev.topic.proxy?groupId=fixed-proxy-consumer-group" />
    
  • 保障偏移量提交可靠性:Camel Kafka组件默认自动提交偏移量,若需更严谨的消费(比如确保消息处理完成后再提交),可开启手动提交:
    <from uri="appKafkaConsumer:dev.topic.proxy?groupId=fixed-proxy-consumer-group&autoCommitEnable=false&allowManualCommit=true" />
    
    随后在路由处理逻辑末尾添加偏移量提交操作,比如调用kafkaManualCommit()处理器。
  • 生产者端增强消息可靠性:当前生产者未处理发送结果,建议添加回调确认消息是否成功写入Kafka,同时配置生产者参数acks=all,确保所有同步副本都确认接收消息:
    // 生产者配置添加
    props.put(ProducerConfig.ACKS_CONFIG, "all");
    
    // 发送消息时添加回调
    producer.send(record, (metadata, exception) -> {
        if (exception != null) {
            System.err.println("消息发送失败: " + exception.getMessage());
        } else {
            System.out.println("消息已发送至分区 " + metadata.partition() + ", 偏移量 " + metadata.offset());
        }
    });
    

二、追踪未消费消息的方法

  • 使用Kafka命令行工具:
    • 查看指定Group的消费进度,计算未消费消息数:
      kafka-consumer-groups.sh --bootstrap-server <kafka-broker地址>:9092 --describe --group fixed-proxy-consumer-group
      
      输出中CURRENT-OFFSET是当前消费到的偏移量,LOG-END-OFFSET是分区最新偏移量,两者差值即为未消费消息数。
    • 查看Topic各分区的总消息数:
      kafka-run-class.sh kafka.tools.GetOffsetShell --bootstrap-server <kafka-broker地址>:9092 --topic dev.topic.proxy --time -1
      
  • Camel路由内监控:在路由中添加日志记录消费偏移量,或者集成Camel Metrics组件,统计消费速率、待处理消息数等指标。
  • Kafka监控指标:通过Prometheus采集Kafka的kafka_consumer_group_current_offset和kafka_topic_partition_current_offset指标,结合Grafana可视化两者差值,实时监控未消费消息量,还可设置告警阈值。

内容的提问来源于stack exchange,提问作者Aayush

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.05 13:08:19