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

