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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 04:52:06