Why Confluent Kafka Avro encoding bothers me!
March 12, 2020
Recently I have been working to get PySpark to read from a Kafka topic. We use a Confluent based deployment of Kafka with Schema Registry and Kafka Connect. If you have tried doing this it has come with quite a bit of pain. But I think this pain partly comes from how Confluent has modified Avro binary encoding with their own header.
Here is the binary format of Confluent Avro messages in Kafka from the documentation.
| Bytes | Area | Description |
|---|---|---|
| 0 | Magic Byte | Confluent serialization format version number; currently always 0. |
| 1-4 | Schema ID | 4-byte schema ID as returned by Schema Registry |
| 5-… | Data | Avro serialized data in Avro’s binary encoding. The only exception is raw bytes, which will be written directly without any special Avro encoding. |
Note that the first 5 bytes are the Confluent Avro header. But being no longer a valid Avro binary encoding (as a whole), this could only be decoded by Confluent Avro decoder. You might point out that this is how encoding/decoding works. But couldn’t Confluent do this while adhering to the Avro binary encoding?
I think they could have! Just imagine that Confluent has mandated that every Confluent Schema registry compatible Avro schema should begin with 2 fields like below.
{
"type": "record",
"fields": [
{
"name": "ser_format",
"type": "int"
},
{
"name": "schema_id",
"type": "int"
},
/* rest of the fields ... */
]
}Schema Registry can enforce that any Avro schema definition conforms to the above format and the extra overhead will be exactly the same (or slightly less - because ints are encoded with variable length encoding in Avro). These fields can be filled in by Confluent Serializer from values in schema registry. This could have made it possible to use any Avro decoder for decoding!
The above document also notes the following.
The serialization format used by Confluent Platform serializers is guaranteed to be stable over major releases. No changes are made without advanced warning. This is critical because the serialization format affects how keys are mapped across partitions.
This implies that the serialized data supplied to the partitioner already has the above headers prepended. But if you change the key schema doesn’t it get partitioned differently to messages before schema change anyway?
Another thing which bothers me is the lack of a schema version field in record. This might be because how Confluent/Avro supports automatic schema evolution (when schemas are compatible). But what about the case when breaking schema changes are present?
I think this is a case of accidental complexity, not essential complexity. Just some thoughts.
Written by Francois Fernando, a software craftsman and tinkerer.