Database Sharding & Partitioning — Keys, Hotspots & Rebalancing
Vertical scaling and a single primary eventually hit CPU, IOPS, storage, or write-throughput walls. Sharding splits data across nodes for independent write capacity. Table partitioning splits one table inside an engine for prune and maintenance. This hub is the map for when each is enough, how a key avoids hotspots, why scatter-gather explodes cost, and how online resharding works.
- 1Gist
- 2Maps
- 3Q&A
- 4Sandbox
Voice readout needs Web Speech Synthesis in this browser.
Shard by user_id with modulo N
Prefer
Virtual shards or a ring, plus a high-cardinality key
Plan membership changes before the first modulo. Isolate mega-tenants. Keep the hot transaction on one shard.
- Adding a node moves a slice of keys, not most of the dataset.
- A whale can leave the shared hash bucket for a dedicated shard.
- The key still has to be stable and match the transaction boundary.
Alternative
hash(user_id) % N
Trivial routing and even spread when IDs are uniform. The bill arrives when N changes or one tenant dominates.
- Changing N reshuffles almost every key. Dual-write or downtime follows.
- A celebrity or whale still lands on one shard.
- Cross-user joins and global uniqueness become application problems.
Overview
Vertical scaling and a single primary eventually hit CPU, IOPS, storage, or write throughput. Sharding is horizontal partitioning across independent nodes. Table partitioning is a split inside one engine so the planner can prune and so maintenance can target one slice. Seniors use both, and they do not start with a custom shard layer because a famous company did.
This hub is the decision map. It does not re-teach page layout, snapshot visibility, or how you add a column under load. Ring placement and vnode theory live on Consistent hashing: rings, virtual nodes, and replicas. Additive schema, online DDL, and migration dual-write live on Zero-downtime database migrations. Use those pages for their subjects. Use this cluster when the question is where rows live and how you move them.
By the end of this hub you should be able to:
- Pick read replicas, in-engine partitioning, or application sharding from the measured ceiling
- Refuse a shard when there is no high-cardinality key or when most queries are global aggregations
- Say why
created_atas the only shard key creates a write hotspot - State scatter-gather as slowest-shard latency plus merge cost
- Describe an online reshard as dual-write, backfill, shadow read, then cutover — or as a range move on a ring
- Keep cross-shard uniqueness and mega-tenants out of a naive hash bucket
Decisions
- 1
Single DB under pressure?
- Reads, lag OKReplicas plus caching
- Bloat or pruningIn-engine partitioning
- Writes or isolationNatural shard key?
- 2
Replicas plus caching
- 3
In-engine partitioning
- ?
Natural shard key?
- High card, evenHash or composite shard
- Need localityRange shards, watch hot
- No good keyDefer shard, fix model
- 5
Hash or composite shard
- nextMinimize cross-shard joins
- 6
Range shards, watch hot
- nextMinimize cross-shard joins
- 7
Defer shard, fix model
- 8
Minimize cross-shard joins
- nextPlan reshard on day one
- 9
Plan reshard on day one
- nextIdempotent multi-shard writes
- 10
Idempotent multi-shard writes
Lesson map
Database Sharding & Partitioning — Keys, Hotspots & Rebalancing
Vertical scaling and a single primary eventually hit CPU, IOPS, storage, or write-throughput walls. Sharding splits data across nodes for independent write capacity. Table partitioning splits one table inside an engine for prune and maintenance. This hub is the map for when each is enough, how a key avoids hotspots, why scatter-gather explodes cost, and how online resharding works.
Architecture. Architecture
Select a node to see why it exists, or an edge to see the protocol, direction, effect, and consequence.
Mermaid export
flowchart TB a["Single DB under pressure?"] b["Replicas plus caching"] c["In-engine partitioning"] d["Natural shard key?"] a -->|Reads, lag OK| b a -->|Bloat or pruning| c a -->|Writes or| d
Default order
Measure the wall first. The key and the reshard plan come before the second cluster exists.
- 1
Name the ceiling
Reads with acceptable lag, table bloat and pruning, or writes, storage, and tenant isolation. Each answer picks a different tool. - 2
Partition inside the engine when that is the wall
Range, hash, list, or composite, so maintenance and pruning do not require a second database. - 3
Shard only with a key
High cardinality, even access, stable, aligned with the transaction. No key means fix the model or split read and write stores first. - 4
Keep the hot path on one shard
A query that must visit every shard is a distributed scan. Design it out of checkout. - 5
Write down the move
Dual-write and backfill, or an online range move. Idempotent writes. A saga or outbox when one transaction must span shards.
Why this shows up in senior interviews
Interviewers are checking whether you can:
- Explain partitioning versus sharding, and when each is enough
- Choose a shard key with high cardinality and even load, and say how you would detect a hotspot
- Say what breaks when a query must scatter-gather across all shards
- Reshard online (dual-write, shadow reads, consistent-hash ring moves) without a maintenance window
- Name the failure modes: orphans, partial commits, dual-write divergence, idempotency gaps
The question underneath is practical. If tomorrow's write QPS doubles and one tenant owns 40 percent of traffic, does the key survive, or do you page at 3am?
Decision rule
Partition inside one engine first when the bottleneck is table bloat or partition pruning. Shard when you need independent write capacity or hard tenant isolation. Do not shard because a case study did. Shard because you measured a write or storage ceiling and you have a key.
| Approach | Strengths | Weaknesses | Prefer when |
|---|---|---|---|
| Vertical scale | Simple ops, ACID intact | Hard ceiling, expensive | Early growth; the hot path still fits |
| Read replicas | Easy read scale | Does not help primary writes | Read-heavy; staleness is acceptable |
| Table partitioning | Pruning, maintenance windows | Still one primary write path | Huge tables, time or tenant prune |
| Application sharding | Independent write capacity | Cross-shard pain, reshard tax | Write wall plus a clear key |
| Distributed SQL | SQL plus auto range splits | Cost, latency, semantics shift | You need SQL and geo without a custom shard layer |
| CQRS / split stores | Isolates write and read models | Dual-write consistency | Analytics and OLTP diverge |
Distributed SQL here means engines such as CockroachDB, Spanner, or Yugabyte that split ranges for you. That does not delete the key decision. It moves split and merge into the database. You still pay for cross-range queries.
Cluster map
| Page | Focus | Previous | Next |
|---|---|---|---|
| Hub (this) | Strategy and when to shard | Failure modes | Partition strategies |
| Partition strategies | Range, hash, list, composite | Hub | Shard keys |
| Shard keys | Hotspots, fan-out, cardinality | Strategies | Cross-shard queries |
| Cross-shard queries | Scatter-gather and avoidance | Keys | Resharding |
| Resharding | Dual-write, online moves, rings | Cross-shard | Failure modes |
| Failure modes | Orphans, partial commits, idempotency | Resharding | Hub |
Adjacent only, and do not duplicate them: consistent hashing is the routing math for rings and vnodes. Zero-downtime migrations is expand/contract dual-write for schema, not for moving shard ranges.
Modulo N is the easy router
shard = hash(user_id) % N is easy to explain. Uniform IDs spread. Three costs show up later:
- Resharding almost always needs dual-write or downtime when N changes, because most keys change buckets.
- Hot users still concentrate on one shard. A celebrity or a whale tenant is a hotspot even when the hash is fair.
- Cross-user joins and global uniqueness become application problems. A unique constraint is per shard.
The better default is virtual shards, or a consistent-hash ring, from day one, a key with high cardinality, and dedicated shards for mega-tenants. The ring itself is the other study. Here you only need the database consequence: membership changes should move a fraction of keys, and the application must copy those keys.
Minimal routers
The Python sketch routes with a stable hash and with numeric ranges, then shows that repeating one key piles writes onto a single shard. The TypeScript sketch maps many virtual slots onto fewer physical shards so a later rebalance can move slots instead of recomputing hash % N. The browser sandbox cannot import Node's crypto module, so the TypeScript hash is a small integer mix. Production code should use one agreed function, such as SHA-256, on every app instance. Do not mix languages or library versions and expect the same bucket.
Press Run. Snippets must be self-contained — no network, files, or native modules.
Press Run. Snippets must be self-contained — no network, files, or native modules.
A hot key stays on one shard in both sketches. Virtual slots do not fix a celebrity. They fix the reshuffle when you add a physical database. Salt, a dedicated shard, or a split is the hotspot tool. That is the shard keys page.
Partition pruning versus an application shard
| Concern | In-engine partition | App-level shard |
|---|---|---|
| Transaction scope | Full ACID across partitions, usually | Single-shard ACID; cross-shard is a saga or 2PC |
| Ops | One cluster to back up | N clusters, N failovers |
| Hotspot fix | Split or merge partitions | Move key ranges or rehash |
| Query planner | Engine can prune | App must route or scatter |
"Usually" on ACID matters. Some engines still cannot enforce a global unique constraint across partitions. Know the engine before you promise one.
Interview Q&A
What is the difference between partitioning and sharding?
Answer
Partitioning splits a table inside one database so you can prune and maintain slices. Sharding splits data across independent databases or nodes so writes and storage scale out. Partitioning does not give you a second primary write path.
When should you not shard?
Answer
When vertical scale, replicas, caching, and better indexes still fit. When you lack a high-cardinality key. When most queries are global aggregations. Sharding those workloads turns every request into a scatter.
What properties make a good shard key?
Answer
High cardinality, even access, alignment with transaction boundaries, and stability. The key should not change under a row, and the hot path should include it so you do not fan out.
Why is created_at a bad sole shard key?
Answer
New writes pile onto the newest range. That range is a write hotspot. Older ranges go cold. Time is a fine archival partition inside a shard. It is a poor way to spread inserts.
What is scatter-gather?
Answer
The query fans out to many or all shards and merges the results. Latency is the slowest shard plus the merge. That dominates p99, and every such query multiplies load by the shard count.
How do you reshard without downtime?
Answer
Dual-write to old and new, backfill history, shadow-read to verify, cut reads over, then stop the dual-write. Or move ranges online on a ring. The steps are on the resharding page. Schema dual-write is a different lesson.
How do cross-shard unique constraints work?
Answer
They usually do not, inside the database. UNIQUE is per shard. Enforce uniqueness in the application or in a global allocator. A directory that maps email to user id is the same idea.
What do you do with a mega-tenant?
Answer
Give them a dedicated shard or an isolation pool. Do not leave one tenant in a shared hash bucket where they drown everyone else who hashed nearby.
Pitfalls
- Sharding a read-heavy database that only needed replicas and a cache.
- Choosing the key after the data is already in one modulo bucket.
- Treating an in-engine partition as extra write capacity.
- Promising a global unique index the engine cannot build.
- Leaving the reshard story as "we will figure it out at 10x."
- Copying ring math into this design instead of linking the hashing study, or copying expand/contract steps into a row move.
One tenant owns 40 percent of writes. The table is already partitioned by month on a single primary. Say whether you add replicas, a composite partition, or shards. Name the key, where the whale goes, and which query you refuse to run on the request path.
Go deeper
PostgreSQL's table partitioning docs are the in-engine vocabulary. Vitess sharding and TiDB horizontal scaling show how a SQL layer spreads writes. The AWS RDS sharding post is the application-router view. Kleppmann's Designing Data-Intensive Applications, chapter 6, is the partitioning chapter behind the tradeoffs on this page. ByteByteGo's sharding talk is a short visual pass.
Next: Partition strategies.