Working with Message Serialization
1. Understanding Serialization Concept
| Concept | Description | Detail |
|---|---|---|
| Serializer | Object → byte[] for the wire | Producer side |
| Deserializer | byte[] → object | Consumer side |
| Serde | Paired (de)serializer (Streams) | Symmetric |
Example: Serializer interface
public interface Serializer<T> {
byte[] serialize(String topic, T data);
}
2. Using String Serializer
| Aspect | Description | Detail |
|---|---|---|
| StringSerializer | UTF-8 encode by default | Most common |
| Encoding | serializer.encoding |
Override charset |
| Pair | StringDeserializer on consume | Match both |
Example: String serializer
props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
3. Using Byte Array Serializer
| Aspect | Description | Detail |
|---|---|---|
| ByteArraySerializer | Pass-through raw bytes | No transform |
| Use Case | Pre-encoded payloads | Images, custom |
| Control | App handles encoding | Full control |
Example: Byte array producer
props.put("value.serializer", "org.apache.kafka.common.serialization.ByteArraySerializer");
producer.send(new ProducerRecord<>("blobs", key, imageBytes));
4. Using Integer Serializer
| Aspect | Description | Detail |
|---|---|---|
| IntegerSerializer | 4-byte big-endian int | Fixed width |
| Key Use | Numeric IDs as keys | Compact |
| Pair | IntegerDeserializer | Match |
Example: Integer key
props.put("key.serializer", "org.apache.kafka.common.serialization.IntegerSerializer");
producer.send(new ProducerRecord<>("metrics", 42, payload));
5. Using Long Serializer
| Aspect | Description | Detail |
|---|---|---|
| LongSerializer | 8-byte big-endian long | Timestamps, IDs |
| Use Case | Epoch millis, sequence | Common |
| Pair | LongDeserializer | Match |
Example: Long timestamp value
props.put("value.serializer", "org.apache.kafka.common.serialization.LongSerializer");
producer.send(new ProducerRecord<>("heartbeats", host, System.currentTimeMillis()));
6. Using JSON Serialization
| Aspect | Description | Detail |
|---|---|---|
| Approach | Object → JSON string → bytes | Jackson |
| Pros | Human-readable, flexible | Easy debug |
| Cons | No enforced schema, larger | Verbose |
Example: Custom JSON serializer
public class JsonSerializer<T> implements Serializer<T> {
private final ObjectMapper mapper = new ObjectMapper();
public byte[] serialize(String topic, T data) {
try { return mapper.writeValueAsBytes(data); }
catch (Exception e) { throw new SerializationException(e); }
}
}
7. Using Avro Serialization
| Aspect | Description | Detail |
|---|---|---|
| Avro | Compact binary + schema | Schema Registry |
| Schema ID | Stored in 5-byte prefix | Magic byte + id |
| Evolution | Backward/forward compatible | Versioned |
Example: Avro producer
props.put("value.serializer", "io.confluent.kafka.serializers.KafkaAvroSerializer");
props.put("schema.registry.url", "http://localhost:8081");
producer.send(new ProducerRecord<>("users", user.getId(), avroUser));
8. Using Protobuf Serialization
| Aspect | Description | Detail |
|---|---|---|
| Protobuf | Google binary format | Strong typing |
| Serializer | KafkaProtobufSerializer |
Registry-backed |
| Performance | Fast, very compact | gRPC-friendly |
Example: Protobuf serializer
props.put("value.serializer",
"io.confluent.kafka.serializers.protobuf.KafkaProtobufSerializer");
props.put("schema.registry.url", "http://localhost:8081");
9. Creating Custom Serializers
| Method | Purpose | Detail |
|---|---|---|
| configure() | Read settings | Optional |
| serialize() | Produce byte[] | Required |
| close() | Release resources | Optional |
Example: Register custom serializer
props.put("value.serializer", "com.example.OrderSerializer");
// OrderSerializer implements Serializer<Order>
10. Handling Serialization Errors
| Aspect | Description | Detail |
|---|---|---|
| SerializationException | Thrown on bad data | Producer/consumer |
| Poison Pill | Undeserializable record blocks consumer | Stuck offset |
| ErrorHandlingDeserializer | Wraps and routes failures | Spring/DLQ |
Example: Safe deserialization wrapper
try {
Order o = deserializer.deserialize(topic, bytes);
} catch (SerializationException e) {
sendToDlq(record, e); // skip poison pill
}
11. Configuring Schema Registry
| Config | Description | Detail |
|---|---|---|
| schema.registry.url | Registry endpoint | http://host:8081 |
| auto.register.schemas | Register on first send | Disable in prod |
| basic.auth. credentials.source |
Registry auth | USER_INFO |