System design practice

Notion Sharding

medium Based on a system published by Notion sharding write-heavy real-world

Notion's Postgres: 480 logical shards by workspace on 32 databases, behind PgBouncer.

Solve it in your browser Read the lesson first

Until 2021 Notion kept every block (each paragraph, heading, to-do and page is a block) in one PostgreSQL database. It held over 20 billion rows and was running out of room: VACUUM stalled, and transaction id wraparound was approaching, after which Postgres stops taking writes. A bigger machine would only buy months.

So Notion split the database. The hard decision is the partition key: the column that decides which shard a row lives on. A good key keeps the data a request needs on one shard and spreads the load evenly.

Functional requirements

The API servers decide which shard a query goes to and send it through pgbouncer, which pools their connections to Postgres. Use these use case names exactly: the traffic, requirements and tests in problem.proschi refer to them.

Scale

Constraints

What is given

problem.proschi declares web, the API servers (the client of this problem), and pgbouncer, ten PgBouncer instances. Add the Postgres fleet as one node, choose how many shards with capacity { db shards N }, and write the two use cases.

What this model leaves out

Based on

How your design is checked

You write the design as text in Proschi. Tests run in your browser: a simulation of the traffic above checks latency, availability, cost and what happens when a machine fails. How the simulation works.

Start designing

Lesson · 12 min read

Learn it: Notion Sharding #

Notion Sharding: when one Postgres primary is not enough #

What you'll learn #

The problem, explained #

Notion stores everything as blocks: a page is a block, and so is every paragraph, heading and to-do inside it. Until 2021 all blocks lived in one PostgreSQL database. It grew past 20 billion rows, and VACUUM (Postgres's background cleanup of dead row versions) could not keep up. The database was also approaching transaction ID wraparound: at that point Postgres stops accepting writes to protect the data. A bigger machine would only buy months, so Notion split the data across many databases.

Two use cases:

The scale is an assumption for the exercise (Notion does not publish request rates): 200k page loads and 100k block edits per second at peak. Both use cases need a p99 under 70 ms and 99.95% availability. Losing any single machine must not break a latency limit. The budget is $30,000 a month, PgBouncer included.

What is given. given.proschi declares web, the API servers, as the client: they decide which shard a query goes to. It also declares pgbouncer, ten PgBouncer instances. PgBouncer is a connection pooler. Postgres runs one process per connection, so thousands of API processes connecting directly would exhaust it. The pooler shares a small number of real connections among all of them. You add the Postgres fleet as one node, choose its shard count with capacity { db shards N } (the only capacity line a solver may write), and write both use cases.

What the tests check: both use cases start at web and go through pgbouncer before Postgres; there is no path from web to any database; page loads read Postgres and edits write Postgres before responding. The requirements add p99, availability, durability, failure survival and cost.

The model leaves some things out, and the statement says so. The simulation cannot see the partition key itself, only its effect: a query that needs every shard is written as a fan-out (x32) and loads every shard. Load is spread evenly over the shards, connection limits are not simulated, and the migration is out of scope.

Back-of-the-envelope #

The simulation's numbers for a PostgreSQL replica are 20k reads and 5k writes per second. The crucial detail: Postgres is single-primary. Every replica serves reads, but every write goes to the one primary of its shard.

QuantityArithmeticResult
Reads (page loads)given200k rps
Writes (edits)given100k rps
Read:write ratio200k : 100k2 : 1
Write capacity of one primarysimulation default5k rps
Primaries needed at 100%100k ÷ 5k20
Read capacity of a shard with a primary and a replica2 × 20k40k rps
PgBouncer load300k ÷ (10 × 50k)60%
Monolith with 4 replicas: write utilisation100k ÷ 5k2,000%

The table answers the first question right away. This is a write problem: 100k writes need at least 20 primaries just to avoid saturation. Adding replicas to a single database does not help at all, because replicas add zero write capacity.

How many shards, then? Twenty primaries at 100% is not a design. The simulation queues each shard's writes on its one primary (one server, so the classic 1 ÷ (1 − ρ) slowdown): at 50% busy a hop takes twice its base latency, at 70% over three times. The p99 limit is 70 ms, and an idle hop's p99 is about 2.8 times its mean, so the primaries must stay well below 70% busy. Write utilisation per shard is 100k ÷ (5k × shards): 20 shards give 100%, 25 give 80%. Keep going until the p99 holds, then check the budget.

Cost. Each replica of each shard costs $400 a month, PgBouncer $100 per instance. So shards × replicas × $400 + $1,000 ≤ $30,000, which caps shards × replicas at 72. Two replicas per shard (a primary and a standby for failover) fits a shard count in the 30s; three replicas per shard cuts the affordable shard count to 24, too few for the write rate.

Availability and failure. With a replica, a Postgres primary that fails is replaced by promotion; the model charges writes 10% of the primary's downtime for the failover, which still clears 99.95%. survive any node failure takes one instance out of one shard: writes keep a primary, reads lose a replica. At 200k reads spread over dozens of shards, one shard's reads fit comfortably on one remaining replica.

Storage, for scale. 20 billion rows at, say, 1 KB each (an assumption) is about 20 TB, which is a lot for one Postgres host to vacuum. Split over a few dozen databases it is well under a terabyte each.

Concepts #

Replicas scale reads, shards scale writes #

Replication copies the same data to several machines. In a single-primary database such as Postgres or MySQL, replicas apply the primary's changes and serve reads. They help read-heavy workloads and give you a standby for failover. They do nothing for writes: every replica must apply every write, so the primary is still the bottleneck, and replicas actually add write work.

Sharding (horizontal partitioning) splits the rows across independent databases, each with its own primary. Ten shards means ten primaries and ten times the write capacity, as long as writes spread evenly. The price is complexity: the application must know where each row lives, cross-shard queries and transactions become hard, and rebalancing is a project.

Use replicas first when reads are the problem; they are cheap and invisible to the application. Shard only when writes, data size or maintenance (like VACUUM) outgrow one primary. Do not shard to fix a slow query or a missing index.

title "Sharded relational store"

app "App"      [Actor]
db  "Store"    [PostgreSQL] x2

capacity {
  db shards 4
}

app -> db : SQL

usecase "Save" {
  app -> db  : UPDATE item SET name WHERE tenant_id = $1
  db --> app : UPDATE 1
}

Choosing a partition key #

The partition key is the column that decides which shard a row lives on, usually through a hash. A good key has two properties:

  1. Locality: the queries you run most need rows from only one shard. Anything else turns one query into many.
  2. Spread: the key has many distinct values with similar load, so no shard is much hotter than the rest.

For Notion, the workspace ID wins on both. Every block, comment and discussion belongs to exactly one workspace, and nearly every request (load a page, save an edit) concerns one workspace, so it touches one shard. There are millions of workspaces, so load spreads well. Compare the alternatives:

When not to shard by tenant: when one tenant can outgrow a shard (a single huge customer becomes a hot shard), or when the core queries cross tenants (a global feed). Then you need a finer key, a dedicated shard for the giant, or a different data model.

Logical shards and the routing layer #

Re-sharding is the expensive part of sharding: moving rows between databases while serving traffic. Notion's trick was to create many more logical shards than physical databases: 480 Postgres schemas, 15 per database. The application maps workspace → logical shard (fixed forever) and logical shard → database (a small table). Growing the fleet means moving whole schemas to new machines; no row is ever re-hashed. 480 was chosen because it divides evenly into many fleet sizes, and Notion later did exactly that, going from 32 to 96 databases.

Routing happens in the application: the API computes the shard from the workspace ID and sends the query to the matching database through PgBouncer. Other designs put routing in a separate layer, such as Vitess for MySQL or the Citus coordinator for Postgres. The trade-off: that layer hides sharding from the application, but it is another moving part on the hot path.

Designing it step by step #

1. Scope the problem #

Ask: what is the read and write rate? (200k and 100k per second.) What do the main queries filter on? (Workspace and parent block.) Must the data stay in Postgres? (Yes.) Is there a connection limit to worry about? (Yes, which is why PgBouncer exists.) Point out that a 2:1 read:write ratio is unusually write-heavy. That alone tells you replicas will not be enough.

2. High-level design #

The flow is short: web → PgBouncer → Postgres, for both use cases. The design question is the shape of the Postgres node.

Walk the interviewer through the alternatives in order. A bigger monolith: the write rate is twenty times what one primary takes. A monolith with read replicas: same primary, same problem. A few shards with many replicas: more primaries, but still not enough for 100k writes, and replicas you pay for sit mostly idle. Many shards with one standby each: write capacity scales with the shard count, and the standby covers failover. That last one wins.

Then the key: partition by workspace ID, so each request names its shard and touches one.

3. Deep dive #

The shard count. Use the formula from the estimates: pick a count that keeps write utilisation per primary under about two-thirds, check the p99 in the analysis, then check that two replicas per shard fit the budget. If you want to reason like Notion did, prefer a count that divides a larger number of logical shards evenly.

Reads. With two replicas per shard, read capacity is far above 200k rps; reads are not the constraint, writes are. Notice that the simulation reports the busier of the two sides for a single-primary store.

PgBouncer. Ten instances at 50k rps each carry 300k requests at 60%. It is given; your job is to route every query through it. The test "Queries go through PgBouncer" enforces that, including its line no path from web to any database.

Failure. Losing a primary promotes the standby; losing a standby costs read capacity on one shard. Run the analysis and confirm no latency limit breaks with an instance missing.

4. Wrap-up #

Summarise: shard by workspace ID across enough primaries to keep writes around two-thirds busy, one standby per shard, routing in the application, pooling in PgBouncer, logical shards for cheap growth. Next steps worth naming: the migration (double writes, backfill, verification, then switching reads), monitoring for hot workspaces, and a plan for the day shards run hot again (Notion's 2023 re-shard).

Common mistakes #

One big database (wrong/one-big-database.proschi). Keep the monolith and add replicas: four replicas, one primary. Every one of the 100k edits still lands on that primary, which takes 5k. In real life this is the "just buy a bigger box" plan, and Notion's story is what happens at its end: maintenance falls behind and wraparound looms. Here it fails the p99 of both use cases, because the saturated database is on both paths (and survive any node failure too).

Replicas instead of shards (wrong/replicas-instead-of-shards.proschi). Four shards with eight replicas each: lots of read capacity, only four primaries. 100k writes on 20k of write capacity saturate the primaries. It fails the p99 of Edit block, and of Load page too, because pages are read from the same overloaded shards. The lesson: count primaries, not machines.

Shard by block ID (wrong/shard-by-block-id.proschi). Writes spread perfectly, but a page's blocks live on every shard, so a page load becomes an x32 fan-out: 200k page loads turn into 6.4 million shard reads a second. It fails the p99 of Load page, and of Edit block too, because edits wait behind those reads on the same shards. In production this is the classic scatter-gather trap: every page load waits for the slowest shard, and a single slow shard slows every page.

Direct connections (wrong/direct-connections.proschi). The API servers connect to Postgres directly. With hundreds of API processes and dozens of databases, that is thousands of connections, each a Postgres backend process with its own memory. It fails "Queries go through PgBouncer".

Three replicas per shard. Tempting for safety, but at a shard count that can handle the writes, it breaks the $30,000 budget. One standby per shard is enough to survive a single failure.

In the interview #

Start with the arithmetic that rules out the easy answers: "100k writes a second against 5k per primary means at least 20 primaries before any headroom. Replicas do not add primaries, so this is a sharding problem." Then spend your time on the partition key, which is what the interviewer most wants to hear about.

Expected follow-ups:

Further reading #

Now design it

More system design problems