Database Sharding
Concept
Read Replicas solve the Read problem. But what if your application is processing 50,000 INSERT statements per second (e.g., a high-frequency trading platform, or IoT sensor data)?
A single Primary Database physically cannot write to its hard drive that fast.
Database Sharding is the ultimate form of Horizontal Scaling.
It solves the Write problem by taking one massive logical table (e.g., 10 Billion Users) and physically slicing it into completely separate databases on completely different servers.
- Shard 1 (Server A): Stores Users with IDs 1 to 5 Billion.
- Shard 2 (Server B): Stores Users with IDs 5 Billion to 10 Billion.
Now, you have two completely independent servers handling INSERTs simultaneously. You have doubled your write capacity.
The Routing Problem
When the Node.js application executes SELECT * FROM users WHERE id = 6000000000, how does it know which database connection string to use?
You must introduce a Routing Layer (either inside your Node.js code, or via a specialized proxy server like Vitess or Citus).
The Routing Layer inspects the incoming query, looks at the id, consults a mathematical map, realizes the data lives on Shard 2, and forwards the SQL query strictly to Server B.
The Massive Drawbacks of Sharding
Sharding is the most complex operation in Software Engineering. It fundamentally breaks relational database theory.
1. The JOIN Problem
Imagine an orders table and a users table.
User 5 lives on Shard 1. Their Order lives on Shard 2.
You cannot execute SELECT * FROM users INNER JOIN orders. The database engine on Server A physically cannot reach across the network to access the hard drive on Server B. Cross-shard JOINs are mathematically impossible in standard SQL.
The Fix: You must fetch the User from Server A in Node.js, fetch the Order from Server B in Node.js, and manually “join” them together using Javascript map() functions in your application code.
2. The Transaction Problem
If you need to deduct 50 to User B (Shard 2), you cannot use a simple BEGIN; ... COMMIT;. Standard ACID transactions do not work across network boundaries. You must implement extremely complex, failure-prone algorithms like the Two-Phase Commit (2PC) or the Saga Pattern.
3. The Rebalancing Problem
If Shard 1 fills up its hard drive, but Shard 2 is mostly empty, you have to move data. Moving millions of active, live-updating rows from one physical server to another without causing downtime is a DevOps nightmare.
Interview Questions
Q: Explain the difference between Table Partitioning and Database Sharding.
A:
- Partitioning splits a table into smaller chunks, but all the chunks remain on the same physical server (same CPU, same RAM). It improves index performance and data management.
- Sharding splits a table and moves the chunks to completely different physical servers. It multiplies total CPU and Disk I/O capacity, but severely breaks SQL functionality (like JOINs and Foreign Keys).
Q: A startup’s database is slow. A junior developer immediately suggests Sharding. How do you respond?
A: You absolutely reject it. Sharding should be avoided at all costs until every other option is exhausted.
Before Sharding, a company should:
- Optimize slow queries with B-Tree Indexes.
- Implement robust RAM Caching (Redis).
- Scale the Database Vertically (buy a massive server).
- Implement Read Replicas.
Only when the Write throughput exceeds the physical limits of the largest available server on AWS should a company accept the architectural nightmare of Sharding.