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 …
LeaderboardSalaryAccount