Notion sharded its Postgres data into 480 logical shards across 32 physical databases in 2021 and then in 2023 spread the same 480 logical shards across 96 databases, with no change to the routing logic. What did the 2021 decision buy, what did the 2023 move still cost, and where would copying the number be a mistake?
Show the full answer Hide the answer
The situation they were in
Notion's workspace content outgrew a single Postgres instance, and the 2021 project split it horizontally. The detail that matters is not that they sharded. It is that they chose a logical shard count far larger than the number of machines they needed: 480 logical schemas spread 15 to a database across 32 databases. By 2023 the same 480 logical schemas were running 5 to a database across 96 databases, described in Notion's "The Great Re-shard" post. The capacity tripled and the application's routing function did not change.
480 is not arbitrary. It divides by 32, 40, 48, 60 and 96, so the fleet can grow in several steps while every database keeps an equal number of logical shards.
What the 2021 choice bought
It converted every future capacity change from an application release into a data movement. A shard key hashed to one of 480 logical shards is stable for the life of the data; only the lookup from logical shard to host changes, and that is configuration. Without the over-provisioning, tripling capacity means rehashing every row, which means the routing function, the application, the caches and every piece of analytics that embedded the old shard count all change together, under a flag, with no way to verify the result except by comparing two full datasets.
The cost of the decision in 2021 was small and paid up front: more schemas to create, more connections to manage, and a per-query cost of nothing, since routing is a hash and a table lookup either way.
What the 2023 move still cost
Over-provisioned logical shards make resharding tractable, not free. The data still had to move. Notion used Postgres logical replication to copy history and then apply changes continuously, and reports cutting the synchronisation time from about three days to about twelve hours by deferring index creation until after the bulk copy — building indexes during ingest is slower than building them once at the end.
The remaining work is the part no shard count removes: a verified cutover per database, a routing configuration change, the connection-pool capacity for three times the hosts, and the operational surface of 96 primaries instead of 32. Day-one over-provisioning buys away the application change, not the migration.
Where copying it would be a mistake
- The number. 480 is an artefact of Notion's factor requirements and of how much data one host could hold. The rule that transfers is "pick a logical shard count divisible by every physical count you might plausibly want, and small enough that one logical shard fits comfortably inside one host's working set". That might be 64 or 4,096.
- The decision to shard at all. A single Postgres primary with read replicas serves a very large number of products, and sharding costs you cross-shard queries, transactions, foreign keys and every join that spans the key. Shard when writes or working set exceed one machine, not when reads do.
- The confidence. This worked because the routing abstraction existed before it was needed. Retrofitting logical shards onto an application that already embeds physical database names in queries is a different and much larger project.
Common weak answers
- "They should have used a distributed SQL engine from the start." Possibly, at a different cost: a new operational model, a new failure catalogue, and a dependency decided years before the evidence existed. The documented path kept the engine the team already ran well.
- "480 is best practice." It is one team's arithmetic. Repeating the number without the working set and the factor requirement behind it is cargo cult.
- "Resharding was zero downtime so it was low risk." Zero downtime describes the user experience, not the risk. The risk sat in verification, and that is where the twelve hours went.