Data Partitioning (Sharding)

Splitting a dataset across multiple machines so no single node has to hold — or serve — all of it.

Intermediate · 18 min read

Why this matters

There is a scale at which "just buy a bigger database server" stops being a plan and becomes a fantasy. CPU, RAM, and disk on one machine are finite — and at some point the cost curve of bigger boxes goes vertical. Every database that serves the modern internet got there by doing the same thing: partitioning (also called sharding) — splitting the dataset across many machines so each one only holds and serves a slice of the whole. This is how your social feed, your order history, and your photo library all survive having billions of rows. The question is never whether to shard at scale; it's how, and which pain you're willing to accept.

Video: Database Optimization in System Design | Replication, CAP Theorem, Partitioning & Sharding Explained — Web Detailed by Mohi
Explains why databases become bottlenecks and how partitioning and sharding unlock performance at scale.

Picture a library system

Our analogy: a beloved city library that got too popular. The central branch is overflowing — shelves groaning, lines out the door. The city opens branch libraries and splits the collection across them. Every book lives at exactly one branch. Now, how do patrons find their books? That question is the entire lesson.

Video: Sharding, Partitioning & Indexing in Database| System Design | Episode 09 — The Tech Intern
Builds intuition with everyday pictures like a phone book at two counters, then walks range vs hash sharding.

Partitioning strategies

Each strategy is a trade-off, and real systems pick their poison deliberately. Cassandra leans on consistent hashing plus tunable replication; MongoDB's sharding is range/hash-based with a config server acting as the directory; DynamoDB hashes partition keys and rebalances behind the scenes. Same toolbox, different choices.

A range-based router sends each request straight to the shard that owns its key range:

flowchart LR
    APP[Application] --> ROUTER{Shard Router}
    ROUTER --> S1[("Shard 1: keys A-M")]
    ROUTER --> S2[("Shard 2: keys N-Z")]

Directory-based partitioning moves the mapping out of a formula and into its own service, which can be repointed without touching clients:

flowchart LR
    APP[Application] --> DIR[Directory Service]
    DIR --> S1[(Shard 1)]
    DIR --> S2[(Shard 2)]
    DIR --> S3[(Shard 3)]

And consistent hashing's payoff: adding a node only steals a slice of keys from its neighbor — every other shard is untouched:

flowchart TD
    N0["Before: Node B owns keys 34-66"] --> N1["Node D is added"]
    N1 --> N2["Node B now owns keys 34-58"]
    N1 --> N3["Node D now owns only keys 59-66"]

Interactive diagram: HashRingPlayground (loads in the app)

Video: Why One Database Is Never Enough — Partitioning Deep Dive (DDIA Ch. 6) — Savant Space
Compares key-range vs hash partitioning, hot spots, and consistent hashing for rebalancing, straight from DDIA.

The hard part: cross-shard operations

Sharding buys horizontal scale at a real price. The library is bigger now, but some things got harder:

  1. Joins across shards — "all orders joined with all users" — now require the application (or a query fan-out layer) to do work the database used to do for you.
  2. Transactions spanning shards need distributed transaction protocols (e.g., two-phase commit) or must be avoided by design. Many teams choose avoidance and model their data so cross-shard writes simply never happen.
  3. Rebalancing — adding a new shard means moving data. With consistent hashing you move a slice; without it, you re-shelve the whole library. Schedule it off-hours and bring snacks. The advanced treatment — virtual nodes, live resharding, and hot-key surgery — is in the sharding-strategies lesson.
  4. Hot shards — a celebrity's profile, a viral post: one key can get a million times the traffic of the average. Even a "perfectly even" hash can't save you from popularity. Teams handle this with caching, key splitting (the celebrity gets celebrity_1…celebrity_N keys), or moving hot data aside.

A scatter-gather query fans out to every shard that might hold a match, then merges the partial results back in the application:

sequenceDiagram
    participant App as Application
    participant S1 as Shard 1
    participant S2 as Shard 2

    App->>S1: Query matching rows
    App->>S2: Query matching rows
    S1-->>App: Partial result
    S2-->>App: Partial result
    Note over App: Merge partial results

Video: System Design | Sharding | System Design: How Distributed Databases Scale (Shards, Joins, & Skew) — CS AI Think Tinker Learn (unverified — confirm channel on YouTube)
Split one recipe book across four kitchens and every dinner becomes a phone tree — that's the cross-shard tax on joins, aggregations, and scatter-gather.

Picking a partition key

The partition key is the single most consequential decision in a sharded design — it's the rule that decides where everything lives. A good partition key:

Get this wrong and no amount of servers fixes it: the classic failure is sharding by country_code and discovering that one country's traffic equals everyone else's combined. Congratulations, you've built a single point of overload with extra steps.

When the constraints start fighting each other — a key with high cardinality that scatters related rows, a whale tenant that dwarfs everyone else — the sharding-strategies lesson picks up this thread with strategy selection under pressure: virtual nodes, salted keys, directory surgery for hot tenants, and the resharding cost math.

Video: What Is a Partition Key? Routing, Pruning, Hot Partitions, and How to Choose One — Kandi Brian
Walks through hash, range, and composite keys, hot partitions, and how cardinality shapes a good choice.

Takeaways

  1. Partition when one node can no longer hold the data or serve the load — not before. Sharding adds real operational complexity, and premature sharding is a tax you pay every day.
  2. The partition key is the most consequential decision in the design; it should match your queries, spread load evenly, and dodge hot keys.
  3. Every strategy trades something: range gives you scans but risks hot spots, hash gives you balance but punishes range queries, directory gives you flexibility but adds a critical dependency, and consistent hashing minimizes the pain of resizing.
  4. Cross-shard joins and transactions are the tax you pay for horizontal scale — design your data model to keep them rare.

Check your understanding

  1. Which partitioning strategy makes range queries expensive because they must fan out to every shard?

    • Directory-based
    • Range-based
    • None of the above
    • Hash-based
  2. What is the main risk of the directory-based (card catalog) strategy?

    • The directory becomes a critical dependency — it must itself be replicated, or its failure takes down all routing
    • Range queries become impossible to run
    • Data can never be rebalanced between shards
    • It only works for alphabetical or numeric keys
  3. What makes a partition key a *good* choice, according to the lesson?

    • It has low cardinality so the cluster needs fewer shards
    • It is the database's auto-incrementing primary key
    • It matches the query pattern, has high cardinality, and avoids monotonically increasing hot keys
    • It is a raw timestamp, keeping every write neatly ordered
  4. What technique minimizes data movement when a sharded cluster is resized?

    • A single central directory server holding the full mapping
    • Consistent hashing
    • Range-based partitioning by alphabetical order
    • Two-phase commit across all shards
  5. Why does a viral post create a 'hot shard' even when hashing is even on average?

    • Hashing silently breaks down once traffic passes a threshold
    • One key gets a wildly disproportionate share of traffic, overloading whichever shard owns it
    • The shard router stops forwarding to healthy shards
    • Two-phase commit locks the shard for the duration

Go deeper

Want to keep pulling this thread? These talks and tutorials go further than we did here:

Sources & further reading