Course Content
System Design Interview
31 sections · 71 lessons
Object Storage: requirements, scale and the architecture
"Design an object storage service like S3." Half of every other system in this course stores its bytes in one of these. This section builds it, and the interesting requirement is not throughput — it is durability, stated as a probability.
This lesson scopes the problem, sizes it and draws the architecture. The second lesson traces the write and read paths and does the durability arithmetic; the third covers consistency, garbage collection and the follow-up questions.
What an object store is, and is not
An object store holds immutable blobs addressed by a key inside a namespace called a bucket. PUT /photos/2027/mar/beach.jpg, then GET it back. There are no directories, only keys that contain slashes; there is no partial update, only replacing the whole object; and there is no file handle you can seek and write into.
That restriction is the whole point. Because objects are immutable and whole, the system never has to coordinate concurrent writers to the middle of a file, which is the single hardest problem in a distributed file system. Give that up and you can build something that scales to exabytes with a simple consistency story.
Compare the three storage families out loud, because interviewers like the distinction:
| Family | Unit | Access | Good for |
|---|---|---|---|
| Block storage | Fixed-size blocks | Attached to one machine, read and write by offset | Databases, virtual machine disks |
| File storage | Files in a directory tree | Shared, with a POSIX interface and locking | Shared home directories, legacy applications |
| Object storage | Whole immutable objects with a key | HTTP, no partial write | Media, backups, data lakes, static assets |
The questions that shape everything after
- What is the object size range? A store optimised for 4 KB thumbnails and one optimised for 5 TB video masters make opposite choices about chunking and erasure coding.
- What durability and availability targets? Ask for numbers, not adjectives. "Eleven nines of durability" is a probability, and the durability lesson turns it into arithmetic.
- What is the access pattern? Write once, read many is the assumption most of these systems make, and it justifies immutability.
- Is versioning required? Keeping old versions on overwrite changes garbage collection and quota accounting.
- Single region or multi-region? Cross-region replication is asynchronous, which means a window where a regional loss loses recent writes. Name the window.
- What consistency guarantee on read after write? Modern object stores offer strong read-after-write for new objects, and that is a design decision with a cost.
The assumptions this section uses
| Question | Assumption |
|---|---|
| Sizes | 1 KB to 5 TB, averaging 1 MB |
| Durability | 99.999999999% — eleven nines — annual, per object |
| Availability | 99.99% for reads |
| Pattern | Write once, read many; overwrites rare |
| Versioning | Supported, off by default per bucket |
| Regions | Single region core, with optional cross-region replication |
Requirements and scale
Functional requirements
- Create and delete buckets, with per-bucket configuration.
- Put, get, and delete objects by key, plus a metadata-only head request.
- List objects in a bucket by key prefix, paginated.
- Multipart upload for large objects, with resume.
- Optional versioning, lifecycle rules, and access control per bucket.
Non-functional requirements
- Durability of eleven nines per object per year — the hard requirement.
- Read availability of four nines.
- Time to first byte in tens of milliseconds for a small object.
- Linear capacity growth by adding storage nodes.
The arithmetic
Logical capacity. 100 billion objects averaging 1 MB = 100 PB of user data.
Raw capacity. With erasure coding at a 1.5× multiplier — the durability lesson justifies the number — that is 150 PB of raw disk. At 16 TB per drive: 150,000 TB ÷ 16 TB = 9,375 drives. At 24 drives per storage node, roughly 400 storage nodes. Add spare capacity for rebuilds and growth and call it 500.
Metadata volume. Each object needs a row: bucket, key, size, content type, checksum, creation time, version, and the list of chunk locations. Keys are long and the location list grows with object size, so call it 1 KB per object. 100 billion × 1 KB = 100 TB of metadata.
That number deserves a pause. The metadata alone is 100 TB — larger than most systems in this course store in total. It will not fit on one machine, it is queried on every single request, and it must support prefix range scans for listing. Metadata is therefore its own distributed database with its own sharding strategy, and the architecture below makes that the central split.
Request rate. Assume 200,000 requests per second at peak: roughly 85% gets, 10% puts, 5% listing and head requests. That is 170,000 gets per second.
Byte throughput. Object sizes are extremely skewed — the median object is a few kilobytes while the mean is 1 MB, because a small number of very large objects carry most of the bytes. Take an aggregate of 100 GB/s read throughput. Across 400 storage nodes that is 250 MB/s per node, or 2 Gbit/s — comfortable on 25 Gbit networking, which means the network is not the constraint. Say that explicitly; it directs attention to where the constraint actually is.
Durability, expressed as objects lost. Eleven nines means an annual per-object loss probability of 10⁻¹¹. With 100 billion objects: 10¹¹ × 10⁻¹¹ = one object lost per year.
The architecture
One split defines this system: metadata and data are different problems and belong in different subsystems.
Why the split is the design
Metadata is small, hot, transactional, and needs range queries by key prefix. Data is enormous, cold per byte, immutable, and needs nothing but "give me these bytes".
Combine them and every property fights. The metadata store would have to hold petabytes; the data store would have to support transactions and listing. Separate them and each becomes a problem with a known solution: a sharded database for one, a fleet of dumb disks for the other.
The components
API service. Terminates HTTP, authenticates the request, authorises it against the bucket policy, and orchestrates everything below. Stateless, behind a load balancer, scaled horizontally.
Metadata store. A sharded distributed database holding bucket records and object records. Sharded by a hash of (bucket, key) for even distribution — with a caveat about listing that the follow-ups return to. Object records point at chunks.
Placement service. Decides which storage nodes hold the chunks of a new object. It knows the cluster topology — which nodes are in which rack, in which power domain — and its job is to spread chunks so that no single failure takes out enough of them to lose the object.
Data nodes. Machines with many drives, running a simple service that stores and returns chunks by identifier and verifies checksums. Deliberately unintelligent: no queries, no transactions, no cross-node coordination.
Garbage collector and scrubber. Background jobs that delete unreferenced chunks and continuously re-read stored data to detect corruption. Consistency and correctness and Follow-ups cover both.
The one ordering rule that makes it correct
Write chunks first. Write the metadata record last, in a single atomic operation.
If the process dies partway through writing chunks, the metadata row does not exist, so the object does not exist. A GET returns not-found, which is correct. The orphaned chunks are cleaned up by the garbage collector.
Reverse the order and a crash leaves a metadata row pointing at chunks that were never written — an object that exists, is listed, and fails on read. That is a much worse failure than a missing object, because it is silent until someone tries to use the data.