Storage 101 and requirements
Three storage models exist. Block storage exposes raw volumes to a server, offering the best performance but no sharing semantics. File storage layers a hierarchical directory namespace on top, easy to share over NFS or SMB. Object storage gives up in-place updates and hierarchy in exchange for virtually unlimited scale, high durability and low cost; data is accessed through a RESTful API, and each object is immutable once written. It is ideal for backups, media, logs and data lakes.
Terminology: a bucket is a globally named container; an object is a payload plus metadata, addressed by a URI like s3://bucket/path/to/object; versioning keeps multiple versions of an object under one name; an SLA formalises promises such as durability and availability. Our system must support bucket creation, object upload and download, versioning and listing, store about 100 PB, achieve six nines of durability and four nines of availability, and keep cost low.
Back-of-the-envelope estimation
Assume 20% of objects are small (under 1 MB, median 0.5 MB), 60% medium (1 to 64 MB, median 32 MB) and 20% large (over 64 MB, median 200 MB). The weighted average size is 0.2 x 0.5 + 0.6 x 32 + 0.2 x 200 = 0.1 + 19.2 + 40 ≈ 59.3 MB. With 100 PB of capacity at 40% utilisation, usable data is 40 PB, i.e. 4e10 MB. Dividing gives 4e10 / 59.3 ≈ 6.8e8, roughly 0.68 billion objects.
If each object needs about 1 KB of metadata, the metadata tier holds about 0.68 TB. The key insight is the asymmetry: data is tens of petabytes while metadata fits in a modest, though still sharded, database. That justifies storing them in completely different systems, each optimised for its own job. The workload is also dominated by object count for metadata operations and by bytes for data operations, which affects how we scale each tier.
High-level design: separate metadata from data
A load balancer spreads HTTP requests across a stateless API service, which consults IAM to authenticate the caller and authorise the action. The API service then coordinates two backends: a metadata store mapping bucket and object names to IDs and attributes, and a data store that stores raw bytes keyed by an opaque object ID. This mirrors UNIX file systems, where an inode holds metadata and pointers while data blocks live elsewhere.
Uploading creates the bucket metadata if needed, then the API service streams the payload to the data store, which persists it and returns a UUID; finally the API service writes a metadata row linking name to UUID. Downloading reverses this: look up the UUID by name, then fetch bytes from the data store. Immutability simplifies everything: since objects are never edited in place, replicas never diverge on partial updates, caching is trivially safe, and versioning is just keeping older UUIDs around.
Deep dive: the data store
The data store has three parts. A stateless data routing service receives read and write requests from the API tier. A placement service maintains a virtual cluster map describing the physical topology (data centers, racks, nodes, disks) and decides which nodes should hold the replicas of a new object, ensuring copies sit in different failure domains. It receives heartbeats from data nodes; a node missing heartbeats for a grace period is marked down. Because the placement service is critical, it runs as a small cluster using consensus such as Raft or Paxos.
Data nodes store the bytes. On write, the routing service asks placement for a primary, sends data to it, and the primary replicates to secondaries. The response timing is a classic trade-off: acknowledge after all replicas persist for strongest consistency but highest latency, after the primary plus one replica for a middle ground, or after the primary alone for lowest latency but weaker durability during the replication window.
- Ack after all replicas: strongest, slowest.
- Ack after a quorum: balanced.
- Ack after primary only: fastest, risk of loss if the primary dies before replicating.
Deep dive: how data is organised on disk
Storing each object as its own file fails at scale. Small objects waste space because file-system blocks have a minimum size, and billions of files exhaust inodes and make metadata operations slow. Instead, data nodes pack many small objects into large files, appending each new object to a read-write file. When that file reaches a threshold, perhaps a few gigabytes, it is sealed as read-only and a new read-write file is started. Writes to one read-write file are serialised, so a node keeps several open, one per CPU core, to parallelise.
To find an object, each data node keeps an object_mapping table with object_id, file_name, start_offset and object_size. A read seeks to the offset and reads the length. Because mapping data is written once and read often, and is local to each node, an embedded store like SQLite or RocksDB on the node is a good fit, avoiding a network hop to a shared database. This is essentially the same idea Facebook's Haystack used for photos.
object_mapping (file, offset, size). When the file reaches its size threshold it is sealed read-only and a fresh one is opened.Deep dive: durability, erasure coding and correctness
Disks fail. With an annual failure rate around 0.81%, three independent replicas give roughly 1 - 0.0081^3 ≈ 0.999999 durability, about six nines, but only if the copies fail independently. That is why replicas are spread across failure domains: different racks protect against a switch or power failure, and different availability zones protect against a site outage. Erasure coding offers an alternative: split data into k chunks and compute m parity chunks, for example 8+4, so any 8 of the 12 can rebuild the original. Overhead is 50% versus 200% for triple replication, and durability is higher, around eleven nines, but writes need parity computation and reads may need to gather chunks from many nodes.
Durability also means detecting silent corruption. Each object gets a checksum such as MD5, SHA-1 or an HMAC, stored alongside it; each sealed file also gets a checksum at its end. On read, the node recomputes and compares; a mismatch triggers recovery from another replica or from parity chunks.
- Replication: 200% overhead, ~6 nines, simple, fast reads.
- Erasure coding 8+4: 50% overhead, ~11 nines, CPU-heavy, slower reads and repairs.
Metadata, listing, versioning and multipart upload
Metadata has a bucket table (owner, name, ID, versioning flag), small enough for one database with read replicas, and an object table (bucket, name, version, object ID) that must be sharded. Sharding by bucket alone creates hotspots for huge buckets; sharding by object name alone breaks name lookups that include the bucket. Sharding by hash(bucket_name, object_name) distributes load evenly and serves the dominant lookup in one shard. The cost appears in listing by prefix, because a bucket's objects are scattered: a naive list scatters to all shards and merges, and paginating across shards requires tracking a cursor per shard. A common fix is a denormalised listing table sharded by bucket, trading write amplification for simple, fast listing.
With versioning enabled, an upload inserts a new row with a new version ID rather than overwriting; a delete inserts a delete marker. Large uploads use multipart upload: the client initiates, uploads parts (say 200 MB each) in parallel with independent retries, then completes the upload so the store assembles parts by ETag. This turns a fragile multi-gigabyte transfer into many small resumable ones.
Garbage collection and wrap-up
Immutability means space is never reclaimed in place. Deletes are lazy: the object is marked deleted in metadata and becomes unreachable, but its bytes stay in the packed file. Orphaned multipart parts and corrupted copies add more garbage. A background garbage collector periodically compacts read-only files, copying live objects into a new file, updating object_mapping atomically, and dropping the old file. Running compaction in batches keeps write amplification manageable. Overall, the recipe is: separate metadata and data, write immutable objects into large append-only files, spread redundancy across failure domains, verify everything with checksums, and reclaim space asynchronously.
Key numbers
Key terms
- Object storage
- Storage exposing immutable objects in a flat namespace over an HTTP API, optimised for scale and durability.
- Bucket
- A globally unique, named container that holds objects.
- Placement service
- The component that tracks cluster topology and chooses which data nodes hold an object's replicas.
- Virtual cluster map
- A logical model of data centers, racks, nodes and disks used to place replicas across failure domains.
- Failure domain
- A set of components that can fail together, such as a rack or availability zone.
- Erasure coding
- Splitting data into k data chunks plus m parity chunks so any k chunks can reconstruct the original.
- Multipart upload
- Uploading a large object as independent parts that are assembled once all parts arrive.
- Compaction
- Rewriting files to drop deleted objects and reclaim disk space.
Common mistakes
- Storing each small object as a separate file, exhausting inodes and wasting block space.
- Placing all replicas in one rack, so a single switch failure loses every copy.
- Sharding the object table by bucket only, creating hotspots for very large buckets.
- Assuming replicas protect against silent bit rot without verifying checksums.
- Forgetting garbage collection, so lazily deleted data accumulates forever.
Further study
- Finding a needle in Haystack: Facebook's photo storage (2010)
- f4: Facebook's Warm BLOB Storage System (2014)
- Ceph: A Scalable, High-Performance Distributed File System (2006)
- The Google File System (2003)
- Reed-Solomon codes