For a cloud-drive blob store, how do three-way replication and Reed-Solomon erasure coding compare on storage overhead and failure tolerance?
answer
- full copies vs fragments
- any k of k plus m
- overhead is (k+m)/k
- repair reads k fragments
- hot replicated, cold coded
basics
~20 sThree-way replication costs 3x raw storage and survives two lost copies, with simple reads and repairs. A Reed-Solomon (6,3) code costs 1.5x and survives any three lost fragments, but reads and repairs must touch several nodes and decode.
solid answer
~40 sReplication writes three full copies on separate failure domains: 3x raw overhead, tolerates two losses, any one replica serves a read, and repair copies one survivor. Reed-Solomon `RS(k, m)` splits an object into `k` data fragments plus `m` parity fragments, and any `k` of them rebuild it, so overhead is `(k + m) / k` and it survives `m` losses. `RS(6, 3)` is 1.5x and tolerates three losses; `RS(10, 4)` is 1.4x and tolerates four. For 1 PB of data, `RS(6, 3)` saves 1.5 PB of disk versus 3x replication. The price is multi-node reads, encode and decode CPU, and expensive repairs that read `k` fragments to rebuild one. Stores therefore usually replicate hot or small objects and erasure-code cold, large ones, placing fragments across racks or zones.
go deeper
Remember the two numbers: three copies cost 3x space, a 6+3 Reed-Solomon code costs 1.5x and survives three losses.
Explain the any-k-of-k-plus-m property, derive the overhead formula, and compare read and repair fan-out for both schemes.
Discuss placement across failure domains, degraded reads, repair traffic after disk failures, and why small or hot objects often stay replicated.
Treat the choice as a cost-versus-operability trade: when transcoding to erasure coding pays off, how repair load shapes the code chosen, and how margins shrink during a zone outage.
## Why blob stores need redundancy A **cloud-drive blob store** runs on thousands of disks, and at that scale disks, machines and whole racks fail every day. **Durability** - the promise that written bytes are never lost - comes from storing **redundant data** across independent **failure domains** (disks, hosts, racks, zones). The two standard schemes are **replication** and **erasure coding**. ## Three-way replication The store writes three full copies of each object on three different failure domains. - **Raw overhead:** 3x - storing 1 PB of user data consumes 3 PB of disk. - **Failure tolerance:** any two copies can be lost; one survivor is enough. - **Reads:** any single replica serves the whole object, so reads are fast and simple. - **Repair:** after a failure, one surviving replica is copied to a new node. ## Reed-Solomon erasure coding **Erasure coding** splits an object into `k` data fragments and computes `m` parity fragments, giving `k + m` fragments in total. A **Reed-Solomon** code has the property that *any* `k` of the `k + m` fragments are enough to rebuild the object. It is written `RS(k, m)`. - **Raw overhead:** `(k + m) / k`. For `RS(6, 3)` that is 9 / 6 = **1.5x**; for `RS(10, 4)` it is 14 / 10 = **1.4x**. - **Failure tolerance:** any `m` fragments can be lost - three for `RS(6, 3)`, four for `RS(10, 4)`. - **Reads:** a full healthy read fetches the `k` data fragments from `k` nodes; if a data fragment is missing, parity must be fetched and the object decoded. - **Repair:** rebuilding one lost fragment requires reading `k` surviving fragments over the network and recomputing - six fragment reads for `RS(6, 3)`, instead of one copy. ## Side by side | | 3x replication | RS(6, 3) | RS(10, 4) | |---|---|---|---| | Raw storage per 1 PB of data | 3 PB | 1.5 PB | 1.4 PB | | Losses survived | 2 | 3 | 4 | | Nodes touched by a full read | 1 | 6 | 10 | | Fragments read to repair one loss | 1 | 6 | 10 | | CPU cost | Negligible | Encode and decode | Encode and decode | The arithmetic is the headline: for 1 PB of data, moving from 3x replication to `RS(6, 3)` frees 1.5 PB of raw disk - half the footprint - while *raising* the number of tolerated losses from two to three. ## Placement matters as much as the code A code's tolerance only helps if fragments fail independently. Fragments of one object must sit on distinct disks and hosts, and ideally be spread across racks or zones. With `RS(6, 3)` placed three fragments per zone across three zones, losing an entire zone removes exactly three fragments - still decodable from the remaining six, though with no margin left until repair restores the missing fragments. Three replicas placed in the same rack, by contrast, share one failure domain and can all disappear together. ## How real systems combine them Because each scheme wins somewhere, blob stores commonly use both: 1. **Write hot data replicated.** Fresh uploads are read often and sometimes deleted soon after; replication keeps those writes and reads cheap. 2. **Transcode cold data to erasure coding.** A background job re-encodes objects that have stopped changing and are read rarely, reclaiming roughly half their raw space. 3. **Treat small objects specially.** Splitting a 4 KB file into nine fragments multiplies per-fragment bookkeeping for almost no saving; tiny objects are often replicated, or packed together into larger containers before encoding. 4. **Reduce repair fan-in.** Some systems use **locally repairable codes**, which add extra local parity so a single lost fragment is rebuilt from a small group rather than from all `k` - trading slightly higher overhead for much cheaper repairs. ## Terms worth defining - **Systematic code:** the data fragments are the original bytes, unchanged, so a healthy read just concatenates them without decoding. - **Degraded read:** a read while some fragments are unavailable; it pays decode cost and extra fetches. - **Repair traffic:** the network load generated when rebuilding lost fragments after a disk or node failure. ## In an interview Give the overhead formula, compute one concrete ratio, state the failure tolerance, and name the price erasure coding pays: multi-node reads, decode CPU and expensive repairs. Close with the hybrid - replicate hot data, erasure-code cold data.
- Why is repair more expensive with erasure coding than with replication?To rebuild one lost fragment of `RS(k, m)` the store reads `k` surviving fragments from `k` nodes and recomputes - six reads for `RS(6, 3)` - whereas replication copies one surviving replica. After a disk failure, that multiplies network and CPU load across the cluster. Locally repairable codes add local parity groups so most single losses are rebuilt from a few fragments, at a small extra storage cost.
- How would you place RS(6, 3) fragments so the store survives losing a whole zone?Put three of the nine fragments in each of three zones, each on a distinct host. A zone outage removes exactly three fragments, which the code tolerates, so data stays readable from the other six. The margin is then zero: a further loss before repair would make objects unreadable, so repair of the missing fragments must be prioritised while the zone is down.
Replication is keeping three photocopies of a document; erasure coding is tearing it into six strips plus three cleverly computed extra strips, where any six strips let you reassemble the page.
saying these in an interview costs you the question
- Erasure coding is just compression that makes files smaller.
- Three replicas cost about the same raw storage as a 6+3 code.
- Erasure coding is always better because it uses less space.
- A 6+3 code survives any three losses, so fragment placement does not matter.
- Rebuilding a lost fragment only needs one surviving fragment, like a replica.