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

如何为Spring Kafka Template指标添加主题信息?版本兼容与方案问询

解决方案:Spring Kafka指标扩展与版本升级指南

一、不升级Spring Kafka的情况下添加主题标签到指标

1. 扩展LoggingProducerListener而非替换

直接实现ProducerListener会覆盖自动配置的LoggingProducerListener,可以通过继承它来保留原有日志功能,同时添加自定义指标逻辑:

@Component
public class MetricsEnhancingProducerListener extends LoggingProducerListener<Object, Object> {

    private final MeterRegistry meterRegistry;

    public MetricsEnhancingProducerListener(MeterRegistry meterRegistry) {
        this.meterRegistry = meterRegistry;
    }

    @Override
    public void onSuccess(ProducerRecord<Object, Object> record, RecordMetadata metadata) {
        super.onSuccess(record, metadata); // 保留原有日志逻辑
        // 自定义成功指标,添加主题标签
        meterRegistry.counter("spring_kafka_template_success_count",
                "topic", record.topic())
                .increment();
        // 记录发送耗时:需结合发送前的时间戳,可通过包装KafkaTemplate实现
    }

    @Override
    public void onError(ProducerRecord<Object, Object> record, Exception exception) {
        super.onError(record, exception); // 保留原有日志逻辑
        // 自定义错误指标,添加主题与错误类型标签
        meterRegistry.counter("spring_kafka_template_error_count",
                "topic", record.topic(),
                "error_type", exception.getClass().getSimpleName())
                .increment();
    }
}

2. 使用ProducerInterceptor拦截器

通过实现ProducerInterceptor在消息发送前后收集数据,自定义带主题标签的指标,不影响原有Listener:

@Component
public class KafkaMetricsInterceptor implements ProducerInterceptor<Object, Object> {

    private final MeterRegistry meterRegistry;
    private final ThreadLocal<Long> sendStartTime = new ThreadLocal<>();

    public KafkaMetricsInterceptor(MeterRegistry meterRegistry) {
        this.meterRegistry = meterRegistry;
    }

    @Override
    public ProducerRecord<Object, Object> onSend(ProducerRecord<Object, Object> record) {
        sendStartTime.set(System.currentTimeMillis());
        return record;
    }

    @Override
    public void onAcknowledgement(RecordMetadata metadata, Exception exception) {
        Long startTime = sendStartTime.get();
        if (startTime != null) {
            long durationMs = System.currentTimeMillis() - startTime;
            String topic = metadata != null ? metadata.topic() : "unknown";
            
            if (exception == null) {
                // 成功耗时指标
                meterRegistry.timer("spring_kafka_template_send_seconds", "topic", topic)
                        .record(durationMs, TimeUnit.MILLISECONDS);
                meterRegistry.counter("spring_kafka_template_success_count", "topic", topic).increment();
            } else {
                // 错误计数指标
                meterRegistry.counter("spring_kafka_template_error_count", "topic", topic, "error", exception.getClass().getSimpleName())
                        .increment();
            }
            sendStartTime.remove();
        }
    }

    @Override
    public void close() {
        sendStartTime.remove();
    }

    @Override
    public void configure(Map<String, ?> configs) {}
}

将拦截器配置到生产者工厂:

@Configuration
public class KafkaConfig {

    @Bean
    public ProducerFactory<Object, Object> producerFactory(KafkaProperties properties, KafkaMetricsInterceptor interceptor) {
        Map<String, Object> configs = properties.buildProducerProperties();
        List<String> interceptors = new ArrayList<>(configs.getOrDefault(ProducerConfig.INTERCEPTOR_CLASSES_CONFIG, Collections.emptyList()));
        interceptors.add(interceptor.getClass().getName());
        configs.put(ProducerConfig.INTERCEPTOR_CLASSES_CONFIG, interceptors);
        return new DefaultKafkaProducerFactory<>(configs);
    }
}

二、Spring Kafka 2.9.8安全升级指南(兼容Spring Boot 2.7.18)

Spring Boot 2.7.x的依赖管理原生支持Spring Kafka 2.8.x~2.9.x,升级无兼容性风险,步骤如下:

  1. 修改依赖版本
    Maven:
    <dependency>
        <groupId>org.springframework.kafka</groupId>
        <artifactId>spring-kafka</artifactId>
        <version>2.9.8</version>
    </dependency>
    
    Gradle:
    implementation 'org.springframework.kafka:spring-kafka:2.9.8'
    
  2. 配置MicrometerTagsProvider
    升级后直接使用官方API添加主题标签:
    @Bean
    public KafkaTemplate<Object, Object> kafkaTemplate(ProducerFactory<Object, Object> producerFactory) {
        KafkaTemplate<Object, Object> template = new KafkaTemplate<>(producerFactory);
        template.setMicrometerTagsProvider(record -> Arrays.asList(
                Tag.of("topic", record.topic())
        ));
        return template;
    }
    
  3. 验证要点
    • 检查核心功能:生产者发送、消费者接收是否正常
    • 确认依赖无冲突:通过mvn dependency:tree或gradle dependencies排查kafka-clients版本(2.9.8对应kafka-clients 2.8.x,与Spring Boot 2.7.18默认版本一致)

三、Spring Boot与Spring Kafka核心版本兼容列表

Spring Boot版本Spring Kafka版本范围对应kafka-clients版本
2.7.x2.8.x ~ 2.9.x2.8.x
3.0.x3.0.x ~ 3.1.x3.0.x
3.1.x3.1.x ~ 3.2.x3.1.x
2.6.x2.7.x ~ 2.8.x2.7.x

内容的提问来源于stack exchange,提问作者dev-007

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.16 16:50:15