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

如何为Spring Boot Kafka指标添加Topic标签?

如何在Spring Boot的Kafka Prometheus指标中添加Topic标签?

Dropwizard生成的Kafka Prometheus指标会将Topic作为生产者和消费者的标签,这对集群特定Pod的负载监控与管控至关重要,但Spring Boot服务的同类指标默认未包含Topic标签,以下是针对你使用版本的实现方法:

示例对比

Dropwizard指标(含Topic标签)

kafka_producer_record_send_rate{application="my-dropwizard-service", client="my-dropwizard-service-20230925-123734.045_my-dropwizard-service-7d659554b4-lgk7r-StreamThread-1-producer", topic="ip-cds-gateway"} 0

Spring Boot默认指标(无Topic标签)

kafka_producer_record_send_rate{application="my-spring-boot-service", client_id="my-message.my-spring-boot-service-2731f24b-1d88-4120-b62f-558059ea9b59-StreamThread-1-producer", kafka_version="3.1.2", spring_id="defaultKafkaStreamsBuilder"} 0

依赖版本

  • spring-boot 2.7.15
  • spring-boot-starter-actuator 2.7.15
  • spring-kafka 2.8.11
  • micrometer-core 1.9.14
  • micrometer-rgistry-prometheus 1.9.14

已有配置

@Bean
MeterRegistryCustomizer<MeterRegistry> metricsCommonTags() {
    return registry -> registry.config().commonTags("application", "my-spring-boot-service");
}

实现方案

1. 自定义生产者指标监听(添加Topic标签)

创建扩展MicrometerProducerListener的自定义监听类,在创建指标时直接注入Topic标签:

import io.micrometer.core.instrument.MeterRegistry;
import org.springframework.kafka.core.MicrometerProducerListener;
import org.springframework.stereotype.Component;

@Component
public class TopicLabelProducerListener<K, V> extends MicrometerProducerListener<K, V> {

    public TopicLabelProducerListener(MeterRegistry meterRegistry) {
        super(meterRegistry);
    }

    @Override
    protected void createMeter(String metricName, String topic, String clientId, MeterRegistry registry) {
        // 创建指标时直接添加topic标签
        registry.counter(metricName, "client_id", clientId, "topic", topic);
    }
}

将该监听注册到生产者工厂:

import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.kafka.core.DefaultKafkaProducerFactory;
import org.springframework.kafka.core.ProducerFactory;

import java.util.Map;

@Configuration
public class KafkaProducerConfig {

    @Bean
    public ProducerFactory<?, ?> producerFactory(Map<String, Object> producerConfigs, TopicLabelProducerListener<?, ?> listener) {
        DefaultKafkaProducerFactory<?, ?> factory = new DefaultKafkaProducerFactory<>(producerConfigs);
        factory.addListener(listener);
        return factory;
    }
}

2. 自定义消费者指标监听(添加Topic标签)

同理,创建消费者端的自定义监听类:

import io.micrometer.core.instrument.MeterRegistry;
import org.springframework.kafka.core.MicrometerConsumerListener;
import org.springframework.stereotype.Component;

@Component
public class TopicLabelConsumerListener<K, V> extends MicrometerConsumerListener<K, V> {

    public TopicLabelConsumerListener(MeterRegistry meterRegistry) {
        super(meterRegistry);
    }

    @Override
    protected void createMeter(String metricName, String topic, String clientId, MeterRegistry registry) {
        registry.counter(metricName, "client_id", clientId, "topic", topic);
    }
}

注册到消费者工厂:

import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.kafka.core.DefaultKafkaConsumerFactory;
import org.springframework.kafka.core.ConsumerFactory;

import java.util.Map;

@Configuration
public class KafkaConsumerConfig {

    @Bean
    public ConsumerFactory<?, ?> consumerFactory(Map<String, Object> consumerConfigs, TopicLabelConsumerListener<?, ?> listener) {
        DefaultKafkaConsumerFactory<?, ?> factory = new DefaultKafkaConsumerFactory<>(consumerConfigs);
        factory.addListener(listener);
        return factory;
    }
}

3. Kafka Streams场景补充配置

如果使用Kafka Streams,需要开启指标并确保标签注入,可通过配置类实现:

import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.kafka.config.StreamsBuilderFactoryBeanConfigurer;

@Configuration
public class KafkaStreamsConfig {

    @Bean
    public StreamsBuilderFactoryBeanConfigurer streamsBuilderFactoryBeanConfigurer() {
        return factoryBean -> {
            // 开启Kafka Streams的Micrometer指标
            factoryBean.setMicrometerEnabled(true);
        };
    }
}

同时在配置文件中开启指标开关(application.yml):

spring:
  kafka:
    producer:
      metrics:
        enabled: true
    consumer:
      metrics:
        enabled: true
    streams:
      metrics:
        enabled: true

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.10 02:02:05