Skip to content

The finished design, decision by decision

How to shard a Live Database Without Downtime

Not the one correct diagram, but a design you can defend under these constraints: the finished architecture, then every stage's question with the reasoning that answers it, the tradeoffs it accepts, and where another engineer could land differently.

Sharding Postgres while it is running

The short answer

7 parts, each with one job. The map below shows how requests and data move between them; the stages after it explain why each part is there.

Users
Read and edit pages.
Application servers
Compute the shard from workspace_id; send each query to its shard; dark-read and compare during migration.
Monolith Postgres
The original database; source of truth until cutover.
Write audit log
Every write made during the migration, in order.
Backfill and catch-up
Copies existing rows, then replays the log; checkpoints per table and shard; never overwrites newer versions.
Connection poolers
PgBouncer clusters that map logical shards to physical hosts and can pause traffic for a cutover.
Shard hosts
480 logical shards (one Postgres schema each), spread across physical hosts.
12345678CLIENTUsersSERVICEApplicationserversDATABASEMonolithPostgresLOG / STREAMWrite audit logWORKERBackfilland catch-upSERVICEConnectionpoolersDATABASEShard hosts

Select a component to see what it is responsible for and which state it owns.

  1. 1Users → Application servers: Requests
  2. 2Application servers → Monolith Postgres: Reads and writes until cutover
  3. 3Application servers → Write audit log: Log every write
  4. 4Backfill and catch-up → Monolith Postgres: Copy existing rows
  5. 5Backfill and catch-up → Write audit log: Replay logged writes
  6. 6Backfill and catch-up → Shard hosts: Write rows by shard
  7. 7Application servers → Connection poolers: Queries by shard
  8. 8Connection poolers → Shard hosts: Schema on host
  • Request / response
  • Asynchronous

Why does this design work?

Sharding divides the one resource that was running out: per-host write volume and table size. The shard key is the workspace, because almost every query and join lives inside one, so most queries still touch one database. A fixed set of 480 logical shards (Postgres schemas) is mapped to physical hosts, so growth is moving whole schemas, never re-hashing rows.

The migration follows the safe order: log writes first, backfill without overwriting newer versions, catch up, verify with sampled comparisons and dark reads, then cut over in a short window, with the monolith authoritative and the switch reversible until then. Later growth uses database replication and a brief pause at the connection poolers. Cross-workspace features are designed explicitly: scatter-gather for a few shards, write-time indexes for many, and a single owner for global constraints.

Invariants, and where they are enforced

  • Every row of a workspace lives on exactly one logical shard.

    Logical shard = hash(workspace_id) mod 480; a map assigns each logical shard (schema) to a host. Enforced by Application servers, Shard hosts.

  • No write made during the migration is lost.

    Writes are logged before the copy starts; catch-up replays them in order; the backfill never overwrites a newer version. Enforced by Write audit log, Backfill and catch-up.

  • Until cutover, the monolith is the source of truth and the switch can be undone.

    Shards receive copies only; reads from shards are dark (compared, not served) until verification passes. Enforced by Application servers, Monolith Postgres.

What does it rely on?

  • Almost every query is scoped to one workspace.
  • Every row's workspace can be determined.
  • Rows carry versions, so concurrent copies can be reconciled.
  • A few minutes of planned maintenance is acceptable for the first cutover.

What tradeoffs does it make?

ChoiceGainsCosts
Workspace as shard keySingle-shard queries and local joins.Large workspaces concentrate load; cross-workspace features get harder.
Many logical shards per hostGrowth without re-sharding.More schemas, and routing that maps shards to hosts.
Application audit log for captureWorks at a scale where logical replication snapshots were too slow.Custom code that must capture every write path.
Scheduled maintenance at cutoverSimple, certain final switch.Minutes of downtime.

What are the reasonable alternatives?

Distributed SQL (CockroachDB, Spanner, YugabyteDB)
Better when starting fresh, or cross-shard transactions are central and the team can adopt a new database.
Vitess or Citus in front of existing databases
Better when you want sharding middleware rather than routing in the application.
Vertical partitioning (separate databases per table group)
Better when A few tables dominate load and they rarely join with the rest. Figma did this first.

When does it stop working?

  • One workspace grows beyond a single host's capacity.
  • Cross-workspace queries become the common case.
  • Global constraints or transactions span many shards routinely.

Every stage, decided and explained

Spoilers, for the whole investigation: each stage's question and its answer, the reasoning behind it, and the tradeoffs it accepts. If you have not worked through the stages yet, you may want to do that first.

Work through the stages

Stage 1 of 9 · Model

What is actually running out?

Postgres tags each transaction with a 32-bit ID. Vacuum must periodically "freeze" old rows so IDs can be reused; if it falls too far behind, Postgres stops accepting writes rather than risk corrupting data.

What you need to know first

Postgres gives each transaction a 32-bit transaction ID. That's about 4 billion IDs, reused in a cycle. To make reuse safe, vacuum must "freeze" old rows so they no longer depend on their original transaction ID.

If vacuum falls too far behind, Postgres stops accepting writes rather than risk misreading old rows as new. That's wraparound, and it's a hard deadline.

Three different fixes target three different limits:

FixWhat it divides
Read replicasread load (every replica still replays every write)
Partitioning tables on one hostmaintenance work per table (same host limits)
Sharding across hostswrite volume, table size and vacuum work

See Partitioning and Replication.

Vacuum can't keep up with the primary's write volume. Do read replicas help?

No: replicas take reads, but the primary still does every write and all the vacuum work.

The limit is per-host writes and table size on the primary. Replicas don't touch it, and each replica replays the same writes.

What the stage asks

Which statements hold?

  1. Fails

    Moving to a bigger instance would fix this for good.

    They are already on the largest instance. And vacuum's work grows with table size, so a faster machine postpones the deadline without changing its direction.

  2. Fails

    Adding read replicas would relieve the pressure.

    Replicas take reads off the primary, but this is a write-volume and table-size problem on the primary. Every replica also replays every write.

  3. Holds

    Splitting the largest tables into partitions on the same host makes each vacuum smaller, but every partition still shares one host's CPU, disk and connections.

    Native partitioning helps maintenance, not capacity. The host's limits are the ceiling.

  4. Holds

    Spreading workspaces across many hosts divides write volume, table sizes and vacuum work between them.

    Each host holds a fraction of the data and receives a fraction of the writes, so every per-host limit is pushed back by roughly the number of hosts.

The reasoning

  1. Name the limit first: here, per-host write volume and table size.
  2. Replicas divide reads; partitioning on one host shrinks maintenance; only sharding divides writes.
  3. Transaction ID wraparound makes vacuum lag a hard deadline.

Name the limit before choosing the fix. Here it is per-host write volume and table size: replicas do not touch it, a bigger machine is not available, and partitioning inside one host only makes maintenance smaller. Horizontal Partitioning is the only option that divides the load itself. That justifies the cost, which is considerable: Notion described sharding as something they would rather have done earlier, while the data was smaller and the migration simpler.

Stage 2 of 9 · Decide

Choose the shard key

Most queries load a page's blocks, a workspace's sidebar or a page's comments. A workspace may have one member or thousands.

What you need to know first

The shard key decides which shard each row lives on. Choose it from the queries: the common query should hit one shard, and rows that are joined together should live together.

An even spread is worth little if the common query has to visit every shard.

Sharding by a hash of each row's own ID spreads data perfectly evenly. What happens to 'load this page's blocks'?

It becomes a scatter-gather across all shards, because the page's blocks are spread everywhere.

Every page load touches every shard. The common query gets slower as you add shards.

Blocks, pages and comments are all sharded by workspace ID. Why does it matter that they share the same key?

So a workspace's blocks, pages and comments land on the same shard, and joins between them stay local. Tables queried together must be split the same way. Figma calls such groups "colos".

What the stage asks

What should determine which shard a row lives on?

  1. Sound

    The workspace ID, for every table that belongs to a workspace

    Almost every query stays inside one workspace, so almost every query hits one shard, and joins between a workspace's blocks, pages and comments stay local. It is the key Notion chose. The cost: very large workspaces are concentrated on one shard, and anything that crosses workspaces becomes a multi-shard query.

  2. Flawed

    The ID of the user who created the row

    A shared page contains blocks written by many users, so loading one page would gather rows from many shards. The key does not match how data is read.

  3. Flawed

    A hash of each row's own ID, for perfectly even distribution

    Data spreads evenly, and every page load becomes a scatter-gather across all shards. Even distribution is worthless if the common query has to touch everything.

  4. Flawed

    Time ranges of creation date, so old data sits on cold shards

    All new writes land on the newest shard, which becomes the hot spot, and a page with old and new blocks spans shards.

What a strong answer covers

  • Common queries and joins stay inside one workspace, so they hit one shard.
  • Every row can be attributed to a workspace, so the key is known for every query.
  • Names the costs: big workspaces concentrate load; cross-workspace queries span shards.

The reasoning

  1. Choose the shard key from the queries: the common query should hit one shard.
  2. Shard tables that are joined together by the same key so joins stay local.
  3. Costs: very large tenants concentrate load, and cross-tenant queries span shards.

The shard key is chosen from the queries, not from the data's size or distribution. A key that sends the common query to one shard, and keeps related rows together for joins, is worth far more than a perfectly even spread. Figma calls groups of tables sharded by the same key "colos" (colocations) for exactly this reason: tables that are queried together must be split the same way.

Notion added one regret afterwards: they wished every table's primary key had included the workspace ID, so the shard could be found from a row's ID alone without looking anything up.

Stage 3 of 9 · Decide

How many shards?

Today 32 hosts are enough. In a few years it may be 96, or more. Moving rows between shards is a migration in its own right; moving a whole shard from one host to another is much easier.

What you need to know first

Separate two layouts:

  • Logical: which shard a row belongs to. Fixed forever, computed from the workspace ID.
  • Physical: which host a shard lives on. Changed whenever you like.

With many small logical shards (say 480), growing from 32 to 96 hosts means moving whole shards between hosts, not re-hashing rows.

Why choose 480 logical shards rather than 500?

480 divides evenly by 32, 40, 48, 60, 80, 96 and more, so every host count stays balanced.

500 divides by 20, 25 and 50 but not 32 or 96, so some hosts would hold more shards than others.

480 logical shards on 96 hosts. How many shards per host?

About 5 shards.

480 ÷ 96 = 5, down from 15 on 32 hosts. Growth is a matter of moving 10 of each host's schemas elsewhere.

What the stage asks

How should shards map to hosts?

  1. Defensible

    32 shards, one per host

    Simple today. But growing to 40 or 96 hosts means splitting shards: rehashing workspaces to new shard numbers and moving their rows, which is another full data migration.

  2. Sound

    480 logical shards (one Postgres schema each), 15 per host on 32 hosts; growing means moving whole schemas to new hosts

    A workspace's logical shard never changes, so growth means moving schemas, not re-hashing rows. 480 divides evenly by 32, 40, 48, 60, 80, 96 and more, so hosts stay balanced at many sizes. This is exactly Notion's layout, and how they later grew to 96 hosts.

  3. Defensible

    16,384 logical shards spread over the 32 hosts

    Very flexible placement, but hundreds of schemas per host multiply catalog size, connection pools and migration tooling. Worth it with good automation; more than this team needs.

  4. Defensible

    Consistent hashing of workspace IDs directly onto hosts, so adding a host moves only a fraction of workspaces

    It limits how many workspaces move, but those that do still have their rows moved individually between live databases. Fixed logical shards turn the same growth into moving whole schemas with database replication, which is much easier to do safely.

What a strong answer covers

  • A workspace's logical shard never changes; only the shard-to-host map does.
  • Growth moves whole logical shards with database tools rather than re-sharding rows.
  • The shard count divides evenly into many host counts.Supporting

The reasoning

  1. Separate the fixed logical layout (row to shard) from the changeable physical one (shard to host).
  2. Growth then moves whole shards with database tools instead of re-sharding rows.
  3. Pick a shard count that divides evenly into many host counts.

Separate the logical layout (which shard a row belongs to, fixed forever) from the physical layout (which host a shard lives on, changed whenever you like). Then growth is a routing change plus a database copy, not a data re-shuffle.

Figma reached the same split by a different route: they sharded "logically" first, using views inside existing databases, so the application could be tested against the sharded layout before any data moved physically.

Tradeoffs

ChoiceGainsCosts
Many logical shards per hostGrowth by moving schemas; balanced at many host counts.More schemas to manage; routing must map shard to host.

Stage 4 of 9 · Change it

Move billions of rows, live

Notion captured writes with an application-level audit log rather than Postgres logical replication, because replicating the initial snapshot of tables this large through logical replication was too slow for them at the time.

What you need to know first

Moving live data has one shape, whatever the tools. See Online data migrations:

  1. Capture every new write (here, an audit log).
  2. Backfill existing rows.
  3. Catch up by replaying captured writes until the copy trails by seconds.
  4. Verify the copy against the original.
  5. Switch reads and writes, briefly.
  6. Keep the old store as a fallback for a while.

Why start capturing writes before the backfill begins?

So every write during the copy is recorded; nothing changes unseen between the copy and the switch.

If you copy first, writes made during the copy are missing, and you can't tell which ones.

What the stage asks

Put the migration steps in a safe order.

In this order

  1. 1Start logging every write to an audit log, alongside the normal write to the monolith
  2. 2Copy every existing row from the monolith to its shard, never overwriting a newer version
  3. 3Replay audit-log entries written since logging began, until the shards trail the monolith by seconds
  4. 4Verify: compare sampled rows, and run dark reads against both and alert on differences
  5. 5In a short maintenance window, stop writes, let catch-up finish, then switch reads and writes to the shards
  6. 6Keep the monolith for a while as a fallback, then retire it

Logging must start before the copy, so that any write the copy misses is in the log. The copy then has a fixed past to cover, and catch-up closes the gap. Verification happens while the monolith is still the truth, so a mismatch costs a fix, not an incident. Only then does the switch happen. See Online data migrations.

Notion's backfill ran on 96 CPUs for about three days; the cutover took five minutes of scheduled maintenance. They noted afterwards that a zero-downtime switch would have been possible with more work, which they did for the 2023 expansion.

The reasoning

  1. Capture writes, backfill, catch up, verify, switch, keep a fallback: in that order.
  2. Each step depends on the previous one being complete.
  3. A brief maintenance window is only for the final switch.

The order is the design. Each step depends on the one before it being complete: capture before copy, copy before catch-up, catch-up before verify, verify before switch. Every safe live migration published by Stripe, GitHub, Notion, Figma and Discord follows this shape, whatever the tools.

Stage 5 of 9 · Break it

Shards that are slightly in the past

Here is the backfill and catch-up code. Find every line that contributes to these three problems.

What you need to know first

Backfill and catch-up run at the same time, against the same rows. The backfill copies a row as it was when scanned; catch-up applies newer writes from the log. Whichever writes last wins, unless writes are conditional.

The rule that works is newest version wins: write only if the row is absent or the stored version is older.

Catch-up writes version 7 of a block to its shard. A moment later the backfill reaches the same block, which it scanned at version 5, and upserts it unconditionally. What's on the shard?

Version 5: the backfill overwrote a newer version with an older one. Nothing will correct it unless that block is edited again. A version check on every write prevents it, and makes replays harmless too.

Catch-up always starts reading the audit log from position 0. A deploy restarts it on day 4. What happens?

It replays four days of entries again: correct only if writes are versioned, and slow either way

Persist the position of the last applied entry and resume from it. Versioned writes make the overlap harmless.

The backfill also competes with production for the primary's CPU and disk. Throttle it, back off when production latency rises, or read from a replica or snapshot. See Backpressure and capacity.

What the stage asks

Select the faulty lines.

Codetypescriptmigrate.ts
  1. 1async function backfill(table: string) {
  2. 2 for await (const batch of monolith.scan(table, { batchSize: 5000 })) {

    Scans the primary as fast as it can, competing with production traffic. Throttle the copy (or read from a replica or snapshot) and back off when the primary is busy.

  3. 3 for (const row of batch) {
  4. 4 const shard = shardFor(row.space_id);
  5. 5 await shard.upsert(table, row, { onConflict: 'id', update: 'all' });

    If catch-up has already written a newer version of this row, the backfill's older copy overwrites it. Only write when the row is absent or the incoming version is newer (WHERE existing.version < incoming.version).

  6. 6 }
  7. 7 }
  8. 8}
  9. 9
  10. 10async function catchUp() {
  11. 11 for await (const entry of auditLog.readFrom(0)) {

    Always starts at the beginning of the log. Persist the position of the last applied entry and resume from it.

  12. 12 await applyToShard(entry); // same versioned write as the backfill
  13. 13 }
  14. 14}

What the fix has to do

  • The backfill can overwrite newer data written by catch-up; writes must be conditional on version.
  • The copy must be throttled so it does not degrade production.
  • Progress must be checkpointed so restarts resume.
  • Notes that replaying log entries must be idempotent, because restarts and overlaps replay some entries twice.Supporting

The reasoning

  1. Backfill and catch-up race on the same rows: write only if the stored version is older.
  2. Versioned writes also make replaying log entries idempotent.
  3. Throttle the copy and checkpoint progress so restarts resume instead of starting over.

Backfill and catch-up run concurrently against the same rows, so they need a rule for who wins, and "the last one to write" is wrong. "The newest version wins" is right, and it makes every write idempotent too: replaying an entry twice, or backfilling a row that catch-up already wrote, changes nothing.

The throttle and checkpoint bugs are about operating the migration rather than its correctness, but they decide whether it finishes in days or drags on for weeks. Discord's migration and Stripe's both describe the same priorities.

Stage 6 of 9 · Decide

Is the copy right?

Before the shards serve a single user, the team wants evidence that they match. Notion compared randomly sampled records and ran dark reads: queries sent to both databases with results compared, while users still got the monolith's answer. Stripe used GitHub's Scientist library for the same kind of comparison.

What you need to know first

Verification is Reconciliation between two stores: compare, explain every difference, fix its cause. Two techniques:

  • Sampled comparison: pick random rows and compare them field by field.
  • Dark reads: send real production queries to both stores, compare results, and still give users the old store's answer.

Row counts match for every table. Is the copy correct?

Not necessarily: counts catch missing rows, not stale ones.

The overwrite bug from the last stage leaves counts identical while content differs.

A comparison runs right after a user edits a block, and reports a mismatch. Is it a bug?

Probably not: catch-up trails by a few seconds, so the shard hasn't received the edit yet. Compare after a short delay, or re-check a mismatch before alerting.

What the stage asks

Which statements hold?

  1. Fails

    If row counts match for every table, the data matches.

    Counts catch missing rows, not stale ones. Every bug in the previous stage would leave counts identical.

  2. Holds

    Dark reads test the sharded data with real query patterns, without risking what users see.

    They run production queries against both stores and report differences, while users get the monolith's result. They also exercise the new routing code.

  3. Holds

    Comparing immediately after a write will report false mismatches while catch-up is a few seconds behind.

    Compare after a short delay, or re-check a mismatch before alerting. Notion's 2023 re-shard sampled with a brief delay for this reason.

  4. Fails

    With thorough unit tests on the migration code, verification against production data is unnecessary.

    Production data contains years of edge cases no fixture has: odd encodings, orphaned rows, rows from old schema versions. Verification is about the data, not the code.

The reasoning

  1. Matching counts catch missing rows, not stale ones; compare content.
  2. Dark reads test real query patterns against both stores without affecting users.
  3. Allow for catch-up lag before reporting a mismatch.

Verification is Reconciliation between two stores: compare, explain every difference, fix the cause. It only has value while the old store is still authoritative, which is why it sits before the cutover in every published migration plan.

Stage 7 of 9 · Change it

From 32 hosts to 96

Each of the 32 hosts holds 15 of the 480 logical shards.

What you need to know first

With fixed logical shards, adding hosts means moving whole schemas. Postgres logical replication can copy a schema to a new host and keep it in sync. Cutover is a brief pause at the connection poolers while replication drains, then a routing change.

No workspace changes shard, so no row-level migration code is needed.

Why not change the hash to spread workspaces over 96 shards?

Changing the hash moves almost every workspace: another full row-by-row migration.

The fixed logical layout exists so growth never requires this.

Each pooler held 50 connections per host to 32 hosts. With 96 hosts, how many connections per pooler?

About 4,800 connections.

50 × 96 = 4,800, up from 1,600. Multiplying hosts multiplies connections. Notion split its poolers into clusters, each in front of a subset of hosts.

What the stage asks

How do you add capacity?

  1. Flawed

    Change the hash so workspaces spread over 96 shards instead

    Changing the hash moves almost every workspace to a different shard: a full row-by-row migration, again. The logical layout exists so you never have to do this.

  2. Sound

    Add 64 hosts; move whole schemas so each host holds 5 instead of 15, using Postgres logical replication; switch each host over by briefly pausing traffic at the poolers

    No workspace changes logical shard; only the schema-to-host map changes. Database replication copies the schemas, a pause of about a second at the poolers lets replication finish, routing flips, and traffic resumes. Notion did this in 2023 and reported CPU and disk utilisation falling to about 20% at peak.

  3. Defensible

    Move each host to a bigger instance

    Quick if bigger instances exist, but it costs more per unit of capacity, can only be done once or twice, and does not let you separate the heaviest shards from each other.

  4. Defensible

    Add read replicas to every shard

    It relieves read load, but the hot resources here include write-driven CPU and disk on the primaries, which replicas do not reduce.

What a strong answer covers

  • The unit of movement is a whole logical shard, so no row changes shard.
  • Database replication can copy whole schemas, replacing custom backfill code.
  • Cutover is a brief pause at the poolers while replication drains, then a routing change.
  • Mentions keeping connection counts bounded as hosts multiply (e.g. splitting the poolers).Supporting

The reasoning

  1. Move whole logical shards between hosts; no row changes shard.
  2. Use database replication to copy schemas, then pause briefly at the poolers to cut over.
  3. Watch connection counts as hosts multiply.

The first migration bought the second one: because the logical layout was fixed, the expansion was "copy schemas with standard replication, flip routing". Notion reported two lessons from it worth stealing: build indexes after the initial sync (it cut their sync from three days to twelve hours), and watch connection counts, because every pooler connecting to three times as many hosts is three times the connections. They split PgBouncer into four clusters, each in front of 24 hosts.

Stage 8 of 9 · Change it

Queries that cross shards

Rows for different workspaces now live in different databases.

What you need to know first

Sharding makes every operation inside the shard key cheaper and every operation across keys harder:

  • Queries across shards become scatter-gathers, or need an index maintained on write.
  • Uniqueness across shards needs one place that owns it.
  • Transactions across shards become multi-step workflows.

Users are sharded by workspace. Each shard has a unique index on email. Are emails globally unique?

No: each index sees only its own shard, so two shards can each accept the same email.

Global uniqueness needs a single owner, such as an unsharded table keyed by email, or users sharded by email.

A user moves a page to a workspace on another host. Why can't that be one local transaction any more?

The two workspaces live in different databases, and a local transaction covers one database. The move becomes copy, switch, delete: a multi-step operation made safe with idempotent steps and a state machine, or a distributed transaction.

What the stage asks

Which statements hold?

  1. Holds

    A query filtered by one workspace ID touches exactly one shard.

    That is the point of the shard key: the router computes the shard from the workspace ID.

  2. Holds

    'Recent pages across all your workspaces' must query every shard holding one of the user's workspaces and merge the results.

    A scatter-gather. It is fine when it touches a few shards; for features that need many, maintain a separate index (per-user recent pages) updated on write.

  3. Fails

    A unique index on email in each shard's users table enforces globally unique emails.

    Each shard's index sees only its own rows, so two shards can each accept the same email. Global uniqueness needs one place that owns it, for example an unsharded table keyed by email, or users sharded by email.

  4. Fails

    Moving a page from one workspace to another can still be a single local transaction.

    The workspaces may be on different hosts. The move becomes a multi-step operation (copy, switch, delete) that must be made safe with idempotent steps and a state machine, or a distributed transaction.

The reasoning

  1. Inside the shard key, everything gets cheaper; across keys, it gets harder.
  2. Global uniqueness needs a single owner, not per-shard indexes.
  3. Cross-shard moves become multi-step, idempotent workflows.

Sharding is a trade: every operation inside the key gets cheaper, and every operation across keys gets harder. Joins, uniqueness and transactions that span shards need explicit designs: an index maintained on write, a separate owner for global constraints, a multi-step workflow. Figma built a query proxy (DBProxy) that parses SQL, routes single-shard queries and scatter-gathers the rest, and still asked teams to avoid cross-shard queries where they could.

Stage 9 of 9 · Defend it

Why not a distributed database?

Your interviewer: "This is months of custom work. Why not move to a distributed SQL database like CockroachDB, Spanner or Vitess, which shards for you?"

What you need to know first

Infrastructure choices are often about risk under a deadline rather than which technology is best in general. Moving to a new database means a full migration plus new operational expertise, while the wraparound deadline approaches.

A strong answer still concedes what the alternative offers.

What does a distributed SQL database (CockroachDB, Spanner) genuinely offer over hand-sharding?

Automatic rebalancing, cross-shard transactions, and no custom routing code

Those are real advantages. They matter most for new projects, bigger teams, or heavy cross-shard transactions.

What the stage asks

Make the case for hand-sharding Postgres here, concede what the alternative offers, and say when you would choose it.

Reference answer

Risk and time. The deadline is wraparound, which stops all writes. Moving to a new database is the same live migration plus learning a new system's failure modes, performance profile and operations at the same moment. Hand-sharding keeps every database a Postgres database the team already knows how to run, tune, back up and debug.

Compatibility. The application's queries, extensions and tooling keep working inside a shard. A distributed SQL database is "Postgres-compatible" up to a point, and finding those points mid-migration is costly.

What I concede. A distributed database rebalances automatically, supports cross-shard transactions, and removes the routing layer and migration tooling we have to own. Those are real; Figma listed the same alternatives and chose sharding Postgres for the same reasons of risk and timeline.

When I would choose it. A new system without legacy data, a team with time to learn it, or a product where cross-shard transactions are central (a ledger spanning many accounts) would make a distributed database the better call.

What a strong answer covers

  • Weighs risk and timeline: a new database means a full migration plus new operational expertise, under a hard deadline.
  • Values the team's existing Postgres expertise, tooling and extensions.
  • Concedes what a distributed database gives: automatic rebalancing, cross-shard transactions, less custom routing code.
  • Names when they would choose the distributed database (greenfield, bigger team, heavy cross-shard transactions).

The reasoning

  1. Weigh risk and timeline: a new database is a migration plus new expertise under a deadline.
  2. Value the team's existing Postgres knowledge and tooling.
  3. Concede what distributed SQL offers, and when you'd choose it.

Good infrastructure decisions are often about risk under a deadline, not about which technology is best in general. Both Notion and Figma made the "boring" choice and said why; that reasoning is what an interviewer wants to hear.

How you did

Now try it as an interview question

  • “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.”

The interview mode mixes stages from this and other investigations with concept recall and questions about your own projects.

Back to the last stage