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

使用Smallrye Kafka CloudEvent构建器时OpenTelemetry追踪传播失败求助

使用SmallRye CloudEvent Builder时OpenTelemetry追踪上下文无法传播的解决方法

问题根因

当用OutgoingCloudEventMetadataBuilder构建CloudEvent消息时,单纯把TracingMetadata作为独立元数据附加到Message上是无效的——SmallRye的CloudEvent处理逻辑会优先处理自身的元数据结构,不会自动把独立的追踪元数据注入到CloudEvent的扩展字段里,导致消费者端读不到生产者的追踪上下文,只能生成新的Trace ID。

而直接用Message.of()添加TracingMetadata时,SmallRye Kafka连接器会自动把追踪信息塞进Kafka消息头,所以能正常传播。

修复方案

核心思路是把OpenTelemetry追踪上下文直接嵌入到CloudEvent的扩展属性中,而不是作为独立元数据存在。修改CloudEventBuilder,在构建CloudEvent元数据时手动注入追踪相关的扩展字段。

修改后的CloudEventBuilder代码

import io.opentelemetry.context.Context;
import io.opentelemetry.context.propagation.TextMapPropagator;
import io.smallrye.reactive.messaging.ce.OutgoingCloudEventMetadata;
import io.smallrye.reactive.messaging.ce.OutgoingCloudEventMetadataBuilder;
import java.net.URI;
import java.util.HashMap;
import java.util.Map;
import org.eclipse.microprofile.reactive.messaging.Message;
import org.eclipse.microprofile.reactive.messaging.Metadata;

public class CloudEventBuilder<T> {

  private static final String DATA_CONTENT_TYPE = "application/json";
  private static final TextMapPropagator PROPAGATOR = TextMapPropagator.composite();
  private T payload;
  private Context otelContext;

  private final OutgoingCloudEventMetadataBuilder<T> metadataBuilder =
      OutgoingCloudEventMetadata.builder();

  public CloudEventBuilder() {
    metadataBuilder
        .withSource(URI.create("my-event-source"))
        .withType("MyEvent")
        .withDataContentType(DATA_CONTENT_TYPE);
  }

  public CloudEventBuilder<T> withId(String id) {
    metadataBuilder.withId(id);
    return this;
  }

  public CloudEventBuilder<T> withPayload(T payload) {
    this.payload = payload;
    return this;
  }

  // 直接接收OpenTelemetry上下文,替代原来的TracingMetadata
  public CloudEventBuilder<T> withOtelContext(Context otelContext) {
    this.otelContext = otelContext;
    return this;
  }

  public Message<T> build() {
    // 把OpenTelemetry上下文注入到CloudEvent扩展字段
    if (otelContext != null) {
      Map<String, String> traceHeaders = new HashMap<>();
      PROPAGATOR.inject(otelContext, traceHeaders, Map::put);
      // 将traceparent、tracestate等追踪字段作为CloudEvent扩展添加
      traceHeaders.forEach(metadataBuilder::withExtension);
    }

    OutgoingCloudEventMetadata<T> ceMetadata = metadataBuilder.build();
    return Message.of(payload, Metadata.of(ceMetadata));
  }
}

同步修改send2方法

@WithSpan("myspan2")
public void send2() {
  var myEvent = new MyEvent(45);
  Message<MyEvent> ceEventMessage =
      new CloudEventBuilder<MyEvent>()
          .withPayload(myEvent)
          .withId("12345")
          .withOtelContext(Context.current())
          .build();

  emitter.send(ceEventMessage);
}

验证步骤

  1. 生产者端:发送消息后,检查Kafka消息的CloudEvent扩展字段,确认存在traceparent和tracestate属性。
  2. 消费者端:接收消息时,SmallRye会自动读取这些扩展字段恢复追踪上下文,此时消费者的Span会作为生产者Span的子Span,保持同一个Trace ID。

额外说明

  • 这种方式符合CloudEvent扩展规范,是跨服务追踪的标准实现方式。
  • 确保Quarkus OpenTelemetry扩展已正确配置,Kafka传播支持默认是开启的,无需额外配置。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.20 19:54:38