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

Example: Registry client config

props.put("schema.registry.url", "https://sr.example.com");
props.put("auto.register.schemas", "false");
props.put("use.latest.version", "true");