System-design interview · Extended interviews
Design a distributed key-value store
Design a store that looks up values by key, rejects conflicting updates, preserves acknowledged writes across replica failures and moves partitions while serving requests.
You will learn to
- Explain the path from a key lookup to a durable conditional update.
- Distinguish replication agreement from merely counting read/write responses.
- Recover a partition leader and move ownership without accepting stale writes.
Practice in this chapter
8 interview questions with model answers and follow-ups.
Go to interview practiceUseful foundations: CAP theorem: consistency, availability, and partition tolerance · Quorums, consensus, leases, and fencing · Consistent hashing and virtual nodes
Workload and timing examples are interview assumptions.
Dotted concept links open the relevant explanation in a new tab.
01Problem and scope
A distributed key-value store maps a key within one tenant's namespace to value bytes and a version; the store does not interpret the value's contents. Before choosing replication or partitioning, specify which updates must succeed together, what a later read must see and which failures an acknowledged write must survive. This design supports exact-key lookup, replacement and deletion, with conditional updates that reject an obsolete expected version. Joins, arbitrary search and transactions across unrelated keys require different contracts. Cart-42 at version 7 illustrates concurrent updates; the store does not interpret its application-level contents.
Interviewer: “Make that store available everywhere.” Candidate: “Must two devices updating cart-42 agree immediately, or may they accept competing versions and merge later?” Interviewer: “Prevent one device silently overwriting a completed edit. Start within one region.” We require each key's successful updates to have one agreed order; we are not offering every operation of a general-purpose database.
Version checks prevent concurrent edits from silently replacing each other. Client A adds an item while client B removes an item, both starting from cart-42 version 7. The store cannot decide the correct shopping meaning of these competing edits, but it can prevent both callers from believing they replaced the same version. The rejected caller rereads the cart, and the application decides how to combine the edits.
We exclude cross-key transactions, arbitrary search, global range scans and active-active cross-region writes. A tenant can have many keys, but one atomic operation addresses one key plus the store’s internal request-result metadata. This scope lets us show which replica group may update a key and when an update is safe to acknowledge. We will use tested consensus and storage engines; an interview sketch describes their required behavior rather than pretending a new database implementation is a weekend project.
02Functional requirements
- Write a value. Create or replace a tenant-scoped value up to 1 MiB. The response returns a version token identifying the committed state.
- Read an exact key. Read one exact key and receive its bytes and version, or an authoritative not-found result. A not-found response must follow the same consistency rule as a returned value.
- Conditionally replace or delete. Replace or delete only if the supplied version matches. A mismatch returns conflict and the current version under the same authority, allowing client A to reconsider the requested edit.
- Retry a mutation. Retry a mutation with the same request ID and canonical payload within the documented retry horizon. The caller receives the original outcome even when the first response disappeared.
- Expand and move partitions. Increase cluster capacity and move partitions without losing committed records or allowing two independent owners to accept conflicting writes.
- Recover replicas and backups. Recover a failed replica and restore historical backups without exposing partially restored state as authoritative.
Acceptance boundaries
Unconditional PUT is allowed only when the application deliberately accepts replacement semantics. It is not a hidden merge operation. DELETE creates a tombstone, a stored deletion marker ordered with updates; recreation receives a fresh version so an old expected-version token cannot accidentally match a new incarnation. Administrative scans for backup and repair are separate privileged interfaces, not an accidental promise of a public global range query.
03Non-functional requirements
- Regional latency. Target p95 of 20 ms per operation under normal conditions. Global low-latency writes are outside this initial promise.
- Availability. Assume 99.95% eligible-operation availability over a month, with brief election pauses and minority-partition unavailability explicitly allowed.
- Durability. Acknowledged mutations survive one replica-node or zone failure using three appropriately placed replicas. A complete regional loss needs a separately defined backup/replication recovery plan.
- Consistency. Provide linearizable operations within each key in its home region: a later read sees a completed write or something newer. Only a group able to establish the required majority accepts these operations.
- Workload and isolation. Assume 100,000 reads/s and 20,000 writes/s, 1 KiB average values and a 1 MiB maximum. Enforce value-size and per-tenant byte quotas so an abusive writer cannot exhaust a partition.
- Retention. Keep committed live values until deletion, retry results for an assumed one-hour horizon, and backup history for an assumed thirty-day retention policy.
Correctness takes priority
- One mutation order. All successful conditional mutations of a key fit one order; at most one succeeds against a particular current version.
- Real-time reads. A linearizable read respects completed operations. An isolated former leader must refuse strong reads even if its disk looks healthy.
- Explicit degradation. Minority-isolated clients receive unavailable, not misleading success. A separately named stale-read API may exist, but cannot silently replace GET. The availability percentage never authorizes dropping successful writes.
04Capacity estimates
Disk capacity depends on retained live values; write bandwidth depends on every replacement, its replicated copies and storage maintenance. Keeping those quantities separate prevents a busy store from looking small merely because it repeatedly updates the same keys.
Assume 100,000 reads/s, 20,000 writes/s and 1 KiB average values.
| Estimate | Arithmetic | Sizing consequence |
|---|---|---|
| Value ingress | 20,000 × 1,024 = 20.48 MB/second, about 1.77 TB/day | Write traffic, not permanent live-data growth |
| Three-copy value writes | 1.77 TB/day × 3 ≈ 5.31 TB/day | Before log and compaction overhead |
| Logical live dataset | Ten billion keys × (1 KiB + 100 B key/version overhead) ≈ 11.24 TB | Replacements do not add a new live value forever |
| Replicated live dataset | 11.24 TB × 3 = 33.72 TB | Before indexes, spare capacity and temporary compaction files |
| Capacity equivalents | 33.72 TB / measured 500 GB usable live data per node ≈ 68 | Not a final server or partition count |
| Value-read bandwidth | 100,000 × 1,024 = 102.4 MB/second | A single hot key can dominate even if bytes fit |
| Leader-partition lower bound | 20,000 / benchmarked 3,000 writes/s ≈ seven | Smaller partitions can move independently, but each replica group requires consensus coordination |
| One-hour retry history | 20,000 mutations/s × one hour = 72 million results | At 100 B each: 7.2 GB logical, before indexes/replicas |
| In-flight operations | 120,000/s × assumed 20 ms average ≈ 2,400 | Little's law requires a stable boundary and an average |
Capacity is not placement
Many logical partitions share nodes. Throughput, failure reserve, distinct-zone replicas and uneven traffic may require more than 68 node equivalents. A p95 latency target is not an average; the 20 ms average in the concurrency calculation is a separate workload assumption.
Retries add load without adding successful writes
Failed conditional results can also need retention. A ten-minute outage does not stop client retries, so attempted ingress may greatly exceed successful mutation rate. Include admission control and client backoff. Retained live bytes and daily rewritten bytes are separate measurements.
05APIs and contracts
The API turns the version rule into a client workflow: read a value and its version, submit the intended change with that version, then handle a conflict or recover the outcome of a timed-out request. The request ID identifies the attempted change; the version identifies the state it expects to replace.
| API | Example | Meaning |
|---|---|---|
| Read | GET /kv/cart-42 |
Return bytes and version=7 |
| Conditional replace | PUT /kv/cart-42 with {"expectedVersion":7,"requestId":"req-a-9","value":{"items":["book","pen"]}} |
Commit only if version is still 7 |
| Delete | DELETE /kv/cart-42?expectedVersion=8 |
Ordered removal, not an untracked disk erase |
Every request is scoped by authenticated tenant, not a caller-selected namespace alone. A mutation includes a requestId scoped to the authenticated tenant and key. Within that scope, its fingerprint identifies the supplied input: expected version, operation type and value. Reusing the ID with different input is rejected. The same requestId on another key is a separate operation, because those keys may have independent partition authorities. Return 409 for a version conflict, 413 for the size limit, 429 for tenant admission limits and a retryable unavailable result when the required owner cannot be reached. A timeout leaves the outcome unknown.
Versions are opaque persistent tokens, not timestamps supplied by client A. A recreated cart must not reuse the deleted cart’s version. For this design, use a partition incarnation plus ordered mutation revision, carrying that identity through migration. The implementation can choose another proven nonrepeating representation.
There is no public listing pagination because range scans are excluded. Administrative snapshots expose a snapshot identifier and continuation cursor under a separate consistency contract. The mutation retry horizon is one hour in this exercise; clients older than that must use an explicit status/reconciliation path or reread before forming a new conditional intent. They cannot expect deduplication records to exist forever.
06Data model and access patterns
The versioned API needs more than stored values: it must remember retry outcomes, route each key to its current owner and recover the agreed update history. A replicated log stores commands in that agreed order; a state machine applies them using fixed rules to produce the next state and result.
The metadata names positions in this process. An owner epoch identifies a placement generation, a log term identifies a leadership period, and a log index identifies a command position. The applied index records how far a replica has installed committed commands into its local state. A snapshot saves state at such a boundary so recovery knows where replay must resume.
| Record | Key and fields | Why it exists |
|---|---|---|
| Value | (tenant,key), bytes, version, tombstone |
Exact lookup and conditional mutation |
| Request result | (tenant,key,requestId), fingerprint, original result |
Distinguishes successful retry from a new conflict |
| Partition metadata | hash range, replica members, owner epoch | Routes to the current authority |
| Replicated log | partition, term, index, ordered command | Recovers agreed state transitions |
| Snapshot manifest | partition incarnation, applied index, checksum | Defines a complete recoverable state |
A storage engine can keep recent ordered entries in a memtable, an in-memory sorted structure backed by durable recovery data, and flush immutable sorted files to disk. A log-structured merge design later combines files through compaction, removing obsolete versions where retention and replication safety permit. Sequential writes are efficient, but reads may inspect several files and compaction rewrites bytes. Bloom filters can cheaply rule out files that definitely lack a key; a positive result is only a possibility.
The replicated consensus log records the commands agreed by the replica group; the storage engine's write-ahead log lets a node recover its local updates after a crash. An implementation may integrate them or avoid redundant logging with a carefully justified protocol; “we have two logs” is not itself a guarantee. Measure write amplification, read amplification, disk space, and fsync latency.
Apply a committed command as one atomic local engine batch: update the value or tombstone, record the request result, and advance applied-index metadata together. On restart, replay only according to the engine’s recovery contract, never expose a value without its matching deduplication result. Snapshotting includes those records and the applied boundary, rather than copying arbitrary files at unrelated moments.
A tombstone suppresses old values still present in immutable files or stale replicas. Retiring it requires proof that the relevant older state cannot reappear under supported repair and snapshot rules. User deletion and physical erasure from backups are distinct policies. The directory owns placement metadata; it does not contain the user value and cannot decide whether client A’s conditional replacement succeeded.
07Basic working design
Begin with one server, a key index, a tested local storage engine and a durable recovery log. An in-memory dictionary alone would lose cart-42 on restart. The API authenticates client A, reads version 7, and accepts req-a-9 with the replacement value only inside the storage engine’s serialized mutation path.
The owner checks the request-result table first. If req-a-9 is new, it compares the expected version with the current cart, constructs version 8 and atomically records both value and result under its durable-write policy. Only then does it reply. A crash after the log becomes durable but before the reply can be recovered: replay reconstructs both records, and a retry returns version 8.
GET consults this owner and the same ordered state; DELETE installs a new tombstone version through the same path. This baseline is already correct for concurrent requests on one machine if its serialization and recovery protocol are correct. It does not survive loss of the only disk or serve traffic during machine repair.
We now have a concrete benchmark target: exact-key read latency, conditional-write latency including log synchronization, live dataset capacity and write amplification under steady-state compaction. Measuring a memory-map microbenchmark would omit the very work that supports the promised acknowledgment.
A local atomic engine batch keeps client A’s value and retry result together.
Read each connection in order
- syncGET / conditional PUTTenant applications → Authenticated KV API
- syncValidated key and req-a-9Authenticated KV API → Single storage owner
- syncAtomic value + result; durable commitSingle storage owner → Engine files + recovery log
- syncReturn committed versionSingle storage owner → Authenticated KV API
08Find the baseline flaws
The live dataset is 11.24 TB before copies. It exceeds the assumed 500 GB usable live budget by more than twentyfold, so one storage server cannot hold it. Rewriting existing keys generates log and compaction traffic even when the number of live keys stays constant. A design counting only retained values can run out of write bandwidth first.
An asynchronous second copy is not enough for the durability contract. At t0, A logs version 8 and acknowledges client A; at t1, A’s disk is destroyed before B receives it. Promoting B loses an acknowledged write. Waiting for a second durable copy improves this interval, but an election protocol must also prevent promoting a history that omits committed entries.
A naive read-then-write conditional check fails independently: client A and client B both read version 7, both compare outside the serialized path, and both write a replacement. The last writer wins while both callers heard success. The check and update must be one ordered state-machine action, not two HTTP calls.
Finally, load-balancing reads across stale copies breaks the selected GET contract. Client A completes version 8, then reads version 7 from B. A replica must also confirm that it can serve current reads; copying writes alone is insufficient. These counterexamples explain why the next changes address ordering, placement and storage behavior separately.
09Improve the design, step by step
Change one: replicate one ordered decision stream. A single-disk loss motivates three replicas across failure zones, with a tested leader-based consensus protocol. The leader proposes commands, waits for the protocol’s durable commit condition and applies them before success. A valid replacement leader preserves committed history. This changes disk-loss recovery from “restore yesterday’s cart” to continuing from committed state. It costs inter-replica bandwidth, synchronization latency and temporary unavailability without a majority. The new danger is a stale leader answering strong reads; read authority must be confirmed. Asynchronous replicas are simpler and may improve availability for a weaker contract, but they are rejected for acknowledged-loss protection here.
Change two: split many keys across independent authorities. The 11.24 TB dataset and measured per-owner throughput trigger hash-based logical partitions with separate replica groups. The router resolves (tenant,cart-42) to partition 18. Moving a small logical partition changes fewer placements than replacing one giant physical-node modulo map. Benefits are aggregate capacity and parallelism across keys. Costs include a replicated directory, more consensus groups and coordinated migrations. A router may use an old map, so storage owners check the placement epoch and redirect requests sent to the wrong owner. Range placement would be preferable for ordered scans, but our exact-key API does not need them. A single very hot key still cannot be split without changing its semantics.
Change three: budget the engine’s deferred work. Steady-state random updates and disk pressure motivate an LSM-style engine with sorted files, Bloom filters and managed compaction. Batching improves sustained ingestion; file filters avoid some absent-key reads. It costs background CPU, rewritten bytes and temporary space. Compaction debt can stall foreground writes, so limit ingestion when maintenance cannot keep up. A B-tree engine remains a reasonable alternative for the measured read/update mix; benchmark both rather than calling an LSM universally faster. Large values may require separate blob placement, but that adds garbage-collection and publication boundaries and is deferred until the 1 MiB workload demonstrates a need.
Change four: isolate operational work. Replica catch-up, backup and tenant bursts can consume the same I/O as client A’s request. Reserve bandwidth and concurrency for each class; throttle migrations and apply per-tenant byte limits before queues grow indefinitely. This protects the 20 ms objective at the cost of slower administrative progress and explicit 429/unavailable responses. Unlimited buffering is rejected because it converts overload into latency and memory exhaustion. If a workload truly needs long asynchronous ingestion, expose a different admission and completion contract rather than silently weakening PUT.
Each step keeps the per-key decision at one authority. Scaling the cluster never changes a successful expected-version check into a best-effort suggestion.
10Detailed architecture
Authenticated routing
Client libraries contact an authenticated gateway or route directly through an equivalent authenticated protocol. The router caches a versioned partition map from a durable metadata quorum. Hashing locates a logical range; metadata identifies its replica group and current routing epoch. The router does not pick a random replica for a strong operation.
Partition replication
Partition 18 has leader A and followers B/C in separate configured failure domains. Each node holds the replicated log and its local state engine. The engine is an implementation boundary inside a storage node, not a fourth independent copy. Leaders order mutations, apply committed commands, and perform a safe read protocol. Followers replicate and catch up; they are eligible for leadership only under the consensus election rules.
Movement and background work
Each request waits for authentication, routing, consensus or a strong-read check, and the response. Compaction, repair, migration, metrics and backup run in the background with limits on their resource use. The final diagram makes these distinctions visible so an arrow to a directory or a backup cannot be mistaken for a committed user-data write.
Implementation option and limits
A coherent implementation uses a proven Raft library for each logical partition and RocksDB for local ordered state, with a small replicated metadata service for placement. RocksDB is an embedded storage engine, not a distributed database; it does not supply ownership, consensus or the retry protocol. An existing distributed database is preferable when its documented operations meet the contract. A small control-plane store such as etcd can hold placement metadata; the ten-billion-key payload estimate is not a recommendation to place the entire dataset in etcd.
Each partition's replica group maintains one ordered update history. Metadata identifies that group, and bandwidth limits keep migration and backup work from blocking client requests.
Read each connection in order
- sync1. Key, expectedVersion, requestIdTenant applications → Authenticated router / cached map
- sync2. Resolve owner + epochAuthenticated router / cached map → Metadata quorum / partition map
- sync3. Strong operation for partition 18Authenticated router / cached map → Partition 18 leader A
- syncRoute other hash partitionsAuthenticated router / cached map → Other partition replica groups
- replication4. Replicate ordered commandsPartition 18 leader A → Partition 18 follower B
- replication4. Replicate ordered commandsPartition 18 leader A → Partition 18 follower C
- sync5. Apply value + result atomicallyPartition 18 leader A → A: local engine and log
- syncLocal durable follower statePartition 18 follower B → B: local engine and log
- syncLocal durable follower statePartition 18 follower C → C: local engine and log
- controlPublish supported ownership transitionPlacement / membership controller → Metadata quorum / partition map
- controlSnapshot/catch-up before cutoverPlacement / membership controller → Partition 18 leader A
- asyncConsistent snapshot + log boundaryPartition 18 leader A → Snapshot / backup worker
- asyncRetain verified historical snapshotSnapshot / backup worker → Protected backup storage
11Write path and acknowledgement
A mutation succeeds only after its ordered command is durably committed and applied together with its request result. The following trace tests two updates against the same version.
- The router hashes
(tenantA,cart-42)and finds partition 18, ownership epoch 6, led by node A with followers B and C. Routing metadata is cached, but a stale epoch receives a redirect or rejection. - A receives request
req-a-9and proposes a command containing the expected version and new value. The replicated log orders it with other commands for partition 18. - The command is durably replicated and committed according to the consensus protocol. When applied in log order, it checks version 7, writes version 8, and records the request result. A competing request based on version 7 cannot also replace version 8.
- A replies with version 8 only after commit and application. Client A's next strong read goes through a leader that confirms its current authority and has applied the necessary committed index; a former isolated leader must not answer stale data as current.
- If the reply is lost, retrying
req-a-9returns version 8. If a different request tries expected version 7, it receives a conflict and must reread before deciding how to merge application data.
The actual acknowledgment includes the result identity, not just a generic 200. If the version check fails when its command is applied, the failure result is also associated with req-a-9 so a repeated request does not change meaning after another cart update. The protocol may optimize known duplicates, but correctness cannot rely on an unreplicated memory cache of request IDs.
Deletes follow the same ordered path and install a fresh tombstone version. A successful deletion does not authorize an old replica to resurrect version 7. During migration, clients may repeat the request through a new owner, so request-result state must move with the key or remain accessible through the owner’s supported retry protocol. Copying only user values would reopen the lost-response ambiguity.
12Read and delivery path
Before serving a strong read, the leader must confirm that it still leads and has applied the required committed commands. Its label alone proves neither.
A block cache inside the engine speeds access without inventing a second authority: cached blocks are interpreted through the current engine state. An application-side value cache would need a separate validated freshness protocol to serve strong GET. We do not quietly add such a cache just to hit a latency target. Clients wanting low-latency stale snapshots can opt into an explicitly weaker operation.
13Correctness deep dive
The consensus protocol orders commands; the deterministic state machine decides their meaning. Suppose committed log positions 101 and 102 contain client A’s add-pen request and client B’s remove-book request, both expecting version 7.
apply(command, committedIndex):
saved = result(command.tenant, command.key, command.requestId)
if saved exists:
outcome = original result if fingerprint matches else invalid-reuse
else:
current = value(command.tenant, command.key)
if current.version != command.expectedVersion:
outcome = conflict(current.version)
else:
outcome = success(freshVersion(partitionIncarnation, committedIndex))
prepare replacement bytes or tombstone with that version
atomic engine batch:
install replacement only for a new successful request
store fingerprint and outcome only when no prior result exists
advance applied index, including duplicate and rejected commands
| Applied position | Before | Decision | After |
|---|---|---|---|
| 101: req-a-9 expects 7 | cart version 7 | Match; return version 8 in this simplified notation | book + pen, version 8 |
| 102: req-b-4 expects 7 | cart version 8 | Conflict; record failed result | Still version 8 |
| Later: req-a-9 repeated | Saved req-a-9 success | Return original result | No extra mutation |
The displayed version 8 is shorthand for the opaque nonrepeating token. The key point is that client B’s check happens after client A’s applied change in the agreed order. Another write cannot run between the version check and its update.
Committed command order and saved results determine the outcome; a lost response does not create another update.
Read each connection in order
- syncreq-a-9: replace if version 7Client A’s device → Partition leader
- syncreq-b-4: replace if version 7Client B’s device → Partition leader
- syncOrder and durably commit 101 then 102Partition leader → Replica quorum / engine
- syncApply 101: value v8 + req-a-9 resultReplica quorum / engine → Replica quorum / engine
- syncApply 102: conflict + req-b-4 resultReplica quorum / engine → Replica quorum / engine
- returnApplied committed resultsReplica quorum / engine → Partition leader
- blockedVersion 8 reply lostPartition leader → Client A’s device
- returnConflict: current version 8Partition leader → Client B’s device
- syncRetry req-a-9 unchangedClient A’s device → Partition leader
- syncRead original committed resultPartition leader → Replica quorum / engine
- returnreq-a-9 succeeded at version 8Replica quorum / engine → Partition leader
- returnOriginal success; no second mutationPartition leader → Client A’s device
14Failure and recovery
Recovery must preserve the same per-key history even when the process, replica group or storage medium changes. For each failure below, identify which owner can still prove that history before allowing more strong reads or writes.
| Failure or race | Required response and boundary |
|---|---|
| Leader loses contact after commit | A commits client A's update with a majority, then loses connectivity. B and C can elect a leader under the protocol; A cannot continue confirming authority alone. Client A may see a temporary timeout, but a committed update must survive a valid election. Merely choosing read count R and write count W with R+W>N does not specify leader fencing, version ordering, failed writes, membership change, or linearizable reads. Overlap is one ingredient, not a complete consistency algorithm. |
| Partition migration | To move partition 18, transfer a consistent snapshot to its new replicas, replay changes after the snapshot index, and switch ownership through a coordinated configuration/epoch transition. Old owners reject epoch-6 writes after epoch 7 is active; new owners must not begin from an incomplete copy. Raft membership changes have their own protocol and must not be replaced with arbitrary simultaneous configuration edits. |
| Majority unavailable | A majority loss leaves partition 18 unavailable even if other partitions work. Report affected-key failure instead of describing the whole cluster as uniformly up or down. Operators restore quorum or recover from a validated snapshot/log history; they must not force two disconnected primaries into existence to remove an alert. |
| Disk stall and retry burst | At 20,000 writes/s, a ten-second disk stall creates 200,000 pending writes if nothing limits admission. Bound the proposal queue and bytes in flight. Slow or reject new work before memory exhaustion; preserve the outcome of already committed commands. Client retries reuse identities with jitter, and each request has a deadline so retry fanout cannot become unlimited. |
| Corruption and historical recovery | Corruption detection uses checksums and comparisons appropriate to the engine; a healthy majority is not proof that every historical backup is clean. Replicas may copy an accidental delete. Test restoring cart-42 at a chosen history point into an isolated environment and verify application-visible versions before any promotion. Recovery still has the stated regional-disaster limits; it does not promise survival of every correlated loss. |
15Operations, security, and cost
Authenticate node-to-node replication and administrative control, encrypt tenant data under the required threat model, and authorize tenant-scoped keys at every serving path. A direct follower endpoint must not bypass namespace checks. Rate-limit bytes as well as operation count: one 1 MiB mutation costs about a thousand average 1 KiB mutations in replication payload.
Monitor p95 and p99 committed-operation latency, unavailable partitions, quorum loss, leader churn, disk synchronization time, compaction debt and hot-key concentration. A low global average can hide one inaccessible customer partition. Distinguish client attempts, proposals, committed mutations and conflicts to diagnose a retry storm accurately.
The 33.72 TB three-copy live estimate excludes transient compaction output and recovery reserve. If operational policy permits only 70% steady storage occupancy, the corresponding capacity budget is about 48.2 TB before additional index/log overhead. This is a planning assumption, not a universal engine threshold. Increasing replica count from three to five raises live copy bytes by two thirds while changing the tolerated failure/latency tradeoff; it does not increase per-key write concurrency proportionally.
Roll out engine formats and protocols with mixed-version compatibility. Exercise leader loss after commit, stale-leader reads, duplicate request IDs, deleted-key recreation and migration during writes. Verify both returned histories and durable state with fault injection. A successful backup command or a green cluster membership page does not prove that the promised conditional operation survives those interleavings.
16Decision ledger and limitations
The central choice was a linearizable, single-key API. The tables relate that promise to its replication, storage and routing costs, and identify the requirements that would justify a different contract.
| Choice | Benefit | Price |
|---|---|---|
| Leader-based strong writes | Simple per-key order and conditional updates | Minority partitions stop |
| Eventually consistent multi-writer store | Writes can continue in more partitions | Conflicts and reconciliation become product work |
| LSM-style local storage | Efficient sustained ingestion | Compaction and read amplification |
| Hash partitions | Good exact-key distribution | No natural global range scan |
| Remaining choice | Consequence | Trigger to reconsider |
|---|---|---|
| One key per atomic command | Can prove one key’s updates are correct; cannot keep several keys consistent atomically | Application actually needs multi-key transactions |
| Hash-based ownership | Spreads many-key lookups across owners; global scans are expensive | Range queries become a first-class requirement |
| One-hour retry history | Bounded result storage; old identities need a new recovery rule | Offline clients need longer safe replay |
| Engine block cache, no unvalidated value cache | Preserves the chosen strong-read path | A documented weaker read mode would meet product needs |
An availability-first multi-writer store can accept updates in more isolated locations, but returns concurrent versions or applies a merge rule. A shopping application may prefer merging independent add-item operations to rejecting a temporary partition; that is a different API and conflict model. It cannot be substituted under the existing expected-version promise without explaining changed outcomes.
The first scaling limit may be one hot cart, not total data size. A storage partition can move to a larger owner, but it cannot create parallel successful mutations against the same prior version. Further scaling requires the application to accept independently updated subkeys or a weaker consistency rule.
17Interview closing
“I start with one durable owner and make value, version and retry outcome one recoverable state transition. The single machine fails our storage and failure targets, so I split many keys into logical partitions and replicate each partition through a tested majority protocol.
“Clients route through versioned placement metadata. The leader commits and applies a mutation before acknowledging success, and confirms current authority before a strong read. Conditional mutations with the same expected version are serialized: once one advances the version, the other conflicts. If a reply is lost, retrying the original request identity recovers its recorded result rather than applying another write.
“I pay for replica bytes, log synchronization, compaction and minority-side unavailability. I protect foreground work from migrations and tenant bursts, and test safe membership changes. The next measurements are hot-key concentration, steady-state write amplification and tail latency during one-node failure.”
Interviewer: “Now allow writes in two disconnected regions.” Candidate: “I cannot keep the same immediate conditional-update promise on both sides. I would discuss routing each key to one authority or introducing application-mergeable operations and explicit conflicts. The user-visible semantics must change before the topology does.”
Practise the interview questions
Say your answer aloud before opening the model answer. Then answer the follow-up and compare the reasoning.
Why does a key-value API need conditional writes?
Reveal a model answer
Conditional writes prevent lost updates by combining the expected-version check with replacement atomically. If two clients read version 7, only one replacement can advance it; the other receives a conflict and rereads before merging application intent.
Interviewer follow-up
Could two conditional writes both succeed?
Reveal the follow-up answer
They cannot both succeed against version 7 if the check and update are one ordered atomic operation. With a separate read and write, another client could change the value between them.
What the answer must demonstrate: State the atomic boundary.
The replicas split into one node and two nodes. Who serves writes?
Reveal a model answer
With a three-node majority protocol, the communicating pair can establish leadership and commit. The isolated node cannot confirm authority and must reject or time out strong operations.
Interviewer follow-up
Can it still serve stale reads?
Reveal the follow-up answer
Only through an explicitly weaker API whose callers accept stale data. I would not label those responses linearizable.
What the answer must demonstrate: Name the client-visible availability cost.
A PUT times out. Did it fail?
Reveal a model answer
I cannot infer failure from a lost response. The command may have committed. A stable request ID lets the client retry and recover the recorded outcome.
Interviewer follow-up
Why is expectedVersion alone not always enough?
Reveal the follow-up answer
A retry may conflict after its own successful update changed the version. Recording the request result distinguishes that success from another writer’s change.
What the answer must demonstrate: Timeout is an unknown outcome.
Why not write every value directly into one disk file?
Reveal a model answer
It can work initially, but frequent random rewrites and index maintenance may limit throughput. A log plus sorted in-memory updates and immutable files supports batching, with compaction paying the cleanup cost later.
Interviewer follow-up
What do you measure before choosing it?
Reveal the follow-up answer
Read/write amplification, tail latency, compaction debt, and temporary disk headroom under the actual value distribution.
What the answer must demonstrate: Describe the cost of the optimization.
How do you move a partition while clients are writing?
Reveal a model answer
Copy a snapshot at a known log index, replay later changes, and use a coordinated ownership epoch transition. Stale routers and old owners are rejected rather than letting both sides independently accept writes.
Interviewer follow-up
Can you just update the routing map first?
Reveal the follow-up answer
No. The new owner may lack committed changes, and an old owner may still be active. First copy all committed changes, then transfer ownership so the old owner can no longer accept writes.
What the answer must demonstrate: Routing is not proof of exclusive ownership.
One key receives half your traffic. Will more virtual nodes solve it?
Reveal a model answer
No. Virtual nodes distribute groups of different keys. This key remains one logical item. I can cache or replicate reads under a clear consistency contract, but serial conditional writes retain a bottleneck.
Interviewer follow-up
When would you split the value?
Reveal the follow-up answer
When the application can update subkeys independently. If all subkeys must change together, they still need coordination after the split.
What the answer must demonstrate: Distinguish many-key balance from one-key contention.
Why is routing GET to the node that says “leader” insufficient?
Reveal a model answer
“That node may be isolated from a newly elected majority. I require the implementation’s safe linearizable-read protocol, such as current quorum confirmation and waiting for the necessary applied index, before reading its local state.”
Interviewer follow-up
Could a lease avoid a quorum round trip?
Reveal the follow-up answer
“Only with the protocol’s timing and validity assumptions enforced. An arbitrary wall-clock timeout is not a proof of current authority.”
What the answer must demonstrate: Name both authority and applied-state requirements.
When a key moves to a new partition owner, what must migrate besides its value?
Reveal a model answer
Its nonrepeating version/incarnation, deletion state and required key-scoped request-result history move with a consistent snapshot index and catch-up log. Otherwise a delayed retry can lose evidence of its original outcome. The destination must finish catch-up before it becomes authoritative.
Interviewer follow-up
Can you delete tombstones to reduce transfer size?
Reveal the follow-up answer
“Only under a proven retention/repair rule that prevents older copies from reappearing. An incomplete snapshot must not become authoritative.”
What the answer must demonstrate: Migrate the correctness metadata, not just payload bytes.
Blank-page exercise · 45 minutes
Build the answer yourself
Design a durable key-value store, then lose the leader after client A’s version-7 replacement commits but before the client receives the reply.
- Clarify per-key strong semantics and minority-side failure.
- Estimate live bytes, write traffic, deduplication retention and headroom.
- Draw baseline and evolved partition ownership, then trace PUT and GET.
- Trace two updates against the same expected version, then show how retrying a lost response returns the original result.
- Test migration, stale-leader reads, compaction pressure and restore.
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 key-value storeWhat does expectedVersion=7 protect?Recall first, then reveal
It prevents replacing a value that another operation has already changed to a later version.
Read version; compare before replace.
Return to lessonDesign a distributed key-value storeDoes R+W>N prove linearizability?Recall first, then reveal
No. It gives set overlap under assumptions, but ordering, failed writes, reads, and reconfiguration still need a protocol.
Overlap is not the whole algorithm.
Return to lessonDesign a distributed key-value storeWhy keep a deletion marker?Recall first, then reveal
Older replicas and disk files must learn that the key was deleted before its history is safely reclaimed.
Keep deletion markers until stale values cannot return.
Return to lessonFinal revision
Summary and interview notes
A strongly consistent key-value store gives each key one ordered mutation history and confirms current authority before reads. Partitioning distributes independent keys; replication preserves committed history through the failures named in the contract.
Remember these points
- Compare expected version and update state in the same committed state-machine operation.
- Store the result under tenant, key and request ID, atomically with the mutation; scope the retry promise explicitly.
- Majority overlap alone does not supply safe elections, strong reads or membership changes.
- Move values, versions, tombstones and retry history together; finish copying and replay before serving from the new owner.
- Budget compaction and failure reserve separately from logical live bytes.
Interview tips
- Trace two updates against the same version, then lose the winning response.
- Explain which operations stop on the minority side and how strong reads prove current authority.
- Distinguish an embedded engine, consensus group and metadata service before naming products.
Important qualifications
- The one-hour retry horizon, 20 ms target and storage capacities are workload assumptions, not engine guarantees.
- Conditional updates to one hot key still require a single agreed order; adding hash partitions only spreads work across different keys.
Technical references
- Raft consensus paperPrimary description of leader election, replicated logs, and safe membership changes.
- etcd API guaranteesConcrete documentation of strong operations and weaker alternatives in a real key-value API.
- RocksDB overviewOfficial storage-engine overview for logs, memtables, sorted files, and compaction.
- Dynamo paperPrimary availability-oriented alternative with version reconciliation.
Practice marks stay in this browser.