如何在Kafka Producer API中发送JSON格式自定义Header
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
RecordHeaderinstances to theList.of()call.
内容的提问来源于stack exchange,提问作者user5702955

