Spring Boot Kafka Listener的TraceId设置及来源判定文件位置查询
Spring Boot Kafka Listener 中 TraceId 的处理方案
1. 自定义 TraceId 工具类
用于存储、获取、生成 TraceId,基于 ThreadLocal 实现线程隔离:
import java.util.UUID; public class TraceIdUtils { private static final ThreadLocal<String> TRACE_ID_THREAD_LOCAL = new ThreadLocal<>(); private static final String TRACE_ID_HEADER = "trace-id"; public static String getTraceId() { return TRACE_ID_THREAD_LOCAL.get(); } public static void setTraceId(String traceId) { TRACE_ID_THREAD_LOCAL.set(traceId); } public static void clear() { TRACE_ID_THREAD_LOCAL.remove(); } // 生成UUID格式的TraceId public static String generateTraceId() { return UUID.randomUUID().toString().replace("-", ""); } public static String getTraceIdHeaderKey() { return TRACE_ID_HEADER; } }
2. 实现 Kafka 消费者拦截器
在消费消息前自动处理 TraceId:优先从消息头读取,读取不到则生成新的 TraceId,并绑定到当前线程
import org.apache.kafka.clients.consumer.ConsumerInterceptor; import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.clients.consumer.ConsumerRecords; import org.apache.kafka.clients.consumer.OffsetAndMetadata; import org.apache.kafka.common.TopicPartition; import java.util.Map; public class TraceIdConsumerInterceptor implements ConsumerInterceptor<String, Object> { @Override public ConsumerRecords<String, Object> onConsume(ConsumerRecords<String, Object> records) { for (TopicPartition partition : records.partitions()) { for (ConsumerRecord<String, Object> record : records.records(partition)) { String traceId = null; // 遍历消息头查找TraceId for (org.apache.kafka.common.header.Header header : record.headers()) { if (TraceIdUtils.getTraceIdHeaderKey().equals(header.key())) { traceId = new String(header.value()); break; } } // 无有效TraceId则生成新的 if (traceId == null || traceId.isBlank()) { traceId = TraceIdUtils.generateTraceId(); } TraceIdUtils.setTraceId(traceId); } } return records; } @Override public void onCommit(Map<TopicPartition, OffsetAndMetadata> offsets) { // 提交消息后清理ThreadLocal,避免内存泄漏 TraceIdUtils.clear(); } @Override public void configure(Map<String, ?> configs) {} @Override public void close() { TraceIdUtils.clear(); } }
3. 注册消费者拦截器
在 application.yml 中配置拦截器全类名:
spring: kafka: consumer: properties: interceptor.classes: com.yourpackage.TraceIdConsumerInterceptor
4. 在 Kafka Listener 中使用 TraceId
直接通过工具类获取当前线程绑定的 TraceId:
import org.springframework.kafka.annotation.KafkaListener; import org.springframework.stereotype.Component; @Component public class DemoKafkaListener { @KafkaListener(topics = "your-topic", groupId = "demo-group") public void handleMessage(String message) { String traceId = TraceIdUtils.getTraceId(); // 业务逻辑中使用TraceId,比如日志打印、链路追踪 System.out.printf("处理消息 | TraceId: %s | 内容: %s%n", traceId, message); } }
补充:生产者端传递 TraceId
如果需要从上游服务(如HTTP接口)传递TraceId到Kafka消息,需在生产者端添加拦截器,将当前线程的TraceId写入消息头:
import org.apache.kafka.clients.producer.ProducerInterceptor; import org.apache.kafka.clients.producer.ProducerRecord; import org.apache.kafka.clients.producer.RecordMetadata; import java.util.Map; public class TraceIdProducerInterceptor implements ProducerInterceptor<String, Object> { @Override public ProducerRecord<String, Object> onSend(ProducerRecord<String, Object> record) { String traceId = TraceIdUtils.getTraceId(); if (traceId != null) { // 将TraceId添加到消息头 return new ProducerRecord<>( record.topic(), record.partition(), record.timestamp(), record.key(), record.value(), record.headers().add(TraceIdUtils.getTraceIdHeaderKey(), traceId.getBytes()) ); } return record; } @Override public void onAcknowledgement(RecordMetadata metadata, Exception exception) {} @Override public void configure(Map<String, ?> configs) {} @Override public void close() {} }
生产者配置添加拦截器:
spring: kafka: producer: properties: interceptor.classes: com.yourpackage.TraceIdProducerInterceptor
内容的提问来源于stack exchange,提问作者Remco Cats
相关产品推荐
相关产品推荐

