Replication vs Sharding: the same data on every node, or a different slice on each
3 min read
The difference between replication and sharding is where a write lands: replication copies every write to every node, so each node holds a complete copy of the data, while sharding routes each write to exactly one node by its shard key, so each node holds a different subset. Both spread data across more than one machine, which is why beginners fold them into a single idea called "scaling the database" — but they do it for opposite reasons. Replication buys availability. Sharding buys capacity. Keeping them straight starts with refusing the one sentence that fits both and distinguishes neither.

Replication: the same copy, everywhere
In replication, writes go to one primary node, and the primary copies every write to each secondary. MongoDB's replica set is the canonical shape — a primary and N secondaries, all holding the same data set. Every node ends up with a complete copy, which is the whole point and also the thing people misread: total storage capacity does not change. Three nodes holding the same 500 GB store 500 GB, not 1.5 TB. Three replicas of a full disk are three full disks.
What you get for that redundancy is availability. Reads can be served from any copy, and if the primary fails a secondary is elected and takes over — with no slice of the data missing, because every node had all of it. Replication answers "what happens when a node dies", not "the data no longer fits".
Sharding: a different slice on each node
Sharding is horizontal partitioning of the data: you pick a shard key — say the first letter of a username — and partition its range across shards, A–H, I–P, Q–Z. Each write is routed to exactly one shard, the one that owns that key's range, so each shard holds a different subset of the whole. Now the storage adds up: three shards holding 500 GB each store 1.5 TB. Total capacity is additive, and that is the point.
The cost is that a shard failing loses that slice of the data rather than nothing — which is why, in production, each shard is itself a replica set. And the shard key is close to irreversible: choose a bad one and you get a hot shard that takes most of the traffic while the others idle, and resharding to fix it is expensive.
Illustrative. Replicas hold the same 500 GB three times; shards hold a different 500 GB each, so capacity adds up.
When to shard a database
Reach for replication first, and shard only when you must. The honest rule for when to shard a database: replicate for availability, and shard when one machine can no longer hold the data or absorb the write load. If your problem is "a node might die" or "reads should come from more than one machine", add replicas — that is the read replicas vs sharding decision, and read replicas win it whenever they can, because they are simpler, reversible, and do not force a shard key on you. Shard only when the data set or the write throughput genuinely exceeds a single node, because sharding is the harder, near-permanent commitment. Most systems replicate for years before they ever need to shard, and the teams that shard early usually regret the shard key they chose before they understood their access pattern.
Replication vs Sharding in a system design interview
An interviewer is checking whether you conflate the two, and the probe is usually a trap: "we're running out of disk — let's add a read replica." The crisp answer is that a replica adds no capacity, because every replica holds the full data set; you add capacity by sharding, which splits the data across nodes. Expect the follow-up, "so which do you use?" — and the strong answer names both. A real cluster shards for capacity and replicates each shard for availability, so it is not a choice but a composition. If you can say "replication is copy-to-all for availability, sharding is route-to-one for capacity, and production does both", you have said the thing they are listening for.
The cluster that does both
The trap in framing this as a versus is that the mature answer is "both". A production sharded cluster partitions the data across shards and makes each shard a replica set, so a single write is first routed to its shard and then copied across that shard's nodes — capacity from the partitioning, availability from the copies. The versus is how you learn the two mechanisms. The composition is how you run them.





