On this page
Tracks

System Design — S3-like Object Storage

Last reviewed 11 Sept 2026

Part of the system design series. See the framework and building blocks first if you haven’t.

1. Requirements

Functional

  • PUT/GET/DELETE an object by key within a bucket; list objects by prefix.
  • Support objects from a few bytes to many terabytes (large objects need multipart/chunked upload).
  • Versioning and basic access control (bucket/object-level permissions) per bucket.
  • Strong read-after-write consistency for new object PUTs (an old S3 limitation, now standard expectation).

Non-functional

  • Extreme durability — “11 nines” (99.999999999%) annual durability is the industry-standard target, meaning with a billion objects you’d expect to go roughly a century without losing one.
  • High availability, but durability and availability are different guarantees and should not be conflated — data can be safe (durable) but briefly unreachable (unavailable) during a partial outage, and that’s an acceptable trade-off; the reverse (available but corrupted/lost) is not.
  • Massive scale — exabytes of data, effectively unbounded object count.
  • Cost efficiency at scale — replication (3x) is simple but expensive; the real system needs a cheaper durability mechanism at this scale.

2. High-level architecture

flowchart TD
C[Client] --> API[API layer: PUT/GET/DELETE]
API --> META[(Metadata store<br/>bucket/key -> chunk locations)]
API -->|PUT: split into chunks| ENC[Erasure encoder]
ENC -->|N data + M parity shards| PLACE[Placement service]
PLACE --> N1[Storage node]
PLACE --> N2[Storage node]
PLACE --> N3[Storage node]
PLACE --> N4[Storage node]
API -->|GET: fetch shards, decode| DEC[Erasure decoder]
N1 --> DEC
N2 --> DEC
N3 --> DEC
BG[Background scrubber/repair] -.checks + rebuilds.-> N1
BG -.-> N2
Metadata and data are split into separate subsystems with very different scaling shapes

This is the same structural split every large object store converges on:

  • Metadata layer (bucket/key → which shards, on which nodes, object size, version, ACLs) is comparatively small and needs to be fast and consistent — it’s a classic indexing problem, often backed by a distributed key-value/table store of its own.
  • Data layer is where the actual bytes live, split into chunks, erasure-coded, and spread across many storage nodes/racks/availability zones — this is the layer that dominates cost and where the durability math lives.
  • Splitting these means they scale independently: metadata QPS and data throughput have very different profiles, and you don’t want a metadata hot spot to be coupled to data-plane capacity, or vice versa.

3. Replication vs. erasure coding

ApproachStorage overheadDurability mechanismRead/write costUsed by
Triple replication (3x)3.0x rawAny 1 of 3 copies survivesCheap — read any full copy, no reconstructionSimpler/smaller systems, hot/recent data tiers
Erasure coding (e.g. 6+3 Reed-Solomon)~1.5x rawAny 6 of 9 shards reconstruct the objectCostlier — reconstruction needs multiple shard reads + decode computeAzure Storage, Backblaze (17+3), Google Cloud Storage
Erasure coding (Backblaze’s 17+3)~1.2x rawAny 17 of 20 shards reconstruct the objectLower overhead, more shards to manage per objectBackblaze (published, open-sourced their Reed-Solomon implementation)

4. Deep dive — erasure coding and multipart upload

Erasure coding mechanics: an object is split into N data shards; a parity-generation algorithm (Reed-Solomon is the standard) computes M additional parity shards; all N+M shards are spread across independent failure domains (different disks, nodes, racks, ideally different availability zones). The object is fully reconstructable from any N of the N+M shards — meaning the system tolerates up to M simultaneous shard losses without data loss. Backblaze’s published scheme splits into 17 data shards + 3 parity shards (20 total, any 17 reconstruct); Azure Storage uses its own “Local Reconstruction Codes,” a variant designed to reduce how many shards need to be read to reconstruct after a failure, trading a bit more overhead for cheaper repair.

Why this beats replication economically: to survive 2 simultaneous failures, replication needs 3 full copies (3x storage) — erasure coding with, say, 6 data + 3 parity shards also survives up to 3 failures but only costs 1.5x storage. The cost is CPU (encode on write, decode on read/repair) and increased request fan-out (a read after a shard loss needs to fetch multiple surviving shards, not just one full copy) — pure storage-vs-compute trade-off, and at hyperscale, storage cost dominates enough that this trade is almost always worth it.

Multipart upload (how you get large objects, tens of GB to TBs, into the system reliably): the client splits the object into parts (e.g. 100MB chunks) client-side, uploads each part independently (parallelizable, and a failed part can retry without restarting the whole upload), and a final “complete multipart upload” call tells the service to assemble/register the parts as one logical object in the metadata layer. This also solves the “what if the upload dies at 90%” problem — only the failed part needs to be retried, not the whole object, and an abandoned multipart upload can be garbage-collected after a TTL to reclaim the partial storage.

5. What real systems do today

  • Backblaze is unusually transparent about their design (most hyperscalers keep this proprietary): their storage “Vaults” split each object into 20 shards — 17 data + 3 parity — using an open-sourced Reed-Solomon implementation, reconstructable from any 17 of 20. They wrote their own Java Reed-Solomon library after finding nothing suitably fast/reliable, and it performs comparably to C implementations.
  • AWS S3’s internal erasure-coding scheme is not publicly documented in detail (S3 discusses user-facing APIs extensively but stays quiet on internal durability engineering), though industry analysis suggests a design in the neighborhood of 5 data + 4 parity shards, plausibly chosen so the shard count divides evenly for distribution across a fixed number of data centers/AZs — AWS discussed some of this design space at re:Invent 2024, but hasn’t published the level of detail Backblaze has.
  • Azure Storage uses Local Reconstruction Codes (LRC), a scheme specifically engineered to reduce the number of shards that must be read during reconstruction after a failure (cheaper, faster repair) compared to standard Reed-Solomon, at a modest additional storage cost.
  • Google Cloud Storage erasure-codes objects and distributes fragments via Colossus (Google’s internal successor to GFS), selecting erasure-coding schemes per application based on required durability and cost; Google states multi-AZ redundancy is in place before a write is even acknowledged as successful.
  • Durability targets converge across vendors: AWS S3 and Google Cloud Storage both publicly target 11 nines (99.999999999%) annual durability; Azure advertises even higher (12, and in some service tiers claimed up to 16 nines) — these numbers are computed from the erasure-coding scheme’s tolerance combined with modeled/observed disk failure rates, not measured empirically (you can’t observe “11 nines” directly in any reasonable timeframe).

6. Scaling & failure

  • Metadata store becomes the bottleneck (billions/trillions of objects, high list/lookup QPS) → shard the metadata store by a hash of the bucket/key, same as any other high-cardinality key-value scaling problem; keep it a separate concern from data-plane scaling entirely.
  • Hot object (a viral file, many concurrent GETs) → this is a caching problem, not a storage-layer problem: put a CDN or edge cache in front of read-heavy objects so the storage nodes aren’t serving every request directly.
  • A storage node fails → this is the expected, designed-for case, not an exception: the placement/repair service detects the missing shard (via background scrubbing/health checks), reads any N surviving shards for affected objects, reconstructs the missing shard, and re-writes it to a healthy node — this background repair loop running continuously is what actually delivers the “11 nines” promise over time, not the erasure coding alone. Without ongoing repair, shard losses would accumulate until objects fall below their reconstruction threshold.
  • A whole rack or AZ goes down → this is why shard placement across independent failure domains matters (see the deep dive callout) — if placement was done correctly, losing an entire AZ still leaves enough surviving shards elsewhere to reconstruct every affected object; if placement clustered shards within that AZ, you’ve silently built a system with much weaker durability than the erasure-coding math implied.

What happens when the metadata store is unavailable: this is a full availability outage even though the actual object bytes are perfectly intact on the storage nodes — without metadata, the system doesn’t know which shards belong to which key. This is precisely the durability-vs-availability distinction worth stating explicitly: the data is durable (safe, not lost) but unavailable (can’t currently be served) — an acceptable, and very different, failure mode from data loss. Fix at the architecture level: replicate/shard the metadata store with its own failover (it’s a much smaller, more tractable consistency problem than the data plane) so a single metadata node loss doesn’t take down access to exabytes of otherwise-healthy data.

Interview follow-ups

  • “Why not just replicate every object 3x like a normal database?” — Works, but costs 3x storage; erasure coding gets similar or better durability at roughly 1.2–1.5x overhead by trading some storage for CPU and read fan-out — the standard move at hyperscale where storage cost dominates.
  • “What actually happens when one storage node dies?” — Nothing user-visible immediately (that’s the point of erasure coding); a background repair process detects the missing shard, reconstructs it from N surviving shards, and re-writes it — this continuous repair loop is what sustains the durability target over time, not a one-time encoding.
  • “Why split metadata and data into separate subsystems instead of one store?” — Very different scaling shapes and cost profiles — metadata is small and needs fast consistent lookups; data is enormous and dominated by storage/durability cost. Coupling them couples unrelated scaling problems.
  • “How do you upload a 500GB file without redoing the whole thing if the network blips at 95%?” — Multipart upload: split client-side into independently-retriable parts, assemble via a final “complete” call in the metadata layer; only the failed part needs to be retried.
  • “Durability vs. availability — what’s the actual difference, and why does it matter here?” — Durability is “is the data still safe/intact”; availability is “can I read it right now.” A well-designed object store can briefly be unavailable (metadata outage, network partition) while remaining fully durable — that’s an acceptable trade-off; losing durability is not.
  • “You erasure-code an object into 9 shards but put them all in the same rack. What did you actually just build?” — A system with much weaker real-world durability than the N+M math implies — a single rack failure (power, top-of-rack switch) can take out enough shards at once to lose the object; failure-domain-aware placement is as load-bearing as the encoding algorithm.
  • “How would you explain ‘11 nines of durability’ — is that measured?” — It’s a modeled number, derived from the erasure-coding scheme’s failure tolerance combined with observed/assumed hardware failure rates, not something empirically observed (you can’t wait a century to verify it) — worth saying this so you don’t overclaim what the number means.

Sources: Backblaze Vaults: Zettabyte-Scale Cloud Storage Architecture · Erasure Coding: Backblaze Open Sources Reed-Solomon Code · Understanding Cloud Storage 11 9s Durability Target — Google Cloud Blog · S3 Architecture — Neo Kim, The System Design Newsletter · Colossus: Google’s Next-Generation Distributed File System — SysTutorials · Cloud Storage Durability vs. Availability — Backblaze