Amazon S3 is designed for eleven nines of annual durability, and from the outside that number looks like a property you either have or do not. From the inside it is the output of a machine that never stops running: data split into fragments, fragments placed so that failures do not line up, every byte checksummed, disks constantly re-read, and lost fragments rebuilt faster than new failures arrive. Designing an object store, or reasoning about the one you depend on, means understanding that machine.
AWS has published parts of how S3 works and kept the rest private, so this article separates the two. Published facts are attributed. Everything else, including every erasure-coding geometry below, is an illustrative design for an S3-like system, not a claim about S3. The request path, the consistency contract, key design and multipart uploads are covered in cloud object storage architecture and S3 architecture; this page goes underneath them, into durability and the metadata layer.
Three layers and one loop
An object store at this scale has stateless front ends, a metadata index that maps each key to its current version and fragment locations, and a data layer of storage nodes. Those layers are described in the linked articles. What they understate is the background loop: scrubbers that re-read stored fragments and verify their checksums, a repair system that rebuilds anything missing or corrupt, and a placement system that moves data to balance load and failure exposure. Durability is a property of that loop, not of the write path alone.
Erasure-coding geometry is a repair-bandwidth decision
With a k+m erasure code, an object stripe becomes k data fragments and m parity fragments, and any k reconstruct it. Storage overhead is (k+m)/k, and the stripe survives any m simultaneous losses. Those two numbers are where most discussions stop, and they hide the cost that actually shapes the design: repair. Rebuilding one lost fragment of a Reed-Solomon stripe means reading k surviving fragments. With 6+3, a lost 1 GB of fragments costs 6 GB of reads; with 12+4, it costs 12 GB. HDFS erasure coding and Google's Colossus show the same storage-versus-repair trade in systems you can read about.
Wide stripes therefore buy cheaper storage and more failure tolerance per byte with more repair traffic and more nodes touched per read of a small object. They also raise the minimum useful object size: a 4 KB object split twelve ways is twelve tiny writes, so practical designs keep small objects whole and replicated, or pack many small objects into one larger erasure-coded unit. Locally repairable codes, which add parity over subgroups so a single loss can be rebuilt from a few fragments, exist precisely to cut this repair cost, at the price of extra parity.
A durability model you can compute
A first-order model makes the loop's role concrete. Assume each fragment's disk fails independently with an annual failure rate (AFR) and that a lost fragment is rebuilt within a repair window T. A stripe is lost only if m more of its fragments fail while the first is still being repaired. The code below estimates the annual loss probability per stripe; it is a teaching model with stated assumptions, not S3's method.
from math import comb, log10
def annual_loss(k, m, afr=0.02, repair_hours=24.0):
# P(lose one stripe in a year): first failure, then m more among the
# remaining n-1 fragments inside the repair window. Independent failures only.
n = k + m
lam_t = afr * repair_hours / 8760.0
return n * afr * comb(n - 1, m) * lam_t ** m
for name, k, m in [("3x replication", 1, 2), ("6+3", 6, 3), ("12+4", 12, 4)]:
for t in (24, 168):
p = annual_loss(k, m, repair_hours=t)
print(f"{name:15s} repair {t:4d}h overhead {(k+m)/k:.2f}x "
f"loss/yr {p:.1e} nines {-log10(p):.1f}")| Scheme | Overhead | Repair in 24 h | Repair in 7 days |
|---|---|---|---|
| 3x replication | 3.00x | 9.7 nines | 8.1 nines |
| 6+3 | 1.50x | 11.8 nines | 9.2 nines |
| 12+4 | 1.33x | 14.4 nines | 11.0 nines |
The model, at a 2 percent AFR, says three things. Erasure coding beats replication on both cost and durability. Repair speed matters as much as geometry: slowing repair from a day to a week costs 6+3 more than two nines. And the absolute numbers are optimistic, because the model assumes independent failures and counts per stripe; a store with billions of stripes multiplies the per-stripe probability by that count to get expected losses. Real durability is set by what the model leaves out, which is correlation.
Correlated failure and placement
Disks fail together. A firmware bug hits every drive of one model, a power event takes a rack, a bad deployment corrupts writes on a set of hosts, an availability zone loses connectivity. Placement is the defence: fragments of one stripe go to different drives, hosts, racks and, for S3 Standard, multiple availability zones, which AWS documents as storing data across a minimum of three. With 9 fragments across 3 zones, three per zone, a 6+3 stripe survives the loss of a whole zone exactly, with nothing left over for a disk failure in the remaining zones; geometry and placement have to be chosen together.
Placement also manages heat. New data is read most, so placing all of today's writes on the newest drives makes them hot. AWS engineers have described spreading each customer's data across very large numbers of drives so that no single drive's bandwidth limits a workload, which also means a single drive failure affects tiny slices of many objects and repair parallelises across the fleet. Drive diversity, staged firmware rollouts and deployment safety are durability measures in their own right.
End-to-end checksums
Silent corruption, bits that change without an error, is the failure that replication cannot see unless every copy is verified. The defence is a checksum computed as close to the source as possible and checked at every hop. Since December 2024, current AWS SDKs compute a CRC-based checksum on upload by default and S3 verifies it before accepting the object, then stores a whole-object checksum in metadata, including for multipart uploads; the AWS CLI v2 defaults to CRC64NVME, while some SDKs default to CRC32. Internally, a design like this keeps per-fragment checksums so scrubbers can verify each fragment independently and repair only the bad one.
# Upload with an explicit checksum, then verify what S3 stored
aws s3api put-object --bucket my-bucket --key data/part-0001.parquet \
--body part-0001.parquet --checksum-algorithm CRC64NVME
aws s3api head-object --bucket my-bucket --key data/part-0001.parquet \
--checksum-mode ENABLED # response includes the stored checksumOn your side, the rule is the same: compute checksums before data leaves the process that produced it, and verify them after reading, so corruption anywhere in between is detected rather than propagated.
The storage node: ShardStore
AWS published the design of ShardStore, a key-value storage node for S3, at SOSP 2021, where the paper won a best-paper award. It is written in Rust, roughly 40,000 lines at the time, and is a log-structured merge tree that keeps shard data outside the tree to reduce write amplification. Crash consistency uses a soft-updates protocol, ordering dependent writes to disk so that a crash at any point leaves a recoverable state, without a journal for every update.
The more transferable lesson is how it was validated. The team used lightweight formal methods: executable reference models that the implementation is checked against, property-based tests that generate operation sequences including crashes, and model checking of concurrent executions with Shuttle, a stateless model checker they open-sourced. Bugs in crash recovery and concurrency were found before production. For anyone building a storage node, the takeaway is that crash points and interleavings must be generated, not hand-written.
Strong consistency and the witness
S3 became strongly consistent for reads after writes in December 2020. Werner Vogels later described how: the metadata layer serves reads from a cache in front of the persistence tier, and the risk was that a write flowing through one part of the cache and a read through another would see different versions. AWS used new replication logic in the persistence tier that orders operations per object, and added a component called a witness that acts as a read barrier: on a read, the cache can learn from the witness whether its view of the object is stale. The witness keeps only in-memory state, so it adds little latency. AWS also reported proofs of the cache coherence algorithm and model checking of the design and code.
For a design exercise, the pattern generalises: keep a fast, cheap service whose only job is to know the latest version number per key, consult it on reads, and fall through to the authoritative store when the cache is behind. The hard part is making the witness itself highly available without making it a bottleneck.
Worked example: designing for nine-nines
Suppose you are building an internal object store on three zones with 2 percent AFR disks. Running the model, 6+3 with 24-hour repair gives about 11.8 nines per stripe, but at a billion stripes the expected annual loss count is about 1.7e-3, or one stripe every few hundred years, before correlation. To survive a zone loss, place three fragments per zone, and size repair so that the fleet can rebuild a failed 20 TB drive's fragments within the window: at six reads per rebuilt byte, that is 120 TB of reads, about 1.4 GB/s sustained over a day, spread across many nodes. Then add scrubbing so that every fragment is re-read within a few weeks, and alert when repair backlog age approaches the window the durability target assumes.
Failure modes
- Repair backlog: the queue of under-replicated stripes grows faster than repair, silently eating durability. Alert on backlog age, not only size.
- Correlated batches: drives of one model or firmware failing together. Diversify and stage firmware.
- Silent corruption: undetected without checksums at every hop and scrubbing of cold data.
- Stale metadata reads: a cache serving an old version. Requires an ordering mechanism such as a witness.
- Orphaned fragments: data written but never committed to the index after a failed upload. Needs garbage collection that is safe against in-flight writes, the kind of issue designing a file upload service also has to handle.
- Hot partitions in the index: concentrated key ranges overload one index partition; the index must split ranges, as in sharding.
Trade-offs
Every lever in this design trades against another. Wider stripes lower storage cost and raise failure tolerance, but they increase repair reads, the number of nodes on every read path, and the tail latency of reads that must wait for the slowest of k fragments. Faster repair buys durability with network and disk bandwidth that customer traffic also wants, so repair needs priorities: a stripe one failure from loss jumps the queue ahead of one with spare parity. Spreading across zones protects against zone loss at the price of cross-zone bandwidth on every write and every repair. A single-zone class gives up that protection for lower latency and cost, which is a reasonable choice only for data you can recreate.
On the metadata side, a cache with a witness keeps reads fast and correct, but adds a component that must be as available as the store itself. The simpler alternative, reading the authoritative index on every request, is correct and slower. Choose by asking which failure you can afford: a slow read or a wrong one.
What to do next
- Run the durability model with your own AFR, geometry and repair time, and then ask what correlated failure would break its assumptions.
- Turn on and verify checksums in your S3 uploads; confirm which algorithm your SDK uses by default.
- If you build storage, choose erasure-coding geometry and placement together, and size repair bandwidth before you size capacity.
- Measure and alert on repair backlog age and scrub coverage, since those are what durability depends on.
- Read the ShardStore paper and Vogels' consistency post for primary-source detail, and test any storage node with generated crash points and interleavings.