Serialization formats and schema evolution
Choosing how data is encoded on the wire and in storage (JSON, Protocol Buffers, Avro, Thrift), forward and backward compatibility, rules for evolving schemas safely, schema registries, and versioning APIs and events.
Reading is half of it. See this used in a real interview: walk through Design a Distributed Message Queue →
Every message between services, every event in Kafka and every record on disk is encoded in some format, and that format will change. Fields get added, renamed and removed while old and new versions of the code run side by side, sometimes for months (mobile apps, stored events). Choosing a format and following compatibility rules keeps those changes from breaking production.
The common formats
| Format | Encoding | Schema | Strengths | Weaknesses |
|---|---|---|---|---|
| JSON | text | optional (JSON Schema) | human-readable, universal, great for public APIs | large, slow to parse, loose types |
| Protocol Buffers | binary, tagged fields | required (.proto) | compact, fast, good evolution rules, gRPC | not human-readable |
| Avro | binary, no tags | required, writer schema needed to read | very compact, great for Kafka and data lakes | needs a schema registry |
| Thrift | binary, tagged fields | required | similar to Protobuf | smaller ecosystem today |
| MessagePack, CBOR | binary JSON-like | none | smaller than JSON, schemaless | same evolution risks as JSON |
| Parquet, ORC | columnar files | embedded | analytics storage, compression | not for messages |
Rule of thumb: JSON at public edges, Protobuf for internal RPC, Avro or Protobuf with a registry for event streams, Parquet for analytics. Binary formats are typically several times smaller and faster than JSON, which matters at millions of messages per second. See Design a Message Queue.
Compatibility
- Backward compatible: new code can read data written by old code.
- Forward compatible: old code can read data written by new code (it ignores what it does not understand).
You usually need both, because during a rolling deploy old and new versions read each other's messages, and because stored data (events, files) is read by code written years later.
Safe and unsafe changes
| Change | Safe? |
|---|---|
| Add an optional field with a default | yes |
| Remove an optional field (and never reuse its tag or name) | yes |
| Rename a field (Protobuf: tags matter, names do not on the wire) | yes in Protobuf binary, no in JSON |
| Change a field's type (int32 to string) | no |
| Make an optional field required | no |
| Reuse a deleted field's tag number | no, corrupts old data |
| Change the meaning of a field | no, even if the type stays the same |
| Add an enum value | careful: old readers must handle unknown values |
In Protobuf, mark removed tags reserved. In Avro, every added field needs a default. In JSON, readers must ignore unknown fields and tolerate missing ones.
Making breaking changes
When a breaking change is unavoidable, use expand and contract:
- Add the new field alongside the old one; writers write both.
- Migrate readers to the new field.
- Backfill old data if needed.
- Stop writing the old field, then remove it.
For events, an alternative is a new event type or topic version (OrderPlaced.v2) with both published during the migration. The same pattern applies to database schemas. See online schema migrations.
Schema registries
In Kafka-style systems, a schema registry stores every version of each topic's schema and enforces a compatibility mode (backward, forward or full) when producers register a new version. Messages carry a small schema id instead of the full schema. A breaking change is rejected at deploy time instead of crashing consumers at 3 a.m. See message queues and streams.
Versioning APIs
- Public REST APIs: additive changes in place; breaking changes behind a new version (
/v2or a version header), with a deprecation period. See API design. - Mobile clients: assume old app versions live for years. The server must keep accepting old request shapes and returning fields old clients need.
- gRPC: follow Protobuf rules; package versions (
payments.v1) for breaking changes.
Stored data lives longest
Events in an event-sourced system, messages in a chat history and files in a data lake are read long after they were written. Store the schema id or version with every record, keep readers able to handle all historical versions (or upcast old events to the new shape when reading), and test against real old data. See event sourcing and CQRS.
Checklist
- JSON at the edge, binary formats with schemas inside.
- Backward and forward compatibility as a default requirement.
- Only additive, optional changes; never reuse tags or change types.
- Expand and contract for breaking changes.
- A schema registry enforcing compatibility for streams.
- Schema versions stored with long-lived data.