System-design interview · Extended interviews
Design a distributed object store
Design resumable file uploads, publish only complete verified versions, read selected byte ranges, recover damaged copies and delete storage only when no upload, retained version or active reader needs it.
You will learn to
- Separate object metadata/visibility from distributed byte placement.
- Trace resumable multipart upload and atomic object-version publication.
- Compare replication/erasure coding, consistency boundaries, checksums, and delegated access.
Practice in this chapter
8 interview questions with model answers and follow-ups.
Go to interview practiceUseful foundations: Replication and durability · Data partitioning and sharding · Quorums, consensus, leases, and fencing · Databases, data models, and ACID transactions
Workload and timing examples are interview assumptions.
Dotted concept links open the relevant explanation in a new tab.
01Problem and scope
An object store maps a tenant-scoped key to one complete immutable version. Keep the metadata mapping a key to a version separate from the storage locations holding its bytes. Parts upload independently, but the key points to the new version only after the database commits its verified manifest: the ordered record of chunks, lengths and integrity information. This design provides resumable multipart upload and strong per-key visibility within one region rather than shared-file mutation. A two-GiB object at T7/report.pdf uses thirty-two 64-MiB parts; completion U31 competes against expected version V4.
Smallest working design
Start with one server: write bytes to a temporary location, verify them, then atomically update the key’s metadata to point at the completed file. The key is a namespace entry; the underlying bytes can be immutable. Distributed object storage extends placement and recovery while preserving that publication idea.
Clarify the object contract
Candidate: “Must readers see a partial overwrite, and can uploads resume?” Interviewer: “Never partial; support multipart resume.” Candidate: “Do we need shared-file mutation or only whole-object versions?” Interviewer: “Whole objects, with strong per-key visibility in one region.” This makes immutable byte chunks plus an atomic namespace pointer a natural starting design.
Protocol cases to prove
The protocol must handle interrupted part uploads, lost acknowledgements and two uploads both trying to replace the same expected version. Each read must retain one manifest throughout the response, called pinning that version, so an overwrite cannot change the remaining chunks mid-read. The hard promise is not that all chunks travel atomically over the network. It is that the key names one complete verified version only after the publication transaction commits.
02Functional requirements
- Operate on private objects. Support PUT, GET, HEAD and DELETE with metadata/checksums; private access is the default.
- Upload in parts. Initiate, upload/list/retry parts, complete or abort, and recover status after a lost response.
- Publish complete versions. Completion supplies an ordered part set with integrity metadata; missing parts fail without changing the current object.
- Read one version. GET/HEAD resolve the current version or an authorized explicit version. Range reads return the requested inclusive byte range and correct metadata; invalid ranges fail clearly.
- List with bounded cursors. Use an opaque tenant/prefix/last-key cursor over lexicographic namespace keys.
- Retain or delete versions. Support optional version retention. Deleting the current key follows that policy and does not automatically erase every historical byte.
- Delegate bounded access. Signed grants can authorize a key/version/upload operation under bounded expiry and headers.
- Control lifecycle and quotas. Bound object size, active uploads and retained bytes; abort abandoned sessions, expire permitted versions and schedule repair/cold-tier moves without exposing incomplete manifests.
Write and retry contract
Choose strong per-key visibility in one region: after a committed overwrite, new current-version reads resolve the new complete version. A request may pin an older explicit version where policy permits. A retry with the same upload/part identity and checksum confirms the same logical work; replacement parts, if allowed before finalization, need an explicit generation. Multipart sessions cannot become unlimited unbilled temporary storage.
Listing is not a snapshot
Each page reads currently committed namespace rows. A whole traversal is not a point-in-time snapshot: concurrent insertions before the cursor may be absent, and later keys can change between pages. An inventory/export needing a fixed view uses a separately retained namespace snapshot.
Scope and provider boundary
In-place byte edits, filesystem locking and multi-key transactions are outside scope. Amazon S3 documents strong read-after-write for PUT/DELETE and offers versioning; these are provider-specific capabilities, not universal object-store axioms. Cross-region replicas, external metadata databases and CDNs may have separate freshness contracts. S3 overview/consistency.
03Non-functional requirements
- Workload assumption. Ten million new objects/day averaging 10 MB; one billion reads/day averaging 4 MB; thirty-day baseline retention.
- Latency. Metadata operations p95 below 100 ms and first-byte reads p95 below 200 ms for healthy hot data within the region. Byte count and network rate dominate whole-transfer time.
- Availability. Target 99.95% API availability under the assumed deployment.
- Durability before publication. For frequently accessed, or hot, writes, require three durable full copies across independent failure domains: groups of storage resources placed so that a single stated failure does not destroy every copy. Metadata commits have their own replicated durability policy.
- Access isolation. Authorize before resolving private content or issuing a grant. Version-keyed immutable caches must enforce the access contract.
- Bounded lifecycle. Bound incomplete-upload lifetime and protect active readers, uploads and retained versions during garbage collection.
Per-key consistency invariants
| Operation | Required result |
|---|---|
| New current-version lookup after V5 commits | Resolve V5 |
| Read already pinned to V4 | May finish V4; never mix V4 and V5 chunks |
| Two conditional overwrites both expecting V4 | At most one succeeds |
| Replay upload completion U31 | Return its one recorded result |
Limits of these promises
The durability promise covers the stated single-node/domain failure; it is not a fabricated universal “eleven nines” claim. Repair restores redundancy. Correlated failures, operator errors and regional loss require further design and testing.
A CDN cached under an unversioned name can have a different freshness policy; do not silently include it in the strong-origin guarantee.
04Capacity estimates
Assume 10M new objects/day averaging 10 MB, thirty-day retention, and 1B daily reads averaging 4 MB.
| Quantity | Calculation | Consequence |
|---|---|---|
| Write requests | 10M / 86,400 ≈ 116/s | Metadata rate differs from byte rate |
| New payload | 10M × 10 MB = 100 TB/day | About 1.16 GB/s ingress |
| Retained raw data | 100 TB × 30 = 3 PB | Storage placement dominates cost |
| Three replicas | 3 PB × 3 = 9 PB | Before headroom/metadata |
| Reads | 1B × 4 MB = 4 PB/day | About 46.3 GB/s outbound payload |
Interpret byte throughput
These are hypothetical sizes. Request rate alone would badly understate the network problem. Cache popular immutable versions and support ranges where clients need only part of a file; account for egress and repair bandwidth separately.
Read and upload throughput
At a fivefold peak, read payload can approach 231.5 GB/s and ingress 5.8 GB/s under these assumptions. Large objects make bandwidth and disk throughput more important than the modest 116 new-object requests/s average. Multipart upload also increases request count: the uploading client's 2 GiB object produces 32 part writes plus initiate/complete operations, not one write call.
Metadata overhead
Thirty days yields roughly 300M objects. At an illustrative 1 KB namespace/manifest header each, metadata is about 300 GB before chunk references, indexes and replicas. A 2 GiB version with 32 chunk references at 64 bytes adds about 2 KB of reference metadata. Small objects may have disproportionate metadata cost and should avoid unnecessary multipart/chunk overhead.
Coding and repair bandwidth
An erasure code stores data fragments plus calculated parity fragments that can reconstruct missing data. A 4+2 code stores four data and two parity fragments, totaling 1.5 times the original bytes. It reduces 3 PB raw from 9 PB at three replicas to about 4.5 PB encoded payload, saving 4.5 PB before extra capacity reserved for placement and repairs. It spends CPU/network during encoding and repair. Reconstructing a missing fragment may read several surviving fragments; budget repair traffic separately so a disk failure does not starve the 46.3 GB/s ordinary read stream. Popular immutable versions can use caches, but cache hit ratio determines actual origin savings.
05APIs and contracts
POST /uploads
{key:"T7/report.pdf",size:2147483648,expectedVersion:"V4"}
→ {uploadId:U31,partSize:67108864,state:"open",expiresAt:...}
PUT /uploads/U31/parts/9
Content-Length:67108864; checksum:<defined algorithm/value>
→ {part:9,generation:1,checksum:...,durabilityStatus:"stored"}
POST /uploads/U31/complete
{parts:[{number:1,generation:1,checksum:...},...32]}
→ {version:V5,state:"completed",length:2147483648}
Completion validation and replay
Completion rejects missing parts, mismatched lengths/checksums, expired/aborted sessions and changed expected key version. The service encodes the ordered completion manifest in a documented, consistent format and records its hash, called the completion fingerprint; a retry with another part set cannot be mistaken for the same request. GET /uploads/U31 recovers the committed V5 outcome after a lost reply. API examples describe this designed service, not literal S3 request syntax.
Version reads and listing
GET/HEAD accepts an explicit version or resolves current; conditional writes use expectedVersion rather than comparing client wall clocks. DELETE can accept an expected version to avoid accidentally deleting a newer overwrite. Listing cursors bind tenant/prefix and chosen namespace view, with bounded expiry. Authorization derives tenant ownership from verified credentials; the string T7/ is not proof that the uploading client may access that key.
Checksum and ETag semantics
A checksum states both its algorithm and which bytes it covers, such as one part or the whole object. An entity tag (ETag) is a provider-defined version validator and not universally the full-object MD5, especially for multipart objects.
06Data model and access patterns
The identities describe different layers: a part identifies a position in one upload, a chunk identifies immutable stored bytes, and a version manifest assembles chunk references into a complete object. The logical key points to one current version. Keeping these layers separate allows part retries and byte-placement repairs without silently changing a published object.
| API/state | Example |
|---|---|
| Begin | POST /uploads {key:T7/report.pdf,size:2147483648,expectedVersion:V4} → U31 |
| Part | PUT /uploads/U31/parts/9 with length/checksum |
| Complete | POST /uploads/U31/complete {parts:[{number:1,generation:1,checksum:...},...32]} |
| Manifest | V5,key,byteLength,contentType,chunkRefs,checksum,owner=T7 |
| Range read | GET /objects/T7/report.pdf?version=V5 plus byte range |
Part identity and replay
At 64 MiB per part, 2 GiB / 64 MiB = 32 parts. Part numbers and upload identity make retries address the same work. An object manifest is the committed map from logical byte ranges to stored chunks. Metadata stores upload state, expected prior version, quotas, retention, and ownership; storage nodes do not decide application authorization.
Upload and chunk records
Add Upload(U31,tenant, key, state, expectedVersion, completionFingerprint, expiry, resultVersion), Part(U31,number, generation, checksum, chunkRef) and KeyHead(T7,key, currentVersion). Version manifests are immutable after publication and include ordered ranges, ownership, retention and integrity data. Chunk placement metadata tracks replica locations, checksums, health and placement generation.
Namespace partitioning
Partition namespace metadata by tenant/key hash or range according to listing needs. All conditional publication state for one key/upload must share a transactional authority or use a specified atomic protocol. Byte nodes own chunk storage/verification; they do not choose which version a logical key names. A replicated metadata leader handles that decision.
Collection roots
Garbage collection (GC) deletes chunks that are no longer needed. A root is a recorded reason to retain bytes, such as a current or retained object version or an active upload. A reader pin records that a download still needs its chunks. Track Chunk(chunkId, state=LIVE|DELETING|DELETED, uploadRefs, versionRefs, readerPins, deletionGeneration) under the same metadata authority that owns the upload and version references. This design does not deduplicate chunks across independent authorities. Create the protected upload reference before granting a byte upload. Every new upload reference, retained-version reference or reader pin is acquired atomically only while the chunk is LIVE; DELETING is irreversible for that chunk identity. Namespace listing reads metadata, not a scan across disk directories. Placement can change during repair while a manifest keeps the same logical chunk identity.
Integrity is not authorization
Keep integrity and access separate: a correct checksum proves bytes match expected content, not that the caller is entitled to receive them. Signed grants and metadata ownership enforce access before storage-node requests are issued.
Enforce immutable bytes
A chunk ID is immutable because the byte service enforces it, not because the name looks unique. Its first accepted write atomically creates the bytes and their checksum; a retry may confirm identical bytes but cannot overwrite that identity. Replacing a multipart part allocates a new chunk ID and part generation. The upload grant binds upload, part generation, chunk ID, allowed size/checksum and expiry; the node checks those restrictions. If managed storage backs the byte plane, retain an exact immutable VersionId or enforce create-only conditional writes and deny bypass writes. A reusable signed PUT to a mutable key does not implement this invariant.
Current head versus historical roots
A current KeyHead is always a version root even when optional historical retention is disabled. Replacing/deleting the head removes that current root under the same authority, but historical retention roots and reader pins may keep the old version alive. This distinction prevents a non-versioned bucket from collecting its still-current object.
07Basic working design
Temporary bytes and metadata
The baseline has an API, a transactional metadata database and one local disk. The uploading client initiates U31 against V4. Each part is written to an immutable temporary file, fsynced under the stated local durability policy and recorded by part identity/checksum. A lost part reply can be retried and compared with the recorded part rather than restarting two GiB.
Freeze, publish and pin reads
- Freeze and verify. Completion first freezes U31's chosen part generations, validates the ordered set and builds a manifest.
- Commit current version and retry result. It then atomically updates KeyHead from V4 to V5 and records U31 completed/result V5.
- Read one immutable manifest. Readers consult KeyHead once, load that immutable manifest and stream its chunks.
- Before/after visibility. Before the transaction they see V4; after it they see V5.
- Never publish through a part rename. No rename of an individual part makes partial V5 visible.
Crash recovery and local durability
If the process crashes before publication, U31 can resume/finalize or eventually abort; V4 remains current. If it crashes after commit before replying, the retry reads U31's recorded result. The baseline demonstrates complete-version publication and safe retries, but its crash recovery depends on the local disk and commit policy actually implemented. It cannot survive losing its sole disk, and the disk/network become clear capacity bottlenecks.
Bytes before namespace publication
The same ordering will apply when storage is distributed: save durable bytes first, then publish metadata that names the complete object.
Readers follow one committed manifest; uncompleted part files are not the current object.
Read each connection in order
- sync1. Begin U31; upload partsUpload / download client → Object API / completion logic
- sync2. Write / verify durable partsObject API / completion logic → Immutable local part files
- sync3. Freeze parts; publish if current = V4Object API / completion logic → KeyHead / uploads / manifests
- sync4. Read key or versionUpload / download client → Object API / completion logic
- sync5. Resolve one manifestObject API / completion logic → KeyHead / uploads / manifests
- sync6. Stream pinned version bytesObject API / completion logic → Immutable local part files
08Find the baseline flaws
| Failure test | What breaks and what must follow |
|---|---|
| One disk is a capacity/failure limit | One disk cannot retain three PB or serve tens of GB/s. A node loss after a part acknowledgment destroys that part unless redundancy exists; an upload-status row does not contain the bytes. The first scaling change must therefore address placement and redundancy, not merely add more stateless API instances. |
| Concurrent conditional overwrites | The correctness counterexample is concurrent overwrites. U31 and U32 both read current V4, independently upload parts, then blindly set current to V5 and V6. Both callers receive success despite an expectedVersion=V4 contract. The final pointer change requires one atomic conditional comparison at the key authority. |
| Part replacement races finalization | Another race occurs within one upload. The completion worker validates part 9 generation 1 while the uploading client replaces part 9 with generation 2. If the manifest later reads an uncontrolled mixture, the published checksum/length may no longer describe the uploaded set. Freeze the exact immutable part-generation list before verification and disallow part mutation for that finalization state. |
| Mixed-version read | Finally, a read that resolves current for every chunk can fetch part 1 from V4 and part 2 from newly published V5. A file that never existed is returned. Pin the version/manifest at request start. Replicating chunks protects their bytes. It does not choose which version a reader uses or coordinate publication and cleanup. |
09Improve the design, step by step
Replication stores full copies; three replicas cost roughly 3× raw bytes and can tolerate selected node/domain failures according to placement. Erasure coding splits data into k data fragments plus m parity fragments; an illustrative 4+2 code uses 6/4=1.5× raw payload and can reconstruct from any four valid fragments under that code’s assumptions.
This figure assumes a suitable maximum-distance-separable 4+2 code. Not every code family offers the same any-four guarantee.
Remember: k data + m parity; sufficient valid fragments reconstruct.
Read the diagram
- Start with data fragments D1 to D4 and parity fragments P1 and P2 under an MDS 4+2 code.
- D2 and P1 are lost. D1, D3, D4 and P2 are four valid survivors.
- Those four reconstruct the original data. Six stored fragments for four source fragments give 1.5x raw payload overhead.
Try from memoryMust the four survivors all be data fragments?
No. In this MDS 4+2 example, any four valid fragments suffice, including the shown mix of data and parity.
| Placement | Benefit | Cost |
|---|---|---|
| Full replicas | Simple low-latency reads/repair | Higher stored-byte overhead |
| Erasure coding | Lower redundancy overhead | Reconstruction/repair CPU and network |
| Hot replicated + cold coded | Matches different access patterns | Lifecycle movement complexity |
| Cross-region copy | Additional disaster-recovery option | Replication lag, cost, and residency policy |
Separate fragment locations across the actual failure domains promised. Six fragments on one failing disk do not survive a disk loss. Acknowledgement must specify which durable placements exist before success.
1. Distribute immutable chunks with replicated hot writes
- Trigger: disk capacity/throughput and single-node loss.
- Mechanism: Placement selects three independent failure domains; publication waits for the policy's verified durable copies. This adds parallel byte capacity and survives the promised failure.
- Benefit, cost and alternative: Costs are roughly 3× stored bytes and write network; correlated placements defeat the benefit. One disk remains suitable only for a weaker local-development contract.
2. Shard and replicate metadata authority
- Trigger: namespace growth and a metadata single point of failure.
- Mechanism: Route each key to a leader-backed shard with atomic expected-version publication and durable upload results.
- Benefit, cost and alternative: This scales independent keys, but introduces routing, rebalance and safe failover. For any one key, one metadata authority still determines the order of conditional publications. A large single metadata database is simpler until measured limits justify sharding.
3. Ranges and immutable-version caching
- Trigger: 46.3 GB/s read payload.
- Mechanism: Versioned caches and range requests avoid repeated or unnecessary origin bytes.
- Benefit, cost and alternative: Costs are cache storage, access-aware keys and stale unversioned-name risks. A download needing an entire infrequently accessed object gains little from many tiny range requests; combine adjacent storage reads to avoid unnecessary request and disk overhead.
4. Background coding, scrubbing and lifecycle GC
- Trigger: three-copy byte cost and latent corruption.
- Mechanism: Infrequently accessed versions can move to verified erasure-coded storage before old replicas are retired. Scrubbers periodically read stored bytes and check their checksums, allowing repair before corruption destroys the last usable copy.
- Benefit, cost and alternative: This saves storage while adding repair CPU/network and transition state. Keep hot small objects replicated when latency and repair simplicity outweigh byte savings. GC retains chunks referenced by committed versions, active uploads or readers until those references can be safely released.
10Detailed architecture
Namespace authority and placement
Clients enter an authenticated metadata/API gateway for namespace and upload operations. A key router resolves the metadata shard. Its leader and replicas own KeyHead, upload state, immutable manifests, quotas and publication results. A placement service maps chunks to storage nodes across failure domains; it can change healthy locations without changing an object's logical version.
Upload and background repair
For uploads, the API can return limited part grants so clients send large bytes directly through the byte path. Storage nodes verify length/checksum and report durable placement evidence to the upload metadata path. A completion coordinator freezes the manifest, validates policy, then performs the key/version transaction. There is no implication that metadata replicas contain the object bytes.
Version-pinned range reads
For reads, the gateway authorizes and resolves one version, then a range/data service obtains the relevant chunks from a versioned cache or healthy placement. Repair/scrub workers operate asynchronously and update placement state only after verified replacement data exists. Lifecycle GC scans metadata roots and staged uploads before deleting unreferenced chunks.
Why background actors matter
The final graph includes these background actors because byte durability depends on continuous repair, not just initial copying. Cross-region copies and CDN behavior are additional boundaries with their own freshness/RPO policies; the regional KeyHead guarantee cannot be casually extended to them.
Per-key publication transfers protected chunk references atomically; GC must win the same metadata guard before physical deletion.
Read each connection in order
- sync1. Begin / complete / authorize readObject clients → Authenticated object / grant API
- sync2. Resolve tenant/key authorityAuthenticated object / grant API → Key metadata router
- sync3. Read or guarded metadata commandKey metadata router → Metadata leader KeyHead / uploads / manifests
- replicationReplicate publication / upload stateMetadata leader KeyHead / uploads / manifests → Metadata durable replicas
- sync4. Request scoped part placementAuthenticated object / grant API → Chunk placement service
- sync5. Upload granted immutable partsObject clients → Chunk nodes across failure domains
- sync6. Record checksums / durable receiptsChunk nodes across failure domains → Metadata leader KeyHead / uploads / manifests
- async7. Frozen finalization workMetadata leader KeyHead / uploads / manifests → Frozen-manifest completion worker
- syncVerify required durable placementsFrozen-manifest completion worker → Chunk placement service
- sync8. Atomic expected-version publicationFrozen-manifest completion worker → Metadata leader KeyHead / uploads / manifests
- sync9. Authorized pinned manifest / rangeAuthenticated object / grant API → Range / byte streaming service
- syncRead versioned cached bytesRange / byte streaming service → Immutable version / chunk cache
- syncResolve healthy chunk locationsRange / byte streaming service → Chunk placement service
- sync10. Fetch / verify requested bytesRange / byte streaming service → Chunk nodes across failure domains
- sync11. Stream one object versionRange / byte streaming service → Object clients
- asyncScrub / reconstruct valid redundancyChecksum scrub / repair workers → Chunk nodes across failure domains
- controlPublish verified placement generationChecksum scrub / repair workers → Chunk placement service
- syncMark DELETING only if no referencesRetention / upload GC workers → Metadata leader KeyHead / uploads / manifests
- asyncDelete claimed generation; no new refsRetention / upload GC workers → Chunk nodes across failure domains
11Write path and acknowledgement
Before transferring a part, the service records an upload reference that prevents cleanup from deleting its chunk. Completion freezes and verifies the manifest, then changes the key pointer and upload result atomically.
Numbered upload and publication trace
- Authorize and create upload state. The API authenticates the uploading client for tenant T7, checks quota, and records U31 pending against current V4.
- Choose independent placements. The placement service assigns part/chunk destinations across chosen failure domains.
- Upload and retry immutable parts. The uploading client sends 32 parts. Part 9’s response is lost; retrying U31/part 9 with the same expected checksum confirms that immutable part generation; different bytes require a new generation, not an overwrite of verified bytes.
- Verify the complete set. Complete verifies the ordered part set, total length, and supported integrity checks; missing/corrupt parts prevent publication.
- Publish conditionally. The metadata transaction conditionally changes the key from V4 to complete V5 and records U31 completed. Readers saw V4 until this commit.
- Recover the recorded result. If the completion reply disappears, querying/retrying U31 returns the recorded V5 outcome rather than publishing another inconsistent version.
Provider-specific multipart option
S3’s documented multipart initiate/upload/complete model is one implementation option. Multipart upload.
Guard each part record
The metadata write checks that U31 is still open and records the exact part generation. The completion request atomically changes open → finalizing and stores its canonical ordered manifest fingerprint. New/replacement part commits are rejected after that transition; in-flight byte uploads may finish as unreferenced objects but cannot alter the frozen manifest.
Verify the frozen manifest
The completion worker verifies total length, each required immutable part generation/checksum and sufficient durable placements. If a storage failure reduced the policy below its publication threshold, repair or request reupload before continuing. Validation failure leaves current V4 unchanged and reports a recoverable or terminal upload state according to the error.
Atomic publication result
The final metadata transaction checks U31 is finalizing with that fingerprint and KeyHead still equals V4, then inserts immutable V5, updates the head and records U31 completed/result V5 atomically. A competing overwrite that already changed the head makes this transaction fail cleanly; its bytes are retained briefly for an explicit retry/rebase policy or GC, not silently published over the winner.
12Read and delivery path
Check access and hold a reference to one immutable version before reading ranges. Cleanup keeps its chunks while a retained version or active reader still needs them.
Resolve one manifest
A reader resolves V5 once, then maps its requested range through the manifest. The half-open interval [128 MiB, 192 MiB) is the third 64-MiB part; the corresponding inclusive HTTP range is bytes=134217728-201326591. Unrelated parts need not be fetched. Pinning the version prevents a concurrent overwrite from mixing chunks from V4 and V5 in one response.
Verify bytes
Numbered range-read flow
- Authorize and resolve one version. Authenticate the requester and validate ownership or the signed grant, including method, key/version and effective expiry. Resolve current KeyHead once if no explicit version is supplied.
- Atomically pin live chunks. In the owning metadata transaction, resolve/load V5 and acquire reader pins for its required chunks only if their state is LIVE. A retained version root already protects its manifest. If a chunk is DELETING, do not stream from a remembered location: fail an explicitly deleted-version read or re-resolve/retry a current-key read. Validate range bounds against byteLength and choose the covered chunks/offsets.
- Use a versioned access-aware cache. Check a cache keyed by tenant/access boundary, V5 and chunk/range identity. Never substitute an unversioned cached V4 body merely because the key name matches.
- Read and verify required bytes. Resolve healthy chunk placements, read required bytes and verify integrity under the supported checksum/range scheme. A whole-part checksum may require verifying a larger chunk than the requested subrange unless subchunk checksums exist.
- Recover from valid redundancy. On corruption/unavailability, try independent valid redundancy, schedule repair and fail clearly if the promised data cannot be reconstructed. Return correct content-range/length/type metadata and stream one version throughout.
Concurrent overwrite behavior
An overwrite committed during step four affects later new reads, not this pinned response. A delete may prevent new current-key resolution while an already-authorized, protected version read finishes under the defined policy.
Release pins and choose ranges
Release the durable reader pins only after the streaming service finishes its use of the chunks. A crashed streaming worker can leave reader pins behind. Recovery removes them only after preventing that worker generation from issuing further reads and waiting for its outstanding reads to finish or terminate; elapsed time alone does not prove that the chunks are unused. This deliberately favors temporary retained bytes over deleting a chunk still in use.
13Correctness deep dive
Freeze before namespace publication
First freeze the upload’s exact part list so it cannot change during validation. Then publish the key’s new version with an atomic update that orders competing uploads. These are separate state changes.
transaction freeze(U31, requestedParts):
lock upload U31
require canonical(requestedParts) matches any recorded completion fingerprint
if COMPLETED: return recorded result
require OPEN and not expired, or matching existing FINALIZING fingerprint
# expiry applies to starting finalization; finalizing recovery uses its stored manifest
store immutable part-generation list and fingerprint
set state=FINALIZING; commit
verify frozen lengths, checksums and durable placements
transaction publish(U31, V5):
lock upload U31; lock KeyHead(T7,report.pdf)
if U31.COMPLETED: return U31.resultVersion
require U31.FINALIZING and verified fingerprint matches
require KeyHead.version == U31.expectedVersion # V4
lock frozen chunk reference rows; require each state == LIVE
insert immutable manifest V5 and its retained-version references
transfer U31 upload references to V5 references atomically
set KeyHead=V5; set U31=COMPLETED,resultVersion=V5
commit; return V5
Guard part metadata
Part updates and freezing use the same upload-state guard. Once the upload is FINALIZING, a new part update sees that state and is rejected before commit. The verified manifest names immutable chunks, so a later write to the same pathname cannot change its bytes. Placement health may change after verification; the redundancy policy is designed to survive the stated failure, while repair maintains it. Do not claim protection against arbitrary simultaneous loss between two instructions.
Competing completions
U31 wins: it locks KeyHead at V4, commits V5 and its completion result. U32 expecting V4 then sees V5 and conflicts. U32 wins: U31 fails the same predicate, leaving V6 current; it cannot publish V5 merely because its upload finished first. U31 reply lost: retry finds COMPLETED/V5 and returns it without another pointer change.
Crash and replay result
Garbage collection shares authority
Garbage collection must use the same metadata transactions as publication and reader-pin creation. If it merely checks for zero references and deletes later, a new reader could acquire a reference between those two steps:
Atomic deletion claim
transaction claimForDeletion(chunk):
lock chunk and associated reference state
require chunk.state == LIVE
require uploadRefs == 0 and versionRefs == 0 and readerPins == 0
require all abandoned owning uploads are durably ABORTED
set state = DELETING; increment deletionGeneration
commit deletion work(chunkId, deletionGeneration)
Abort versus publish
An abort and publication contend on the upload state: ABORTED prevents publication, while COMPLETED transfers protection to the retained version. GC cannot reclaim an OPEN or FINALIZING upload by just observing an old timestamp. Physical workers delete only the immutable identity claimed by that DELETING generation and retry until removal is recorded. No new upload, repair publication, version reference or read pin can resurrect that identity; a later upload uses a new chunk ID.
Reference versus deletion outcomes
KeyHead and upload result change atomically, so concurrent completion and a lost reply do not create inconsistent versions.
Read each connection in order
- syncFreeze manifest; verified partsCompletion U31 → Metadata authority
- syncFreeze alternate manifest; verified partsCompletion U32 → Metadata authority
- syncIf current V4: commit V5 + U31 result + LIVE refsCompletion U31 → Metadata authority
- blockedCommit succeeds; response lostMetadata authority → Completion U31
- syncPublish V6 only if current = V4Completion U32 → Metadata authority
- returnConflict: current is V5Metadata authority → Completion U32
- syncResolve current keyNew reader → Metadata authority
- returnComplete immutable V5 manifestMetadata authority → New reader
- syncRetry complete U31Completion U31 → Metadata authority
- returnReturn stored V5; no new publishMetadata authority → Completion U31
14Failure and recovery
| Failure or condition | Surviving state, response and recovery |
|---|---|
| Part acknowledgement then node loss | A node fails after acknowledging part 9 but before completion. The service verifies enough valid durable redundancy or repairs/reuploads before it publishes V5; an acknowledged part token alone cannot substitute for the promised durability. A metadata leader failover must preserve committed U31/V5 state and reject stale writers. Two concurrent overwrites using expectedVersion V4 cannot both succeed as the same conditional update. |
| Read-time replica loss | If a chunk node fails during a read, select another verified replica or reconstruct an encoded stripe. Repair writes a new copy first, verifies it, then atomically updates placement generation; it does not remove the last healthy copy before replacement succeeds. Throttle repair separately so a fleet failure does not consume all customer read bandwidth. |
| Metadata leader partition | If the metadata leader partitions, new publication pauses until a safely fenced leader can recover committed heads/uploads. A stale leader must not accept another expected-V4 write after V5 is committed elsewhere. If only the API process fails, callers recover through U31 or an explicit version; no retransmission of already recorded parts is necessary. |
| GC races a new reference | If GC races a read or publication, its atomic LIVE-to-DELETING transition competes with reference acquisition under the same metadata authority. The reference winner blocks collection; the deletion winner blocks new pins and publication. Workers act only on the committed deletion generation, so there is no unguarded check-to-delete window. Abort an abandoned upload before removing its upload roots, and reject completion after abort. Reader recovery releases leaked pins only after fencing and draining that serving generation. Version deletion includes old versions, retention and backup treatment; a delete marker alone is not physical erasure. |
Observe repair and recovery
Measure byte throughput, first-byte/range latency, checksum failures, incomplete-upload age, repair backlog, replication lag, metadata conflicts, and storage overhead. Test interrupted parts, lost completion replies, corrupt replicas, concurrent overwrite, stale CDN content, and expired signed URLs. Report each boundary’s guarantees rather than saying all storage is simply consistent.
15Operations, security, and cost
Scoped grants and access enforcement
An authorized service can issue a short-lived signed upload/download URL tied to the operation, key/version, and permitted headers. Anyone possessing it may exercise that capability until its effective expiry; it is not inherently single-use. S3’s presigned URL documentation also explains credential-lifetime effects. Presigned URLs. Avoid logging grants, scope the signer narrowly, and do not expose storage credentials to clients.
Retention and abuse limits
Version retention helps recover overwrites but costs bytes and complicates deletion policy. A delete marker can hide the current name while older versions remain retrievable to authorized callers. Garbage collection removes only unreferenced expired versions/parts after respecting active uploads, retention, and recovery policy; abort abandoned uploads explicitly.
Service and repair metrics
Measure successful-byte throughput, first-byte p95, range read amplification (storage bytes fetched divided by bytes requested by the client), checksum failures, desired-versus-actual replica count, repair age, metadata conflicts and abandoned-upload bytes. A healthy PUT success rate can hide a growing repair backlog that reduces failure tolerance. Track bytes by tenant and lifecycle state so incomplete uploads cannot quietly consume unlimited capacity.
Stored bytes and recovery cost
The 3 PB retained workload costs 9 PB under three replicas versus about 4.5 PB for illustrative 4+2 coding before overhead. The comparison excludes encoding and repair bandwidth; benchmark reads and one-domain recovery before moving cold data. If a popular immutable version has 90% cache hit rate, its origin read bytes fall tenfold, but cache egress and authorization still cost resources.
Compatible rollout and fault drills
Roll out a new manifest or checksum format with readers that understand both before enabling writers. Test interrupted parts, completion after concurrent overwrite, corruption with one replica unavailable, lost completion replies and GC while a reader is pinned. For cold-tier conversion, publish a new verified placement only after every needed fragment meets policy, then retire old replicas gradually. Never use a storage-node filename or ETag as a universal authorization or full-object-integrity proof.
Reusable grants and immutable versions
| Access mechanism | What the byte service must enforce | Limit |
|---|---|---|
| Upload part grant | Exact immutable chunk/generation and required integrity/size conditions | Reusing the grant must not mutate verified bytes |
| Versioned download grant | Signed method, immutable version and effective expiry | Bearer possession permits use; it is not requester identity |
| Identity-bound download | Authenticate the actual requester, compare principal with the grant, check current policy | Adds an online authorization dependency |
| One-use application grant | Atomically mark the grant’s unique server-side token used before allowing the request | Retries/range requests need an explicit session policy |
Expiry normally governs admitting a request; an already admitted transfer can continue under the service's stated policy. A promise to terminate bytes immediately on revocation requires a serving-layer cancellation protocol and cannot be inferred from a presigned URL.
16Decision ledger and limitations
Storage and publication choices
| Decision | Benefit | Cost/limit | Change trigger |
|---|---|---|---|
| Immutable chunks plus atomic KeyHead | No partial published versions | Manifest/GC lifecycle complexity | Mutable-file semantics require another interface |
| Three-copy hot publication | Simple reads and selected failure tolerance | 3× raw bytes and write traffic | Cold objects justify encoding overhead |
| Conditional expected-version overwrite | Prevents lost concurrent updates | Clients handle conflicts | Last-writer-wins is explicitly the desired contract |
| Multipart resume | Retries only missing parts | Session/part metadata and orphan cleanup | Small objects use a simpler single-part path |
| Version-pinned reads/caches | No mixed-version response | Retention/pins and versioned cache keys | Product chooses different cache freshness semantics |
Cold-data and locality limits
Replication protects against selected live failures; backups/version retention protect different mistakes. Erasure coding reduces bytes but adds reconstruction work, and its failure tolerance depends on independent placement of fragments. Cross-region replication introduces latency, cost and possibly a nonzero recovery-point gap. For any chosen policy, state which failures an acknowledged object survives.
What remains outside this design
Strong per-key origin visibility does not promise atomic transactions across multiple objects, immediate CDN invalidation or a consistent snapshot of a huge listing unless separately implemented. Signed URLs delegate a bounded capability and can often be reused until effective expiry. Their convenience does not make them single-use or instantly revocable without additional enforcement.
17Interview closing
Rehearse the architecture and contract
“I separate the logical object name from immutable bytes. Multipart uploads make parts independently retriable under a stable upload identity. Completion freezes the exact part generations, validates integrity and durable placement, then atomically replaces the expected key version with a complete new version and records the upload result. Two uploads expecting the same prior version cannot both win, and a lost completion reply returns the recorded version. Reads resolve one manifest and pin it, so an overwrite cannot mix chunks from different versions.
Defend the critical boundary
“I scale the byte plane across failure domains, keep metadata in replicated key authorities and use ranges/versioned caches for read bandwidth. Hot replicas simplify low-latency reads; cold erasure coding saves storage with repair cost. Garbage collection atomically marks chunks DELETING only after retained-version, upload and read references are absent. Reference acquisition and publication use the same metadata guard, so a new reader or publisher cannot race a deletion check. The main bottleneck here is tens of GB/s of reads and petabytes of retained bytes, not just request QPS. My next tests are completion races, corrupt-chunk recovery and one-domain repair under live load.”
Answer the follow-up
If the interviewer asks for a shared mutable filesystem, explain that byte-range mutation, locks and namespace semantics need another contract. If they ask for globally immediate reads after a regional write, revisit replication/coordination and latency rather than assuming the single-region KeyHead guarantee extends across asynchronous replicas and caches.
Practise the interview questions
Say your answer aloud before opening the model answer. Then answer the follow-up and compare the reasoning.
Why must an object key remain unpublished while its parts are still uploading?
Reveal a model answer
The key promises one complete object version. Exposing incomplete parts would make reads depend on upload timing and could combine missing or unverified data. I store parts privately and publish the manifest only after the full required set is durable and validated.
Interviewer follow-up
What do existing readers see during overwrite?
Reveal the follow-up answer
They can pin V4 while a new request resolves V5 after publication. A read must not resolve individual chunks against a changing current-version pointer.
What the answer must demonstrate: State which metadata transaction makes the complete version visible to new readers.
A part-upload acknowledgement is lost during a two-GiB multipart upload. How does the client resume without restarting the whole object?
Reveal a model answer
The upload session U31 and part number identify reusable work. The client retries or queries that part with the expected checksum and continues the remaining parts. Completion names the verified ordered part set; the retry does not create a second whole report.
Interviewer follow-up
What if completion’s response is lost too?
Reveal the follow-up answer
U31 retains its committed completion result V5. Retry or status lookup resolves that outcome instead of blindly creating another publication. A completion replay must match the recorded manifest fingerprint before returning that success; reusing U31 with a different part set is a conflict.
What the answer must demonstrate: Both part and completion operations need stable identities.
What does a 4+2 erasure code buy compared with three replicas?
Reveal a model answer
It uses six fragments for four fragments’ worth of original data, about 1.5× payload rather than 3×. Under the code/placement assumptions, any four valid fragments reconstruct the data. It trades stored bytes for more complex reconstruction and repair work.
Interviewer follow-up
Does that promise survive two rack failures?
Reveal the follow-up answer
Only if fragment placement across racks actually preserves four valid fragments after that failure pattern. The code’s fragment tolerance is not automatically a rack or region guarantee.
What the answer must demonstrate: Connect mathematical redundancy to physical failure domains.
Can you verify this multipart file by treating its ETag as MD5?
Reveal a model answer
Not universally. ETag semantics depend on the provider and upload method; multipart ETags need not be the MD5 of the complete bytes. I choose supported explicit checksum algorithms and verify part/object integrity under that documented contract.
Interviewer follow-up
Does a matching checksum prove the uploading client had permission?
Reveal the follow-up answer
No. It establishes integrity relative to the expected checksum, not authority to read/write the object. Authentication and tenant ownership checks are separate.
What the answer must demonstrate: Integrity identifiers are not access credentials.
Can a presigned download link be used twice?
Reveal a model answer
Generally yes within its effective validity; it is a bearer capability, not inherently a one-use token. I scope key/version, operation, and expiry, and protect it from logs/leaks. If a link must work only once, the serving application must atomically record its first use and reject later uses, with an explicit policy for retries and range requests.
Interviewer follow-up
Why might it expire before the written deadline?
Reveal the follow-up answer
The signer’s temporary credentials or another policy may end earlier. The effective authorization includes those dependencies, so the app should handle refresh through an authenticated request.
What the answer must demonstrate: Describe delegated authority and its lifetime accurately.
If object PUT is strongly consistent, is my application database/CDN automatically current?
Reveal a model answer
No. The provider’s per-object contract does not atomically update an external application row or invalidate every cache/region replica. I link those changes through an explicit workflow and pin immutable versions where possible, then state each boundary’s freshness.
Interviewer follow-up
Can deleting the current object prove every old byte is erased?
What the answer must demonstrate: Do not extend one subsystem’s guarantee across independent stores.
Why freeze the part-generation list before validating completion?
Reveal a model answer
A replacement could change part 9 after validation, leaving the final manifest inconsistent with the checked bytes. Finalization freezes the ordered part list and fingerprint and rejects further part-record changes; unreferenced upload bytes can be collected later. Storage must also reject overwriting those chunk identities: frozen metadata is insufficient if a reusable upload URL can still replace the bytes.
Interviewer follow-up
What if the completion reply is lost?
Reveal the follow-up answer
The upload’s completed state and resultVersion are stored atomically with KeyHead publication. Retrying U31 returns V5 rather than publishing again.
What the answer must demonstrate: Name both the freeze and publication boundaries.
V6 overwrites the object while a range read of V5 is streaming. What should the response contain?
Reveal a model answer
Only V5. Resolve and protect one immutable manifest at request start, map the range to its chunks and keep that version throughout. New current-key reads can resolve V6, but existing reads must not re-resolve current per chunk.
Interviewer follow-up
Can GC delete V5 immediately after the overwrite?
Reveal the follow-up answer
No while a retained-version root or reader pin exists. New references/pins require LIVE, while GC atomically changes LIVE to DELETING only with zero roots/pins and aborted abandoned uploads. If the pin wins, GC is blocked; if DELETING wins, the pin is rejected and the reader retries or returns the defined deleted-version error. Grace time is not the guard.
What the answer must demonstrate: Separate name visibility from in-flight version lifetime.
Blank-page exercise · 45 minutes
Build the answer yourself
Store the uploading client’s two-GiB report as 32 resumable parts. Lose part 9 and completion responses, fail a storage node, overwrite V4 concurrently, and share only authorized V5 access.
- Separate key metadata, upload session, manifest, and bytes.
- Calculate byte throughput and redundancy overhead.
- Trace multipart retry and atomic publication.
- Read a pinned byte range and validate integrity.
- Explain failure-domain placement and repair.
- Scope signed grants and retained-version deletion.
Check that each component and design decision follows from your requirements and workload.
Recall the key ideas
Answer from memory before opening each card. Explain why the choice works and what it costs. Revisit missed cards tomorrow.
Design a distributed object storeWhat is an object key?Recall first, then reveal
A name inside a bucket/namespace, resolved to object metadata and bytes; it is not automatically a filesystem path.
Namespace + key → object version.
Return to lessonDesign a distributed object storeWhen should a new upload become visible?Recall first, then reveal
After its required parts and integrity checks are complete and the metadata manifest is committed.
Bytes ready, then publish the name.
Return to lessonDesign a distributed object storeWhat is a presigned URL?Recall first, then reveal
A link that lets anyone possessing it perform its specified operation on the specified resource before effective expiry, within the signer’s permissions. It is normally reusable.
Possession grants the scoped operation.
Return to lessonFinal revision
Summary and interview notes
Store verified durable bytes before publishing a complete immutable version. Freeze the chosen parts, then atomically update the key and completion result. Readers keep one version; cleanup deletes chunks only after every protected upload, version and reader has released them.
Remember these points
- Immutable chunk identity requires enforced write protection or an exact provider version, not a unique-looking path.
- Freeze exact part generations, then conditionally publish KeyHead and the replay result atomically.
- The current object version, retained historical versions and active readers each create references that prevent their chunks from being deleted.
- GC marks a chunk DELETING only when no upload, version or reader reference remains. New references must then fail, preventing a later reader or upload from reusing bytes scheduled for deletion.
- Three replicas and 4+2 coding trade stored bytes for read/repair complexity under explicit failure-domain placement.
Interview tips
- Compute both object QPS and byte throughput; this workload is dominated by petabytes and read bandwidth.
- Show competing completion transactions, then reverse a reader-versus-GC race under the same metadata guard.
Important qualifications
- S3 syntax and provider guarantees are distinct from this custom service's upload-expiry, fingerprint and reference protocol.
- Strong per-key origin visibility does not create a multi-page listing snapshot or invalidate a CDN.
- Anyone holding a presigned URL may normally reuse it until effective expiry; additional serving-side checks are required to restrict identity, use count or immediate revocation.
Technical references
- Amazon S3 overview and consistencyProvider-specific object, versioning, permissions, and strong read-after-write capabilities.
- Amazon S3 multipart uploadDocuments upload sessions, ordered part completion, integrity checks, and multipart ETag limitations.
- Amazon S3 presigned URLsDocuments delegated method/key access, reuse, expiry, and credential dependencies.
- Amazon S3 conditional writesProvider-specific create-only and conditional-update options; enforce conditions at storage, not through key naming alone.
Practice marks stay in this browser.