Notion: sharding postgres without downtime
One Postgres database split by workspace while the product stayed online, then split further two years later.
The idea
Notion's first post covers the choice most teams face once one database is not enough: what to shard by, and into how many pieces. They chose the workspace, because nearly every query is scoped to one, and created far more logical shards than machines, so that adding machines later would mean moving whole shards rather than re-hashing rows.
The migration follows a sequence worth knowing by heart: capture new writes first, copy the history, check that old and new agree (including by reading from both and comparing), and only then switch. The second post shows why the shard count mattered: when they needed three times the capacity, they moved existing shards to new hosts with replication and switched each with a brief pause.
Read the originals
Written by the engineers who built it.
- The Great Re-shard: adding Postgres capacity (again) with zero downtime
Arka Ganguli and others · Post, Jul 2023
The payoff of fixed logical shards: growing to more machines by moving shards whole.
- Herding elephants: lessons learned from sharding Postgres at Notion
Garrett Fidalgo · Post, Oct 2021
Picking a shard key and a shard count, then moving a live database across with no data lost. Includes an honest list of what they would do differently.
Practise it
Make the decisions yourself, then compare.