Spring Integration(Spring Boot3.0.9)链路追踪:如何保留传入traceId
问题描述
我正在用Spring Integration实现一条消息流转流程:从JMS接收消息,再发送到Kafka,代码如下:
@Override protected IntegrationFlowDefinition<?> buildFlow() { return from(Jms.messageDrivenChannelAdapter(connectionFactory).destination(mqInQueue)) .log(LoggingHandler.Level.DEBUG, message -> "Received JMS message: " + message.getPayload()) .channel(channels -> MessageChannelFactory.create(channels, "request-channel-1")) .handle(Kafka.outboundChannelAdapter(kafkaTemplate).topic(kafkaOutTopic)); }
为了测试,我通过POST请求配合JmsTemplate向MQ发送消息:
@PostMapping("/mq") public String sentToMq(@RequestBody final String body) { jmsTemplate.convertAndSend(mqRequestQueue, body, m -> { final var span = tracer.startScopedSpan("jms-send"); try { final var context = span.context(); m.setStringProperty("b3", "%s-%s-%s".formatted(context.traceId(), context.spanId(), Boolean.TRUE.equals(context.sampled()) ? "1" : "0")); } finally { span.end(); } return m; }); return "done"; }
目前除链路追踪外一切正常。我在发送前手动设置了b3消息属性,但Spring Integration接收后会将其覆盖,请问如何保留传入的traceId?
我的相关配置如下:
@EnableIntegrationManagement(observationPatterns = "*")
@Bean @GlobalChannelInterceptor(order = Ordered.HIGHEST_PRECEDENCE) public ChannelInterceptor observationPropagationChannelInterceptor(final ObservationRegistry observationRegistry) { return new ObservationPropagationChannelInterceptor(observationRegistry); }
解决方案
问题根源在于ObservationPropagationChannelInterceptor和Spring Integration的默认观测机制:JMS消息驱动适配器会自动创建新的观测上下文,覆盖传入的b3头信息。可以通过以下方式解决:
- 自定义JMS适配器的观测约定,优先使用传入的trace信息
修改JMS消息驱动适配器配置,通过自定义observationConvention从JMS消息属性中提取b3头,构建观测上下文:
from(Jms.messageDrivenChannelAdapter(connectionFactory) .destination(mqInQueue) .observationConvention(new JmsMessageDrivenChannelAdapterObservationConvention() { @Override public KeyValues getLowCardinalityKeyValues(JmsMessageDrivenContext context) { String b3 = context.getMessage().getStringProperty("b3"); if (b3 != null) { String[] b3Parts = b3.split("-"); if (b3Parts.length >= 2) { return KeyValues.of( TraceContext.TRACE_ID_KEY, b3Parts[0], TraceContext.SPAN_ID_KEY, b3Parts[1] ); } } return super.getLowCardinalityKeyValues(context); } }) )
- 调整全局拦截器逻辑或优先级
当前ObservationPropagationChannelInterceptor优先级最高,会自动传播观测上下文覆盖传入信息:
- 降低该拦截器优先级,让自定义trace处理逻辑先执行;
- 或重写拦截器逻辑,在传播前先检查消息中的
b3属性,存在则优先用它构建观测上下文。
- 禁用Spring Integration自动观测,手动处理trace传播
如果不需要自动生成观测信息,可缩小观测范围或直接关闭:
@EnableIntegrationManagement(observationPatterns = "") // 禁用所有自动观测
之后在消息接收后手动解析b3头,设置到当前线程的trace上下文,再传递到Kafka。
注意:不同Spring Boot/Spring Integration版本的观测API可能略有差异,请根据实际版本调整代码。
内容的提问来源于stack exchange,提问作者Roman
相关产品推荐
相关产品推荐

