What Is Database Sharding?
Problem What is sharding, and what are its benefits and trade-offs?
Be ready to discuss
- Definition: horizontally partitioning a dataset across multiple database instances ("shards"), each holding a disjoint subset of the rows rather than a copy of all of them.
- Shard key strategies: range-based (good for range scans, prone to hotspots on sequential keys), hash-based (even distribution, kills range queries), and directory/lookup-based (flexible, adds a lookup service and a single point of failure).
- Benefits: distributes read and write load across machines, scales storage past a single node's capacity, and shrinks the per-shard working set so indexes stay in memory.
- Trade-offs: cross-shard joins and transactions become expensive or need two-phase commit; global uniqueness and secondary indexes get hard; rebalancing after adding nodes is operationally painful (consistent hashing and virtual nodes help).
- Hot shards: a poorly chosen key concentrates traffic on one shard — celebrity users, a popular city, a monotonically increasing timestamp — and the fix usually means resharding.
- Sharding versus replication: sharding splits data for scale; replication copies the same data for availability and read scale. They are complementary, and production systems run both.
- Operational reality: sharding is a last resort after read replicas, caching, and vertical scale — worth saying out loud, since knowing when not to shard is part of the answer.
asked …