System Design Guide

Fundamentals · 11

Sharding (partitioning)

Split a dataset across many machines so writes and storage scale out. Covers range, hash and directory sharding, choosing a shard key, hotspots and resharding.

3 min read · 6 flashcards

When one database can’t hold all the data or absorb all the writes, you shard: split the data horizontally into pieces (shards or partitions), each living on a different machine. Each shard is usually replicated as well.

App serverShard router(or client library)Shard 1users A–HShard 2users I–QShard 3users R–Zuser_id
A router sends each request to the shard that owns its key

Sharding strategies

Range-based

Each shard owns a contiguous range of keys (A–H, I–Q…, or by date).

Pros

  • Range queries stay on one shard (orders in March)
  • Easy to reason about and to split a range

Cons

  • Hotspots: sequential keys (timestamps, auto-increment IDs) send all new writes to the last shard
  • Uneven sizes unless ranges are rebalanced

Hash-based

Shard = hash(key) mod N, or better, consistent hashing.

Pros

  • Spreads keys and load evenly
  • No hotspot from sequential keys

Cons

  • Range queries must ask every shard
  • With plain mod N, changing N remaps almost every key, so use consistent hashing

Directory-based

A lookup service maps each key (or tenant) to a shard.

Pros

  • Fully flexible: move a big tenant to its own shard
  • Easy rebalancing by updating the map

Cons

  • The directory is an extra hop and must be highly available
  • Another component to keep consistent

Geo-sharding, splitting by region, is a variant: users’ data lives near them, which helps latency and data-residency laws.

Choosing a shard key

A good shard key:

  • Has high cardinality: many distinct values, so data can spread out.
  • Spreads load evenly: no single value gets most of the traffic.
  • Matches the main query: most requests can be answered by one shard.
Example Good key Why
Chat messages conversation_id A conversation is read together; there are many conversations
Multi-tenant SaaS tenant_id Queries stay within a tenant (watch for one giant tenant)
Time-series metrics (metric_id, time bucket) Hashing the metric avoids an “all writes to now” hotspot

What gets harder

  • Cross-shard queries: “top 10 posts overall” must ask every shard and merge (scatter–gather).
  • Joins across shards are slow, so denormalise, or co-locate related data with the same key.
  • Transactions across shards need two-phase commit or a saga, so design so that most operations touch one shard.
  • Resharding: moving data while serving traffic. Plan for it: use many small virtual shards mapped onto fewer machines, so growing means moving whole virtual shards rather than rehashing everything.
  • Unique IDs: auto-increment no longer works across shards. See unique IDs.

Test yourself

Answer in your head, then click a card to check. All cards are in the Anki deck.