如何为Spring Kafka消息添加自定义Header值?
Hey there! Since you're working with Spring Kafka 1.1.6 (a solid older version that still has all the tools you need for custom headers), let's break down exactly how to add fileName and fileId headers to your messages, and how to read them on the consumer side.
Producer: Attaching Custom Headers
You have two straightforward ways to add custom headers when sending messages—using Spring's Message abstraction (the most Spring-like approach) or working directly with Kafka's ProducerRecord.
Approach 1: Using Spring MessageBuilder
This is the cleanest method if you're leveraging Spring's messaging framework:
import org.springframework.kafka.core.KafkaTemplate; import org.springframework.messaging.Message; import org.springframework.messaging.support.MessageBuilder; import org.springframework.stereotype.Component; @Component public class CustomHeaderProducer { private final KafkaTemplate<String, String> kafkaTemplate; // Constructor injection works with Spring Boot's auto-configuration public CustomHeaderProducer(KafkaTemplate<String, String> kafkaTemplate) { this.kafkaTemplate = kafkaTemplate; } public void sendMessageWithHeaders(String topic, String payload, String fileName, String fileId) { // Build a message with your payload and custom headers Message<String> message = MessageBuilder .withPayload(payload) .setHeader("fileName", fileName) .setHeader("fileId", fileId) .build(); // Send the message to your Kafka topic kafkaTemplate.send(topic, message); } }
Approach 2: Using Kafka's ProducerRecord (Low-Level API)
If you prefer working directly with Kafka's native API, you can construct a ProducerRecord with custom headers:
import org.apache.kafka.clients.producer.ProducerRecord; import org.apache.kafka.common.header.Header; import org.apache.kafka.common.header.internals.RecordHeader; import org.springframework.kafka.core.KafkaTemplate; import org.springframework.stereotype.Component; import java.util.Arrays; import java.util.List; @Component public class RawProducerRecordProducer { private final KafkaTemplate<String, String> kafkaTemplate; public RawProducerRecordProducer(KafkaTemplate<String, String> kafkaTemplate) { this.kafkaTemplate = kafkaTemplate; } public void sendWithRawHeaders(String topic, String payload, String fileName, String fileId) { // Create header objects (Kafka stores headers as byte arrays) Header fileNameHeader = new RecordHeader("fileName", fileName.getBytes()); Header fileIdHeader = new RecordHeader("fileId", fileId.getBytes()); List<Header> headers = Arrays.asList(fileNameHeader, fileIdHeader); // Build the ProducerRecord with payload and headers ProducerRecord<String, String> record = new ProducerRecord<>( topic, null, // Let Kafka assign partition automatically null, // Optional message key payload, headers ); kafkaTemplate.send(record); } }
Consumer: Reading Custom Headers
Reading the headers is just as simple—you can either extract them directly with annotations or access the full Message object.
Option 1: Use @Header Annotations (Cleanest for Specific Headers)
Annotate your listener method parameters to pull out individual headers:
import org.springframework.kafka.annotation.KafkaListener; import org.springframework.messaging.handler.annotation.Header; import org.springframework.stereotype.Component; @Component public class CustomHeaderConsumer { @KafkaListener(topics = "${your.kafka.topic.name}") public void consumeMessage(String payload, @Header("fileName") String fileName, @Header("fileId") String fileId) { // Process your message and headers here System.out.println("Received payload: " + payload); System.out.println("File Name: " + fileName); System.out.println("File ID: " + fileId); } }
Option 2: Receive the Full Message Object
If you need access to all headers or additional message metadata, accept the full Message in your listener:
import org.springframework.kafka.annotation.KafkaListener; import org.springframework.messaging.Message; import org.springframework.stereotype.Component; @Component public class FullMessageConsumer { @KafkaListener(topics = "${your.kafka.topic.name}") public void consumeFullMessage(Message<String> message) { String payload = message.getPayload(); // Retrieve headers with type safety String fileName = message.getHeaders().get("fileName", String.class); String fileId = message.getHeaders().get("fileId", String.class); System.out.println("Payload: " + payload); System.out.println("File Name from Header: " + fileName); System.out.println("File ID from Header: " + fileId); } }
Quick Notes for Spring Kafka 1.1.6
- Spring will automatically handle converting
Stringheaders to/from byte arrays, so you don't need to manually encode/decode unless you're using non-string types. - Make sure your header names match exactly between producer and consumer (case-sensitive!).
- Spring Boot's auto-configuration for
KafkaTemplateworks out of the box for this—no extra config needed unless you have custom producer/consumer properties.
内容的提问来源于stack exchange,提问作者BlueStar

