SysDesignPrep.com
Study guide 116 of 183

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

FormatEncodingSchemaStrengthsWeaknesses
JSONtextoptional (JSON Schema)human-readable, universal, great for public APIslarge, slow to parse, loose types
Protocol Buffersbinary, tagged fieldsrequired (.proto)compact, fast, good evolution rules, gRPCnot human-readable
Avrobinary, no tagsrequired, writer schema needed to readvery compact, great for Kafka and data lakesneeds a schema registry
Thriftbinary, tagged fieldsrequiredsimilar to Protobufsmaller ecosystem today
MessagePack, CBORbinary JSON-likenonesmaller than JSON, schemalesssame evolution risks as JSON
Parquet, ORCcolumnar filesembeddedanalytics storage, compressionnot 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

ChangeSafe?
Add an optional field with a defaultyes
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 requiredno
Reuse a deleted field's tag numberno, corrupts old data
Change the meaning of a fieldno, even if the type stays the same
Add an enum valuecareful: 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:

  1. Add the new field alongside the old one; writers write both.
  2. Migrate readers to the new field.
  3. Backfill old data if needed.
  4. 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 (/v2 or 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.

Open in your browser to sign in

Google does not allow sign-in inside this app's built-in browser. Open this page in Safari and sign in there. The link opens this same page.

Tap the ⋯ or share button at the top or bottom of the screen, then Open in browser. Or copy the link and paste it into Safari.