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

如何在Kafka生产者回调中添加Spring Sleuth的trace_id

在Spring Sleuth中为Kafka生产者回调添加Trace ID

好问题!Spring Sleuth的Trace上下文默认绑定到当前线程,但Kafka生产者的回调是在Kafka客户端的异步线程中执行的,这就导致原来的Trace上下文丢失,日志里自然看不到traceId。不过有两种简单可行的方法解决这个问题:

方法一:手动传递并恢复Trace上下文

如果你需要直接在代码中获取traceId,或者想完全控制上下文的传递,可以手动保存当前的TraceContext,在回调执行时恢复它:

首先注入Sleuth的Tracer bean:

@Autowired
private Tracer tracer;

然后修改回调代码:

// 保存当前的Trace上下文
TraceContext traceContext = TraceContextHolder.getCurrentContext();

ListenableFuture<SendResult<K, V>> send = kafkaTemplate.send(topic, data);
send.addCallback(new ListenableFutureCallback<>() {
    @Override
    public void onFailure(Throwable throwable) {
        // 恢复Trace上下文,try-with-resources会自动关闭scope
        try (TraceContextScope scope = tracer.withSpan(traceContext)) {
            log.error(
                MessageFormat.format(
                    "Error when sending message {0} to Kafka, traceId: {1}", 
                    data.getGlobalUUID(), 
                    traceContext.traceIdString()),
                throwable);
            deferredResult.setErrorResult(
                new MarkusKafkaInputException(
                    MessageFormat.format(
                        "Error proceed {0} with message: {1}, traceId: {2}", 
                        data.getGlobalUUID(), 
                        throwable.getMessage(),
                        traceContext.traceIdString())));
        }
    }

    @Override
    public void onSuccess(SendResult<K, V> result) {
        try (TraceContextScope scope = tracer.withSpan(traceContext)) {
            log.trace("Message sent successfully, traceId: {}, result: {}", 
                traceContext.traceIdString(), result.toString());
        }
    }
});

方法二:使用Sleuth的TraceableListenableFutureCallback(推荐)

Spring Cloud Sleuth提供了专门的工具类来处理异步回调的上下文传递,不需要手动管理TraceContext,代码更简洁:

同样需要注入Tracer,然后用TraceableListenableFutureCallback包装你的回调:

import org.springframework.cloud.sleuth.instrument.async.TraceableListenableFutureCallback;

@Autowired
private Tracer tracer;

// ...

ListenableFuture<SendResult<K, V>> send = kafkaTemplate.send(topic, data);
send.addCallback(new TraceableListenableFutureCallback<>(tracer, new ListenableFutureCallback<>() {
    @Override
    public void onFailure(Throwable throwable) {
        log.error(
            MessageFormat.format(
                "Error when sending message {0} to Kafka", data.getGlobalUUID()),
            throwable);
        deferredResult.setErrorResult(
            new MarkusKafkaInputException(
                MessageFormat.format(
                    "Error proceed {0} with message: {1}", 
                    data.getGlobalUUID(), 
                    throwable.getMessage())));
    }

    @Override
    public void onSuccess(SendResult<K, V> result) {
        log.trace(result.toString());
    }
}));

这种方式下,Sleuth会自动帮你保存和恢复Trace上下文,只要你的日志配置里已经开启了traceId的打印(比如Logback的pattern中加入%X{traceId}),日志就会自动带上traceId,不需要手动拼接。

额外提示:确保日志配置输出Trace ID

为了让日志自动显示traceId,需要在日志框架的配置中添加对应的变量。比如Logback的配置:

<appender name="CONSOLE" class="ch.qos.logback.core.ConsoleAppender">
    <encoder>
        <pattern>%d{yyyy-MM-dd HH:mm:ss.SSS} [%thread] %-5level %logger{36} - traceId: %X{traceId} - %msg%n</pattern>
    </encoder>
</appender>

这样无论用哪种方法,日志里都会自动包含traceId,方便追踪链路。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.14 07:37:03