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,需排查配置问题。
问题分析与解决方案
核心问题点
- Kafka消息未携带Baggage元数据:Service A仅在本地Context设置Baggage,但Quarkus的Kafka emitter默认不会自动将Baggage注入消息Headers,导致Service B无法获取到Baggage。
- Service B拦截器类型错误:使用
OutgoingInterceptor是处理向外发送的消息,而Service B是接收Kafka消息,应使用IncomingInterceptor处理进入的消息,才能正确读取Baggage并注入MDC。 - 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
相关产品推荐
相关产品推荐

