Database Sharding: Horizontal Scaling Guide
588 words · Reviewed for accuracy

Sharding is what you reach for when one database server genuinely cannot hold your data or absorb your writes any more — and not a moment before. It splits one logical database across many machines, each holding a slice. The power is real; so is the pain. A good interview answer shows you understand both halves of that sentence.
Sharding defined: horizontal partitioning where each server (shard) owns a distinct subset of the rows, chosen by a shard key. Every query that can be answered from one shard stays fast; every query that spans shards gets expensive.
The shard key is the whole game
Pick the key that matches your dominant query pattern, because it determines whether most queries hit one shard or fan out to all of them. Designing a messaging app? Shard by conversation ID, because reads and writes cluster around conversations. Designing a multi-tenant SaaS? Shard by tenant ID. The wrong key turns every request into a scatter-gather across the fleet.
Three strategies, three trade-offs
| Strategy | How it routes | Watch out for |
|---|---|---|
| Range-based | IDs 1–1M on shard A, 1M–2M on B | Hot shards: new data all lands on the last range |
| Hash-based | hash(key) mod N decides the shard | Even spread, but adding a shard reshuffles everything — use consistent hashing to soften it |
| Directory-based | A lookup table maps key → shard | Flexible rebalancing, but the directory is a single point of failure to protect |
A worked example
Say you're designing an orders service and — purely for the sake of the exercise — assume 200 million orders a year with each order row at roughly 1 KB. That's about 200 GB a year of new data: one beefy server can hold it, but growth plus indexes plus replication gets uncomfortable fast, so sharding by customer ID is a reasonable call. Reads like "my order history" hit one shard. The cost appears when finance asks "total revenue by region yesterday" — that query now fans out to every shard and aggregates. In the interview, you'd propose an analytics pipeline for exactly that class of question rather than pretending the problem away.
What sharding breaks (say this part out loud)
- Joins across shards vanish. You denormalise or join in application code.
- Transactions spanning shards become distributed transactions — slow and fragile, so you design to avoid them.
- Rebalancing when a shard gets hot is an operation, not a config change. Consistent hashing and virtual shards make it survivable.
This is exactly the kind of trade-off naming that pairs well with the SQL vs NoSQL decision — many NoSQL stores shard natively, which is a legitimate point in their favour. One more operational note worth raising: once sharded, conveniences like auto-increment IDs stop working globally, so you need a separate ID-generation strategy. Interviewers probe exactly there.
Common mistakes
- Proposing sharding before exhausting the simpler options: indexes, read replicas, caching, vertical scaling.
- Choosing a shard key that doesn't match the query pattern, then hand-waving the fan-out.
- Forgetting resharding entirely. "What happens when shard 3 fills up?" is a guaranteed follow-up.
FAQ
Sharding vs replication? Replication copies the same data for read throughput and failover; sharding splits different data for capacity. Big systems use both.
When is it truly time to shard? When a single write-primary can't absorb the write rate, or data size outgrows one machine's practical limits — after caching and replicas are already in place.
Sharding slots into the scale-out phase of the system design spine. Rehearse the trade-off narration with Aissence mock interviews.
Put this into practice
Continue with the Aissence workflow this guide supports.