Sharding & Distributed SQL

Reviewed & published by Brayan K

By the end of this lesson you'll be able to explain how a database is split across many servers, choose a shard key that spreads load evenly instead of creating a hotspot, work out which shard any given row lands on, and recognise when a managed distributed SQL engine is the right call instead of rolling your own. This is how systems serve billions of rows when one machine simply isn't enough.

Part of the free SQL course at LearnCodingFast — hands-on lessons with examples you run in your browser, plus practice exercises and a quick quiz.

What You'll Learn

1. Sharding vs Partitioning vs Replication

People use these three words interchangeably, but they solve different problems. Partitioning splits one table within a single server. Sharding splits rows across multiple servers. Replication copies the same data to several servers. Getting these straight is the whole foundation of the lesson.

Think of your data as paperwork. Partitioning is adding more drawers to one filing cabinet — still one cabinet, just better organised. Sharding is buying many cabinets and putting them in different rooms, each holding a slice of the files. Replication is making photocopies so every room has the same files for safety and faster reading. You shard when one cabinet physically can't hold (or be opened fast enough for) all the paper.

-- Three different ways to "split" a database. They are NOT the same thing.

-- 1) VERTICAL PARTITIONING — split COLUMNS, still one server.
--    Move rarely-used / huge columns into their own table.
CREATE TABLE users      (id, email, name);            -- hot, read constantly
CREATE TABLE user_blobs (id, avatar, bio_long_text);  -- cold, read rarely
-- Result: the hot table is smaller, so it caches better. One machine still.

-- 2) HORIZONTAL PARTITIONING (a.k.a. table partitioning) — split ROWS,
--    still one server. The DB stores ranges of rows in separate files.
--    e.g. PostgreSQL: orders_2024, orders_2025 under one "orders" table.
-- Result: one machine, but queries can skip whole partitions. One machine still.

-- 3) SHARDING — split ROWS across DIFFERENT SERVERS (machines).
--    Each server (a "shard") holds a slice of the rows and runs independently.
-- Server A: users 0–999      Server B: users 1000–1999    ...
-- Result: more machines = more CPU, RAM, disk, and write throughput.

-- Mental model:
-- Partitioning  = more drawers in ONE filing cabinet.
-- Sharding      = MANY filing cabinets, in different rooms.
-- Replication   = PHOTOCOPIES of the same cabinet (see below).

-- The concept sketch above is not runnable SQL. This is the real orders
-- table the later blocks query, so each snippet runs on its own.
CREATE TABLE orders (
    id          INTEGER PRIMARY KEY,
    user_id INTEGER,
    product_id  INTEGER,
    order_date  TEXT,
    quantity    INTEGER,
    total       REAL,
    status      TEXT
);

INSERT INTO orders (id, user_id, product_id, order_date, quantity, total, status) VALUES
    (1, 101, 1, '2026-01-14', 2,  49.98, 'shipped'),
    (2, 102, 3, '2026-01-22', 1,  79.00, 'shipped'),
    (3, 101, 4, '2026-02-03', 5,  16.25, 'pending'),
    (4, 103, 5, '2026-02-17', 1,  32.00, 'shipped'),
    (5, 104, 6, '2026-02-28', 3,  38.97, 'cancelled'),
    (6, 102, 2, '2026-03-05', 4,  38.00, 'pending'),
    (7, 105, 3, '2026-03-19', 2, 158.00, 'shipped'),
    (8, 101, 6, '2026-03-30', 1,  12.99, 'shipped');

Notice that the first two stay on one machine — they make queries cheaper but they can't add CPU, RAM, or write throughput. Sharding is the only one of the three that adds capacity, because each shard is a separate server doing its own work.

2. Replication Is a Different Axis

Before you reach for sharding, be sure you actually need it. Replication copies all your data to extra servers so more machines can answer reads and so you survive a server dying. It is far simpler than sharding — but it does nothing for write load or storage, because every copy holds everything. Sharding is what you use when the data itself is too big or too write-heavy for one machine.

-- REPLICATION is a different axis from sharding — don't confuse them.

-- REPLICATION = copy the SAME data to multiple servers.
-- Primary (writes) ──► Replica 1 (reads)
--                 └──► Replica 2 (reads)
-- Goal: more READ capacity + high availability (a replica can take over).
-- It does NOT add write capacity or storage — every copy holds ALL the data.

-- SHARDING = split DIFFERENT data across servers.
-- Shard 0 (users 0–999)   Shard 1 (users 1000–1999)   Shard 2 ...
-- Goal: more WRITE capacity + more storage — no single server holds it all.

-- Real systems use BOTH at once:
-- Shard 0 → Primary + 2 replicas   (this slice, copied 3x for safety)
-- Shard 1 → Primary + 2 replicas   (a different slice, also copied 3x)
-- So you scale writes with shards AND survive failures with replicas.

-- Rule of thumb:
-- Read-heavy and fits on one machine?  → add REPLICAS first (much simpler).
-- Write-heavy or too big for one disk? → you need SHARDING.

3. Hash vs Range Sharding

Once you've decided to shard, you need a rule that maps each row to a shard. That rule reads one column — the shard key — and decides the destination. There are two classic rules. Hash sharding runs the key through a hash function and takes the remainder modulo the number of shards, giving an even spread. Range sharding assigns contiguous ranges of the key to each shard, which keeps range scans fast but risks a hotspot on whichever shard holds the newest data.

-- Two ways to decide which shard a row lives on. Both use a "shard key".
-- A shard key is the column whose value picks the shard, e.g. user_id.

-- HASH SHARDING: shard = hash(shard_key) % number_of_shards
-- Spreads rows evenly, even if the keys are sequential.
-- shard = hash(user_id) % 4   (4 shards: 0, 1, 2, 3)
--   user_id 101 → hash(101)=...  % 4 = 1  → Shard 1
--   user_id 102 → hash(102)=...  % 4 = 2  → Shard 2
--   user_id 103 → hash(103)=...  % 4 = 3  → Shard 3
--   user_id 104 → hash(104)=...  % 4 = 0  → Shard 0
-- 👍 Even spread.   👎 "users 100–200" is scattered across ALL shards.

-- RANGE SHARDING: contiguous ranges of the key map to a shard.
--   user_id      1 –   999  → Shard 0
--   user_id   1000 –  1999  → Shard 1
--   user_id   2000 –  2999  → Shard 2
-- 👍 "users 1000–1500" all live on Shard 1 — fast range scans.
-- 👎 New sign-ups always get the HIGHEST id, so every new row hits the
--    LAST shard. That shard becomes a "hotspot" doing all the writes.

-- Quick test: pick HASH for even write spread; pick RANGE when you mostly
-- run "give me everything between X and Y" queries on the shard key.

Diagram: where does each user_id land?

Four shards, hash sharding with shard = hash(user_id) % 4. To keep the maths readable we pretend hash(user_id) = user_id, so the destination is just user_id % 4. Trace a couple of rows yourself before reading on.

See how consecutive IDs (100, 101, 102, 103) scatter across all four shards? That's exactly what gives hash sharding its even write spread — and exactly why "give me users 100–200" has to visit every shard.

Your Turn: compute the shard

Eight shards, hash sharding. Work out which shard user_id = 42 lands on using shard = hash(user_id) % 8 (pretend hash(x) = x). The answer is in the comments — and it deliberately shows you how to double-check your modulo arithmetic.

-- 🎯 YOUR TURN — work out the destination shard by hand (this is pseudo-SQL/maths).
-- Setup: 8 shards, numbered 0–7. Hash sharding: shard = hash(user_id) % 8.
-- For this exercise assume hash(user_id) = user_id (a "perfect" hash) to keep
-- the arithmetic simple.

-- user_id 42 lands on which shard?
-- shard = 42 % ___        -- 👉 fill in the number of shards
-- shard = ___             -- 👉 fill in the remainder of 42 divided by 8

-- ✅ Expected: 42 % 8 = 5, so user_id 42 → Shard 5.
--    (Check: 8 × 5 = 40, and 42 − 40 = 2 left over... wait — recompute!
--     8 × 5 = 40, remainder 2, so 42 % 8 = 2 → Shard 2. Always re-do the maths!)

4. Choosing a Good Shard Key

The shard key is the single most important decision you'll make, and it's painful to change later. A good key does two things at once: it spreads load evenly (so no shard becomes a bottleneck) and it co-locates data that's queried together (so your common queries stay on one shard). A key with too few distinct values — like status with only 'active'/'inactive' — can't spread across many shards. A key that always increases — like created_at or an auto-increment id under range sharding — sends every new write to the same shard.

Data that's queried together should live on the same shard. In a multi-tenant SaaS, shard by tenant_id so all of one company's rows sit on one server and "show me everything for tenant 7" is a single-shard query. Cross-tenant analytics then runs on a separate analytical copy, not on the live shards.

Your Turn: pick the shard key

A multi-tenant SaaS where almost every query is "for one tenant". Choose the shard key and justify why it keeps the hot query single-shard. The expected answer is in the comments.

-- 🎯 YOUR TURN — choose a shard key and justify it.
-- Workload: a multi-tenant SaaS. Every company ("tenant") only ever sees its
-- OWN data. 99% of queries are "give me <something> for tenant 7". You almost
-- never query across tenants in the live app.

-- A good shard key keeps data that is queried together ON THE SAME SHARD,
-- and spreads load EVENLY so no single shard is a hotspot.

-- Pick the shard key:
SELECT shard_for('___');   -- 👉 which column? (tenant_id / created_at / status)

-- Why is this the right choice? (fill in the blank in the comment)
-- Because all of a tenant's rows hash to ONE shard, so the common
-- "for tenant 7" query stays ___-shard (single / cross) — the fast kind.

-- ✅ Expected: shard key = tenant_id.  It co-locates each tenant's data on one
--    shard, so per-tenant queries are SINGLE-shard (fast). created_at would
--    hotspot today's writes; status has too few values to spread evenly.

5. The Hard Parts: Joins, Transactions, Fan-Out & Resharding

Sharding buys scale but charges a tax in complexity. Once rows live on different servers, four things that were trivial on one machine get hard:

-- The trade-off: some queries can no longer touch just one shard.

-- SINGLE-SHARD query (great): the WHERE filters on the shard key, so the
-- router knows EXACTLY which shard to ask.
SELECT * FROM orders WHERE user_id = 42;   -- → only Shard (hash(42) % N)

-- CROSS-SHARD / FAN-OUT query (expensive): no shard key in the WHERE, so the
-- router must ask EVERY shard and merge the answers. This is "scatter-gather".
SELECT COUNT(*) FROM orders WHERE status = 'pending';
-- Step 1 (scatter): send the query to all N shards in parallel.
-- Step 2: each shard returns its partial count.
-- Step 3 (gather): the coordinator adds the partials into one number.
-- With 16 shards that is 16 round-trips' worth of work for one answer.

-- Takeaway: design so your HOT queries include the shard key. Push rare
-- cross-shard reporting to a separate analytics copy, not the live shards.

-- ✅ Expected result:
-- COUNT(*)
-- 2

6. Distributed SQL Engines (So You Don't Build This Yourself)

Almost nobody hand-rolls sharding from scratch any more — the routing, resharding, and distributed transactions are too easy to get wrong. Instead you reach for a system that does it for you. At a high level there are two families:

The example below is real CockroachDB-flavoured SQL. The CREATE TABLE is standard and runs anywhere; the geo-pinning and follower-read lines are CockroachDB-specific and show what these engines hand you for free.

-- CockroachDB / YugabyteDB style: auto-sharded, auto-replicated, PostgreSQL-compatible.

-- A perfectly normal CREATE TABLE — the engine splits & distributes it for you.
CREATE TABLE products (
    id       UUID DEFAULT gen_random_uuid() PRIMARY KEY,  -- UUID: globally unique, no central counter
    name     VARCHAR(200) NOT NULL,
    category VARCHAR(50),
    price    DECIMAL(10,2),
    region   VARCHAR(20)
);
-- Rows are auto-split into "ranges" (~512 MB each) spread across nodes.

-- Pin EU rows to EU nodes (data-residency / GDPR), engine-specific syntax:
ALTER TABLE products
    CONFIGURE ZONE USING constraints = '[+region=eu]'
    WHERE region = 'eu';

-- "Follower reads": read from the NEAREST replica, accepting slightly stale data
-- in exchange for much lower latency. Great for read-heavy dashboards.
SET TRANSACTION AS OF SYSTEM TIME '-10s';
SELECT * FROM products WHERE category = 'electronics';

-- Distributed transaction — looks exactly like single-node PostgreSQL.
-- Under the hood it uses two-phase commit to stay ACID across nodes.
BEGIN;
UPDATE accounts SET balance = balance - 100 WHERE id = 1;  -- may live on node A
UPDATE accounts SET balance = balance + 100 WHERE id = 2;  -- may live on node B
COMMIT;

-- ✅ Expected: the CREATE TABLE runs in any SQL playground; the CONFIGURE ZONE
--    and "AS OF SYSTEM TIME" lines only work on CockroachDB-class engines.

Common Errors (and the fix)

Frequently Asked Questions

Q: When should I shard versus just adding a read replica?

If your problem is read load and the data still fits on one machine, add a replica — it's a far smaller change. Shard only when write throughput or total size genuinely exceeds a single server. Sharding is the heavier hammer; don't pick it up first.

Q: Hash or range sharding — which is the safe default?

Hash is the safer default because it spreads load evenly and resists hotspots. Choose range only when your dominant queries are "everything between X and Y" on the shard key (e.g. time-series scans) and you've accepted the newest-shard hotspot risk.

Q: Why is changing the number of shards (resharding) such a big deal?

With plain modulo, almost every row's destination changes when N changes — hash(id) % 4 is unrelated to hash(id) % 8 — so you'd move most of your data. Consistent hashing moves far less, and managed engines reshard online without downtime, which is a big reason to use them.

Q: Do I have to give up JOINs and transactions when I shard?

Not entirely, but they get harder. Keep related rows on the same shard (co-locate by shard key) and most joins and transactions stay single-shard and fast. For the rest, denormalise, use replicated reference tables, or pick a distributed SQL engine that handles cross-node joins and ACID for you.

Mini-Challenge: Shard a Chat App

Put it all together — a brief, a blank canvas, and the expected answer in the comments. Reason it out, then check yourself against the solution.

-- 🎯 MINI-CHALLENGE: design a sharding plan for a chat app
-- The app stores messages. The #1 query is "load the last 50 messages for
-- conversation X". You have 4 shards (0–3) and expect billions of messages.
--
-- Decide and write down (as SQL comments) FOUR things:
--   1. The shard key   (hint: what does every hot query filter on?)
--   2. hash or range   (hint: do you need even spread, or range scans?)
--   3. Which shard conversation_id = 1000 lands on, using hash % 4
--      (assume hash(1000) = 1000 for the maths)
--   4. ONE query that would be a SLOW cross-shard fan-out, and why
--
-- ✅ Expected (one valid answer):
--    1. shard key = conversation_id (every hot query filters on it)
--    2. hash (spreads busy and quiet conversations evenly; no hotspot)
--    3. 1000 % 4 = 0 → Shard 0
--    4. "count all messages sent today across the whole app" — no
--       conversation_id in the WHERE, so it must fan out to all 4 shards.

-- your plan here (as comments)

🎉 Lesson Complete

Practice quiz

What distinguishes sharding from partitioning?

  • They are the same thing
  • Sharding copies all data; partitioning splits it
  • Partitioning splits one table within a single server; sharding splits rows across multiple servers
  • Partitioning needs more servers than sharding

Answer: Partitioning splits one table within a single server; sharding splits rows across multiple servers. Partitioning stays on one machine; sharding splits rows across different servers, which is what adds capacity.

What does replication do, as opposed to sharding?

  • Copies the same data to multiple servers for read capacity and availability
  • Splits different data across servers
  • Adds write capacity
  • Removes the primary

Answer: Copies the same data to multiple servers for read capacity and availability. Replication copies all data to extra servers (reads + safety); sharding splits different data to add write capacity and storage.

What is a shard key?

  • A password for the shard
  • The number of shards
  • An encryption key
  • The column whose value decides which shard a row lives on

Answer: The column whose value decides which shard a row lives on. The shard key is the column (e.g. user_id) whose value the routing rule reads to pick the destination shard.

How does hash sharding decide a row's shard?

  • By contiguous ranges of the key
  • shard = hash(shard_key) % number_of_shards
  • Alphabetically by name
  • By insertion time

Answer: shard = hash(shard_key) % number_of_shards. Hash sharding runs the key through a hash and takes it modulo the shard count, giving an even spread.

What is the main downside of range sharding on a sequential key?

  • New sign-ups always hit the highest range, creating a hotspot on the last shard
  • It spreads rows too evenly
  • It can't do range scans
  • It requires hashing

Answer: New sign-ups always hit the highest range, creating a hotspot on the last shard. Range sharding keeps range scans fast but sends every new (highest-id) row to the same shard, making it a write hotspot.

What are the two goals of a good shard key?

  • Be short and unique
  • Be a timestamp and an integer
  • Spread load evenly and co-locate data that's queried together
  • Match the primary key exactly

Answer: Spread load evenly and co-locate data that's queried together. A good shard key spreads load so no shard is a bottleneck and keeps rows queried together on one shard.

Why is a low-cardinality column like status a poor shard key?

  • It changes too often
  • With only a few distinct values it can't spread across many shards
  • It is always NULL
  • It can't be hashed

Answer: With only a few distinct values it can't spread across many shards. A key with too few distinct values (e.g. active/inactive) gives a handful of giant shards instead of an even spread.

What is a cross-shard 'fan-out' (scatter-gather) query?

  • A query on one shard
  • A backup operation
  • A type of index
  • A query with no shard key in its WHERE that must ask every shard and merge results

Answer: A query with no shard key in its WHERE that must ask every shard and merge results. Without the shard key in the WHERE, the router must scatter the query to all shards and gather the partial results, getting slower as shards grow.

Why is resharding from 4 to 8 shards expensive with plain modulo?

  • Modulo is not allowed
  • hash(id) % 4 is unrelated to hash(id) % 8, so most rows must move
  • It deletes all data
  • It only works with range sharding

Answer: hash(id) % 4 is unrelated to hash(id) % 8, so most rows must move. Changing N reshuffles almost every row's destination; consistent hashing minimises movement and managed engines reshard online.

What protocol keeps an all-or-nothing transaction correct across two shards?

  • A single COMMIT
  • Read-your-writes
  • Two-phase commit
  • Consistent hashing

Answer: Two-phase commit. A cross-shard transaction needs two-phase commit: phase 1 every shard confirms it can commit, phase 2 the coordinator commits or all roll back.

Continue this course

Related lessons