Skip to content
diff/reel
All reels
Databases

Replication vs Sharding

Replication vs Sharding — opening frame

sandboxed iframe · 45s loop · 27 KB

Made with Diffreel — draw your own →

Copy every write to every node, or send each write to exactly one.

By · Posted Aug 15, 2026 · 7 views

Both put more than one database node behind your application, for opposite reasons.

Replication copies the same write to every node. Every node holds the same data, so any of them can serve a read — reads scale, and a lost node loses no data. Writes do not get cheaper: every node still does all of them.

Sharding routes each write to one node by its shard key. Each node holds a slice, so total capacity grows with the node count — and a query that does not name the key has to ask every shard, while a lost shard takes its slice with it.

Watch the labels in the animation: identical on every node under replication, disjoint under sharding.

Source

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.

The reel's fan: one write on the left, three data nodes on the right.
The same three nodes in both scenes. Under replication every node is labelled A–Z and the write fires down all three lines at once; under sharding the labels read A–H / I–P / Q–Z and the write picks exactly one line.

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.

Three 500 GB nodes, total usable capacity
Replication500 GB
Sharding1,500 GB

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.

Sources

Related reels