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

