Distributed Systems · advanced

Sharding

Range/hash/geo partitioning.

distributed-systemsdatabasespartitioning

Mental model

Sharding splits one dataset across many machines so it scales past a single box. The shard key is the whole game: hash keys spread load evenly but kill range scans; range keys keep scans but invite hot spots. Re-sharding later is painful, so choose deliberately.

How to study Sharding

Begin by restating the mental model in your own words, then connect it to a concrete system you have built or operated. Name the mechanism, the constraint it addresses, and the trade-off it introduces. Use Designing Data-Intensive Applications (Kleppmann) — book site, Notion — Herding elephants: sharding Postgres, Vitess — Sharding to check details, but close the source before writing your explanation. Retrieval is the learning step; rereading is only preparation.

Next, compare Sharding with Replication. Ask what changes in correctness, latency, resource use, operability, and failure recovery. Complete Pick a shard key and preserve the command, input, output, and one failed attempt as evidence. Finish by explaining the idea without jargon to someone who has not studied the track.

Proof of understanding

  • Explain the mechanism from first principles and identify the state it reads or changes.
  • Give one situation where the concept is the right choice and one where it is not.
  • Predict a realistic failure mode before running the drill, then compare the prediction with evidence.
  • Connect the result to a roadmap or build artifact instead of treating the concept as isolated trivia.

Common mistakes

  • A shard key that creates a hot shard (e.g. sharding by date)
  • Cross-shard joins and transactions, which are slow or impossible
  • No plan for re-sharding as data grows

Learn from primary sources

Practice and explain it back

Pick a shard key

Multi-tenant SaaS: shard by tenant_id vs user_id vs hash(id). Which avoids hot tenant and cross-shard admin queries?

Expected evidence: tenant_id co-locates tenant data; hash spreads load; user_id splits tenant across shards.

Open the interactive drill →

Review prompts

  • Why is changing the shard key later so much harder than choosing a different index?

Build evidence

Synthesize: Distributed Systems

Reason about coordination, data placement, logs, durable workflows, consistency, and recovery under partial failure. Produce one working system, benchmark, or evidence-backed design that integrates the path.

  • Implements or precisely models the core mechanisms from all three milestones
  • Includes at least one injected failure or adversarial case and demonstrates recovery
  • Reports quality, latency, resource, reliability, or usability measurements relevant to the domain
  • Ships a concise architecture note explaining decisions, trade-offs, and remaining risks

Prerequisites

Related concepts

Learning paths