System Design Interview

Course Content

System Design Interview

31 sections · 71 lessons

Object Storage: consistency, garbage collection and follow-ups


The previous lesson was about not losing bytes. This one starts with what a reader sees, and when, then works through the follow-up questions that fill the last part of the interview: storage classes, cross-region copies, bit rot and listing.

Durability is about not losing bytes. Consistency is about what a reader sees, and when.

Creates are easy; overwrites are notA new object• Metadata written after the chunks• There is nothing to be stale against• Read-after-write comes for freeAn overwrite• Two versions exist for a moment• Readers may hold the old chunk list• Old chunks need garbage collection
Versioning turns an overwrite back into a create, which is why the honest answer is to avoid overwriting at all.

Strong read-after-write for new objects

A client writes an object and immediately reads it. Should the read succeed?

Older object stores said "eventually" — the write might not be visible for a moment, because metadata was replicated asynchronously. That produced a genuinely awkward class of bug: a job writes a file, a downstream step reads it, and the read fails intermittently in a way that looks like a race in the user's code.

Strong read-after-write for new objects is achievable because of the ordering rule from The architecture. The metadata write is the commit point, and if that write is a single atomic operation on one metadata shard, then any read routed to that shard's current leader sees it immediately. The cost is that the metadata store must offer strong consistency on that shard — leader reads or a quorum — rather than serving reads from asynchronous replicas.

Say the cost out loud: strong read-after-write means metadata reads pay a coordination round trip. At 200,000 requests per second that is real load, and it is why the metadata store gets its own capacity planning.

Overwrites are harder than creates

Two clients overwrite the same key at the same instant. Both write chunks, both write metadata. One wins the metadata write, the other loses.

Whose bytes are stored? Both sets are on disk. Whose metadata row survives? The last writer's. The loser's chunks are orphaned and collected later. That is last-writer-wins, and it is acceptable only because the object store makes no promise about concurrent overwrites of the same key.

Versioning is the honest fix. With versioning enabled, an overwrite creates a new version rather than replacing anything, and both writes survive with distinct version identifiers. A plain GET returns the newest. Nothing is lost and nothing is ambiguous.

The costs are storage — every overwritten version is retained until a lifecycle rule removes it — and a subtlety about deletion: deleting a versioned object writes a delete marker rather than removing data, so GET returns not-found while the bytes remain until the versions are explicitly purged. That surprises people and is worth naming.

Garbage collection of orphaned chunks

Chunks with no metadata reference accumulate from three sources: failed writes that never committed, overwrites whose losing chunks are unreferenced, and deleted objects.

The naive collector — scan every chunk, check whether any metadata row references it — is a join across 100 TB of metadata and 150 PB of chunks. Not feasible per run.

Two workable designs:

Reference-counted deletion. When metadata is deleted or superseded, enqueue the specific chunk identifiers for deletion. The collector consumes that queue. Fast and exact for the common cases, and it misses chunks whose metadata write never happened.

Periodic mark-and-sweep with a generation stamp. Every chunk carries a creation timestamp. Periodically, sweep the metadata store and build a set of referenced chunk identifiers — a Bloom filter keeps this tractable — then delete unreferenced chunks older than the longest possible in-flight write, typically a day. The age threshold is what stops the collector racing a write that has placed chunks but not yet committed metadata.

Run both: the queue for the common path, the sweep as the safety net.

Storage classes and lifecycle transitions

The rest of this lesson is the follow-ups. The first is cost. Not all objects deserve the same media. A typical ladder: a standard class on fast disks with millisecond access; an infrequent-access class at lower storage cost with a retrieval fee; and an archive class on very cheap media where retrieval takes minutes to hours.

Not every object deserves fast mediaStandard: millisecondsInfrequent: retrieval feeArchive: minutes to hoursScrubbing checks all tiers
Bit rot is silent, so checksums verified on a schedule are the only reason the durability figure is still true in year five.

Lifecycle rules automate the movement: "transition to infrequent access after 30 days, to archive after 180, expire after 7 years."

The arithmetic that decides it. Suppose archive storage costs roughly a fifth of standard per byte, with a retrieval fee. If 70% of your 100 PB is older than 180 days and read a few times a year, moving it saves about 80% of the storage cost on 70% of the data — roughly 56% of the total storage bill. If that data is read weekly, retrieval fees erase the saving. Lifecycle policy is an access-frequency bet, and the right advice is to measure access before setting the rule rather than assuming old means cold.

Cross-region replication

Copy objects asynchronously to a bucket in another region for disaster recovery or for read locality.

Asynchronous is the important word. There is a replication lag — seconds to minutes depending on object size and queue depth — during which the second region does not have the newest objects. If the primary region is lost in that window, those writes are gone. Name the window and call it the recovery point objective; pretending replication is instant is the mistake.

Synchronous cross-region writes would remove the window and add the inter-region round trip to every write: tens of milliseconds within a continent, well over a hundred across oceans. For a store whose selling point is fast writes, that is usually the wrong trade — but say it is a per-bucket choice, because for a small volume of critical data it can be right.

Bit rot, checksums, and scrubbing

Data on disk degrades silently. A bit flips; the drive reports success; the object is now subtly wrong and nobody knows.

Three defences, all needed:

  • Checksum on write, stored in metadata, covering the whole object and each chunk.
  • Verify on every read. A corrupt fragment is detected, the object is reconstructed from the others, and the bad fragment is repaired.
  • Background scrubbing. A job that continuously re-reads stored data and re-verifies checksums, because objects that are never read would otherwise never be checked.

Size the scrubber: re-reading 150 PB once per quarter requires 150 PB ÷ 7.9 million seconds ≈ 19 GB/s of sustained background read across the fleet, about 47 MB/s per node out of a 250 MB/s working budget. Affordable, and worth computing in the interview because it turns "we scrub in the background" into a capacity decision.

Listing, and why it is the awkward operation

Every other operation addresses one key. Listing addresses a prefix, in sorted order, across what may be a billion keys in one bucket.

If metadata is sharded by hash(bucket, key), keys with the same prefix land on different shards, so a prefix scan becomes a scatter-gather across every shard followed by a merge — expensive, and it gets more expensive as the cluster grows.

Two options. Range-partition by key within a bucket, so a prefix scan touches a small number of contiguous shards; this makes listing cheap and creates hot shards for buckets whose writes cluster at the end of the key space, such as timestamp-prefixed keys. Or keep hash partitioning and maintain a separate ordered index per bucket for listing, which keeps writes even and adds a second structure to keep consistent.

Recommendation: range-partition within a bucket, and tell users to randomise or reverse their key prefixes for write-heavy buckets. That is what production systems document, and recommending it shows you know the guidance exists because of this exact trade-off.