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

如何在Kafka Producer API中发送JSON格式自定义Header

How to Send a JSON Payload as a Kafka Header via Producer API

No problem, let's walk through exactly how to get that JSON payload into a Kafka header when using the producer.send() method. Here's a step-by-step breakdown with practical code examples:

Step 1: Prepare Your JSON Payload

First, you need to convert your JSON object into a byte array (since Kafka headers expect byte values). You have two solid options here:

Option A: Use a POJO + JSON Library (Cleaner for Dynamic Data)

If your JSON values change dynamically, create a simple POJO to match your structure and use a library like Jackson to serialize it to a string. This avoids manual JSON syntax errors:

import com.fasterxml.jackson.databind.ObjectMapper;
import java.nio.charset.StandardCharsets;

// POJO matching your JSON schema
class EventMetadata {
    private String type;
    private int version;
    private String orgId;
    private String userId;
    private String userName;
    private String correlationId;
    private String jobId;

    // Add getters, setters, and a constructor (or use Lombok for brevity)
}

// Initialize and populate your metadata
EventMetadata metadata = new EventMetadata();
metadata.setType("record_created");
metadata.setVersion(1);
metadata.setOrgId("your-org-id"); // Fill in your actual values
metadata.setUserId("your-user-id");
metadata.setUserName("your-username");
metadata.setCorrelationId("your-correlation-id");
metadata.setJobId("your-job-id");

// Serialize to JSON string
ObjectMapper objectMapper = new ObjectMapper();
String jsonHeaderValue = objectMapper.writeValueAsString(metadata);

Option B: Raw JSON String (For Static Values)

If your JSON is static, you can skip the POJO and use a direct string:

import java.nio.charset.StandardCharsets;

String jsonHeaderValue = "{ \"type\": \"record_created\", \"version\": 1, \"orgId\": \"your-org-id\", \"userId\": \"your-user-id\", \"userName\": \"your-username\", \"correlationId\": \"your-correlation-id\", \"jobId\": \"your-job-id\" }";

Step 2: Create a Kafka Header Instance

Kafka's RecordHeader class defines a header with a key and byte array value. Convert your JSON string to bytes using UTF-8 to ensure consistent parsing on the consumer side:

import org.apache.kafka.common.header.Header;
import org.apache.kafka.common.header.internals.RecordHeader;

Header eventMetadataHeader = new RecordHeader(
    "event-metadata", // Choose a descriptive header key
    jsonHeaderValue.getBytes(StandardCharsets.UTF_8)
);

Step 3: Attach Header to Your ProducerRecord

Use the ProducerRecord constructor that accepts an Iterable<Header> to pass your header(s). Wrap the header in a list (you can add multiple headers here if needed):

import org.apache.kafka.clients.producer.ProducerRecord;
import java.util.List;

// Assume you have these variables defined
String topic = "your-target-topic";
Integer messageNo = 1;
String records = "your-main-payload-content";

// Build the ProducerRecord with headers
ProducerRecord<Integer, String> producerRecord = new ProducerRecord<>(
    topic,
    messageNo,
    records,
    List.of(eventMetadataHeader)
);

// Send the record
producer.send(producerRecord);

Quick Tips

  • Error Handling: If using Jackson, wrap the serialization in a try-catch block to handle JsonProcessingException (in case of invalid data).
  • Consumer Side: When reading the header, remember to convert the byte array back to a UTF-8 string before parsing the JSON.
  • Multiple Headers: To send more than one header, just add additional RecordHeader instances to the List.of() call.

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.26 10:21:48