如何为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
相关产品推荐
相关产品推荐

