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

如何在SQS通信中实现Sleuth traceId的自动发送与消费

解决方案:Spring Cloud Sleuth在SQS场景下的TraceId自动透传

可以实现SQS通信过程中Sleuth TraceId的自动收发,你当前的组件版本无需升级,通过扩展现有组件的拦截扩展点即可完成适配,具体实现如下:

实现逻辑

Spring Cloud Sleuth默认没有适配spring-cloud-aws-messaging的SQS链路传播,核心实现思路是利用SQS的消息属性存储B3格式的链路信息:

  • 发送端:发消息前从当前Sleuth上下文中取出traceId、spanId等信息,写入SQS消息属性
  • 消费端:收到消息后从消息属性中读取链路信息,注入到当前线程的Sleuth上下文,完成链路恢复

具体代码实现

1. 发送端链路信息自动注入

先实现SQS请求拦截器,在发送消息前自动写入链路属性:

@Component
public class SleuthSQSSendInterceptor implements RequestHandler {
    private final Tracer tracer;

    public SleuthSQSSendInterceptor(Tracer tracer) {
        this.tracer = tracer;
    }

    @Override
    public AmazonWebServiceRequest beforeExecution(AmazonWebServiceRequest request) {
        if (request instanceof SendMessageRequest) {
            SendMessageRequest sendRequest = (SendMessageRequest) request;
            Span currentSpan = tracer.currentSpan();
            if (currentSpan != null) {
                TraceContext context = currentSpan.context();
                sendRequest.addMessageAttributesEntry("X-B3-TraceId",
                    new MessageAttributeValue().withDataType("String").withStringValue(context.traceId()));
                sendRequest.addMessageAttributesEntry("X-B3-SpanId",
                    new MessageAttributeValue().withDataType("String").withStringValue(context.spanId()));
                sendRequest.addMessageAttributesEntry("X-B3-Sampled",
                    new MessageAttributeValue().withDataType("String").withStringValue(context.sampled() ? "1" : "0"));
            }
        }
        return request;
    }

    @Override
    public void afterResponse(AmazonWebServiceRequest request, Object response, TimingInfo timingInfo) {}

    @Override
    public void afterError(AmazonWebServiceRequest request, Exception e, TimingInfo timingInfo) {}
}

将拦截器注册到SQS客户端:

@Bean
public AmazonSQSAsync amazonSQSAsync(
    AWSCredentialsProvider credentialsProvider,
    RegionProvider regionProvider,
    SleuthSQSSendInterceptor sendInterceptor
) {
    return AmazonSQSAsyncClientBuilder.standard()
        .withCredentials(credentialsProvider)
        .withRegion(regionProvider.getRegion().getName())
        .withRequestHandlers(sendInterceptor)
        .build();
}

2. 消费端链路上下文自动恢复

自定义消息处理拦截器,在执行业务逻辑前恢复链路上下文:

@Component
public class SleuthSQSReceiveEnhancer implements QueueMessageHandlerFactory {
    private final Tracer tracer;
    private final QueueMessageHandlerFactory defaultFactory;

    public SleuthSQSReceiveEnhancer(Tracer tracer, QueueMessageHandlerFactory defaultFactory) {
        this.tracer = tracer;
        this.defaultFactory = defaultFactory;
    }

    @Override
    public QueueMessageHandler createQueueMessageHandler() {
        QueueMessageHandler originalHandler = defaultFactory.createQueueMessageHandler();
        // 扩展参数解析器,在解析参数前注入链路上下文
        originalHandler.setArgumentResolvers(Collections.singletonList(new HandlerMethodArgumentResolver() {
            @Override
            public boolean supportsParameter(MethodParameter parameter) {
                return originalHandler.getArgumentResolvers().stream().anyMatch(r -> r.supportsParameter(parameter));
            }

            @Override
            public Object resolveArgument(MethodParameter parameter, ModelAndViewContainer mavContainer, NativeWebRequest webRequest, WebDataBinderFactory binderFactory) throws Exception {
                Message sqsMessage = (Message) webRequest.getAttribute("message", NativeWebRequest.SCOPE_REQUEST);
                if (sqsMessage != null && sqsMessage.getMessageAttributes() != null) {
                    Map<String, MessageAttributeValue> attrs = sqsMessage.getMessageAttributes();
                    if (attrs.containsKey("X-B3-TraceId")) {
                        TraceContext context = TraceContext.newBuilder()
                            .traceId(attrs.get("X-B3-TraceId").getStringValue())
                            .spanId(attrs.get("X-B3-SpanId").getStringValue())
                            .sampled("1".equals(attrs.get("X-B3-Sampled").getStringValue()))
                            .build();
                        // 开启新的span并绑定到当前线程
                        tracer.nextSpan(context).start().scope();
                    }
                }
                // 走原有参数解析逻辑
                return originalHandler.getArgumentResolvers().stream()
                    .filter(r -> r.supportsParameter(parameter))
                    .findFirst()
                    .orElseThrow(() -> new IllegalArgumentException("No matched argument resolver"))
                    .resolveArgument(parameter, mavContainer, webRequest, binderFactory);
            }
        }));
        return originalHandler;
    }
}

注意事项

  • 消费端处理完消息后可在@SqsListener方法末尾手动调用tracer.currentSpan().finish()关闭span,避免内存泄漏
  • SQS单条消息的消息属性最多支持10个,链路传播仅占用3个属性,不会超出限制
  • 上述代码完全适配你当前使用的spring-cloud-sleuth.version:3.0.4和spring-cloud-aws-messaging:2.2.6.RELEASE版本,无需调整API

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.25 15:45:04