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

SpringBoot3+Kafka+Otel+Micrometer微服务链路追踪异常求助

问题:Kafka消息消费时无法继承TraceId(Spring Boot 3 + Micrometer Tracing)

我们正将多个微服务升级至Spring Boot 3,同时从Spring Sleuth切换为Micrometer处理链路追踪,微服务间通过Kafka通信。微服务A生成并记录TraceId和SpanId,写入Kafka的消息带有traceparent、X_INSTANA_C、X_INSTANA_L、X_INSTANA_L_S、X_INSTANA_S和X_INSTANA_T头(traceparent包含已记录的TraceId),但微服务B消费消息时未使用该TraceId,无任何TraceId和SpanId记录,预期消费时应通过traceparent设置TraceId。

环境配置信息

共用POM依赖

<dependencies>
        <dependency>
            <groupId>org.springframework</groupId>
            <artifactId>spring-context</artifactId>
        </dependency>
        <dependency>
            <groupId>io.micrometer</groupId>
            <artifactId>micrometer-tracing</artifactId>
        </dependency>
        <dependency>
            <groupId>io.micrometer</groupId>
            <artifactId>micrometer-observation</artifactId>
        </dependency>
        <dependency>
            <groupId>io.micrometer</groupId>
            <artifactId>micrometer-tracing-bridge-otel</artifactId>
        </dependency>
        <dependency>
            <groupId>io.opentelemetry</groupId>
            <artifactId>opentelemetry-exporter-otlp</artifactId>
        </dependency>
        <dependency>
            <groupId>io.micrometer</groupId>
            <artifactId>micrometer-registry-otlp</artifactId>
        </dependency>
        <dependency>
            <groupId>org.springframework.boot</groupId>
            <artifactId>spring-boot-starter-aop</artifactId>
        </dependency>
</dependencies>

追踪配置类

@Configuration(proxyBeanMethods = false)
@ComponentScan
@PropertySource("classpath:tracing-settings.properties")
public class CommonTracingConfig {

  @Bean
  OtlpHttpSpanExporter otlpHttpSpanExporter(@Value("${management.otel.tracing.endpoint}") final String uri) {
    return OtlpHttpSpanExporter.builder()
      .setEndpoint(uri)
      .setCompression("gzip")
      .build();
  }

  @Bean
  ObservedAspect observedAspect(final ObservationRegistry observationRegistry) {
    return new ObservedAspect(observationRegistry);
  }
}

tracing-settings.properties

management.tracing.sampling.probability=1.0
management.tracing.propagation.type=W3C
management.otel.tracing.endpoint=http://${INSTANA_AGENT_HOST:localhost}:4318/v1/traces
management.otlp.metrics.export.resource-attributes.service.name=${spring.application.name}-${STAGE:local}

通用Kafka配置类

@Configuration
@ComponentScan
@RequiredArgsConstructor
@PropertySource("classpath:kafka-settings.properties")
@PropertySource("classpath:kafka-settings-${STAGE:local}.properties")
public class CommonKafkaConfig {

  @Bean
  public KafkaTemplate<?, ?> kafkaTemplate(final ProducerFactory<?, ?> producerFactory) {
    final KafkaTemplate<?, ?> kafkaTemplate = new KafkaTemplate<>(producerFactory);
    kafkaTemplate.setObservationEnabled(true);
    return kafkaTemplate;
  }

  @Bean
  public KafkaListenerContainerFactory<ConcurrentMessageListenerContainer<Object, Object>> kafkaBatchListenerContainerFactory(final ConcurrentKafkaListenerContainerFactoryConfigurer configurer,
                                                                                                                              final ConsumerFactory<Object, Object> consumerFactory,
                                                                                                                              final KafkaBatchRetrier kafkaBatchRetrier) {
    final ConcurrentKafkaListenerContainerFactory<Object, Object> factory = createKafkaListenerContainerFactory(configurer, consumerFactory);
    factory.setBatchListener(true);
    factory.setBatchInterceptor((records, consumer) -> BatchRecordCompactor.compactRecords(records));
    factory.getContainerProperties().setAdviceChain(kafkaBatchRetrier);
    return factory;
  }

  private ConcurrentKafkaListenerContainerFactory<Object, Object> createKafkaListenerContainerFactory(final ConcurrentKafkaListenerContainerFactoryConfigurer configurer,
                                                                                                      final ConsumerFactory<Object, Object> consumerFactory) {
    final ConcurrentKafkaListenerContainerFactory<Object, Object> factory = new ConcurrentKafkaListenerContainerFactory<>();
    configurer.configure(factory, consumerFactory);
    factory.getContainerProperties().setTransactionManager(null);
    factory.getContainerProperties().setObservationEnabled(true);
    return factory;
  }
}

application.yml日志格式

logging:
  pattern:
    console: timestamp=%d thread=%t loglevel=%-5p class=%c appname=@artifactId@ traceid=%X{traceId:-} spanid=%X{spanId:-} message="%m"%n

缺失的配置及修复方案

1. 补充Spring Kafka Starter依赖

当前POM中缺少spring-boot-starter-kafka,它包含了Micrometer与Kafka集成所需的自动配置类,确保追踪上下文能在消息消费时自动传播。添加以下依赖:

<dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-kafka</artifactId>
</dependency>

2. 为Kafka容器配置ObservationRegistry

在CommonKafkaConfig中,需要将ObservationRegistry注入到容器配置中,确保消费时能从消息头提取追踪上下文:

  • 修改createKafkaListenerContainerFactory方法,添加ObservationRegistry参数并配置到容器属性:
private ConcurrentKafkaListenerContainerFactory<Object, Object> createKafkaListenerContainerFactory(final ConcurrentKafkaListenerContainerFactoryConfigurer configurer,
                                                                                                      final ConsumerFactory<Object, Object> consumerFactory,
                                                                                                      final ObservationRegistry observationRegistry) {
    final ConcurrentKafkaListenerContainerFactory<Object, Object> factory = new ConcurrentKafkaListenerContainerFactory<>();
    configurer.configure(factory, consumerFactory);
    factory.getContainerProperties().setTransactionManager(null);
    factory.getContainerProperties().setObservationEnabled(true);
    // 绑定ObservationRegistry,启用上下文提取
    factory.getContainerProperties().setObservationRegistry(observationRegistry);
    return factory;
}
  • 更新kafkaBatchListenerContainerFactory方法的参数列表,注入ObservationRegistry:
public KafkaListenerContainerFactory<ConcurrentMessageListenerContainer<Object, Object>> kafkaBatchListenerContainerFactory(final ConcurrentKafkaListenerContainerFactoryConfigurer configurer,
                                                                                                                              final ConsumerFactory<Object, Object> consumerFactory,
                                                                                                                              final KafkaBatchRetrier kafkaBatchRetrier,
                                                                                                                              final ObservationRegistry observationRegistry) {
    final ConcurrentKafkaListenerContainerFactory<Object, Object> factory = createKafkaListenerContainerFactory(configurer, consumerFactory, observationRegistry);
    // 原有逻辑保持不变
    factory.setBatchListener(true);
    factory.setBatchInterceptor((records, consumer) -> BatchRecordCompactor.compactRecords(records));
    factory.getContainerProperties().setAdviceChain(kafkaBatchRetrier);
    return factory;
}

3. 确保MDC追踪上下文传播

Micrometer Tracing需要将TraceId/SpanId同步到MDC才能被日志格式捕获。如果自定义了ObservationRegistry,需要手动添加MdcObservationHandler:
在CommonTracingConfig中添加:

@Bean
ObservationRegistry observationRegistry() {
    ObservationRegistry registry = ObservationRegistry.create();
    registry.observationConfig()
            // 启用MDC传播
            .observationHandler(new MdcObservationHandler())
            // 绑定追踪处理器
            .observationHandler(new DefaultTracingObservationHandler(
                    new PropagatingSenderTracingObservationHandler<>(
                            W3CTraceContextPropagator.getInstance(),
                            (carrier, key, value) -> {}
                    )
            ));
    return registry;
}

注:如果使用Spring Boot自动配置的ObservationRegistry,这一步可省略,但自定义时必须显式添加。

4. 验证批量拦截器是否保留消息头

检查BatchRecordCompactor.compactRecords(records)逻辑,确保压缩后的ConsumerRecords完整保留了原始消息的traceparent头。如果头信息丢失,消费端无法提取追踪上下文,需要修改压缩逻辑保证头的完整性。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.12 00:29:51