如何配置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
相关产品推荐
相关产品推荐

