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

如何将Spring-Kafka的Consumer Lag暴露为Prometheus指标?

如何在Spring-Kafka中暴露Consumer Lag到Prometheus

1. 确保依赖齐全

首先确认项目中包含所需依赖(以Maven为例),这些依赖是支撑指标采集和暴露的基础:

<dependencies>
    <dependency>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-starter-actuator</artifactId>
    </dependency>
    <dependency>
        <groupId>org.springframework.kafka</groupId>
        <artifactId>spring-kafka</artifactId>
    </dependency>
    <dependency>
        <groupId>io.micrometer</groupId>
        <artifactId>micrometer-registry-prometheus</artifactId>
    </dependency>
    <!-- 若为非Web应用可省略spring-boot-starter-web,仅需保证Actuator能正常提供端点 -->
    <dependency>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-starter-web</artifactId>
    </dependency>
</dependencies>

2. 配置Actuator与指标开关

在application.properties或application.yml中开启Prometheus端点,并启用Kafka相关指标采集:

# 暴露Prometheus端点
management.endpoints.web.exposure.include=prometheus

# 显式开启Kafka指标采集(部分版本默认开启,明确配置更稳妥)
management.metrics.enable.kafka=true

# 可选:为指标添加全局标签,方便后续筛选
management.metrics.tags.kafka.consumer.group=your-base-group-id

如果每个监听器对应不同的消费组或Topic,建议在@KafkaListener注解中单独指定groupId,这样指标会自动按消费组、Topic维度拆分。

3. 配置Kafka容器的指标采集

Spring Kafka的ConcurrentKafkaListenerContainerFactory默认支持Micrometer指标,但可以显式配置确保Lag指标被采集:

@Configuration
public class KafkaConsumerConfig {

    @Bean
    public ConcurrentKafkaListenerContainerFactory<String, String> kafkaListenerContainerFactory(
            ConsumerFactory<String, String> consumerFactory,
            MeterRegistry meterRegistry) {
        ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>();
        factory.setConsumerFactory(consumerFactory);
        // 设置并发数为3
        factory.setConcurrency(3);
        
        // 启用Micrometer指标采集
        ContainerProperties containerProperties = factory.getContainerProperties();
        containerProperties.setMicrometerEnabled(true);
        // 添加Micrometer监听器,用于收集包括Lag在内的消费者指标
        containerProperties.addListener(new MicrometerConsumerListener<>(meterRegistry));
        
        return factory;
    }
}

4. 验证指标是否生成

启动应用后,访问你的myUrl/prometheus端点,搜索以下核心Lag相关指标:

  • kafka_consumer_records_lag:单分区的消费者Lag值
  • kafka_consumer_records_lag_sum:当前消费组下所有分区的Lag总和
  • kafka_consumer_records_lag_max:当前消费组下的最大分区Lag

这些指标会附带topic、consumer_group、partition等标签,可直接区分不同监听器、Topic的Lag情况。

常见排查点

  • 若未找到Lag指标,检查management.metrics.enable.kafka是否设为true
  • 确认消费者已成功连接Broker且有消息消费(无消息时Lag为0,部分版本可能不生成对应指标)
  • 优先使用Spring Boot管理的依赖版本,避免Micrometer与Spring Kafka的版本冲突

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.08 08:25:21