Database Sharding & Partition Keys
One Customer, One Shard, One Outage
Test your architecture intuition: Pitch a 7-axis solution, survive two aggressive reviewer objections, and inspect the staff-level Teacher Gold Answer.
1. What It Is & Why It Exists
The Core Problem: The Vertical Scaling Wall
A single relational or document database instance is bound by physical hardware constraints: CPU core limits, memory bus bandwidth, NVMe IOPS limits, and operating system file descriptor maximums. When an application scales to tens of terabytes of data or hundreds of thousands of write queries per second:
- Financial Prohibitions: Renting extreme-tier instances (e.g., AWS
u-12tb1.112xlargewith 448 vCPUs and 12 TB RAM) costs about per node on demand ($109.20 an hour in us-east-1, Linux). - Connection & Thread Contention: Thousands of concurrent application worker threads exhaust database connection pools, spending more CPU time on lock arbitration and context switching than on query execution.
- Maintenance Paralysis: Executing standard database maintenance (e.g., PostgreSQL
VACUUM, B-Tree index rebuilds, schema migrations, or restoring point-in-time backups) on a monolithic 50 TB database can run for hours or days while competing with live traffic.
The Breakdown: Monolithic Database vs. Shared-Nothing Sharding
Concrete Proof: The Vertical Scaling Wall vs. Horizontal Sharding
Consider a global e-commerce datastore scaling from to and to :
| Architectural Dimension | Monolithic Single-Instance DB (Vertical Scale) | Shared-Nothing Sharded Fleet (16 Commodity Shards) |
|---|---|---|
| Physical Topology | Single massive instance (m6i.32xlarge, 128 vCPU, 512 GB RAM) | 16 mid-sized instances (m6i.2xlarge, 8 vCPU, 32 GB RAM) |
| Write IOPS Limit | 🚨 Capped by one instance's EBS limit and the volumes attached to it | 🛡️ Adds up across shards: each shard has its own volumes and its own instance EBS limit |
| Storage Capacity | 🚨 Bounded by what one instance can attach and one engine can maintain | 🛡️ Grows with the shard count |
| Blast Radius on Crash | 🚨 Total Global Blackout: 100% of users offline | 🟢 : of users unaffected |
| P99 Write Latency | 🚨 Rises with contention: one buffer pool, lock manager and log absorb all 120,000 writes/s | ⚡ Each shard absorbs about 1/16: writes/s |
| Monthly Cloud Cost | on-demand EC2 ($6.144/hour, us-east-1 Linux) | \approx \4,485/\text{month}16 \times $0.384/\text{hour}$): the same vCPU and RAM at the same price, because m6i prices scale linearly with size. Sharding buys headroom beyond the largest instance and fault isolation, not a discount |
The First-Principles Solution: Shared-Nothing Horizontal Partitioning
Database Sharding partitions a massive monolithic dataset horizontally across independent database servers (shards), where each shard runs on completely separate physical or virtual compute:
- Shared-Nothing Architecture: Shards share no physical RAM, CPU, or disk volumes, completely eliminating cross-node hardware contention.
- Deterministic Route Mapping: Incoming queries pass through a routing intelligence layer (smart client driver or gateway proxy) that inspects the query's Partition Key and routes the operation directly to the specific shard owning that key.
- Fault Domain Isolation: A catastrophic hardware crash or kernel panic on Shard 4 impacts only of keys, allowing shards to continue processing user traffic without interruption.
Synthesizing vector architecture diagram...
Both panels start with 100k workers sending about 120k QPS. In the "Broken Baseline" panel, all of it lands on one database: connections run out, CPU hits 100%, latency climbs, and even maintenance cannot run. In the "Production Standard" panel, a router computes hash(PK) mod 4 for each request and sends it to the one shard that owns that key; each shard has its own storage and CPU and shares nothing with the others, so each holds about a quarter of the keys (and carries about a quarter of the load only if no single key is hot), and one shard failing affects only a quarter of the keys. The catch: queries that do not include the partition key must ask every shard, so choose the key to match your most frequent query.
2. Core Mechanics & Algorithmic Architecture
Sharding Strategies Compared
Synthesizing vector architecture diagram...
1. Hash-Based Sharding
- Mathematical Mechanics: The partition key is hashed using a uniform hash function (non-cryptographic MurmurHash3 or xxHash; MD5, a cryptographic hash, also works and is what Kinesis uses), and mapped to a shard:
- Pros: Uniform key distribution; eliminates temporal clustering. It does not even out load: every request for one hot key still lands on the one shard that owns it.
- Cons: Range queries (
SELECT * WHERE timestamp BETWEEN t1 AND t2) cannot be routed to a single node and must be executed as Scatter-Gather operations across all shards. And changing undermod Nmoves most keys: going from 4 to 5 shards moves 80% of them, and even doubling moves 50%. Consistent hashing with many virtual nodes, or fixed slots with a directory, moves only about the new shard's share (the section Adding capacity: rings, slots and splits of Sharding, Hot Keys & Rebalancing).
2. Range-Based Sharding
- Mathematical Mechanics: Contiguous lexicographical or numerical key intervals map to dedicated shards:
- Pros: Range scans within an interval execute on a single physical shard.
- Cons: Monotonic Write Hotspotting. Using auto-increment IDs or timestamps routes 100% of concurrent writes to the newest shard, defeating horizontal scaling.
3. Directory-Based (Lookup Table) Sharding
- Mechanics: A centralized distributed coordinator (e.g., DynamoDB, etcd, ZooKeeper) maintains an explicit routing dictionary:
- Pros: Extreme flexibility. Individual massive enterprise tenants can be isolated onto dedicated physical hardware without altering cluster-wide hash functions.
- Cons: Introduces an extra network hop to query the directory service (mitigated by local client LRU caching).
Formal Mathematical Sizing Formulations
1. Minimum Cluster Shard Sizing Formula
To dimension the required number of physical shards ():
Where:
- is the total dataset size (e.g., 20 TB), and is the safe operational storage ceiling per shard (e.g., 2 TB).
- is peak write QPS, and is write throughput ceiling per node.
- is peak read QPS, and is read capacity per node.
2. Load Imbalance Factor (Coefficient of Variation)
To quantify partition key skew across shards:
Where is the request load on shard . Lower is better, and the alert threshold is your team's choice (for example as a target). CV is a spread, not the worst case, so also track the busiest shard's load divided by the average: that ratio is what hits a shard's ceiling first.
Step-by-Step Deterministic Routing Pipeline
Synthesizing vector architecture diagram...
In-Memory Shard Mapping Representation
In production proxies (like Vitess VTGate or Citus Coordinator), routing topologies are maintained in memory as sorted slot boundary tables:
| Keyspace ID Range | Physical Shard Host | Role | Network Address | Maximum Capacity |
|---|---|---|---|---|
0x0000 - 0x3FFF | shard_00 | Primary | 10.0.10.10:3306 | () |
0x4000 - 0x7FFF | shard_01 | Primary | 10.0.10.11:3306 | () |
0x8000 - 0xBFFF | shard_02 | Primary | 10.0.10.12:3306 | () |
0xC000 - 0xFFFF | shard_03 | Primary | 10.0.10.13:3306 | () |
Query Execution & Routing Trace Matrix
Under cluster configuration shards, partition key customer_id, and routing rule :
| Step # | Event / Input | In-Memory / Distributed State | Evaluation & Transition | Outcome / Query Latency |
|---|---|---|---|---|
| 1 | Targeted Point QuerySELECT * WHERE customer_id = 'c_102' AND order_id = 9918 | 4-shard cluster (); Routing table: {0: S0, 1: S1, 2: S2, 3: S3} | Hash evaluation (illustrative value): Modulo: | Single-Shard Hop () Dispatched directly to Shard 2 ( 10.0.10.12:3306); 0 impact on other 3 shards |
| 2 | Cross-Shard QuerySELECT * WHERE status = 'PENDING' ORDER BY created_at LIMIT 10 | 4-shard cluster; Predicate lacks partition key ( customer_id) | Router broadcasts query to all 4 shards in parallel; Shards 0, 1, 2 finish in , Shard 3 finishes in | Scatter-Gather Bounded () Router merges partial streams in RAM; latency capped by slowest tail node |
| 3 | Whale Tenant Mutationwhale_corp generates | Hotspot tenant threatens to exhaust Shard 1 storage and network IOPS | Salting applied: Generates 8 sub-keys, each hashed to a shard | Even only if placement is checked 8 sub-keys of give per shard only if they land exactly two per shard. Placed at random, 96.15% of the time some shard gets 3 or more (), so check where each lands or use a larger S. Every read of the whale's value now fans out to all 8 sub-keys (the section Hot writes, and every fix on equal terms of Sharding, Hot Keys & Rebalancing) |
3. Data Migration, Anti-Entropy & Consistency Protocols
1. Online Live Resharding Protocol (Zero-Downtime Cluster Expansion)
When a shard reaches capacity, splitting it (e.g., from ) without taking the database offline requires a strict 4-phase protocol (modeled after Vitess VReplication and DynamoDB partition splitting):
Synthesizing vector architecture diagram...
2. Cross-Shard Transactions: Two-Phase Commit (2PC) vs. Entity Co-Location
Executing atomic transactions spanning multiple shards requires Two-Phase Commit (2PC):
- Phase 1 (Prepare): Coordinator asks all participating shards if they can commit. Shards acquire local locks and write prepare records to their WAL.
- Phase 2 (Commit): If all shards answer "YES", coordinator writes commit to its log and commands shards to commit. If any shard fails, coordinator commands all shards to rollback.
- The Performance Cost: 2PC adds a second round of messages (prepare, then commit), each with a durable log write on the coordinator and the participants, and holds locks across those round trips; a coordinator crash between the phases leaves the participants holding their locks until it recovers.
- The Architectural Safeguard (Entity Co-Location): Design schemas so that related entities (e.g.,
User,Orders,Payments) share the exact same partition key (customer_id). All related mutations execute within a single shard, completely avoiding cross-shard 2PC!
3. Shard Draining Lifecycles & Anti-Entropy Resharding Verification
- Anti-Entropy Verification Before Cutover: Prior to advancing the routing epoch in Phase 4, background verification jobs execute cryptographic hash sampling across the child shards ( on Shard 1A vs. source shard). Only when CDC stream lag drops to zero and sample checksums match does the coordinator fence the source (it refuses reads and writes for the moving range with a retryable error), wait until the children have applied everything up to the fence, and then advance the cluster routing epoch.
- Graceful Shard Draining Protocol: When decommissioning a retired source shard:
- Keep the source shard refusing the moved range (the fence set before the epoch advanced), so a router with an old map is redirected instead of writing to a copy nobody reads.
- Maintain a graceful connection draining window () allowing in-flight transactions to conclude cleanly.
- Terminate idle client TCP connections and archive a final storage snapshot to Amazon S3 before de-provisioning instance hardware.
- Stale Topology Cache Invalidation: When an application pod caches an obsolete routing epoch and dispatches writes to a retired shard, the node returns
SHARD_MOVED <new_shard_id> <epoch>from a record it keeps permanently, which is what makes the move safe (timers and polls only bound how long a stale client bounces; see the section Moving live data of Sharding, Hot Keys & Rebalancing). Smart client drivers intercept this error, invalidate local memory routing maps, back off with jitter, and fetch the authoritative epoch topology from the coordinator.
Unlock Complete Architecture & Production Runbooks
You have explored the free architectural preview (~45%). Spend 1 Coin to unlock the remaining 7 production deep-dive sections for a full 24 hours.