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

如何为Kafka入站消息添加请求作用域UUID并通过Graylog实现追踪?求带详细解释的Kafka示例代码(面向初学者)

Hey there! Since you're new to Apache Kafka, let's break this down step by step with concrete, runnable code examples that fit your exact requirements. We'll build out everything you need to add a request-scoped UUID to Kafka messages for Graylog tracing.

解决方案:为Kafka消息添加请求作用域UUID以实现Graylog追踪

整体思路

We need three core components to make this work:

  • A request-scoped UUID generator: Ensures all Kafka messages in the same request share one unique ID
  • A message wrapper class: Holds the original Kafka message and its associated trace ID
  • Enhanced Kafka producers/listeners: Automatically attach and consume the trace ID without extra boilerplate

1. Request-Scoped UUID Generator

This class creates a single UUID per HTTP request, and returns the same ID every time it's called during that request's lifecycle.

import org.springframework.context.annotation.Scope;
import org.springframework.stereotype.Component;
import java.util.UUID;

// Mark as request-scoped: Spring creates a new instance for every HTTP request
@Component
@Scope("request")
public class RequestTraceIdGenerator {
    private final String traceId;

    // Generate UUID once when the instance is created (per request)
    public RequestTraceIdGenerator() {
        this.traceId = UUID.randomUUID().toString();
    }

    // Return the fixed trace ID for this request
    public String getTraceId() {
        return traceId;
    }
}

Quick Explanation:

  • @Scope("request") guarantees a fresh instance per HTTP request, so each request gets a unique ID
  • The UUID is generated in the constructor, so it never changes during the request's lifetime

2. Kafka Message Wrapper Class

This class wraps your original message content with the trace ID, making it easy to pass both together through Kafka.

public class TracedKafkaMessage<T> {
    // Your original message payload (works with any type: String, POJO, etc.)
    private final T payload;
    // The request-scoped trace ID for Graylog tracing
    private final String traceId;

    public TracedKafkaMessage(T payload, String traceId) {
        this.payload = payload;
        this.traceId = traceId;
    }

    // Getters for accessing the payload and trace ID
    public T getPayload() {
        return payload;
    }

    public String getTraceId() {
        return traceId;
    }

    // Optional: Make logging easier with a readable toString
    @Override
    public String toString() {
        return "TracedKafkaMessage{" +
                "payload=" + payload +
                ", traceId='" + traceId + '\'' +
                '}';
    }
}

Quick Explanation:

  • The generic <T> lets you wrap any type of message without rewriting the class
  • We keep the original payload intact so your business logic doesn't need to change

3. Producer Service: Auto-Attach Trace ID

This service handles sending Kafka messages, automatically wrapping them with the request-scoped trace ID.

import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.stereotype.Service;

@Service
public class TracedKafkaProducer {
    private final KafkaTemplate<String, Object> kafkaTemplate;
    private final RequestTraceIdGenerator traceIdGenerator;

    // Inject KafkaTemplate and our trace ID generator via constructor
    public TracedKafkaProducer(KafkaTemplate<String, Object> kafkaTemplate, RequestTraceIdGenerator traceIdGenerator) {
        this.kafkaTemplate = kafkaTemplate;
        this.traceIdGenerator = traceIdGenerator;
    }

    // Send a message with auto-attached trace ID
    public void sendMessage(String topic, Object payload) {
        String traceId = traceIdGenerator.getTraceId();
        TracedKafkaMessage<Object> wrappedMessage = new TracedKafkaMessage<>(payload, traceId);
        
        kafkaTemplate.send(topic, wrappedMessage);
        
        // Log with your structured logger here (include the trace ID!)
        // Example: structuredLog.info("Sent Kafka message", "traceId", traceId, "topic", topic);
    }
}

4. Listener: Consume and Use Trace ID

For your @KafkaListener, we'll receive the wrapped message directly, then extract the trace ID for your structured logging.

import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.stereotype.Component;

@Component
public class TracedKafkaListener {

    // Listen to your target topic, and receive the wrapped message
    @KafkaListener(topics = "your-target-topic", groupId = "your-consumer-group-id")
    public void handleMessage(TracedKafkaMessage<?> wrappedMessage) {
        String traceId = wrappedMessage.getTraceId();
        Object originalPayload = wrappedMessage.getPayload();

        // Add the trace ID to your structured log here
        // Example: structuredLog.info("Received Kafka message", "traceId", traceId, "payload", originalPayload);

        // Pass the original payload to your business logic
        processBusinessLogic(originalPayload);
    }

    private void processBusinessLogic(Object payload) {
        // Your existing business code goes here
    }
}

Note for Non-Web Consumers:

If your consumer runs as a standalone service (not part of a web app), request scope won't work. Swap @Scope("request") in the trace ID generator with @Scope("thread")—this gives each consumer thread its own trace ID, so all messages processed by the same thread share one ID.


5. Spring Kafka Configuration

Add these properties to your application.properties (or application.yml) to set up Kafka serialization/deserialization for our wrapped messages:

# Kafka Producer Settings
spring.kafka.producer.bootstrap-servers=localhost:9092
spring.kafka.producer.key-serializer=org.apache.kafka.common.serialization.StringSerializer
spring.kafka.producer.value-serializer=org.springframework.kafka.support.serializer.JsonSerializer

# Kafka Consumer Settings
spring.kafka.consumer.bootstrap-servers=localhost:9092
spring.kafka.consumer.group-id=your-consumer-group-id
spring.kafka.consumer.auto-offset-reset=earliest
spring.kafka.consumer.key-deserializer=org.apache.kafka.common.serialization.StringDeserializer
spring.kafka.consumer.value-deserializer=org.springframework.kafka.support.serializer.JsonDeserializer
# Allow deserialization of our wrapper class (use specific packages in production!)
spring.kafka.consumer.properties.spring.json.trusted.packages=*

Quick Explanation:

  • We use JSON serialization/deserialization to easily pass our wrapped messages between producer and consumer
  • spring.json.trusted.packages=* lets Spring deserialize any class (restrict to your app's packages in production for security)

6. Test It Out!

Here's a simple controller to trigger message sending and test the flow:

import org.springframework.web.bind.annotation.PostMapping;
import org.springframework.web.bind.annotation.RequestBody;
import org.springframework.web.bind.annotation.RestController;

@RestController
public class TestMessageController {
    private final TracedKafkaProducer kafkaProducer;

    public TestMessageController(TracedKafkaProducer kafkaProducer) {
        this.kafkaProducer = kafkaProducer;
    }

    @PostMapping("/send-kafka-message")
    public String sendMessage(@RequestBody String messageContent) {
        kafkaProducer.sendMessage("your-target-topic", messageContent);
        return "Message sent with trace ID!";
    }
}

When you call POST /send-kafka-message:

  1. Spring creates a new RequestTraceIdGenerator instance with a unique UUID
  2. The producer wraps your message with this UUID and sends it to Kafka
  3. The listener receives the wrapped message, extracts the UUID, and logs it with your structured logger
  4. You can search Graylog using this UUID to find all logs related to this request

Alternative: Use Kafka Headers (No Message Wrapping)

If you don't want to modify your message structure, you can pass the trace ID in Kafka message headers instead:

Producer Side:

public void sendMessageWithHeader(String topic, Object payload) {
    String traceId = traceIdGenerator.getTraceId();
    kafkaTemplate.send(topic, payload)
            .addHeader("traceId", traceId.getBytes());
}

Listener Side:

@KafkaListener(topics = "your-target-topic")
public void handleMessageWithHeader(ConsumerRecord<String, Object> record) {
    String traceId = new String(record.headers().lastHeader("traceId").value());
    Object payload = record.value();
    
    // Log with trace ID and process payload
}

This is lighter weight if you want to keep your original message format intact.

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.29 19:02:47