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

Spring Cloud Sleuth Kafka自定义Baggage无法传播问题咨询

问题:Spring Cloud Sleuth远程Baggage在Kafka消息传递中失效

我有两个基于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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.11 22:05:57