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

Spring Boot3中DefaultKafkaProducerFactory观测/B3头传播配置及多集群示例问询

Spring Boot 3 Kafka B3头转发与多集群配置方案

一、B3头转发逻辑的实现原理

Spring Boot 3中Kafka的B3传播并非直接在DefaultKafkaProducerFactory中硬编码注入,而是通过Micrometer观测拦截器完成:

  • 当引入spring-boot-starter-actuator、spring-boot-starter-micrometer和spring-kafka依赖后,自动配置会向DefaultKafkaProducerFactory注册ObservationProducerInterceptor。
  • 这个拦截器会在发送ProducerRecord前,从当前线程的Observation上下文(由Spring MVC/WebFlux或其他观测入口初始化)中提取B3格式的traceId、spanId等信息,添加到消息头中。
  • 你看到的MicrometerProducerListener负责观测指标的收集,而B3头的注入是由拦截器完成的,这部分逻辑在KafkaObservationDocumentation和ObservationProducerInterceptor的源码中可以找到。

二、多Kafka服务器(单主题对应单集群)配置示例

假设我们有两个Kafka集群,分别对应topic-a和topic-b,需要为每个集群单独配置ProducerFactory和KafkaTemplate,同时保留观测与B3传播能力。

1. 依赖配置(pom.xml)

<dependencies>
    <dependency>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-starter-web</artifactId>
    </dependency>
    <dependency>
        <groupId>org.springframework.kafka</groupId>
        <artifactId>spring-kafka</artifactId>
    </dependency>
    <dependency>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-starter-actuator</artifactId>
    </dependency>
    <dependency>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-starter-micrometer</artifactId>
    </dependency>
    <!-- B3传播依赖 -->
    <dependency>
        <groupId>io.micrometer</groupId>
        <artifactId>micrometer-tracing-bridge-brave</artifactId>
    </dependency>
</dependencies>

2. 配置文件(application.yml)

# 第一个Kafka集群(对应topic-a)
kafka-cluster-a:
  bootstrap-servers: localhost:9092
  producer:
    key-serializer: org.apache.kafka.common.serialization.StringSerializer
    value-serializer: org.apache.kafka.common.serialization.StringSerializer

# 第二个Kafka集群(对应topic-b)
kafka-cluster-b:
  bootstrap-servers: localhost:9093
  producer:
    key-serializer: org.apache.kafka.common.serialization.StringSerializer
    value-serializer: org.apache.kafka.common.serialization.StringSerializer

# 观测与B3配置
management:
  tracing:
    sampling:
      probability: 1.0 # 全采样便于测试
    propagation:
      type: B3 # 指定传播格式为B3

3. Java配置类

手动创建两个集群对应的ProducerFactory和KafkaTemplate,自动绑定观测组件以启用B3传播:

import org.apache.kafka.clients.producer.ProducerConfig;
import org.springframework.boot.context.properties.ConfigurationProperties;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.kafka.core.DefaultKafkaProducerFactory;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.kafka.core.ProducerFactory;
import io.micrometer.observation.ObservationRegistry;
import org.springframework.kafka.support.observation.KafkaProducerObservationConvention;

import java.util.Map;

@Configuration
public class MultiKafkaClusterConfig {

    // 绑定第一个集群的生产者配置
    @Bean
    @ConfigurationProperties(prefix = "kafka-cluster-a.producer")
    public Map<String, Object> kafkaClusterAProducerProps() {
        return Map.of();
    }

    @Bean
    public ProducerFactory<String, String> kafkaClusterAProducerFactory(Map<String, Object> kafkaClusterAProducerProps,
                                                                       ObservationRegistry observationRegistry,
                                                                       KafkaProducerObservationConvention observationConvention) {
        DefaultKafkaProducerFactory<String, String> factory = new DefaultKafkaProducerFactory<>(kafkaClusterAProducerProps);
        // 启用观测,自动注入B3拦截器
        factory.setObservationRegistry(observationRegistry);
        factory.setProducerObservationConvention(observationConvention);
        return factory;
    }

    @Bean
    public KafkaTemplate<String, String> kafkaClusterAKafkaTemplate(ProducerFactory<String, String> kafkaClusterAProducerFactory) {
        return new KafkaTemplate<>(kafkaClusterAProducerFactory);
    }

    // 第二个集群配置,逻辑与第一个一致
    @Bean
    @ConfigurationProperties(prefix = "kafka-cluster-b.producer")
    public Map<String, Object> kafkaClusterBProducerProps() {
        return Map.of();
    }

    @Bean
    public ProducerFactory<String, String> kafkaClusterBProducerFactory(Map<String, Object> kafkaClusterBProducerProps,
                                                                       ObservationRegistry observationRegistry,
                                                                       KafkaProducerObservationConvention observationConvention) {
        DefaultKafkaProducerFactory<String, String> factory = new DefaultKafkaProducerFactory<>(kafkaClusterBProducerProps);
        factory.setObservationRegistry(observationRegistry);
        factory.setProducerObservationConvention(observationConvention);
        return factory;
    }

    @Bean
    public KafkaTemplate<String, String> kafkaClusterBKafkaTemplate(ProducerFactory<String, String> kafkaClusterBProducerFactory) {
        return new KafkaTemplate<>(kafkaClusterBProducerFactory);
    }
}

4. 测试代码

创建Controller验证消息发送与B3头注入:

import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.PathVariable;
import org.springframework.web.bind.annotation.RestController;

@RestController
public class KafkaMessageController {

    private final KafkaTemplate<String, String> kafkaClusterAKafkaTemplate;
    private final KafkaTemplate<String, String> kafkaClusterBKafkaTemplate;

    public KafkaMessageController(KafkaTemplate<String, String> kafkaClusterAKafkaTemplate,
                                  KafkaTemplate<String, String> kafkaClusterBKafkaTemplate) {
        this.kafkaClusterAKafkaTemplate = kafkaClusterAKafkaTemplate;
        this.kafkaClusterBKafkaTemplate = kafkaClusterBKafkaTemplate;
    }

    @GetMapping("/send/a/{message}")
    public String sendToClusterA(@PathVariable String message) {
        kafkaClusterAKafkaTemplate.send("topic-a", message);
        return "Sent to cluster A: " + message;
    }

    @GetMapping("/send/b/{message}")
    public String sendToClusterB(@PathVariable String message) {
        kafkaClusterBKafkaTemplate.send("topic-b", message);
        return "Sent to cluster B: " + message;
    }
}

三、验证B3传播

启动项目后调用接口,通过Kafka客户端工具(如kafka-console-consumer)查看消息头,会发现b3头已自动添加,格式类似traceId-spanId-sampled。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.12 07:13:16