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

如何用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并正确关联上下文:

核心思路

  1. 创建驻留时长Span:用record.timestamp()作为Span的开始时间,消费者获取消息的当前时间作为结束时间。
  2. 关联生产者追踪上下文:将驻留Span通过Link关联到生产者的Trace Context(从traceparent头提取)。
  3. 创建消费者处理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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.12 18:54:52