Lesson 02

Sharding

Consistent hashing, range sharding, rebalancing, virtual nodes — splitting data across many machines.

Plays in the sticky player at the bottom of the page

Transcript

Lesson 2: Sharding

A reference-style deep dive into consistent hashing, range sharding, rebalancing, virtual nodes. Read this as an article, not a transcript. The accompanying audio is the spoken companion; the article below is the canonical written reference.

Audience. Engineers designing, building, or operating distributed systems who want a clear mental model rather than a checklist of tools.

Prerequisites. Working knowledge of HTTP, basic SQL, and the idea of running more than one server behind a load balancer.

Table of contents

  1. Why sharding matters
  2. The core mental model and where to start 3-12. In-depth sections below (see "Lesson body")

Lesson diagram

Lesson 2 diagram — Sharding

Figure 1. The canonical sharding topology and control flow covered in this lesson.


1. Why sharding matters

A single database cannot scale beyond the size of one machine. Sharding splits one logical dataset into many physical pieces so the system can grow horizontally. Sharding is the answer to "where does this row live?". Get it right and reads/writes distribute evenly; get it wrong and one shard becomes the bottleneck for the whole system.

Three motivations drive sharding decisions:

  • Dataset size — when a single machine's disk or memory cannot hold the data.
  • Write throughput — when a single machine's write IOPS cannot keep up.
  • Read throughput — when a single machine's read QPS cannot keep up (replicas often handle this).

The trade-off: sharding adds enormous complexity. Cross-shard queries are slow. Cross-shard transactions are very slow. Schema migrations become coordination nightmares. Do not shard until the single-database ceiling is provably hit.

2. The sharding key — the most important decision

The sharding key determines which shard each row lives on. Pick it wrong and the system falls over. Pick it right and the system scales linearly.

Properties of a good sharding key:

  • High cardinality — many distinct values, so rows spread across shards.
  • Uniform distribution — values are evenly distributed, so no shard becomes hot.
  • Stable — the value rarely changes (a sharding key change means moving the row).
  • Query-aligned — most queries filter by the sharding key, so the query hits one shard.

Examples:

-- Good: user_id
-- Cardinality = users, distribution = uniform (with hash), stable, most queries filter by user_id
SELECT * FROM orders WHERE user_id = 42;

-- Bad: status
-- Cardinality = 3-5 values, distribution = wildly skewed (most orders are "pending"), changes often
SELECT * FROM orders WHERE status = 'pending';

-- Acceptable: (tenant_id, user_id)
-- Composite key that ensures tenant isolation. Best for SaaS.

3. Sharding strategies

Four strategies, each with a different geometry:

Hash sharding

shard = hash(key) mod N

def shard_for(key, n_shards):
    return hash(key) % n_shards

Simple, uniform distribution. Worst property: adding a shard reshuffles nearly every row. With N=4 shards, adding a 5th moves ~80% of rows to new shards. Online resharding is impossible without downtime.

Range sharding

shard = which range does key fall in?

ranges = [
    (0, 10000, 0),       # keys 0-10000 → shard 0
    (10001, 20000, 1),   # keys 10001-20000 → shard 1
    (20001, 30000, 2),
    (30001, None, 3),    # keys 30001+ → shard 3
]
def shard_for(key):
    for low, high, shard in ranges:
        if (high is None or key <= high) and key >= low:
            return shard

Great for range queries (WHERE key BETWEEN 100 AND 200 hits one shard). Bad for uniform distribution: if your keys are timestamps and you write the latest, one shard gets all the writes.

Consistent hashing

Map both keys and shards to a ring; the key belongs to the next shard clockwise.

import hashlib
def ring_hash(s):
    return int(hashlib.md5(s.encode()).hexdigest(), 16)

def build_ring(shards, vnodes=100):
    ring = {}
    for shard in shards:
        for v in range(vnodes):
            ring[ring_hash(f"{shard}-{v}")] = shard
    return ring

def shard_for(key, ring):
    h = ring_hash(key)
    for node in sorted(ring.keys()):
        if h <= node:
            return ring[node]
    return ring[min(ring.keys())]

Adding a shard only moves ~1/N of keys. Removing one only moves the keys that hashed there. The cost: more complex routing code, harder to reason about which shard a key is on.

Directory-based

A lookup table that maps key ranges to shards.

CREATE TABLE shard_directory (
    range_start BIGINT,
    range_end   BIGINT,
    shard_id    INT
);

Most flexible: you can re-shard by editing the directory. Most operational overhead: the directory itself becomes a single point of failure, and updates to it must be coordinated.

4. Virtual nodes — making consistent hashing balanced

Plain consistent hashing has a flaw: with few shards, the distribution is uneven. One shard might own 40% of the ring while another owns 10%.

Solution: give each physical shard multiple slots on the ring (virtual nodes).

def build_ring_vnodes(shards, vnodes_per_shard=200):
    ring = {}
    for shard in shards:
        for v in range(vnodes_per_shard):
            ring[ring_hash(f"{shard}-{v}")] = shard
    return ring

With 200 virtual nodes per shard, the distribution is uniform to within 1-2%. The cost: the ring is 200x larger, but lookups are still O(log N).

5. Rebalancing — when one shard gets too big

Even with good sharding, some shards will get bigger than others (skewed data, hot keys). Rebalancing options:

Split a shard

When shard S grows past threshold, split it into S1 and S2. Move half the rows. Update the directory.

Add shards to an underutilized ring

With consistent hashing, just add the new shard. ~1/N of keys move automatically.

Online rebalancing

The hard part. Most production systems do it like this:

  1. Provision new shard, copy 1/N of data from existing shards.
  2. Start dual-writing (write to old + new shard).
  3. Migrate reads gradually (10%, 50%, 100%).
  4. Stop writing to old shard, delete it.

This takes hours to days. The system serves traffic the entire time. The failure modes are well known: missed dual-writes cause data loss; partial migrations cause inconsistent reads.

6. Cross-shard queries

The single biggest cost of sharding: a query that needs data from multiple shards.

-- On a single database: one query, one round-trip
SELECT * FROM users WHERE email LIKE '%@example.com';

-- Sharded: scatter-gather to every shard, then aggregate
results = []
for shard in all_shards:
    results += shard.query("SELECT * FROM users WHERE email LIKE '%@example.com'")
return merge(results)

Latency goes from O(1) to O(N). Throughput goes from "unlimited" to "bounded by the slowest shard". Mitigation strategies:

  • Avoid cross-shard queries. Denormalize, pre-aggregate, or restructure the schema so the common queries are shard-local.
  • Use a search index (Elasticsearch, Solr) that indexes data from all shards and serves the cross-shard queries.
  • Accept the cost. Some workloads genuinely need cross-shard queries; just budget for the latency.

7. Cross-shard transactions

Worse than queries. A transaction touching K shards needs distributed commit (2PC, Paxos, Raft), with K network round-trips and a coordinator that must survive failures.

# 2PC pseudocode
def transfer(from_user, to_user, amount):
    coordinator.begin()
    shard_A.prepare(from_user, debit=amount)   # round-trip
    shard_B.prepare(to_user, credit=amount)    # round-trip
    if all_prepared:
        coordinator.commit()
        shard_A.commit()
        shard_B.commit()
    else:
        coordinator.abort()
        shard_A.abort()
        shard_B.abort()

The cost is high enough that most sharded systems forbid cross-shard transactions entirely. The workarounds:

  • Saga pattern — break the transaction into local steps; compensate failures by running compensating actions.
  • Eventual consistency — accept that the multi-shard state will be temporarily inconsistent.
  • Co-location — put related data on the same shard so transactions stay local.

8. Schema migrations on a sharded database

Changing the schema on one database is a coordinated operation; on N shards, it is a nightmare. The standard playbook:

  1. Add the new column as nullable. Deploy the application code that writes both old and new columns.
  2. Backfill the new column for existing rows (a batch job that iterates every shard).
  3. Deploy the application code that reads the new column.
  4. Drop the old column.

Each step is a separate deploy. The whole migration takes weeks. The temptation is to skip steps — don't. Skipping steps is how production data gets corrupted.

9. Real-world sharding examples

MongoDB sharded cluster

Three components: mongos (router), config servers (store the shard map), shards (the data nodes). Sharding is hash-based or range-based. Rebalancing runs in the background.

sh.enableSharding("mydb")
sh.shardCollection("mydb.users", { _id: "hashed" })

PostgreSQL sharding (Citus)

Citus is an extension that turns Postgres into a distributed database. Sharding is hash-based; queries are distributed automatically when possible. Reference tables are replicated to all worker nodes for join locality.

SELECT create_distributed_table('events', 'user_id');

Vitess (YouTube's sharding layer for MySQL)

Vitess sits in front of MySQL and handles sharding, routing, and rebalancing. Used by YouTube, Slack, GitHub. Vindexes (Vitess indexes) define how rows map to shards.

10. When to shard — the decision

Rule of thumb: shard when one database can no longer handle the load, not before. The signals:

  • CPU saturation on the primary — can't handle write throughput.
  • Memory pressure — working set no longer fits in RAM.
  • Disk IOPS saturation — writes are queueing.
  • Replication lag — replicas cannot keep up with the master.

When you see these, you have three options before sharding:

  1. Vertical scale — bigger machine. Often the right first step.
  2. Read replicas — offload reads to async replicas. Cheaper than sharding.
  3. Functional split — split by table, not by row. Users go to one database, orders to another.

Only after all three are exhausted do you reach for sharding.

11. Operational concerns

Running a sharded system requires operational discipline:

  • Per-shard monitoring — each shard is its own failure domain; alert on each.
  • Per-shard capacity planning — uneven growth can leave some shards at 90% full and others at 20%.
  • Backups per shard — independent backup windows, independent restore procedures.
  • Failover per shard — each shard needs its own replica + failover story.
  • Migration tooling — scripts that move data between shards safely, with dry-run and verification.

12. Key takeaways

  • Sharding is the answer to "where does this row live?". Pick the sharding key carefully.
  • Hash sharding is uniform but reshuffles everything on rebalance.
  • Consistent hashing rebalances smoothly but is harder to reason about.
  • Range sharding is great for range queries but terrible if one range is hot.
  • Cross-shard queries and transactions are very expensive — avoid them when possible.
  • Operational discipline is the cost; budget for it before you shard.

Appendix: terms

  • Sharding key — the column used to determine which shard a row lives on.
  • Virtual node — a hash slot on the consistent-hash ring, owned by a physical shard.
  • Rebalancing — redistributing rows across shards to keep the system balanced.
  • Scatter-gather — a query that fans out to every shard and aggregates the results.
  • 2PC (two-phase commit) — a protocol for committing a transaction across multiple nodes.
  • Saga — a sequence of local transactions with compensating actions for failures.

Appendix: source dialogue excerpt

The audio for this lesson was synthesized from the following Cantonese dialogue (verbatim, not translated):

  • M: 各位同學早晨, 我係子謙。歡迎收聽系統架構課程第二課。今日嘅主題係 Sharding, 即係水平分區。…
  • F: 大家好, 我係曉晴。Sharding 係當一個 database 或者 storage 嘅容量或者 throughput 到咗 bottleneck 嘅時候, 將 data 拆分去多個 node 嘅技術。今日我哋會拆解 shard 嘅策略, 同埋 operation 上面嘅 rebalancing 同 hot spot…
  • M: 首先講解基本概念。Sharding 嘅核心係 partition key, 即係將 key 經過 hash function 之後, 對 shard count 取模, 然後 route 去對應嘅 shard。例如十六個 shard, key user 嘅 hash mod 16, 就決定去邊個 shard。Parti…
  • F: Sharding 同 replication 嘅分別, replication 係 copy 同一份 data 去多個 node 提升 availability 同 read scalability, sharding 係 split 唔同嘅 data 去唔同嘅 node 提升 write scalability 同 …
  • M: 好, 第一個 pattern 係 hash partitioning。即係 hash key 然後 modulo。簡單直接, 但係當 shard count 改變嘅時候, 即係 add shard 或者 remove shard, 大部分 key 嘅 mapping 會改變。例如由四個 shard 加到五個 shard…
  • F: Hash partitioning 嘅 rehash cost 係 O total keys, 即係每次 shard count 改變都要重新搬大部分 data。對於 production system 嘅 capacity planning 係 major problem, 因為 growth 過程入面 shard …

Full dialogue contains 33 segments; see script_raw.json in the source folder.

Lesson quiz · 30 questions

Question 1 of 30Answered 0 / 30
Question 1 of 30

Why is consistent hashing better than modulo sharding for rebalancing?

Pick an answer to lock it in. We'll tell you immediately whether you got it right and show an explanation. Then press Enter or click Next to continue.

Shortcuts:ABCDpick answer on current questionEntergo to next unanswered
30 unanswered