System designby Learnastra

System-design interview · Extended interviews

Design a distributed key-value store

By Anup Rai

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 practice

Useful 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.

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

  1. Write a value. Create or replace a tenant-scoped value up to 1 MiB. The response returns a version token identifying the committed state.
  2. 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.
  3. 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.
  4. 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.
  5. Expand and move partitions. Increase cluster capacity and move partitions without losing committed records or allowing two independent owners to accept conflicting writes.
  6. 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

  1. Regional latency. Target p95 of 20 ms per operation under normal conditions. Global low-latency writes are outside this initial promise.
  2. Availability. Assume 99.95% eligible-operation availability over a month, with brief election pauses and minority-partition unavailability explicitly allowed.
  3. 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.
  4. 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.
  5. 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.
  6. 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

  1. One mutation order. All successful conditional mutations of a key fit one order; at most one succeeds against a particular current version.
  2. Real-time reads. A linearizable read respects completed operations. An isolated former leader must refuse strong reads even if its disk looks healthy.
  3. 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.

architecture · baselineOne durable owner

A local atomic engine batch keeps client A’s value and retry result together.

One durable ownerA local atomic engine batch keeps client A’s value and retry result together. client to api: GET / conditional PUT; api to owner: Validated key and req-a-9; owner to disk: Atomic value + result; durable commit; owner to api: Return committed versionGET / conditionalPUTValidated key and req-a-9Atomic value + result; durablecommitReturn committed versionACTORTenant applicationsSERVICEAuthenticated KV APISERVICESingle storage ownerSTOREEngine files +recovery logsync
Read each connection in order
  1. syncGET / conditional PUTTenant applications → Authenticated KV API
  2. syncValidated key and req-a-9Authenticated KV API → Single storage owner
  3. syncAtomic value + result; durable commitSingle storage owner → Engine files + recovery log
  4. 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.

architecture · finalPartitioned strong store

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.

Partitioned strong storeEach 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. client to router: 1. Key, expectedVersion, requestId; router to meta: 2. Resolve owner + epoch; router to leader: 3. Strong operation for partition 18; router to other: Route other hash partitions; leader to b: 4. Replicate ordered commands; leader to c: 4. Replicate ordered commands; leader to disk-a: 5. Apply value + result atomically; b to disk-b: Local durable follower state; c to disk-c: Local durable follower state; controller to meta: Publish supported ownership transition; controller to leader: Snapshot/catch-up before cutover; leader to backup: Consistent snapshot + log boundary; backup to archive: Retain verified historical snapshot1. Key, expectedVersion,requestId2. Resolve owner + epoch3. Strong operation for partition18Route other hash partitions4. Replicate orderedcommands4. Replicate orderedcommands5. Apply value + resultatomicallyLocal durable follower stateLocal durable follower statePublish supported ownershiptransitionSnapshot/catch-up beforecutoverConsistent snapshot + logboundaryRetain verified historicalsnapshotACTORTenant applicationsG1SERVICEAuthenticated router/ cached mapG1STOREMetadata quorum /partition mapG1SERVICEPartition 18 leader AG2SERVICEPartition 18 followerBG2SERVICEPartition 18 followerCG2STOREOther partitionreplica groupsG3SERVICEPlacement /membershipcontrollerG3WORKERSnapshot / backupworkerG4STOREProtected backupstorageG4STOREA: local engine andlogG2STOREB: local engine andlogG2STOREC: local engine andlogG2syncreplicationcontrolasyncG1 Tenant and routing boundaryG2 Partition ownership; per-node local enginesG3 Other owners and placement controlG4 Independent history recovery
Read each connection in order
  1. sync1. Key, expectedVersion, requestIdTenant applications → Authenticated router / cached map
  2. sync2. Resolve owner + epochAuthenticated router / cached map → Metadata quorum / partition map
  3. sync3. Strong operation for partition 18Authenticated router / cached map → Partition 18 leader A
  4. syncRoute other hash partitionsAuthenticated router / cached map → Other partition replica groups
  5. replication4. Replicate ordered commandsPartition 18 leader A → Partition 18 follower B
  6. replication4. Replicate ordered commandsPartition 18 leader A → Partition 18 follower C
  7. sync5. Apply value + result atomicallyPartition 18 leader A → A: local engine and log
  8. syncLocal durable follower statePartition 18 follower B → B: local engine and log
  9. syncLocal durable follower statePartition 18 follower C → C: local engine and log
  10. controlPublish supported ownership transitionPlacement / membership controller → Metadata quorum / partition map
  11. controlSnapshot/catch-up before cutoverPlacement / membership controller → Partition 18 leader A
  12. asyncConsistent snapshot + log boundaryPartition 18 leader A → Snapshot / backup worker
  13. 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.

  1. 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.
  2. A receives request req-a-9 and proposes a command containing the expected version and new value. The replicated log orders it with other commands for partition 18.
  3. 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.
  4. 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.
  5. If the reply is lost, retrying req-a-9 returns 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.

sequence · raceTwo updates from version 7

Committed command order and saved results determine the outcome; a lost response does not create another update.

Two updates from version 7Committed command order and saved results determine the outcome; a lost response does not create another update. tara to leader: req-a-9: replace if version 7; lee to leader: req-b-4: replace if version 7; leader to quorum: Order and durably commit 101 then 102; quorum to quorum: Apply 101: value v8 + req-a-9 result; quorum to quorum: Apply 102: conflict + req-b-4 result; quorum to leader: Applied committed results; leader to tara: Version 8 reply lost; leader to lee: Conflict: current version 8; tara to leader: Retry req-a-9 unchanged; leader to quorum: Read original committed result; quorum to leader: req-a-9 succeeded at version 8; leader to tara: Original success; no second mutationPARTICIPANTClient A’s devicePARTICIPANTClient B’s devicePARTICIPANTPartition leaderPARTICIPANTReplica quorum /engine1. req-a-9: replace if version 72. req-b-4: replace if version73. Order and durably commit101 then 1024. Apply 101: value v8 +req-a-9 result5. Apply 102: conflict +req-b-4 result6. Applied committed results7. Version 8 reply lost8. Conflict: current version 89. Retry req-a-9 unchanged10. Read original committedresult11. req-a-9 succeeded atversion 812. Original success; no second mutationsyncreturnblocked
Read each connection in order
  1. syncreq-a-9: replace if version 7Client A’s device → Partition leader
  2. syncreq-b-4: replace if version 7Client B’s device → Partition leader
  3. syncOrder and durably commit 101 then 102Partition leader → Replica quorum / engine
  4. syncApply 101: value v8 + req-a-9 resultReplica quorum / engine → Replica quorum / engine
  5. syncApply 102: conflict + req-b-4 resultReplica quorum / engine → Replica quorum / engine
  6. returnApplied committed resultsReplica quorum / engine → Partition leader
  7. blockedVersion 8 reply lostPartition leader → Client A’s device
  8. returnConflict: current version 8Partition leader → Client B’s device
  9. syncRetry req-a-9 unchangedClient A’s device → Partition leader
  10. syncRead original committed resultPartition leader → Replica quorum / engine
  11. returnreq-a-9 succeeded at version 8Replica quorum / engine → Partition leader
  12. 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.

Foundation · Question 1

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.

What the answer must demonstrate: State the atomic boundary.

Applied · Question 2

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.

What the answer must demonstrate: Name the client-visible availability cost.

Applied · Question 3

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.

What the answer must demonstrate: Timeout is an unknown outcome.

Foundation · Question 4

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.

What the answer must demonstrate: Describe the cost of the optimization.

Follow-up · Question 5

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.

What the answer must demonstrate: Routing is not proof of exclusive ownership.

Follow-up · Question 6

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.

What the answer must demonstrate: Distinguish many-key balance from one-key contention.

Applied · Question 7

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.”

What the answer must demonstrate: Name both authority and applied-state requirements.

Follow-up · Question 8

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.

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 lesson
Design 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 lesson
Design 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 lesson

Final 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.