如何在RabbitMQ监听器中启动新的追踪Span?
解决方案:RabbitMQ消费者链路追踪实现(响应式场景)
一、依赖选型与基础配置
优先采用OpenTelemetry实现链路追踪,它对Reactor Netty这类响应式框架支持更完善,且兼容W3C TraceContext标准,无需手动传递traceId/spanId。
1. 添加核心依赖(Maven)
<!-- OpenTelemetry 核心API与SDK --> <dependency> <groupId>io.opentelemetry</groupId> <artifactId>opentelemetry-api</artifactId> <version>1.32.0</version> </dependency> <dependency> <groupId>io.opentelemetry</groupId> <artifactId>opentelemetry-sdk</artifactId> <version>1.32.0</version> </dependency> <!-- RabbitMQ 自动链路注入 --> <dependency> <groupId>io.opentelemetry.instrumentation</groupId> <artifactId>opentelemetry-amqp-client-5.1</artifactId> <version>1.32.0</version> <scope>runtime</scope> </dependency> <!-- Reactor Netty 自动链路适配 --> <dependency> <groupId>io.opentelemetry.instrumentation</groupId> <artifactId>opentelemetry-reactor-netty-1.1</artifactId> <version>1.32.0</version> <scope>runtime</scope> </dependency> <!-- Spring Boot 自动配置(可选,Spring生态下推荐) --> <dependency> <groupId>io.opentelemetry.instrumentation</groupId> <artifactId>opentelemetry-spring-boot-starter</artifactId> <version>1.32.0</version> </dependency>
2. 初始化OpenTelemetry
创建配置类完成追踪系统的基础配置(以导出到OTLP为例):
@Configuration public class OtelConfig { @Bean public OpenTelemetry openTelemetry() { return OpenTelemetrySdk.builder() .setTracerProvider(SdkTracerProvider.builder() .addSpanProcessor(BatchSpanProcessor.builder(OtlpGrpcSpanExporter.builder().build()).build()) .build()) .setPropagators(ContextPropagators.create(W3CTraceContextPropagator.getInstance())) .buildAndRegisterGlobal(); } }
二、消息发送端:自动传递追踪上下文
无需手动往消息头写入traceId/spanId,OpenTelemetry的RabbitMQ instrumentation会自动将当前链路上下文注入到消息的traceparent标准头中。响应式场景下(如ReactorRabbitMQ),上下文会随Reactor Context自动传递:
@Autowired private Sender reactorRabbitSender; public Mono<Void> sendToQueue(String payload) { OutboundMessage message = OutboundMessage.create("exchange", "routing-key", payload.getBytes()); return reactorRabbitSender.send(Mono.just(message)); }
三、消费者端:链路延续与全链路追踪
1. 自动链路追踪(推荐)
OpenTelemetry的消费者instrumentation会自动从消息头提取traceparent,创建与发送端链路关联的子Span,覆盖从消息消费到处理完成的全流程。后续调用Microsoft API、数据库操作等,只要对应组件有OpenTelemetry instrumentation支持,都会自动生成子Span,形成完整链路。
示例代码(ReactorRabbitMQ消费者):
@Autowired private Receiver reactorRabbitReceiver; @Autowired private WebClient microsoftGraphClient; @Autowired private MailRepository mailRepo; public void startConsumer() { reactorRabbitReceiver.consumeAutoAck("queue-name") .flatMap(message -> { String messageId = extractMessageIdFromPayload(new String(message.getBody())); // 调用Microsoft Graph API return microsoftGraphClient.get() .uri("https://graph.microsoft.com/v1.0/me/messages/{id}", messageId) .retrieve() .bodyToMono(MailContent.class) // 存储到数据库 .flatMap(mailContent -> mailRepo.save(mailContent)); }) .subscribe(); }
2. 手动链路延续(自动方式失效时)
如果自动instrumentation不生效,可手动解析消息头的追踪上下文,创建关联Span:
发送端手动注入上下文
public Mono<Void> sendToQueue(String payload) { // 提取当前链路上下文 Map<String, String> headers = new HashMap<>(); W3CTraceContextPropagator.getInstance().inject(Context.current(), headers, (carrier, key, value) -> carrier.put(key, value)); OutboundMessage message = OutboundMessage.create("exchange", "routing-key", payload.getBytes()); message.getProperties().getHeaders().putAll(headers); return reactorRabbitSender.send(Mono.just(message)); }
消费者端手动延续链路
public void startConsumer() { reactorRabbitReceiver.consumeAutoAck("queue-name") .flatMap(message -> { // 从消息头提取追踪上下文 Context traceContext = W3CTraceContextPropagator.getInstance() .extract(Context.current(), message.getProperties().getHeaders(), (carrier, key) -> carrier.get(key)); return Mono.deferContextual(ctx -> { return Context.of(ctx).putAll(traceContext) .call(() -> { // 创建关联子Span Span consumeSpan = OpenTelemetry.getGlobalTracer("inbox-service-consumer") .spanBuilder("consume-rabbitmq-message") .setParent(traceContext) .startSpan(); try (Scope scope = consumeSpan.makeCurrent()) { // 业务逻辑:调用API+存储数据库 String messageId = extractMessageIdFromPayload(new String(message.getBody())); return microsoftGraphClient.get() .uri("https://graph.microsoft.com/v1.0/me/messages/{id}", messageId) .retrieve() .bodyToMono(MailContent.class) .flatMap(mailRepo::save) .doOnSuccess(v -> consumeSpan.end()) .doOnError(e -> { consumeSpan.recordException(e); consumeSpan.end(); }); } }); }); }) .subscribe(); }
四、验证与排查
- 确保Span导出到Jaeger/Zipkin等追踪系统,查看链路是否包含:
- 网关→inbox-service发送消息的Span
- 消费者消费消息的子Span
- 调用Microsoft API的子Span
- 数据库存储的子Span
所有Span应共享同一个traceId,形成完整调用链。
- 若自动instrumentation失效,检查依赖版本是否匹配、runtime scope的instrumentation依赖是否被正确加载。
- 响应式场景下禁止用ThreadLocal存储Span,必须通过Reactor Context传递上下文。
内容的提问来源于stack exchange,提问作者Shrey Soni
相关产品推荐
相关产品推荐

