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

如何在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();
}

四、验证与排查

  1. 确保Span导出到Jaeger/Zipkin等追踪系统,查看链路是否包含:
    • 网关→inbox-service发送消息的Span
    • 消费者消费消息的子Span
    • 调用Microsoft API的子Span
    • 数据库存储的子Span
      所有Span应共享同一个traceId,形成完整调用链。
  2. 若自动instrumentation失效,检查依赖版本是否匹配、runtime scope的instrumentation依赖是否被正确加载。
  3. 响应式场景下禁止用ThreadLocal存储Span,必须通过Reactor Context传递上下文。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.24 03:27:39