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

OpenTelemetry Baggage在Quarkus微服务栈中无法传播的问题排查

问题描述

在Quarkus 2.16.5.Final微服务栈中,向OpenTelemetry追踪(trace)添加Baggage后,Baggage无法跨服务传播,但traceId与spanId传播正常。具体场景与配置如下:

Service A 代码(发送Kafka消息)

try (var scope =
        Context.current()
            .with(
                Baggage.builder()
                    .put("RequestId", String.valueOf(message.getRequestId()))
                    .build())
            .makeCurrent()) {
            //......业务逻辑
            emitter.send(KafkaRecord.of(key, message));
           }

Service B 代码(接收Kafka消息并处理)

@ApplicationScoped
public class MDCEnricher implements OutgoingInterceptor {
  @Override
  public Message<?> onMessage(Message<?> message) {
    Baggage.current()
        .forEach(
            (s, baggageEntry) -> {
              MDC.put(s, baggageEntry.getValue());
            });

    return message;
  }
}

Service C 代码(接收REST请求)

@Slf4j
public class MDCFilter {
  @ServerRequestFilter
  public Uni<RestResponse<Void>> copyBaggageIntoMDC(UriInfo uriInfo, HttpHeaders httpHeaders, ContainerRequestContext requestContext){
    SpanContext spanContext = Span.current().getSpanContext();
    Baggage.current()
        .forEach(
            (s, baggageEntry) -> {
              MDC.put(s, baggageEntry.getValue());
            });
    return Uni.createFrom().item((RestResponse<Void>) null);
  }
}

所有服务的Quarkus日志格式配置

quarkus.log.console.format=%d{HH:mm:ss} %-5p [%X] [%c{2.}] (%t) %s%e%n

预期所有应用日志中能显示Baggage中的RequestId,但实际仅输出traceId和spanId,需排查配置问题。


问题分析与解决方案

核心问题点

  1. Kafka消息未携带Baggage元数据:Service A仅在本地Context设置Baggage,但Quarkus的Kafka emitter默认不会自动将Baggage注入消息Headers,导致Service B无法获取到Baggage。
  2. Service B拦截器类型错误:使用OutgoingInterceptor是处理向外发送的消息,而Service B是接收Kafka消息,应使用IncomingInterceptor处理进入的消息,才能正确读取Baggage并注入MDC。
  3. REST调用未配置Baggage自动传播:Quarkus需要额外配置,才能让OpenTelemetry通过HTTP Headers自动传播Baggage到下游服务。

具体修复步骤

1. Service A:修改Kafka发送逻辑,注入Baggage到消息Headers

手动将当前Context中的Baggage转为Kafka Headers,确保消息携带Baggage信息:

try (var scope = Context.current()
        .with(Baggage.builder().put("RequestId", String.valueOf(message.getRequestId())).build())
        .makeCurrent()) {
    // 将Baggage转为Kafka Headers
    Headers headers = new RecordHeaders();
    Baggage.current().forEach((key, entry) -> {
        headers.add("baggage-" + key, entry.getValue().getBytes(StandardCharsets.UTF_8));
    });
    // 发送携带Headers的消息
    emitter.send(KafkaRecord.of(key, message).withHeaders(headers));
}

2. Service B:替换拦截器类型为IncomingInterceptor,从Headers恢复Baggage

@ApplicationScoped
public class MDCEnricher implements IncomingInterceptor {
    @Override
    public Message<?> onMessage(Message<?> message) {
        // 从Kafka Headers读取Baggage并恢复到当前Context
        Baggage.Builder baggageBuilder = Baggage.builder();
        message.getHeaders().forEach(header -> {
            if (header.getKey().startsWith("baggage-")) {
                String key = header.getKey().substring("baggage-".length());
                String value = new String(header.getValue(), StandardCharsets.UTF_8);
                baggageBuilder.put(key, value);
            }
        });
        Baggage baggage = baggageBuilder.build();
        
        // 将Baggage注入当前Context和MDC
        try (var scope = Context.current().with(baggage).makeCurrent()) {
            baggage.forEach((s, baggageEntry) -> MDC.put(s, baggageEntry.getValue()));
            return message;
        }
    }
}

3. 配置Quarkus自动传播Baggage到REST调用

在所有服务的application.properties中添加以下配置:

# 启用Baggage传播
quarkus.opentelemetry.baggage.propagation.enabled=true
# 指定要传播的Baggage键
quarkus.opentelemetry.baggage.propagation.include=RequestId

4. Service C:调整过滤器执行优先级(可选)

为确保过滤器在其他逻辑前执行,添加@Priority注解:

@Slf4j
@Priority(Priorities.AUTHENTICATION)
public class MDCFilter {
  @ServerRequestFilter
  public Uni<RestResponse<Void>> copyBaggageIntoMDC(UriInfo uriInfo, HttpHeaders httpHeaders, ContainerRequestContext requestContext){
    Baggage.current()
        .forEach((s, baggageEntry) -> MDC.put(s, baggageEntry.getValue()));
    return Uni.createFrom().item((RestResponse<Void>) null);
  }
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.16 13:32:43