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

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头信息。可以通过以下方式解决:

  1. 自定义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);
            }
        })
)
  1. 调整全局拦截器逻辑或优先级
    当前ObservationPropagationChannelInterceptor优先级最高,会自动传播观测上下文覆盖传入信息:
  • 降低该拦截器优先级,让自定义trace处理逻辑先执行;
  • 或重写拦截器逻辑,在传播前先检查消息中的b3属性,存在则优先用它构建观测上下文。
  1. 禁用Spring Integration自动观测,手动处理trace传播
    如果不需要自动生成观测信息,可缩小观测范围或直接关闭:
@EnableIntegrationManagement(observationPatterns = "") // 禁用所有自动观测

之后在消息接收后手动解析b3头,设置到当前线程的trace上下文,再传递到Kafka。

注意:不同Spring Boot/Spring Integration版本的观测API可能略有差异,请根据实际版本调整代码。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.11 00:33:40