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

Spring Boot如何为ReactiveKafkaProducerTemplate配置Micrometer指标监控

为ReactiveKafkaProducerTemplate配置Actuator指标

针对Spring Boot 2.6.13 + Spring Kafka 2.8.10的场景,要让ReactiveKafkaProducerTemplate的指标能在Actuator端点展示,可通过以下两种方式实现:

方法一:添加Micrometer Kafka生产者拦截器

这种方式无需依赖Spring Kafka的ProducerListener,直接通过Kafka原生拦截器实现指标收集,对普通KafkaTemplate和ReactiveKafkaProducerTemplate都适用。

  1. 配置生产者参数并添加Micrometer拦截器:
@Bean
public ProducerSettings reactiveProducerSettings(MeterRegistry meterRegistry) {
    Map<String, Object> configProps = new HashMap<>();
    configProps.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "127.0.0.1:9092");
    configProps.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, Serdes.String().serializer().getClass().getName());
    configProps.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, Serdes.String().serializer().getClass().getName());
    // 配置Micrometer生产者拦截器
    configProps.put(ProducerConfig.INTERCEPTOR_CLASSES_CONFIG, 
        Collections.singletonList(MicrometerProducerInterceptor.class.getName()));
    // 绑定MeterRegistry实例
    configProps.put(MicrometerProducerInterceptor.METER_REGISTRY_BEAN_NAME, meterRegistry);
    
    return ProducerSettings.create(configProps);
}

@Bean
public ReactiveKafkaProducerTemplate<String, String> reactiveKafkaProducerTemplate(ProducerSettings producerSettings) {
    return new ReactiveKafkaProducerTemplate<>(producerSettings);
}

方法二:通过Spring Kafka的ProducerFactory关联监听器

虽然ReactiveKafkaProducerTemplate没有直接接收ProducerFactory的构造器,但可以借助ProducerFactory生成的配置,结合Spring的MicrometerProducerListener来收集指标:

  1. 配置带有Micrometer监听器的ProducerFactory:
@Bean
public ProducerFactory<String, String> customProducerFactory(MeterRegistry meterRegistry) {
    Map<String, Object> configProps = new HashMap<>();
    configProps.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "127.0.0.1:9092");
    configProps.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, Serdes.String().serializer().getClass().getName());
    configProps.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, Serdes.String().serializer().getClass().getName());
    
    DefaultKafkaProducerFactory<String, String> producerFactory = new DefaultKafkaProducerFactory<>(configProps);
    producerFactory.addListener(new MicrometerProducerListener<>(meterRegistry));
    return producerFactory;
}
  1. 基于ProducerFactory的配置创建ReactiveKafkaProducerTemplate:
@Bean
public ReactiveKafkaProducerTemplate<String, String> reactiveKafkaProducerTemplate(ProducerFactory<String, String> producerFactory) {
    // 从ProducerFactory中提取配置,构建ProducerSettings
    ProducerSettings producerSettings = ProducerSettings.create(producerFactory.getConfigurationProperties());
    return new ReactiveKafkaProducerTemplate<>(producerSettings);
}

验证配置

配置完成后启动应用,访问Actuator的/actuator/metrics端点,即可看到以kafka.producer.为前缀的指标,比如kafka.producer.record.send.total、kafka.producer.record.error.total等。

注意:确保项目已引入micrometer-registry-*依赖(如micrometer-registry-prometheus),且Actuator的metrics端点已在配置中启用。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.07 06:20:47