使用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); }
验证步骤
- 生产者端:发送消息后,检查Kafka消息的CloudEvent扩展字段,确认存在
traceparent和tracestate属性。 - 消费者端:接收消息时,SmallRye会自动读取这些扩展字段恢复追踪上下文,此时消费者的Span会作为生产者Span的子Span,保持同一个Trace ID。
额外说明
- 这种方式符合CloudEvent扩展规范,是跨服务追踪的标准实现方式。
- 确保Quarkus OpenTelemetry扩展已正确配置,Kafka传播支持默认是开启的,无需额外配置。
内容的提问来源于stack exchange,提问作者Weller
相关产品推荐
相关产品推荐

