Learn Labs
5. Encoding and Evolution

5.5 Avro

Started 2009 as a Hadoop subproject, because Protocol Buffers was not a good fit for Hadoop's use cases.

Started 2009 as a Hadoop subproject, because Protocol Buffers was not a good fit for Hadoop's use cases.

Two schema languages: Avro IDL (for human editing) and a JSON-based one (machine-readable). Like protobuf, they specify only fields and types — no complex validation rules.

record Person {
    string               userName;
    union { null, long } favoriteNumber = null;
    array<string>        interests;
}
{ "type": "record", "name": "Person",
  "fields": [
    {"name": "userName",       "type": "string"},
    {"name": "favoriteNumber", "type": ["null", "long"], "default": null},
    {"name": "interests",      "type": {"type": "array", "items": "string"}}
  ] }

5.1 The encoding — 32 bytes, and NOTHING self-describing

Notice: the schema has NO TAG NUMBERS.

0c  M a r t i n        length 6 (varint), then UTF-8 bytes
02  f2 14              union branch 1 (long), then varint 1337
04  16 daydreaming  0e hacking  00    array: count, items…, terminator

TOTAL: 32 bytes — the most compact of all the encodings.

Nothing in the encoded data says the first field is a string — it could just as well be an integer. That is only workable because the reader always has a schema.

The encoding is simply VALUES CONCATENATED TOGETHER. To parse it, you go through the fields in the order they appear in the schema and use the schema to determine each field's datatype.

⇒ The binary data can be decoded correctly ONLY IF the reading code uses the EXACT SAME SCHEMA as the writing code. Any mismatch means incorrectly decoded data.

So how does Avro evolve at all? This is the clever part.

5.2 Writer's schema and reader's schema

Protocol Buffers
encodeschema v1bytesdecodeschema v2Encoder and decoder may use different versions — the tags carry the identity.
Avro
writer's schemabytesreader's schemaolder or newerdecode using bothThe writer's schema must be identical to the one used for encoding; the reader's may be a different version.
Figure 5.5.25.2 Writer's schema and reader's schema

Schema resolution — how differences are reconciled:

Writer's schemaReader's schemaResolution
userNameuserName✔ matched by name
favoriteNumberfavoriteNumber✔ matched by name — field order may differ
interests—Not in the reader's schema ⇒ ignored
—photoURLReader expects it, writer does not have it ⇒ filled in with the reader's declared default

5.3 Evolution rules

Meaning in Avro
Forward compatibilityThe WRITER can use a NEWER schema than the reader
Backward compatibilityThe WRITER can use an OLDER schema than the reader

THE RULE: you may add or remove only a field THAT HAS A DEFAULT VALUE.

  • Add a field with a default → a reader on the new schema reading old data fills in the default. ✔
  • Add a field without a default → new readers can't read data written by old writers → breaks BACKWARD compatibility. ✗
  • Remove a field without a default → old readers can't read data written by new writers → breaks FORWARD compatibility. ✗

Nulls are explicit. In some languages null is an acceptable default for any variable — not in Avro. To allow null you must use a union type: union { null, long, string } field;. You can use null as a default only if it is the FIRST branch of the union. More verbose than nullable-by-default, but it helps prevent bugs by being explicit about what can and cannot be null.

Two asymmetric cases worth remembering:

ChangeBackwardForward
Change a field name (reader's schema declares an alias matching the old writer's name)✔✗
Add a branch to a union type✔✗
Change a datatype (where Avro can convert)possiblepossible

5.4 But how does the reader know the writer's schema?

You can't include the whole schema with every record — it would likely be much bigger than the encoded data, negating all the space savings. Three answers depending on context:

ContextMechanism
Large file with lots of records (millions, all one schema)Include the schema once at the beginning of the file. Avro specifies a file format for this: object container files
Database with individually written records (different records written at different times with different schemas)Include a version number at the start of every encoded record and keep a list of schema versions. Reader extracts the version, fetches the corresponding writer's schema, decodes. This is how Confluent's Schema Registry for Kafka and LinkedIn's Espresso work
Records over a network connectionNegotiate the schema version at connection setup, use it for the connection's lifetime. The Avro RPC protocol works like this

A database of schema versions is useful in any case: it acts as documentation and gives you a chance to CHECK SCHEMA COMPATIBILITY. Version number = a simple incrementing integer or a hash of the schema.

5.5 Why "no tag numbers" matters: dynamically generated schemas

The scenario: dump a relational database's contents to a file in a binary format.

AvroProtocol Buffers
Generate the schema from the relational schema automatically — one record schema per table, one field per column, column name becomes field name.Field tags would have to be assigned by hand.
A column added or removed? Regenerate the Avro schema and export in the new one. The export process pays no attention to the change.Every schema change would need an administrator to update the column-name → tag mapping manually, and automating it means being very careful never to reassign a used tag.
Readers see the fields changed, but fields are identified by name, so the new writer's schema still matches the old reader's schema.This simply was not a design goal of Protocol Buffers — but it was for Avro.

On this page