Spring Cloud Sleuth Kafka自定义Baggage无法传播问题咨询
我有两个基于Spring Boot 2.7.2的微服务,使用Kafka作为通信broker,借助Spring Cloud Sleuth实现链路追踪,目前链路追踪的连续性在Kafka生产者与消费者间正常工作。但我尝试通过添加远程Baggage(如correlation-id)传递额外信息时失败。
服务A(生产者)示例代码
Span newSpan = sleuthTracer.nextSpan().start(); try (SpanInScope spanInScope = sleuthTracer.withSpan(newSpan);) { sleuthTracer.createBaggage("correlation-id").set(UUID.randomUUID().toString()); String key = product.getData().getMarketplaceGroupId(); var listenableFuture = kafkaTemplate.send(topic, key, message); futures.add(listenableFuture.completable()); processThePublishingResults( listenableFuture, topic, key, resource); } finally { newSpan.end(); }
服务B(消费者)示例代码
String correlationIdBaggage = sleuthTracer.getBaggage("correlation-id").get(); log.info("Received Sleuth baggage with correlation-id = " + correlationIdBaggage);
两个服务的application.yaml配置
spring: sleuth: traceId128: true baggage: remote-fields: - correlation-id tag-fields: - correlation-id correlation-enabled: true correlation-fields: - correlation-id
请问该功能是否仅支持REST通信,而不适用于Kafka?
Sleuth的Baggage功能完全支持Kafka,并非仅适用于REST通信。你的问题出在Baggage的绑定方式和代码细节上,以下是具体问题和解决办法:
Baggage未正确绑定到当前Span上下文:
你当前的代码只是创建了Baggage但未将其绑定到当前活跃的Span上下文中,导致Kafka消息拦截器无法读取到该Baggage并传递到消息头中。需要使用makeCurrent()将Baggage绑定到Span上下文,修改服务A代码如下:Span newSpan = sleuthTracer.nextSpan().start(); try (SpanInScope spanInScope = sleuthTracer.withSpan(newSpan); BaggageInScope baggageInScope = sleuthTracer.createBaggage("correlation-id", UUID.randomUUID().toString()).makeCurrent()) { String key = product.getData().getMarketplaceGroupId(); var listenableFuture = kafkaTemplate.send(topic, key, message); futures.add(listenableFuture.completable()); processThePublishingResults( listenableFuture, topic, key, resource); } finally { newSpan.end(); }消费者端需安全获取Baggage:
直接调用get()可能会因Baggage不存在抛出异常,建议先判断是否存在:Optional<String> correlationId = sleuthTracer.getBaggage("correlation-id").value(); if (correlationId.isPresent()) { log.info("Received Sleuth baggage with correlation-id = " + correlationId.get()); } else { log.warn("correlation-id baggage not found in current span"); }确认配置一致性:
确保两个服务的spring.sleuth.baggage.remote-fields配置完全一致,这样消费者端才能正确解析Kafka消息头中的Baggage信息。验证Kafka消息头:
可以检查生产者发送的Kafka消息头,确认是否包含baggage头(格式为baggage: correlation-id=xxx),如果不存在,说明生产者端的Baggage注入逻辑仍有问题。
总结:只要正确绑定Baggage到Span上下文并配置正确,Sleuth就能通过Kafka实现远程Baggage的传递。
内容的提问来源于stack exchange,提问作者Ahmed Ali

