Fundamentals · 11
Sharding (partitioning)
Split a dataset across many machines so writes and storage scale out. Covers range, hash and directory sharding, choosing a shard key, hotspots and resharding.
When one database can’t hold all the data or absorb all the writes, you shard: split the data horizontally into pieces (shards or partitions), each living on a different machine. Each shard is usually replicated as well.
Sharding strategies
Range-based
Each shard owns a contiguous range of keys (A–H, I–Q…, or by date).
Pros
- Range queries stay on one shard (
orders in March) - Easy to reason about and to split a range
Cons
- Hotspots: sequential keys (timestamps, auto-increment IDs) send all new writes to the last shard
- Uneven sizes unless ranges are rebalanced
Hash-based
Shard = hash(key) mod N, or better, consistent hashing.
Pros
- Spreads keys and load evenly
- No hotspot from sequential keys
Cons
- Range queries must ask every shard
- With plain
mod N, changing N remaps almost every key, so use consistent hashing
Directory-based
A lookup service maps each key (or tenant) to a shard.
Pros
- Fully flexible: move a big tenant to its own shard
- Easy rebalancing by updating the map
Cons
- The directory is an extra hop and must be highly available
- Another component to keep consistent
Geo-sharding, splitting by region, is a variant: users’ data lives near them, which helps latency and data-residency laws.
Choosing a shard key
A good shard key:
- Has high cardinality: many distinct values, so data can spread out.
- Spreads load evenly: no single value gets most of the traffic.
- Matches the main query: most requests can be answered by one shard.
| Example | Good key | Why |
|---|---|---|
| Chat messages | conversation_id |
A conversation is read together; there are many conversations |
| Multi-tenant SaaS | tenant_id |
Queries stay within a tenant (watch for one giant tenant) |
| Time-series metrics | (metric_id, time bucket) |
Hashing the metric avoids an “all writes to now” hotspot |
What gets harder
- Cross-shard queries: “top 10 posts overall” must ask every shard and merge (scatter–gather).
- Joins across shards are slow, so denormalise, or co-locate related data with the same key.
- Transactions across shards need two-phase commit or a saga, so design so that most operations touch one shard.
- Resharding: moving data while serving traffic. Plan for it: use many small virtual shards mapped onto fewer machines, so growing means moving whole virtual shards rather than rehashing everything.
- Unique IDs: auto-increment no longer works across shards. See unique IDs.
Test yourself
Answer in your head, then click a card to check. All cards are in the Anki deck.