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

SpringBoot 3 Kafka消费者日志无法打印TraceId-SpanId问题

Kafka消费时无法打印TraceId和SpanId问题

我有两个基于SpringBoot 3的Producer和Consumer应用,通过Kafka通信,均引入micrometer brave依赖。调用Consumer的REST接口时日志能正常打印TraceId和SpanId,但消费Kafka消息时无法打印。

Producer已开启kafkaTemplate.setObservationEnabled(true);,通过Kafka控制台消费者可看到消息头中包含traceparent信息;Consumer能获取到该消息头,但日志中无TraceId和SpanId。相关代码配置如下:

Producer依赖

<dependency>
    <artifactId>spring-boot-starter-actuator</artifactId>
    <groupId>org.springframework.boot</groupId>
</dependency>
<dependency>
    <artifactId>micrometer-tracing-bridge-brave</artifactId>
    <groupId>io.micrometer</groupId>
</dependency>

Consumer监听器代码

@KafkaListener(containerFactory = "tracingKafkaConsumerFactory", topics = "tracingTopic3")
public void kafkaTracingListener(ConsumerRecords<String, Object> consumerRecords){
    consumerRecords.forEach(consumerRecord ->{
        try {
            if(consumerRecord.headers()!=null) {
                consumerRecord.headers().forEach(header -> {
                    LOGGER.info("Header key {}, Header value {}", header.key(), new String(header.value()));
                });
            }
            TraceDTO traceDto = objectMapper.readValue((String) consumerRecord.value(), TraceDTO.class);
            LOGGER.info("consumer record {}", consumerRecord);
        } catch (Exception e) {
            LOGGER.error("Exception :", e);
        }
    });
    LOGGER.info("Is trace printed here? ");
}

Consumer日志示例

2024-08-02T11:42:07.483+05:30  INFO 30481 --- [trace-demo] [ntainer#0-1-C-1] [               ] c.e.trace.demo.TracingDemoKafkaConsumer  : Header key traceparent, Header value 00-66ac78b618581c0549d46cbed2dfa9dc-9675e35b152ee946-01
2024-08-02T11:42:07.484+05:30  INFO 30481 --- [trace-demo] [ntainer#0-1-C-1] [               ] c.e.trace.demo.TracingDemoKafkaConsumer  : Header key __TypeId__, Header value com.srs.chargesessionmanagement.tracetest.TraceDTO
2024-08-02T11:42:07.537+05:30  INFO 30481 --- [trace-demo] [ntainer#0-1-C-1] [               ] c.e.trace.demo.TracingDemoKafkaConsumer  : data TraceDTO(id=11, data=eleven) 
2024-08-02T11:42:07.571+05:30  INFO 30481 --- [trace-demo] [ntainer#0-1-C-1] [               ] c.e.trace.demo.TracingDemoKafkaConsumer  : consumer record ConsumerRecord(topic = tracingTopic3, partition = 1, leaderEpoch = 0, offset = 2, CreateTime = 1722579126973, serialized key size = 3, serialized value size = 25, headers = RecordHeaders(headers = [RecordHeader(key = traceparent, value = [48, 48, 45, 54, 54, 97, 99, 55, 56, 98, 54, 49, 56, 53, 56, 49, 99, 48, 53, 52, 57, 100, 52, 54, 99, 98, 101, 100, 50, 100, 102, 97, 57, 100, 99, 45, 57, 54, 55, 53, 101, 51, 53, 98, 49, 53, 50, 101, 101, 57, 52, 54, 45, 48, 49]), RecordHeader(key = __TypeId__, value = [99, 111, 109, 46, 115, 114, 115, 46, 99, 104, 97, 114, 103, 101, 115, 101, 115, 115, 105, 111, 110, 109, 97, 110, 97, 103, 101, 109, 101, 110, 116, 46, 116, 114, 97, 99, 101, 116, 101, 115, 116, 46, 84, 114, 97, 99, 101, 68, 84, 79])], isReadOnly = false), key = key, value = {"id":11,"data":"eleven"})
2024-08-02T11:42:07.571+05:30  INFO 30481 --- [trace-demo] [ntainer#0-1-C-1] [               ] c.e.trace.demo.TracingDemoKafkaConsumer  : Is trace printed here? 

Consumer配置类

@EnableKafka
@Configuration
public class TraceDemoKafkaConsumerConfig {

    @Bean
    public Map<String, Object> tracingKafkaConsumerConfigs() {
        Map<String, Object> props = new HashMap<>();
        props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG,
                "localhost:9094");
        props.put(ConsumerConfig.GROUP_ID_CONFIG, "tracing-group-3");
        props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, 10);
        props.put(ConsumerConfig.MAX_POLL_INTERVAL_MS_CONFIG, "300000");
        props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false);
        props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
        props.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, "15000");
        props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        return props;
    }

    @Bean
    public DefaultKafkaConsumerFactory<String, Object> tracingKafkaConsumerDefaultFactory() {
        return new DefaultKafkaConsumerFactory<>(tracingKafkaConsumerConfigs());
    }

    @Bean
    public ConcurrentKafkaListenerContainerFactory<String, Object> tracingKafkaConsumerFactory(){
        ConcurrentKafkaListenerContainerFactory<String, Object> factory = new ConcurrentKafkaListenerContainerFactory<>();
        factory.setConsumerFactory(tracingKafkaConsumerDefaultFactory());
        factory.setConcurrency(2);
        factory.setBatchListener(true);
        factory.getContainerProperties().setObservationEnabled(true);
        factory.getContainerProperties().setSyncCommits(true);
        factory.getContainerProperties()
                .setAckMode(ContainerProperties.AckMode.valueOf("BATCH"));
        return factory;
    }
}

application.yaml配置

management:
  tracing:
    sampling:
      probability: 1

解决方法

1. 补全必要依赖

确认项目中包含micrometer-observation依赖(Spring Boot 3部分场景需手动引入):

<dependency>
    <groupId>io.micrometer</groupId>
    <artifactId>micrometer-observation</artifactId>
</dependency>

2. 修正日志格式配置

检查日志框架(如Logback)的配置,确保格式中包含traceId和spanId占位符,示例Logback配置:

<encoder>
    <pattern>%d{yyyy-MM-dd HH:mm:ss.SSS} [%thread] %-5level %logger{36} - %msg [%X{traceId:-}, %X{spanId:-}]%n</pattern>
</encoder>

3. 为Consumer工厂绑定观测注册中心

手动创建的DefaultKafkaConsumerFactory需注入ObservationRegistry才能正确处理trace上下文,修改配置类:

@EnableKafka
@Configuration
public class TraceDemoKafkaConsumerConfig {

    private final ObservationRegistry observationRegistry;

    // 构造注入观测注册中心
    public TraceDemoKafkaConsumerConfig(ObservationRegistry observationRegistry) {
        this.observationRegistry = observationRegistry;
    }

    // 原配置方法不变...

    @Bean
    public DefaultKafkaConsumerFactory<String, Object> tracingKafkaConsumerDefaultFactory() {
        DefaultKafkaConsumerFactory<String, Object> factory = new DefaultKafkaConsumerFactory<>(tracingKafkaConsumerConfigs());
        // 绑定观测注册中心
        factory.setObservationRegistry(this.observationRegistry);
        return factory;
    }

    // 原容器工厂配置不变...
}

4. 批量监听模式下的上下文激活(可选)

如果使用批量监听,可手动为每条消息激活trace上下文,确保日志正确关联:

@Autowired
private ObservationRegistry observationRegistry;

@KafkaListener(containerFactory = "tracingKafkaConsumerFactory", topics = "tracingTopic3")
public void kafkaTracingListener(ConsumerRecords<String, Object> consumerRecords){
    consumerRecords.forEach(record -> {
        Header traceHeader = record.headers().lastHeader("traceparent");
        if (traceHeader != null) {
            String traceParent = new String(traceHeader.value());
            // 基于traceparent创建并激活观测上下文
            try (Observation.Scope scope = Observation.start("kafka-batch-record", observationRegistry)
                    .setParent(TraceContextParser.parseTraceParent(traceParent))
                    .start()
                    .openScope()) {
                processRecord(record);
            }
        } else {
            processRecord(record);
        }
    });
    LOGGER.info("Is trace printed here? ");
}

private void processRecord(ConsumerRecord<String, Object> record) {
    // 原消息处理逻辑
}

5. 优先使用Spring Boot自动配置(推荐)

放弃手动创建容器工厂,通过application.yaml配置Kafka属性,让Spring Boot自动生成带观测支持的容器:

spring:
  kafka:
    consumer:
      bootstrap-servers: localhost:9094
      group-id: tracing-group-3
      max-poll-records: 10
      max-poll-interval-ms: 300000
      enable-auto-commit: false
      auto-offset-reset: earliest
      session-timeout-ms: 15000
      key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
      value-deserializer: org.apache.kafka.common.serialization.StringDeserializer
    listener:
      observation-enabled: true
      ack-mode: BATCH
      concurrency: 2
      sync-commits: true

这样无需手动配置工厂,Spring会自动处理观测上下文的传递与激活。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.20 02:17:03