System Design Interview

Course Content

System Design Interview

31 sections · 71 lessons

Object Storage deep dive: write and read paths, and the durability arithmetic


The first lesson drew the architecture: an API service, a sharded metadata store, a placement service and a fleet of simple data nodes. This lesson puts requests through it, then does the arithmetic that justifies the 1.5× storage multiplier the capacity estimate assumed.

Trace both paths end to end, because "walk me through a put" is the most likely follow-up.

A put, end to endPUT withkey and bytesAPI serviceauthorisesSplitinto chunksErasurecode and placeMetadatarow committedMetadata is written last, so a half-written object is never visible.
Committing metadata after the data is what gives read-after-write on new objects with no coordination at all.

The write path

  1. Authenticate and authorise. Verify the request signature, then check the bucket policy. Rejecting here costs nothing.
  2. Ask the placement service where the chunks go. For a 1 MB object with erasure coding at 6 data plus 3 parity fragments, that is nine placements, chosen across distinct racks and power domains.
  3. Split and encode. Chunk the object — typically a few megabytes per chunk for large objects, a single chunk for small ones — and compute the parity fragments.
  4. Write fragments in parallel to the nine data nodes. Each node writes to disk, computes a checksum, and acknowledges.
  5. Wait for enough acknowledgements. With 6 data plus 3 parity you can technically declare success at 6, but that leaves no redundancy. Requiring 8 of 9 tolerates one slow node while keeping a spare, and a background repair completes the ninth.
  6. Write the metadata record, including the object's checksum and the fragment locations, as one atomic write. This is the commit point.
  7. Return 200 with the object's entity tag, which is its content checksum.

Failure at any step before 6 leaves no visible object.

The read path

  1. Authenticate and authorise.
  2. Look up the metadata record to obtain size, checksum, and fragment locations.
  3. Fetch fragments. With erasure coding, request more than the minimum — say 7 of 9 — and use the first 6 that arrive, cancelling the rest. This is hedged reading, and it converts a slow node from a latency spike into a non-event.
  4. Reconstruct if necessary, verify the checksum against the metadata, and stream bytes to the client.

Verifying the checksum on read is what turns silent corruption into a loud error. It is cheap and it should never be skipped.

Multipart upload

A 5 TB object cannot be a single HTTP request. At 1 Gbit/s it would take about 11 hours, and any interruption would restart it.

The multipart flow:

  1. InitiateMultipartUpload returns an upload identifier.
  2. The client uploads parts independently — say 100 MB each, so 5 TB is 50,000 parts — each returning a checksum. Parts can go in parallel and in any order, and a failed part is retried alone.
  3. CompleteMultipartUpload sends the list of part numbers and checksums. The service validates and writes the final metadata record, making the object visible atomically.

Three properties fall out. Parallelism: 20 concurrent part uploads at 1 Gbit/s each cuts 11 hours to about 33 minutes. Resumability: a failure retries one 100 MB part, not 5 TB. Atomic visibility: the object appears only at completion.

The cost is abandoned uploads. Parts of an upload that is never completed occupy space forever unless something removes them, which is why a lifecycle rule to abort incomplete multipart uploads after a few days is standard, and worth mentioning unprompted.

Durability: replication versus erasure coding

This is the section's centrepiece. The write path above placed nine fragments; here is why nine, and why not three copies. "Eleven nines" is arithmetic, and here is the arithmetic.

The inputs

Two numbers, both of which you should state as assumptions rather than facts.

Annualised drive failure rate. Take 2%. Backblaze publishes drive-failure statistics from its own fleet and the overall figure has generally sat in the low single digits of per cent — cite it as a public reference point and as an order of magnitude, not a constant.

Repair window. When a drive dies, how long until its data is fully re-replicated elsewhere? A 16 TB drive rebuilt at an aggregate 1 GB/s, with the work spread across many source and destination nodes, takes 16,000 seconds — about 4 hours. Call it T = 4 h.

From those: hourly failure rate per drive λ = 0.02 ÷ 8,760 h = 2.28 × 10⁻⁶ per hour. The probability a given drive fails within a 4-hour window is λT = 9.1 × 10⁻⁶.

Three-way replication, computed

A placement group is 3 drives holding the same data. Loss requires all three to be down at once.

  • Expected first failures per group per year: 3 × 0.02 = 0.06
  • Given one down, chance one of the remaining 2 fails within the repair window: 2λT = 1.83 × 10⁻⁵
  • Given two down, chance the last fails within the window: λT = 9.1 × 10⁻⁶

Annual loss probability ≈ 0.06 × 1.83 × 10⁻⁵ × 9.1 × 10⁻⁶ ≈ 1.0 × 10⁻¹¹

That is exactly eleven nines, at a 200% storage overhead. Three copies of 100 PB is 300 PB.

Reed–Solomon 6+3, computed

Nine fragments, loss requires four failures within the repair window.

  • First failures per group per year: 9 × 0.02 = 0.18
  • Second, from 8 remaining: 8λT = 7.3 × 10⁻⁵
  • Third, from 7 remaining: 7λT = 6.4 × 10⁻⁵
  • Fourth, from 6 remaining: 6λT = 5.5 × 10⁻⁵

Annual loss probability ≈ 0.18 × 7.3 × 10⁻⁵ × 6.4 × 10⁻⁵ × 5.5 × 10⁻⁵ ≈ 4.6 × 10⁻¹⁴

Better than thirteen nines, at a 50% storage overhead. Two orders of magnitude more durable, and it stores 150 PB instead of 300 PB.

A wider 10+4 code runs the same way — fourteen fragments, loss needs five failures: 0.28 × 13λT × 12λT × 11λT × 10λT ≈ 3 × 10⁻¹⁷, about sixteen nines, at a 40% overhead.

Storage overhead against durability, by redundancy scheme0501001502002-way replication3-way replicationReed-Solomon 6+3Reed-Solomon 10+4Storage overhead (%)Durability (nines,independent-failure model)
Storage overhead against durability, by redundancy scheme

Notice the two series move in opposite directions: erasure coding costs less storage and delivers more nines. If that were the whole story, nobody would replicate.

What erasure coding costs, which the chart does not show

Reconstruction traffic. Repairing one lost copy under replication means reading one copy — one chunk of network traffic. Repairing one lost fragment under 6+3 means reading 6 fragments to rebuild one. Repair bandwidth is six times higher, and repair is exactly what happens during a failure, when the cluster is already stressed.

Read latency. A replicated read contacts one node. A 6+3 read contacts at least 6, so its latency is the slowest of 6 rather than of 1. At a 1% chance of any node being slow, the probability that at least one of 6 is slow is 1 − 0.99⁶ ≈ 5.9%, against 1%. Hedged reads mitigate this but do not remove it.

Small objects. A 4 KB object split 6 ways gives 667-byte fragments. Per-fragment overhead — metadata, checksums, filesystem block minimums — can exceed the data itself.

Recommendation: erasure code large and cold objects, replicate small and hot ones. A common cut is to replicate below a few hundred kilobytes and erasure-code above. Say the rule and say why, because "use erasure coding" without the small-object caveat is a partial answer.