Skip to content

Sharding Postgres while it is running

Shard a Live Database Without Downtime

Built from how Notion sharded Postgres while millions of people were using it: choose a shard key, choose a shard count you can live with for years, move every row while writes continue, prove the copy is right, and grow again later without starting over.

Advanced, about 50 minutes, 9 stages

The situation

A collaborative workspace app stores everything (pages, blocks, comments) in one Postgres primary, already on the largest instance available. The biggest table has billions of rows. Writes keep climbing, autovacuum can no longer keep up, and the database is drifting toward transaction ID wraparound: the point at which Postgres stops accepting writes to protect itself.

Each block lives in a single workspace, and almost every query stays inside one workspace. The infrastructure team is small, and it knows Postgres well.

Notion was in this position in 2020 and wrote up the whole migration, then wrote again in 2023 about growing from 32 to 96 database hosts. Figma described a similar journey in 2024, and Stripe and GitHub have published the patterns for moving live data safely. This investigation follows their decisions.

What it has to do

Functional

  • Route every query to the shard that holds its data.
  • Move all existing data from the monolith into shards.
  • Keep serving reads and writes during the migration.
  • Add database hosts later without changing how the application finds data.

Non-functional

  • No acknowledged write is ever lost.
  • At most a few minutes of planned disruption at cutover.
  • Until cutover, every step can be reversed.
  • The sharded copy is verified to match before it serves users.

Constraints and assumptions

  • One Postgres primary on the largest instance size; billions of rows in the largest tables.
  • Vacuum is falling behind, and transaction ID wraparound is a hard deadline.
  • Stay on Postgres: the team's expertise and tooling are built around it.
  • Every row can be attributed to a workspace (directly, or through its parent).
  • IDs are UUIDs, so rows can be copied between databases without ID conflicts.
  • Rows carry an updated-at version that increases on every write.

Interview questions it prepares you for

  • “How would you shard a database that is already in production?”
  • “Migrate a large table to a new schema or store with zero downtime.”
  • “How do you choose a shard key?”
  • “Your database is running out of capacity. Walk me through your options.”

Read and practise next

How Notion built it · Sharding Postgres without downtime, in their engineers' own words

Concepts to know first: Partitioning, Replication.

Similar systems: Design Discord's Message Storage, Design a Payment System.