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

Spring WebFlux+Reactor Kafka分布式追踪traceId/spanId丢失问题

解决Reactor Kafka与Micrometer Tracing的上下文传播问题

1. 确认核心依赖

确保项目包含以下Micrometer Tracing与Kafka集成的关键依赖(根据你的构建工具选择):

Maven

<dependency>
    <groupId>io.micrometer</groupId>
    <artifactId>micrometer-tracing-bridge-brave</artifactId>
</dependency>
<dependency>
    <groupId>io.micrometer</groupId>
    <artifactId>micrometer-tracing-reactor</artifactId>
</dependency>
<dependency>
    <groupId>io.micrometer</groupId>
    <artifactId>micrometer-tracing-instrumentation-kafka</artifactId>
</dependency>

Gradle

implementation 'io.micrometer:micrometer-tracing-bridge-brave'
implementation 'io.micrometer:micrometer-tracing-reactor'
implementation 'io.micrometer:micrometer-tracing-instrumentation-kafka'

2. 配置Reactor Kafka生产者上下文传播

在创建SenderOptions时,添加Micrometer提供的生产者拦截器,将trace上下文注入Kafka消息头:

import io.micrometer.tracing.instrument.kafka.TracingProducerInterceptor;
import reactor.kafka.sender.SenderOptions;
import org.springframework.boot.autoconfigure.kafka.KafkaProperties;
import java.util.ArrayList;
import java.util.List;

@Bean
public SenderOptions<String, Object> kafkaSenderOptions(KafkaProperties properties) {
    Map<String, Object> producerProps = properties.buildProducerProperties();
    
    // 添加Micrometer追踪拦截器
    List<String> interceptors = new ArrayList<>();
    interceptors.add(TracingProducerInterceptor.class.getName());
    producerProps.put(ProducerConfig.INTERCEPTOR_CLASSES_CONFIG, interceptors);
    
    return SenderOptions.create(producerProps);
}

@Bean
public ReactiveKafkaProducerTemplate<String, Object> reactiveKafkaProducerTemplate(SenderOptions<String, Object> senderOptions) {
    return new ReactiveKafkaProducerTemplate<>(senderOptions);
}

3. 配置Reactor Kafka消费者上下文传播

同样为消费者添加Micrometer拦截器,从Kafka消息头中解析trace上下文并绑定到Reactor上下文:

import io.micrometer.tracing.instrument.kafka.TracingConsumerInterceptor;
import reactor.kafka.receiver.ReceiverOptions;
import org.springframework.boot.autoconfigure.kafka.KafkaProperties;
import java.util.ArrayList;
import java.util.List;
import java.util.Collections;

@Bean
public ReceiverOptions<String, Object> kafkaReceiverOptions(KafkaProperties properties) {
    Map<String, Object> consumerProps = properties.buildConsumerProperties();
    
    // 添加Micrometer追踪拦截器
    List<String> interceptors = new ArrayList<>();
    interceptors.add(TracingConsumerInterceptor.class.getName());
    consumerProps.put(ConsumerConfig.INTERCEPTOR_CLASSES_CONFIG, interceptors);
    
    return ReceiverOptions.create(consumerProps)
            .subscription(Collections.singletonList("your-target-topic"));
}

@Bean
public KafkaReceiver<String, Object> kafkaReceiver(ReceiverOptions<String, Object> receiverOptions) {
    return KafkaReceiver.create(receiverOptions);
}

4. 关键配置验证

  • 确保上下文传播钩子早启动:在应用主类的main方法中最先执行Hooks.enableAutomaticContextPropagation():
@SpringBootApplication
public class KafkaTracingApplication {
    public static void main(String[] args) {
        Hooks.enableAutomaticContextPropagation();
        SpringApplication.run(KafkaTracingApplication.class, args);
    }
}
  • 日志格式包含traceId/spanId:在日志配置(如logback-spring.xml)中添加%X{traceId:-}和%X{spanId:-}:
<pattern>%d{yyyy-MM-dd HH:mm:ss.SSS} [%thread] %-5level %logger{36} - %msg %X{traceId:-} %X{spanId:-}%n</pattern>
  • WebFlux中保持上下文链:发送Kafka消息时不要脱离WebFlux的Reactor上下文,直接返回Mono而非手动调用subscribe():
@RestController
public class KafkaController {
    private final ReactiveKafkaProducerTemplate<String, Object> producerTemplate;

    public KafkaController(ReactiveKafkaProducerTemplate<String, Object> producerTemplate) {
        this.producerTemplate = producerTemplate;
    }

    @PostMapping("/send")
    public Mono<Void> sendMessage(@RequestBody String message) {
        return producerTemplate.send("your-target-topic", message).then();
    }
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.25 17:13:19