Skip to content

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.

Practise it

Make the decisions yourself, then compare.