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

如何为Spring Kafka消息添加自定义Header值?

Adding Custom Headers to Spring Kafka Messages (1.1.6 Version)

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 String headers 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 KafkaTemplate works out of the box for this—no extra config needed unless you have custom producer/consumer properties.

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.19 09:45:03