What is database sharding, and why would a team split one big table across multiple database servers instead of just running a bigger single machine?
answer
- horizontal split of rows across servers
- shard key routes to a shard
- scales writes/storage ~linearly
- vertical scaling has a ceiling
- routing layer / router computes shard
basics
~20 sSharding splits a big table's rows across several separate database servers so no single machine has to hold or serve all the data. Each server (shard) keeps a slice of rows, chosen by a shard key like user_id.
solid answer
~50 sSharding is horizontal partitioning: instead of scaling one database server vertically (bigger CPU/RAM/disk), you split a table's rows across N independent database instances, each with its own storage and often its own compute. A shard key (e.g., user_id, tenant_id) decides which shard a row lives on. This lets write throughput and total storage scale roughly linearly with the number of shards, because writes for different keys hit different machines in parallel, and each shard only indexes and caches its own slice. The cost is operational complexity: the application or a routing layer has to know how to find the right shard, cross-shard queries and transactions stop being simple, and rebalancing as shards fill up is nontrivial. Teams reach for it when a single instance can't handle write volume or dataset size even after vertical scaling and read replicas.
go deeper
Should be able to state that sharding splits data across multiple database instances by some key, and that it's different from just buying a bigger server.
Should explain why vertical scaling and read replicas fall short for write throughput, and mention that a shard key and routing layer are required.
Should discuss when sharding is actually warranted versus premature, and name at least one concrete cost (cross-shard joins/transactions, migration difficulty).
Should frame it as an organizational and roadmap decision — the multi-week migration cost, the tooling investment, and how it compares to alternatives like better indexing, caching, or read replicas before committing a team to it.
## What sharding is Sharding is a data-management pattern for horizontally partitioning a dataset: instead of one database instance owning every row of a table, you split the rows across N separate instances called **shards**, and route each row to exactly one shard using a **shard key** — a column or combination of columns (like `user_id` or `tenant_id`) whose value determines placement. A **routing layer**, either built into the application, a driver, or a dedicated proxy, computes "which shard owns this key" on every read and write so the request lands on the correct instance. Mechanically this usually looks like: 1. the client sends a request with a key, 2. the router applies a mapping function (a hash, a range lookup, or a directory lookup — the specific strategies are covered elsewhere), 3. and the request is forwarded to the shard that function names. Each shard is otherwise a normal, independent database with its own storage engine, its own indexes, and its own transaction log; it has no built-in awareness of its sibling shards. ## Why the pattern exists The reason this pattern exists is that **vertical scaling** — buying a bigger box — has a ceiling. A single database instance can only accept so many writes per second before it's bottlenecked: - on **CPU** for query execution, - on **disk IOPS** for the write-ahead log, - or on **lock contention** as more concurrent transactions compete for the same rows and indexes. **Read replicas** solve read scaling by copying the full dataset to additional read-only instances, but they don't help at all with writes, because every replica still has to apply the same stream of writes that originated on the primary. When a product's write volume or total dataset size outgrows what one primary instance can hold or serve — a common story once you get into the hundreds of millions of rows or thousands of writes per second — sharding is one of the few remaining levers, because it multiplies both write capacity and storage capacity by spreading the load across genuinely separate machines rather than layering more reads onto one writer. ## The trade-off The trade-off is real and shows up immediately in day-to-day engineering. Anything that used to be a single-instance operation now potentially has to reach across shard boundaries, and cross-shard operations are slower, harder to make atomic, and sometimes simply not supported by the database engine at all: - a `JOIN` across two tables, - a multi-row ACID transaction, - an `ORDER BY` across the whole dataset, - a `COUNT(*)`. Schema migrations must run on every shard. Operational tooling (backups, monitoring, capacity planning) has to be shard-aware. And because shard assignment is baked into how data is stored, changing the sharding scheme later (adding shards, changing the key) requires migrating live data, which is a delicate, often multi-week project. None of this is free, which is why experienced teams treat sharding as a scaling tool of last resort, not a default architecture. ## Failure modes Failure modes in production tend to cluster around two things: uneven load and routing correctness. - **Uneven load.** If the shard key doesn't distribute traffic evenly, one shard becomes a "hot" shard that saturates while its siblings sit idle — this shows up as elevated latency or throttling on a subset of requests that all happen to hit the same key range or hash bucket, while dashboards for the overall cluster look healthy on average. - **Routing correctness.** The other common failure is a routing bug: if the application computes the wrong shard for a key (a bug in the hash function, a stale routing table after a resharding operation, or a race during a shard migration), writes can silently land on the wrong shard, and subsequent reads for that key return empty or stale results even though the data "exists" somewhere in the cluster. ## Where it shows up A concrete, widely cited example is how large social and e-commerce platforms partition their primary user or post tables once a single MySQL/PostgreSQL instance can no longer keep up — Instagram's early Postgres architecture famously sharded by a custom ID scheme embedding a logical shard number directly in generated primary keys, so the shard for any given row could be computed from its ID alone without a separate lookup. MongoDB and Vitess (used by YouTube and Slack for MySQL) provide sharding as a built-in cluster feature, where a router process (mongos, or Vitess's VTGate) transparently forwards each query to the shard(s) that can answer it. In every one of these systems the underlying motivation is the same: a single database process hit a hard ceiling on write throughput or storage, and horizontal partitioning was the mechanism chosen to push that ceiling out, in exchange for giving up the convenience of treating the whole dataset as one logically simple database.
- If sharding scales writes so well, why don't teams shard from day one?Because it adds real complexity before it's needed: routing logic, harder joins and transactions, and a schema that's expensive to change later. Most applications can absorb years of growth with vertical scaling, indexing, and read replicas alone, so sharding early is usually premature optimization that pays interest in engineering time long before it pays off in capacity.
- What's the difference between sharding and adding read replicas?Read replicas copy the entire dataset to extra instances that only serve reads, which scales read throughput but does nothing for write throughput since every replica still applies the same write stream. Sharding instead partitions the data itself across instances that each handle a subset of both reads and writes, which is what actually relieves a write bottleneck.
Like splitting one giant library into several branch libraries by the first letter of the author's last name — each branch only has to hold and manage its own slice of the collection, so no single building groans under the entire catalog, but finding a book when you don't know the author means checking every branch.
saying these in an interview costs you the question
- Says sharding and replication solve the same problem
- Thinks sharding just means adding more RAM/CPU to one box
- Never mentions a shard key or routing
- Assumes cross-shard joins work exactly like single-database joins
- No acknowledgment of the added operational complexity