如何用Spring Kafka/OpenTelemetry创建Kafka驻留时长的额外追踪Span
如何为Kafka消息驻留时长创建独立的OpenTelemetry Span?
我正在开发一个基于Spring Boot的Kafka消费者,使用OpenTelemetry实现链路追踪。目前只有一个覆盖消息全处理流程的Span,但我需要新增一个独立的专用Span,用来统计消息被消费前在Kafka中的驻留时长。
当前实现代码
package org.example; import org.springframework.kafka.annotation.KafkaListener; import org.springframework.stereotype.Service; import org.apache.kafka.clients.consumer.ConsumerRecord; import io.opentelemetry.api.trace.Span; import io.opentelemetry.api.trace.Tracer; import io.opentelemetry.api.trace.StatusCode; import io.opentelemetry.context.Context; import io.opentelemetry.context.propagation.TextMapGetter; import io.opentelemetry.context.propagation.TextMapPropagator; import io.opentelemetry.api.OpenTelemetry; import org.springframework.beans.factory.annotation.Autowired; import java.time.Instant; import java.util.Map; import java.util.HashMap; @Service public class ConsumerService { @Autowired private Tracer tracer; @Autowired private OpenTelemetry openTelemetry; @KafkaListener(topics = "sample-topic", groupId = "group-2") public void consumeLinksWithRecord(ConsumerRecord<String, String> record) { String word = record.value(); String traceparent = record.headers().lastHeader("traceparent") != null ? new String(record.headers().lastHeader("traceparent").value()) : "No traceparentid header found"; // Create a span for Kafka message processing with ConsumerRecord Span span = tracer.spanBuilder("inside-kafka") .setAttribute("service.name", "insidekafkaservice") .setAttribute("kafka.topic", record.topic()) .setAttribute("kafka.partition", record.partition()) .setAttribute("kafka.offset", record.offset()) .setAttribute("kafka.key", record.key() != null ? record.key() : "null") .setAttribute("kafka.message", word) .setAttribute("kafka.traceparent", traceparent) .setStartTimestamp(Instant.ofEpochMilli(record.timestamp())) .startSpan(); try (var scope = span.makeCurrent()) { extractTraceContext(traceparent, span); System.out.println("Received Message: " + word + " from partition: " + record.partition()); System.out.println("Trace Parent ID: " + traceparent); // Add custom attributes to the span span.setAttribute("message.length", word.length()); span.setAttribute("header.count", record.headers().spliterator().getExactSizeIfKnown()); try { Thread.sleep(80); } catch (InterruptedException e) { Thread.currentThread().interrupt(); span.setStatus(StatusCode.ERROR, "Processing interrupted"); return; } span.setStatus(StatusCode.OK); } catch (Exception e) { span.setStatus(StatusCode.ERROR, e.getMessage()); span.recordException(e); throw e; } finally { span.end(); } } private void extractTraceContext(String traceparentId, Span span) { try { // Create a carrier map with the traceparentid Map<String, String> carrier = new HashMap<>(); carrier.put("traceparent", traceparentId); // Extract the context using the W3C trace context propagator TextMapPropagator propagator = openTelemetry.getPropagators().getTextMapPropagator(); Context extractedContext = propagator.extract(Context.current(), carrier, new TextMapGetter<Map<String, String>>() { @Override public String get(Map<String, String> carrier, String key) { return carrier.get(key); } @Override public Iterable<String> keys(Map<String, String> carrier) { return carrier.keySet(); } }); // Link the extracted context to the current span if (extractedContext != Context.current()) { span.addLink(Span.fromContext(extractedContext).getSpanContext()); } } catch (Exception e) { // Log the error but don't fail the processing System.err.println("Failed to extract trace context: " + e.getMessage()); } } }
时间线示例
- 生产者于00:00开始业务逻辑
- 生产者完成逻辑并将消息存入Kafka的时间为00:01
- 消费者获取消息的时间为00:04
- 消费者完成消息业务逻辑的时间为00:05
消息在Kafka中的驻留时长为00:01到00:04
追踪效果对比
- 期望效果:包含两个关联的Span,一个代表消息在Kafka中的驻留时长(从写入到被消费),另一个代表消费者处理消息的时长,两者属于同一条追踪链路。
- 实际效果:仅显示一个覆盖从消息写入到处理完成全流程的Span,无法区分Kafka驻留和消费处理的时间阶段。
解决方案
要实现独立的Kafka驻留时长Span,需要分别创建两个Span并正确关联上下文:
核心思路
- 创建驻留时长Span:用
record.timestamp()作为Span的开始时间,消费者获取消息的当前时间作为结束时间。 - 关联生产者追踪上下文:将驻留Span通过Link关联到生产者的Trace Context(从
traceparent头提取)。 - 创建消费者处理Span:以驻留Span的上下文作为父上下文,开始时间设为消费者获取消息的当前时间,结束时间为处理完成时间。
修改后的代码实现
package org.example; import org.springframework.kafka.annotation.KafkaListener; import org.springframework.stereotype.Service; import org.apache.kafka.clients.consumer.ConsumerRecord; import io.opentelemetry.api.trace.Span; import io.opentelemetry.api.trace.Tracer; import io.opentelemetry.api.trace.StatusCode; import io.opentelemetry.context.Context; import io.opentelemetry.context.propagation.TextMapGetter; import io.opentelemetry.context.propagation.TextMapPropagator; import io.opentelemetry.api.OpenTelemetry; import org.springframework.beans.factory.annotation.Autowired; import java.time.Instant; import java.util.Map; import java.util.HashMap; @Service public class ConsumerService { @Autowired private Tracer tracer; @Autowired private OpenTelemetry openTelemetry; @KafkaListener(topics = "sample-topic", groupId = "group-2") public void consumeLinksWithRecord(ConsumerRecord<String, String> record) { String word = record.value(); Instant messageWriteTime = Instant.ofEpochMilli(record.timestamp()); Instant consumerReceiveTime = Instant.now(); String traceparent = record.headers().lastHeader("traceparent") != null ? new String(record.headers().lastHeader("traceparent").value()) : "No traceparentid header found"; // 1. 创建Kafka驻留时长Span Span kafkaResidenceSpan = tracer.spanBuilder("kafka-message-residence") .setAttribute("service.name", "kafka-residence-service") .setAttribute("kafka.topic", record.topic()) .setAttribute("kafka.partition", record.partition()) .setAttribute("kafka.offset", record.offset()) .setStartTimestamp(messageWriteTime) .startSpan(); // 关联生产者的追踪上下文 extractTraceContext(traceparent, kafkaResidenceSpan); // 立即结束驻留Span,因为驻留时间到消费者接收消息时已结束 kafkaResidenceSpan.end(consumerReceiveTime); // 2. 创建消费者处理Span,以驻留Span的上下文作为父上下文 Span processingSpan = tracer.spanBuilder("kafka-message-processing") .setAttribute("service.name", "insidekafkaservice") .setAttribute("kafka.topic", record.topic()) .setAttribute("kafka.partition", record.partition()) .setAttribute("kafka.offset", record.offset()) .setAttribute("kafka.key", record.key() != null ? record.key() : "null") .setAttribute("kafka.message", word) .setAttribute("kafka.traceparent", traceparent) .setParent(Context.current().with(kafkaResidenceSpan)) // 设置父Span .setStartTimestamp(consumerReceiveTime) .startSpan(); try (var scope = processingSpan.makeCurrent()) { System.out.println("Received Message: " + word + " from partition: " + record.partition()); System.out.println("Trace Parent ID: " + traceparent); // 添加自定义属性 processingSpan.setAttribute("message.length", word.length()); processingSpan.setAttribute("header.count", record.headers().spliterator().getExactSizeIfKnown()); try { Thread.sleep(80); } catch (InterruptedException e) { Thread.currentThread().interrupt(); processingSpan.setStatus(StatusCode.ERROR, "Processing interrupted"); return; } processingSpan.setStatus(StatusCode.OK); } catch (Exception e) { processingSpan.setStatus(StatusCode.ERROR, e.getMessage()); processingSpan.recordException(e); throw e; } finally { processingSpan.end(); } } private void extractTraceContext(String traceparentId, Span span) { try { if ("No traceparentid header found".equals(traceparentId)) { return; } Map<String, String> carrier = new HashMap<>(); carrier.put("traceparent", traceparentId); TextMapPropagator propagator = openTelemetry.getPropagators().getTextMapPropagator(); Context extractedContext = propagator.extract(Context.current(), carrier, new TextMapGetter<Map<String, String>>() { @Override public String get(Map<String, String> carrier, String key) { return carrier.get(key); } @Override public Iterable<String> keys(Map<String, String> carrier) { return carrier.keySet(); } }); if (!extractedContext.equals(Context.current())) { span.addLink(Span.fromContext(extractedContext).getSpanContext()); } } catch (Exception e) { System.err.println("Failed to extract trace context: " + e.getMessage()); } } }
关键说明
- 驻留Span:
kafka-message-residence的开始时间是消息写入Kafka的时间,结束时间是消费者接收消息的时间,直接结束该Span即可统计驻留时长。 - 处理Span:
kafka-message-processing以驻留Span为父Span,确保链路层级正确,开始时间为消费者接收消息的时间,结束时间为处理完成时间。 - 上下文关联:通过
extractTraceContext方法将驻留Span与生产者的追踪上下文关联,保证整条链路的连贯性。
内容的提问来源于stack exchange,提问作者PatPanda
相关产品推荐
相关产品推荐

