3.8 Serializers
That last point is the organizational killer — the same lockstep-deploy problem Ch. 1 identified, now embedded in your byte layout.
8.1 Why you should not write your own — demonstrated by writing one
public class Customer {
private int customerID;
private String customerName;
public Customer(int ID, String name) { this.customerID = ID; this.customerName = name; }
public int getID() { return customerID; }
public String getName() { return customerName; }
}public class CustomerSerializer implements Serializer<Customer> {
@Override
public void configure(Map configs, boolean isKey) { /* nothing to configure */ }
@Override
/**
We are serializing Customer as:
4 byte int — customerId
4 byte int — length of customerName in UTF-8 bytes (0 if name is Null)
N bytes — customerName in UTF-8
**/
public byte[] serialize(String topic, Customer data) {
try {
byte[] serializedName;
int stringSize;
if (data == null) return null;
else {
if (data.getName() != null) {
serializedName = data.getName().getBytes("UTF-8");
stringSize = serializedName.length;
} else {
serializedName = new byte[0];
stringSize = 0;
}
}
ByteBuffer buffer = ByteBuffer.allocate(4 + 4 + stringSize);
buffer.putInt(data.getID());
buffer.putInt(stringSize);
buffer.put(serializedName);
return buffer.array();
} catch (Exception e) {
throw new SerializationException("Error when serializing Customer to byte[] " + e);
}
}
@Override
public void close() { /* nothing to close */ }
}Now count the ways this fails:
*"This example is pretty simple, but you can see how fragile the code is.
- If we ever have too many customers and need to change
customerIDtoLong...- or if we ever decide to add a
startDatefield...- ...we will have a serious issue maintaining compatibility between old and new messages.
Debugging compatibility issues between different versions of serializers and deserializers is fairly challenging: you need to compare arrays of raw bytes.
To make matters even worse, if multiple teams in the same company end up writing
Customerdata to Kafka, they will all need to use the same serializers and modify the code at the exact same time."*
That last point is the organizational killer — the same lockstep-deploy problem Ch. 1 identified, now embedded in your byte layout.
Recommendation: use existing serializers — JSON, Apache Avro, Thrift, or Protobuf.
8.2 Apache Avro
What it is: a language-neutral data serialization format, created by Doug Cutting to provide a way to share data files with a large audience.
- Data described in a language-independent schema, usually in JSON.
- Serialization usually to binary (JSON output also supported).
- Avro assumes the schema is present when reading and writing files — usually by embedding the schema in the files themselves.
The killer feature for messaging:
"When the application writing messages switches to a new but compatible schema, the applications reading the data can continue processing messages without requiring any change or update."
Worked schema evolution example
Original schema (used for months, a few terabytes of data generated):
{"namespace": "customerManagement.avro",
"type": "record",
"name": "Customer",
"fields": [
{"name": "id", "type": "int"},
{"name": "name", "type": "string"},
{"name": "faxNumber", "type": ["null", "string"], "default": "null"}
]
}id and name are mandatory; faxNumber is optional, defaults to null.
New schema — "upgrade to the 21st century": drop fax, add email:
{"namespace": "customerManagement.avro",
"type": "record",
"name": "Customer",
"fields": [
{"name": "id", "type": "int"},
{"name": "name", "type": "string"},
{"name": "email", "type": ["null", "string"], "default": "null"}
]
}The real-world constraint: "In many organizations, upgrades are done slowly and over many months." So old records have faxNumber and new records have email, simultaneously, and both old and new readers must cope.
reads OLD record (has faxNumber) | reads NEW record (has email) | |
|---|---|---|
OLD app — getName / getId / getFaxNumber() | getName() ✓ · getId() ✓ · getFaxNumber() ✓ | getName() ✓ · getId() ✓ · getFaxNumber() → null |
NEW app — getName / getId / getEmail() | getName() ✓ · getId() ✓ · getEmail() → null | getName() ✓ · getId() ✓ · getEmail() ✓ |
No exceptions. No breaking errors. No expensive updates of existing data.
Two caveats — do not skip these
- "The schema used for writing the data and the schema expected by the reading application must be compatible." (The Avro documentation includes the compatibility rules.)
- "The deserializer will need access to the schema that was used when writing the data, even when it is different from the schema expected by the application." In Avro files, the writing schema is included in the file itself — but there is a better way for Kafka messages.
8.3 Avro + Schema Registry
The problem with embedding schemas per record:
"Unlike Avro files, where storing the entire schema in the data file is associated with a fairly reasonable overhead, storing the entire schema in each record will usually more than double the record size. However, Avro still requires the entire schema to be present when reading the record, so we need to locate the schema elsewhere."
The pattern:
KEY POINT: all of this — storing the schema in the registry and pulling it up when required — happens INSIDE the serializers and deserializers. Your producing code “simply uses the Avro serializer just like it would any other serializer.”
The Schema Registry is NOT part of Apache Kafka, but there are several open source options. The book uses the Confluent Schema Registry (on GitHub, or as part of the Confluent Platform).
Producing generated Avro objects
Properties props = new Properties();
props.put("bootstrap.servers", "localhost:9092");
props.put("key.serializer",
"io.confluent.kafka.serializers.KafkaAvroSerializer");
props.put("value.serializer",
"io.confluent.kafka.serializers.KafkaAvroSerializer");
props.put("schema.registry.url", schemaUrl); // where schemas live
String topic = "customerContacts";
Producer<String, Customer> producer = new KafkaProducer<>(props);
// keep producing new events until someone ctrl-c
while (true) {
Customer customer = CustomerGenerator.getNext();
System.out.println("Generated customer " + customer.toString());
ProducerRecord<String, Customer> record =
new ProducerRecord<>(topic, customer.getName(), customer);
producer.send(record);
}Notes:
KafkaAvroSerializercan also handle primitives — that's whyStringworks as the key whileCustomeris the value.Customeris NOT a POJO. It is "a specialized Avro object, generated from a schema using Avro code generation." "The Avro serializer can only serialize Avro objects, not POJO." Generate classes withavro-tools.jaror the Avro Maven plug-in (both part of Apache Avro).
Producing generic Avro objects (no codegen)
Properties props = new Properties();
props.put("bootstrap.servers", "localhost:9092");
props.put("key.serializer", "io.confluent.kafka.serializers.KafkaAvroSerializer");
props.put("value.serializer", "io.confluent.kafka.serializers.KafkaAvroSerializer");
props.put("schema.registry.url", url);
String schemaString =
"{\"namespace\": \"customerManagement.avro\","
+ "\"type\": \"record\", "
+ "\"name\": \"Customer\","
+ "\"fields\": ["
+ "{\"name\": \"id\", \"type\": \"int\"},"
+ "{\"name\": \"name\", \"type\": \"string\"},"
+ "{\"name\": \"email\", \"type\": [\"null\",\"string\"], \"default\":\"null\" }"
+ "]}";
Producer<String, GenericRecord> producer = new KafkaProducer<String, GenericRecord>(props);
Schema.Parser parser = new Schema.Parser();
Schema schema = parser.parse(schemaString);
for (int nCustomers = 0; nCustomers < customers; nCustomers++) {
String name = "exampleCustomer" + nCustomers;
String email = "example " + nCustomers + "@example.com";
GenericRecord customer = new GenericData.Record(schema);
customer.put("id", nCustomers);
customer.put("name", name);
customer.put("email", email);
ProducerRecord<String, GenericRecord> data =
new ProducerRecord<>("customerContacts", name, customer);
producer.send(data);
}Generic Avro objects are used as key-value maps rather than generated objects with getters/setters. You must provide the schema yourself, since no generated object carries it. "The serializer will know how to get the schema from this record, store it in the Schema Registry, and serialize the object data."