Sharding & partitioning · Key takeaways

1 min read
Senior15 min read
Rapid overview

Key takeaways

  • Shard for write throughput and storage; scale reads with replicas and caches first.
  • Hash for even spread, range for scans, directory for flexible placement.
  • Never hash mod N if N will change: consistent hashing with virtual nodes, or many fixed partitions.
  • Shard key = the field most queries filter on, with high cardinality and even load.
  • Hot keys: cache, replicate, split with suffixes, or isolate.
  • Resharding: copy, catch up, fence and cut over, verify, clean up.
  • Cross-shard queries: avoid by key choice, else scatter-gather, global indexes or denormalised copies.

See also