← All posts
ConceptsSeptember 21, 202614 min read

Database Sharding Explained: Shard Keys, Resharding & the Limits

Sharding is the only way to raise a database's write ceiling, and the only change on the scaling menu you cannot walk back. Once your rows live on eight machines that know nothing about each other, every join, transaction, unique constraint, secondary index and schema migration is permanently harder. That asymmetry is why the answer to "how would you scale this database?" is usually something else — and why the candidates who reach for sharding in the first thirty seconds of a system design interview lose points rather than gain them.

This post covers when the numbers actually justify sharding, why the shard key is the whole decision, the three ways to map a key to a shard and when each one is right, the logical-shard trick that makes resharding survivable, everything that breaks the day you split, and the alternatives that are often the better answer.

What one machine does before you shard

Start with the ceiling you are claiming to have hit, because an interviewer will ask for the number and most candidates do not have one.

A single well-tuned Postgres or MySQL instance on modern NVMe — say 64 vCPUs and 512 GB of RAM — handles somewhere in the region of 10,000 to 30,000 small write transactions per second, and holds several terabytes without drama. Managed instances top out in roughly the same place. That is a real ceiling, but it is much higher than most systems ever reach, and the workloads that hit it usually hit the storage limit first, not the throughput one.

Work an example. A messaging product takes 8,000 messages per second at peak, each row about 400 bytes with indexes:

8,000 writes/s × 400 B      = 3.2 MB/s
3.2 MB/s × 86,400 s         = 276 GB/day
276 GB/day × 365            ≈ 100 TB/year

Eight thousand writes per second is within a single instance's throughput. A hundred terabytes a year is not within any single instance's disk, and long before the disk fills, the operational facts get ugly: an index rebuild takes days, a restore takes longer, and VACUUM or its equivalent never finishes. That is the honest trigger for sharding — data volume and the operations that scale with it, far more often than raw write throughput.

Before sharding, five cheaper things exist, and it is worth knowing which problem each one actually solves:

RemedyFixes readsFixes writesFixes data volume
Fix the query or add an indexYesIndirectlyNo
Bigger instanceYesYesSomewhat
Read replicasYesNoNo
Cache in frontYesNoNo
Archive or tier cold rowsSomewhatNoYes
ShardYesYesYes

Read replicas are the one people get wrong, and the reason is structural rather than a tuning detail.

app tier 8k w/s · 12k r/s writes primary 8,000 writes/s all 8k, to each replica 1 8k w/s + 4k r/s replica 2 8k w/s + 4k r/s replica 3 8k w/s + 4k r/s reads · 12k/s split three ways Reads divide. 12k/s over three machines is 4k/s each. Writes do not. Every machine replays the whole 8k/s write stream. A fourth replica lowers reads per machine to 3k/s and lowers writes by zero.
Replication divides reads and multiplies writes. A replica is a full copy of the primary, so it applies every write the primary took. Adding machines to a replica set raises read capacity and total storage cost while leaving the write ceiling exactly where it was. Sharding is the only entry in the table above that moves it.

Partitioning, sharding, and the words people mix up

Three distinct things get called the same name, and defining them in one sentence at the start of an answer buys you credibility cheaply.

Vertical partitioning splits by column or by table group: move the 2 KB profile_bio out of the hot users row, or move the events table onto its own database. It is cheap, reversible, and buys a year or two. Figma's engineering write-ups describe doing exactly this for years before splitting a single table horizontally, and that ordering is the right instinct.

Horizontal partitioning splits by row, within one database — Postgres declarative partitioning, MySQL PARTITION BY. The rows land in different physical files managed by the same instance. Queries stay simple, planning gets faster, dropping old data becomes an instant DETACH PARTITION rather than a month-long DELETE. But one machine still serves all of it, so it fixes maintenance pain and nothing about capacity.

Sharding is horizontal partitioning across independent database servers that do not know about each other. Separate CPU, separate disk, separate connection pool, separate failure domain — and no query planner spanning them, which is where all the pain comes from.

If the only problem is that a table has grown unwieldy, partition it. Sharding is for when you need more machines.

The shard key is the whole decision

Everything downstream — query fan-out, hot spots, whether transactions stay local, how bad resharding will be — is determined by the column you shard on. Choose it badly and you get the cost of a distributed system with the capacity of a single one.

Four properties, in the order they matter:

  1. It appears in your highest-volume queries. This is the one that decides fan-out. If the key is not in the WHERE clause, the router cannot tell which shard holds the row and must ask all of them.
  2. High cardinality. user_id has millions of values; country has about two hundred and status has four. A low-cardinality key puts a floor under how finely you can ever split.
  3. Even access distribution, not just even row distribution. These come apart constantly: tenant_id may distribute rows evenly across ten thousand customers while one enterprise account generates 40% of the traffic.
  4. Stable. Changing a row's shard key means deleting it from one machine and inserting it on another — a cross-shard write with no transaction to protect it. Never shard on anything a user can edit.

Take the chat system from design a chat system. Three plausible keys, three very different systems:

SHARD KEY = message_id shard 0 conv 42: 3 shard 1 conv 42: 4 shard 2 conv 42: 2 shard 3 conv 42: 3 open thread conv 42 four round trips, merged and re-sorted in the app Rows spread perfectly evenly. Every read of the most common query touches every shard, so the latency of a thread open is the latency of the slowest of four. SHARD KEY = conversation_id shard 0 conv 42: 0 shard 1 conv 42: 12 shard 2 conv 42: 0 shard 3 conv 42: 0 open thread conv 42 one round trip, already ordered Same data, same machines, one column changed — and the hot path goes from four round trips to one.
The shard key decides your fan-out, and fan-out decides your p99. Sharding by message_id gives a textbook-even row distribution and a scatter-gather on the query that runs a thousand times a second. Sharding by conversation_id makes the hot path a single-shard lookup, at the cost of unevenness when one group chat is a thousand times busier than the median.

conversation_id wins, and the reasoning generalises: shard on the entity that your hot query already names. For a chat system that is the conversation; for an e-commerce order service it is usually the customer; for a B2B SaaS product it is the tenant. The candidates who arrive at message_id are optimising for the metric that is easy to measure (row balance) rather than the one that hurts (fan-out on the common path).

Be honest about the cost you just accepted. conversation_id means one enormous group chat lives entirely on one machine, and a busy tenant can saturate a shard by itself. The mitigations are worth naming: give outsized entities their own dedicated shard through a lookup table, or use a composite key such as (conversation_id, bucket) that splits only the top 0.1% of conversations into a handful of sub-partitions. Both are standard and both are things interviewers are hoping you will volunteer.

Three ways to map a key to a shard

StrategyHowBest atFails at
RangeContiguous key ranges per shard (A–F, G–M)Range scans, time-window queries, cheap splittingSequential keys: all new writes land on the newest shard
Hashhash(key) decides the shardEven write distribution, no coordinationRange scans and ORDER BY become scatter-gather
DirectoryAn explicit key-to-shard lookup tablePer-entity placement, moving one big tenantAn extra hop, and a lookup service you must keep available

The range trap is the one that shows up in production most often. Shard on created_at or a monotonic ID and every insert in the system goes to a single shard while the others sit idle — you have built a sharded system with the write capacity of one machine. If you need time-ordered data, shard on something else and partition within the shard by time.

The hash trap is quieter: shard = hash(key) % N bakes the node count into the routing function, so adding a node remaps nearly every key. That is the exact failure consistent hashing exists to solve, and it is worth naming explicitly rather than hand-waving — for a cache a remap is a miss storm, but for a database it means physically moving almost the entire dataset.

The directory approach is what large real migrations converge on, because it decouples the two decisions. Vitess calls the mapping a vindex and supports lookup vindexes for exactly this; the routing table is small, changes rarely, and every client caches it.

My default: hash into a fixed set of logical shards, with a directory mapping those logical shards to physical nodes. You get hash's even distribution and directory's operational freedom at the cost of one indirection.

Logical shards make resharding survivable

The trick is to decouple "which bucket does this key belong to" (fixed forever) from "which machine holds that bucket" (a small table you edit).

from hashlib import blake2b

BUCKETS = 1024  # chosen once, at design time, and never changed

def bucket_of(key: str) -> int:
    digest = blake2b(key.encode(), digest_size=8).digest()
    return int.from_bytes(digest, "big") % BUCKETS

def node_of(key: str, routing: dict[int, str]) -> str:
    # `routing` is a versioned bucket -> node table, cached by every client
    # and refreshed on change. Moving data edits this, not the hash above.
    return routing[bucket_of(key)]

With BUCKETS = 1024 on eight nodes, each node owns 128 buckets. Doubling to sixteen nodes moves 64 buckets off each old node — about half the data, unavoidably, because half of it has to end up somewhere new — but no key ever changes bucket, so nothing needs rehashing and no client needs new logic. You are copying rows and editing 512 rows in a routing table, not rewriting the addressing scheme while serving traffic.

Notion's 2021 write-up on sharding Postgres describes landing on 480 logical shards spread across 32 physical databases for exactly this reason: when they later needed more machines, they moved logical shards rather than recomputing anything.

Pick the bucket count generously. It caps your maximum node count forever, and 1,024 or 4,096 costs nothing at small scale.

bucket = hash(key) % 12 — the routing function, identical before and after node A b0 b3 b6 b9 node B b1 b4 b7 b10 node C b2 b5 b8 b11 not yet provisioned b6, b7 and b8 move — one bucket from each node node A b0 b3 b9 node B b1 b4 b10 node C b2 b5 b11 node D b6 b7 b8 Rows are copied and one table is edited. No key is rehashed, and no client changes.
The indirection is the whole point. Keys address buckets; buckets address machines. Scaling out becomes a data copy plus a routing-table edit, which is a job you can run gradually, verify, and roll back — rather than a cluster-wide rehash you have to get right on the first attempt.

The move itself is a four-step dance per bucket, and being able to describe it is what separates an answer that has thought about operations from one that has not:

  1. Backfill. Copy the bucket's rows to the target node while tailing the change stream (logical replication or CDC) so the copy keeps catching up.
  2. Verify. Compare row counts and checksums while both copies are live. This is the step people skip and regret.
  3. Cut over. Briefly reject or buffer writes for that one bucket — seconds, affecting 1/1024th of users — flip the routing entry, bump its version, and let clients pick it up.
  4. Soak, then drop. Keep the source copy read-only for a day so a rollback is a second routing edit, then delete it.

Vitess automates this pipeline with VReplication; MongoDB's balancer does the equivalent with chunks. If you are building it yourself, budget a quarter, not a sprint.

What breaks the day you shard

Cross-shard joins stop existing. There is no planner above the shards. You either denormalise so the joined data lives on the same shard, or you fetch from two shards and join in the application — which is a join without an optimiser, so watch the row counts.

Transactions stop being free. Two-phase commit across shards works and is slow: it holds locks across a network round trip and blocks if the coordinator dies mid-commit. The better move is to choose the shard key so your invariants live inside one shard. If money must move between two accounts on different shards, use a saga with compensating actions and an idempotency key, and say out loud that you have accepted eventual consistency in exchange for availability — which is the CAP theorem trade made concrete.

Secondary indexes become a design decision. An index on email when you shard by user_id is either local (present on every shard, so a lookup by email fans out to all of them) or global (a separate table, itself sharded by email, that maps to the owning shard — one extra hop, and eventually consistent with the base row). DynamoDB makes this choice explicit as LSI versus GSI; on a hand-rolled shard set you build it yourself.

Unique constraints and auto-increment IDs stop working. UNIQUE(email) is only unique within a shard. You need a separate uniqueness table on a single shard, or you fold the constraint into the shard key. For IDs, Snowflake-style generation is the standard answer: 41 bits of millisecond timestamp, 10 bits of node ID, 12 bits of per-millisecond sequence — 4,096 IDs per node per millisecond, roughly time-ordered, no coordination.

Pagination, sorting and aggregates get expensive. ORDER BY created_at LIMIT 20 across 10 shards means fetching 20 rows from each, merging 200, and discarding 180. Deep offsets are worse: OFFSET 10000 requires 10,020 rows from every shard.

Tail latency amplifies. This is the arithmetic worth memorising. If each shard serves a query under 50 ms 99% of the time, a request that touches one shard is slow 1% of the time. A request that fans out to ten shards and waits for all of them is fast only if every shard is fast: 0.99¹⁰ ≈ 0.904, so roughly 10% of requests exceed the single-shard p99. Fan-out turns a good p99 into a bad one, which is the real reason to fight for single-shard queries.

Availability multiplies down too. Ten shards at 99.9% each, on a request that needs all of them, gives 0.999¹⁰ ≈ 99.0% — from nine hours of downtime a year to eighty-eight. Each shard needs its own replicas, so your machine count is shards × replication factor, and your backup and restore story now has to be consistent across all of them.

Schema migrations multiply. Every ALTER TABLE runs N times, at different speeds, with a window where shards disagree about the schema. Every migration needs to be backward compatible for the duration.

What sharding does not solve

It does not reduce total work. A query doing a sequential scan still scans; you have just bought N machines to do it in parallel. Profile before you shard, or you will pay ten times as much for the same bad plan.

It does nothing for a single hot row. One viral post, one celebrity account, one global counter — all requests for that key land on one shard regardless of how you partition. The fixes are replication of that key, a cache in front, or request coalescing, exactly as with consistent hashing's hot-key problem.

It does not make single-row reads faster. Looking up one row by primary key was already an index seek in single-digit milliseconds. After sharding it is an index seek plus a routing decision. Sharding buys throughput and capacity, not latency.

It is not high availability. A shard is a single point of failure for its slice of the data. Without replicas underneath, sharding makes availability worse, not better — more machines, each of which can take out a portion of your users.

It does not fix analytics. SELECT count(*) ... GROUP BY across every shard is the worst possible query shape here. Ship changes to a column store such as ClickHouse or a warehouse, and leave the shards serving point queries.

The alternatives worth naming

Distributed SQL. CockroachDB, Spanner, TiDB, YugabyteDB and Aurora's distributed variants shard automatically by splitting key ranges, rebalance without you, and offer genuine cross-shard transactions. The cost is real — distributed commits mean extra round trips, and some query patterns get surprisingly slow — but you get to keep joins, transactions and a single schema. For a greenfield system that clearly will outgrow one machine, this is where I would start in 2026.

A database that was sharded from day one. Cassandra, DynamoDB and ScyllaDB make the partition key a schema decision you cannot avoid, which forces the shard-key thinking above before you have written any code. The trade you accept is the one described in SQL vs NoSQL: you design for your access patterns and you give up ad hoc querying. Note the ceiling still exists per partition — DynamoDB caps a single partition at roughly 3,000 read units and 1,000 write units per second, so a hot partition key is a hard limit rather than a slowdown.

Managed sharding on the database you already run. Vitess, and PlanetScale on top of it, give you MySQL semantics with a routing layer, online resharding and query rewriting. If you have an existing MySQL application, this is a far smaller step than a rewrite.

Not sharding. Genuinely the most underrated option, and the one interviewers respect when it is argued with numbers. Most large tables are 90% cold: move rows older than 90 days to object storage or a separate archive database and the active set fits on one machine again. Split the single hottest table onto its own instance. Put reads behind a cache and absorb write bursts with a queue. Buy the bigger instance — an engineer-quarter of sharding work costs more than several years of a larger box.

Answering this in an interview

The sequence that lands:

  1. Justify it with a number. "At 8,000 writes a second and 400-byte rows, we add 100 TB a year — that is a volume problem, not a throughput one, and it is why replicas and caching do not help."
  2. Name the shard key and the query that chose it. "Shard by conversation_id, because opening a thread is our highest-volume query and it names a conversation."
  3. State the mapping. "Hash into 1,024 logical buckets, with a routing table mapping buckets to physical nodes."
  4. Volunteer the skew. "Large group chats break this, so the top 0.1% get a composite key or a dedicated shard."
  5. Say what you gave up. "Cross-conversation search now fans out, so it goes to a separate search index. Cross-shard transactions do not exist, so anything spanning conversations becomes a saga."
  6. Have a resharding story. Backfill, verify, flip the routing entry, soak.

Steps 4 through 6 are where the level gets decided. Anyone can say "shard by user ID". The signal an interviewer is looking for is whether you know what the decision costs, and whether you would have tried to avoid it first.

If you want to hear how this sounds out loud, including the follow-ups about hot tenants and cross-shard transactions that always come next, you can run a voice mock interview on Whitepad and get scored on it.

Practice this out loud

Reading is the easy part. Sit across from a senior AI interviewer that talks, watches your whiteboard, and scores you like the real thing — your first mock is free.

Start a free mock →