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

如何配置Sleuth在KafkaMessageListenerContainer而非监听器侧启动Span

Kafka+Sleuth链路ID在ErrorHandler中丢失问题解决方案

首先明确:完全可以通过自定义配置将Span的启动时机前置到KafkaMessageListenerContainer层级,解决ErrorHandler丢失spanId、traceId的问题。


问题根因

默认Spring Cloud Sleuth对Kafka监听器的链路埋点是在MessageListenerMethodInterceptor(方法级切面)中实现的,该切面的执行顺序晚于Kafka容器的异常拦截逻辑:当消费逻辑抛出异常时,切面会先结束Span、清空MDC中的链路信息,再将异常抛给ErrorHandler,最终导致ErrorHandler中无法获取链路ID。


具体实现步骤

  • 第一步:关闭Sleuth默认的Kafka消费侧自动埋点,在配置文件中添加如下配置:
spring:
  sleuth:
    messaging:
      kafka:
        consumer:
          enabled: false
  • 第二步:自定义KafkaMessageListenerContainer的Bean后置处理器,在容器初始化时包装原生消息监听器,将Span启动逻辑前置到监听器执行的最外层。
  • 第三步:实现链路追踪的监听器包装类,确保整个消费生命周期(包括异常抛出到ErrorHandler的阶段)都在Span的生效范围内,示例代码如下:
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.springframework.beans.BeansException;
import org.springframework.beans.factory.config.BeanPostProcessor;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.kafka.listener.KafkaMessageListenerContainer;
import org.springframework.kafka.listener.MessageListener;
import brave.Span;
import brave.Tracer;
import brave.Tracing;
import brave.kafka.clients.KafkaTracing;

@Configuration
public class SleuthKafkaCustomConfig {

    private final Tracing tracing;

    public SleuthKafkaCustomConfig(Tracing tracing) {
        this.tracing = tracing;
    }

    @Bean
    public BeanPostProcessor kafkaContainerTracingPostProcessor() {
        return new BeanPostProcessor() {
            @Override
            public Object postProcessAfterInitialization(Object bean, String beanName) throws BeansException {
                if (bean instanceof KafkaMessageListenerContainer<?, ?> container) {
                    Object rawListener = container.getContainerProperties().getMessageListener();
                    if (rawListener instanceof MessageListener<?, ?> rawMessageListener) {
                        // 替换为带链路追踪的包装类
                        container.getContainerProperties().setMessageListener(
                            new TracingMessageListenerWrapper(rawMessageListener, tracing)
                        );
                    }
                    return container;
                }
                return bean;
            }
        };
    }

    static class TracingMessageListenerWrapper<K, V> implements MessageListener<K, V> {
        private final MessageListener<K, V> delegate;
        private final KafkaTracing kafkaTracing;
        private final Tracer tracer;

        public TracingMessageListenerWrapper(MessageListener<K, V> delegate, Tracing tracing) {
            this.delegate = delegate;
            this.kafkaTracing = KafkaTracing.newBuilder(tracing).build();
            this.tracer = tracing.tracer();
        }

        @Override
        public void onMessage(ConsumerRecord<K, V> record) {
            // 从消息头提取链路信息,启动新Span
            Span consumeSpan = kafkaTracing.nextSpan(record).name("kafka.consume").start();
            try (Tracer.SpanInScope ignored = tracer.withSpanInScope(consumeSpan)) {
                delegate.onMessage(record);
            } catch (Exception e) {
                consumeSpan.error(e);
                // 抛出异常给ErrorHandler处理,此时Span尚未结束,MDC链路信息依然存在
                throw e;
            } finally {
                // ErrorHandler执行完成后才结束Span
                consumeSpan.finish();
            }
        }
    }
}

效果验证

上述配置完成后,Span的启动和关闭逻辑都在KafkaMessageListenerContainer的消息消费最外层,异常抛出到ErrorHandler时,Span还未关闭,MDC中的traceId、spanId都正常保留,可直接在ErrorHandler中获取打印。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.24 02:15:03