In this chapter
GreenMart splits its Orders table by city — one table per shard — and watches a single-city lookup stay effortless while a company-wide report reveals its real cost: EXPLAIN QUERY PLAN shows three separate scans and three separate sorts, merged only at the very end.
The Problem in Real Life
GreenMart is opening in city after city now — ten of them, and counting. Sarah watches the single Orders table grow past anything the original design ever imagined, and starts wondering whether one table can really hold an entire country's worth of orders forever.
Mike asks the obvious question: "Could we just... split it up? Give each city its own slice?" Sarah realizes that's exactly the right instinct — and exactly the tradeoff every large system eventually has to make.
Could we just give each city its own slice of the data?
Mike
One Table for Everyone vs. One Slice Per City
One giant table isn't the only option
Splitting data by a natural key — like city — is a standard way to keep individual queries fast as a dataset grows.
Single-shard queries stay genuinely isolated
A query about one city never has to know or care how large any other city's data has grown.
A view can hide sharding, not eliminate its cost
UNION ALL across every shard makes cross-shard queries convenient to write, but the underlying work still happens.
Cross-shard reports do real extra work
Each shard needs its own scan and its own sort before the results can be merged together.
What Sharding Actually Trades Away
Horizontal partitioning (often called sharding when the pieces live on entirely separate servers) means splitting one logical dataset across multiple physical tables or databases, usually by some natural key — here, by city. Each shard is smaller, faster to query on its own, and can even live on completely separate hardware in a real large-scale system.
The tradeoff is real: a question about one city gets faster and simpler, because it only ever touches that city's own small slice. A question about every city gets harder, because now something has to reach into every shard and combine the results — work that a single unsplit table would have handled with one straightforward query.
With each city sharded into its own table, a query about Mumbai only ever touches Mumbai's own small table — completely unaffected by how large Delhi's or Bangalore's tables get.
SELECT * FROM Orders_Mumbai WHERE CustomerID = 501;
This is the entire appeal of sharding: most real queries in a multi-city business are about one city at a time, and those stay simple and fast no matter how many other cities exist.
"How many orders has GreenMart taken, total?" used to be one COUNT(*) against one table. Now it means reaching into every single shard and adding the results together by hand.
SELECT(SELECT COUNT(*) FROM Orders_Mumbai) +(SELECT COUNT(*) FROM Orders_Delhi) +(SELECT COUNT(*) FROM Orders_Bangalore) AS TotalOrdersAllCities;
This is the real cost of sharding: it optimizes for the common case (one city) at the expense of the less common one (everywhere at once).
A UNION ALL view can combine every shard into one queryable name, so most reporting code doesn't need to know or care how many cities actually exist underneath it.
CREATE VIEW AllCitiesOrders ASSELECT 'Mumbai' AS City, OrderID, CustomerID, OrderDate, Total FROM Orders_MumbaiUNION ALLSELECT 'Delhi', OrderID, CustomerID, OrderDate, Total FROM Orders_DelhiUNION ALLSELECT 'Bangalore', OrderID, CustomerID, OrderDate, Total FROM Orders_Bangalore;
Convenient — but the view doesn't make the underlying work disappear, it just hides where it's happening.
"The 3 most recent orders, company-wide" looks like a simple ORDER BY ... LIMIT. Underneath the view, it isn't.
EXPLAIN QUERY PLANSELECT City, OrderID, OrderDate, Total FROM AllCitiesOrdersORDER BY OrderDate DESCLIMIT 3;
The plan shows a MERGE (UNION ALL) combining all three shards — and each shard runs its own separate SCAN plus its own USE TEMP B-TREE FOR ORDER BY, independently, before the results are merged together at the very end. One query became three scans and three sorts.
Notice exactly what sharding bought GreenMart and what it cost: single-city lookups are now genuinely isolated and fast, but a report spanning every city does real extra work — one scan and one sort per shard, merged only at the end, instead of one scan and one sort against a single table.
This is a fundamental, unavoidable tradeoff in real distributed systems, not something a cleverer query can fully undo. The right shard key is the one that matches how the business actually asks questions most of the time — and for GreenMart, most real questions genuinely are about one city.
Key Takeaway
Sharding trades a harder company-wide question for an easier single-shard one — a cross-shard report has to scan and sort every shard independently before merging the results, real extra work that a single unsplit table never needed.
Why This Matters
Every large-scale system eventually faces this exact tradeoff — the right partitioning key can make routine queries dramatically faster, but the wrong one (or ignoring the cost of cross-shard reporting entirely) can quietly turn simple company-wide questions into expensive, ever-slower operations.
GreenMart now has a real plan for splitting data by city, with eyes open about what gets faster and what gets harder. The checkpoint ahead puts this to work: a new city joins the map, without disturbing anyone already on it.
