Level 2 · The Database Is the Bottleneck
Session 20: When One Database Isn't Enough: Sharding, Shard Keys, & Cross-Shard Queries
Imagine a single customer ledger that has grown 5 meters tall and is far too thick to flip through, so we split it into 10 separate books in 10 different cabinets. What suddenly becomes very easy, and what suddenly becomes an operational nightmare?
1. Imagine If
A bank records every account transaction for 100 million customers in one giant steel cabinet. The cabinet is packed to the brim, and the floor underneath is starting to crack from the sheer weight.
The bank's director decides to split the data across 4 separate cabinets: - Cabinet 1: Customers with names A–F - Cabinet 2: Customers with names G–L - Cabinet 3: Customers with names M–R - Cabinet 4: Customers with names S–Z
Looking up the balance of a customer named "Budi" is now 4 times faster. But what happens when "Budi" (Cabinet 1) transfers Rp 100,000 to "Zainal" (Cabinet 4)? The clerk has to lock both cabinets at once, write in book 1, run across the room to book 4, and if the lights go out halfway there, Budi's balance has gone down but Zainal's balance hasn't gone up yet!
2. What Actually Happens
When the write workload (Write Workload) pushes past the hardware limits of a single machine (disk I/O saturated, tables tens of TB in size), we're forced into Horizontal Partitioning / Sharding:
-
Choosing a Shard Key (Partition Key): - The shard key is the reference column that decides which database a given row is stored in:
\text{Node ID} = \text{Hash}(\text{user\_id}) \pmod N- Good Shard Key: Evenly distributed (for example: a UUIDuser_idortenant_id), so the IOPS load and storage capacity are balanced across all shards. -
The Deadly Danger: Hot Shard: - If the shard key is chosen by country (
country_code) or date (created_at):- 85% of users are from Indonesia \rightarrow the Indonesia Shard sits at 100% CPU/Disk, while the Singapore Shard is 2% busy and idle.
- A date shard key \rightarrow every write today slams into Today's Shard, leaving yesterday's shards cold with no load at all.
-
Scatter-Gather (Cross-Shard Query): - A query with the shard key (
WHERE user_id = 123) uses Single-Shard Routing (very fast, straight to 1 node). - A query without the shard key (SELECT * FROM orders WHERE status = 'PENDING') forces the proxy to send the query to ALL shards (scatter), wait for every shard to reply, then merge and sort in memory (gather). Latency shoots up!
3. What If We Try...
"Let's use Two-Phase Commit (2PC) for cross-shard transactions!"
Why is Two-Phase Commit (2PC) avoided at large scale? - Phase 1 (Prepare): The coordinator locks the rows in Shard A and Shard B, then asks "Are you ready to commit?". - Phase 2 (Commit): If both shards answer "Ready", the coordinator orders the permanent write. - Fatal Weakness: Throughout the 2PC process, data locks (locks) are held across the network. If network latency hits or the coordinator node crashes midway, the rows on both shards stay locked forever (blocking protocol), jamming up every other transaction. - Modern Solution: Replace 2PC with the Saga Pattern (Eventual Consistency & Compensating Transaction).
4. The Official Name
- Database Sharding: Splitting one logical dataset across many independent physical database nodes.
- Shard Key: The attribute/column that decides where a data partition is stored.
- Scatter-Gather: A query execution pattern that runs in parallel across all shards and merges the results.
- Hot Shard: An imbalance of data/traffic concentrated on one particular partition.
- Saga Pattern: A distributed transaction pattern built from a series of local transactions plus compensating actions (rollback logic).
5. In Our World
- NewSQL / Distributed SQL: Databases like CockroachDB, TiDB, YugabyteDB, Google Spanner that automate sharding at the engine level (Raft consensus), so developers don't have to manage routing by hand.
- Application-Level Sharding: Frameworks like Vitess (for MySQL) or Citus (for PostgreSQL) that act as a transparent coordinator proxy.
6. The Performance Tester's Lens
Crucial metrics and tests: 1. Shard Data & Traffic Skew: Measure the standard deviation of disk capacity and QPS across shards. A skew > 15% between shards signals a bad shard key. 2. Scatter-Gather Latency Penalty: Compare the p99 response time of queries with the Shard Key vs Cross-Shard queries. 3. Saga Compensating Latency: Test a failure scenario: make step 3 of a multi-shard transaction fail. How many seconds does it take for the compensating event to restore the user's balance to its original state?
7. Question for the Next Round
We now know how to split a database across 10 servers. But what if the data shape that's most efficient for writing transactions (Normalized OLTP) turns out to be the slowest, most painful shape for displaying analytical history (Denormalized Query)?
The answer is in Session 21: Separating the Read Path and the Write Path (CQRS & Materialized Views).