Learn Labs
3. Kafka Producers: Writing Messages to Kafka

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 customerID to Long...
  • or if we ever decide to add a startDate field...
  • ...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 Customer data 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() → nullgetName() ✓ · getId() ✓ · getEmail() ✓

No exceptions. No breaking errors. No expensive updates of existing data.

Two caveats — do not skip these
  1. "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.)
  2. "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:

Avro objectregister, get an idfetch by idAvro objectSCHEMA REGISTRYid 42 → {schema…} · id 43 → {schema…}Producer codeSerializerDeserializerConsumer codeKAFKA record[schema-id][serialized data] — only the ID travels

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.”

Figure 3.8.2The pattern

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:

  • KafkaAvroSerializer can also handle primitives — that's why String works as the key while Customer is the value.
  • Customer is 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 with avro-tools.jar or 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."


On this page