Design a Database Sharding Strategy
Problem Explain database sharding — what it is, when it is warranted — and design a sharding strategy for a dataset that has outgrown a single database instance.
Functional requirements
- Horizontally partition a dataset across multiple database instances by a shard key.
- Route every read and write to the correct shard transparently to the application.
- Support adding shards (and rebalancing) as the dataset grows.
- Handle queries that must span shards.
- Preserve per-entity transactional guarantees.
Non-functional requirements
- ~10 TB and growing ~1 TB/month; ~80k writes/sec, beyond a single primary's ~10-15k/sec ceiling.
- Target ~500 GB and ~10k writes/sec per shard → ~20 shards initially, planned to 100.
- Single-shard query p99 < 20 ms; resharding online with < 1 min of write unavailability per key range.
- No more than 20% load skew between the hottest and median shard.
Areas to go deep
- Shard-key selection and its effect on skew, cross-shard query rate, and transaction locality.
- Range vs hash vs directory strategies, and how to rebalance online.
- Cross-shard joins/transactions and hot-key mitigation.
asked …