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

