System designby Learnastra

System-design interview · Extended interviews

Design a distributed key-value store

By Anup Rai

Build a durable exact-key store with conditional updates, safe retries and strong regional reads, then partition independent keys and replicate each partition.

You will learn to

  • Explain value, version and request identity using two concurrent edits.
  • Trace a committed mutation and a current-authority read through a replica group.
  • Account for storage maintenance, hot keys, migration and failure without overstating availability.

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.

01Choose the meaning of a successful operation

A key-value store maps a tenant’s key to value bytes whose contents it does not interpret. Offer exact-key GET, PUT and DELETE, with optional expected-version conditions. Exclude joins, arbitrary search, global range scans and multi-key transactions. Those operations require additional indexes or coordination and should not appear accidentally through an underspecified API.

Use cart-42 at version 7. Client A adds a pen while client B removes a book, both from version 7. The store cannot decide how the shopping application should merge those intentions. It can ensure only one conditional replacement succeeds; the other client rereads and decides what to do next.

Choose strong per-key operations in one home region: a read begun after a successful write returns that write or something newer. This is linearizability, a real-time ordering contract. An isolated replica cannot simply keep accepting conditional updates or serving old values as current. A separate stale-read option can exist, but must be named explicitly.

Assume 1 KiB average values, a 1 MiB maximum and a 20 ms normal-operation p95 target. Acknowledged mutations should survive one replica-node or zone failure, with replica placement chosen accordingly. Regional disaster recovery is a separate promise.

Clarify whether callers need strong current reads or may accept stale replicas, whether operations cover one key or several, and what failure must preserve an acknowledged write. This answer chooses exact-key linearizable operations in one home region with one-zone failure tolerance.

02Functional requirements

  1. Read and mutate exact keys. Provide authenticated GET, PUT and DELETE for opaque value bytes under a tenant’s key.

  2. Update conditionally. Accept an expected version so clients can reject a replacement based on an outdated value.

  3. Recover mutation retries. Return the stored outcome for the same request identity and payload within the supported retry horizon, including after a lost response.

  4. Distribute and recover data. Route keys to their current owner and support safe replica recovery and partition movement without changing the client’s key or operation meaning.

03Non-functional requirements

These are illustrative interview assumptions, not product facts or measured benchmarks. Confirm them before choosing components, then validate the completed design under the stated workload. Here p95 means the 95th-percentile latency: 95% of measured requests take no longer than that value. Report errors and rejected work alongside latency; a fast failure is not a successful outcome.

  1. Capacity and size limits. Plan for ten billion live keys, 100,000 reads/s and 20,000 writes/s, with 1 KiB average values and a 1 MiB maximum. Budget three copies plus indexes, maintenance and retry history.

  2. Latency. Target 20 ms p95 for admitted same-region exact-key operations under normal load. Benchmark durable writes, strong reads and ongoing compaction together; timeouts are misses, not fast successful operations.

  3. Consistency. Provide linearizable per-key reads and mutations: a read started after a successful write sees it or something newer. Only one conditional replacement of the same current version may succeed.

  4. Durability and failure boundary. Use three consensus replicas placed across zones so acknowledged mutations survive one replica or zone failure. Stop strong operations in a partition that lacks a safe majority; regional disaster survival is outside this promise.

  5. Retry and tenant isolation. Retain mutation outcomes for the selected one-hour retry horizon. Authenticate tenant scope on every path, reject conflicting request-ID reuse and enforce byte as well as operation quotas.

04Begin with one durable storage owner

The baseline is one server using a tested local storage engine. The server authenticates the tenant, identifies the key and processes its changes in order. Before changing cart-42, the owner checks whether this request identity already has a saved result. If so, it returns that result after confirming the payload matches.

For a new conditional PUT:

  1. Compare expected version 7 with the current value inside the same atomic operation that installs the replacement.

  2. On success, write the new value, a fresh version and the saved request result together.

  3. Make the operation durable before replying.

A crash after commit but before response can then recover both the value and what the caller should receive.

GET reads the current committed value. DELETE writes an ordered deletion marker rather than relying on an untracked physical erase. Recreating the key receives a new version so an old expected-version token cannot accidentally match a later incarnation.

This single-owner version is complete for local concurrency and process restart under the engine's guarantees. It cannot tolerate losing the only disk or serve while that machine is unavailable. Replication is the next change because of those specific limits, not because a dictionary automatically becomes distributed when copied.

Design diagramA durable conditional operation

Value, version and request result are one atomic local change.

A durable conditional operationValue, version and request result are one atomic local change. client to api: Key, expected version, request ID; api to owner: Authorized operation; owner to disk: Durable atomic value + result; owner to client: Committed version or conflictKey, expected version, requestIDAuthorized operationDurable atomic value + resultCommitted version or conflictACTORTenant applicationSERVICEAuthenticatedkey-value APISERVICESerialized storageownerSTORELocal engine andrecovery datasyncreturn
Read each connection in order
  1. syncKey, expected version, request IDTenant application → Authenticated key-value API
  2. syncAuthorized operationAuthenticated key-value API → Serialized storage owner
  3. syncDurable atomic value + resultSerialized storage owner → Local engine and recovery data
  4. returnCommitted version or conflictSerialized storage owner → Tenant application

05Separate stored data from rewritten bytes

Assume 100,000 reads/s and 20,000 writes/s. With 1,024-byte average values, value ingress is 20.48 MB/s, about 1.77 TB/day. That is write traffic, not necessarily new live storage: repeatedly replacing the same cart rewrites bytes without adding a new live cart each time.

Ten billion live keys at 1 KiB plus 100 bytes of key/version overhead require about 11.24 TB logically. Three copies require 33.72 TB before engine indexes, logs, temporary compaction files and failure reserve. If a tested node has 500 GB of usable live-data capacity, the ratio suggests about 68 node equivalents, not a final fleet layout. Placement, throughput and failure headroom still determine the actual topology.

Retaining one hour of results at 20,000 mutations/s means 72 million results. At an illustrative 100 bytes each, that adds 7.2 GB logically before replication and indexing; rejected conditions may also need saved outcomes.

Read payload is roughly 102.4 MB/s before overhead. A hot key can dominate one owner despite comfortable total bandwidth. Benchmark durable writes and steady-state storage maintenance, not only an in-memory map. Client retries increase attempted traffic during outages without increasing useful completed writes.

06Make versions and request identities distinct

Interfaces

Request or message Contract
GET /kv/cart-42 Returns bytes and an opaque version, or an authoritative not-found result.
PUT /kv/cart-42 Replaces only the named current version and saves the result.
DELETE /kv/cart-42 Records an ordered deletion under the same retry contract.

Conditional replacement request

PUT /kv/cart-42

Request fields for client A’s update

expectedVersion: 7
requestId:       the stable identity reused for this attempted update
value:           replacement cart bytes containing A’s added pen

Conditional delete request

DELETE /kv/cart-42

Delete request fields

expectedVersion: the version being deleted
requestId:       the stable identity reused for this attempted deletion

Stored records

Record Fields or identity Purpose
Value tenant, key, bytes, version, tombstone Current application-visible state.
RequestResult tenant, key, requestId, fingerprint, result Distinguishes a retry from a new mutation.
Partition metadata range, members, ownershipVersion Routes keys to their current replica group.

A version identifies the state being replaced. A request ID identifies the attempted operation. Reusing an ID with another value or condition must fail; otherwise a caller could accidentally retrieve the result of unrelated work. Tenant identity comes from authentication on every access path.

Specify a retry horizon, such as one hour, rather than retaining results forever by implication. Older uncertain operations require status recovery or an application-level decision after rereading; issuing a new unconditional write blindly can repeat the original change. Return distinct conflicts, size-limit errors, admission limits and unavailable outcomes. A network timeout alone cannot establish whether the mutation committed.

07Give each partition one agreed command history

Use a proven leader-based consensus implementation for each replica group, with three appropriately placed replicas in this exercise. The leader proposes commands to an ordered log. The protocol determines when a command is durably committed and ensures a valid replacement leader preserves committed history. Apply a committed command before acknowledging its application result.

Suppose two commands expect version 7. A's command is applied first, finds version 7 and creates version 8. B's command is applied next, sees version 8 and records a conflict. The comparison occurs while applying the agreed order, not in an earlier unprotected read. Two callers therefore cannot both successfully replace the same current version.

The local engine applies the value or deletion marker, saved result and applied-log position as one recoverable batch. After a crash, recovery knows which commands are installed and which must be replayed. An in-memory duplicate-request cache alone would lose the result when leadership changes.

Three copies are not sufficient without the protocol. If A acknowledges before another durable copy receives the update, losing A can lose a successful write. If a replacement leader can choose an obsolete history, waiting for extra copies still does not establish safe failover. Use the library's supported election and membership rules rather than inventing them during the interview.

08Explain why a strong read cannot trust an old leader

A client routes GET to the partition leader, but the process still needs to establish that its authority is current. Imagine A loses contact with B and C. They elect a new leader and commit version 8 while A's disk still contains version 7. A must not answer a later strong GET from its isolated local copy.

One supported approach uses the consensus protocol to confirm leadership with a quorum and obtain a safe read position. The leader waits until its local state has applied through that position, then returns the key. A correctly implemented lease can be an alternative, but it adds timing assumptions that the baseline need not introduce.

A storage-engine block cache speeds disk access without changing which committed state is read. An independent application value cache cannot silently serve this strong API unless it has an appropriate freshness-validation protocol. Followers may support explicitly stale reads for callers who accept them, with different documented semantics.

Not-found responses also matter: after a successful creation, an old replica saying “absent” violates the same promise as returning an old value. During majority loss, the affected partition becomes unavailable for strong operations even if its surviving node can still answer network requests.

09Partition independent keys and preserve ownership on moves

Hash the authenticated tenant and key into logical partitions. A routing directory maps each partition to its replica group. Many logical partitions may share physical nodes; choosing many manageable units allows capacity changes without moving an entire server's data at once. A cached directory reduces lookup work, while stale owners reject or redirect requests using a placement version.

Different partitions can process independent keys in parallel. A single hot cart still has one ordered mutation history. Moving that partition to a larger owner can help resource pressure, but adding hash partitions does not create parallel successful replacements of the same version. Splitting that value into subkeys would change which updates can commit together.

To move a partition, transfer a consistent snapshot, replay subsequent committed changes and verify the destination is caught up before activating its ownership. Include values, deletion markers, versions and request-result history. Copying only live values would lose retry safety and might resurrect deleted state.

Coordinate cutover through the replication system's supported configuration protocol and a current routing generation. The previous owner must lose write authority before an independent new owner accepts conflicting commands. Updating the directory before the copy is ready creates a missing or stale partition, not a completed migration.

Design diagramIndependent replicated partitions

Routing chooses a group; that group commits mutations and validates strong reads.

Independent replicated partitionsRouting chooses a group; that group commits mutations and validates strong reads. client to map: Resolve partition and ownership; client to leader: Strong operation for cart-42; leader to b: Ordered durable commands; leader to c: Ordered durable commands; client to other: Independent keysResolve partition andownershipStrong operation for cart-42Ordered durable commandsOrdered durable commandsIndependent keysSERVICEClient/routerSTOREPartition directorySERVICEPartition leaderSTOREFollower BSTOREFollower CSTOREOther partitiongroupssyncreplication
Read each connection in order
  1. syncResolve partition and ownershipClient/router → Partition directory
  2. syncStrong operation for cart-42Client/router → Partition leader
  3. replicationOrdered durable commandsPartition leader → Follower B
  4. replicationOrdered durable commandsPartition leader → Follower C
  5. syncIndependent keysClient/router → Other partition groups

10Explain storage maintenance without designing a new database

A log-structured merge storage engine keeps recent entries in an in-memory sorted structure backed by recovery data, then writes sorted immutable files. Background compaction combines files and discards obsolete versions when safe. This can support efficient sustained ingestion, but it consumes CPU, I/O and temporary disk space after the original write returns.

Reads may consult several structures. Bloom filters help skip files that definitely lack a key; a positive result only means the file may contain it. A block cache helps repeated reads. Neither feature changes the store's consistency contract or replaces the authoritative version decision.

A deletion marker suppresses older values still present in files or replicas. Remove it only when the supported repair, snapshot and retention rules ensure that old state cannot reappear. Physical removal from backups follows a separate retention policy. Immediate API deletion does not imply that every historical byte vanished immediately.

A B-tree engine is a valid alternative for another measured read/write mix. Explain which work the local engine performs, then discuss write amplification, read amplification and resources reserved for maintenance. A local engine such as RocksDB does not independently provide distributed routing, consensus or request deduplication; those belong to the system around it.

11Recover while protecting foreground capacity

Failure Behavior
Leader dies after commit but before reply A valid successor retains the command and result; the same request ID recovers the original outcome.
Majority becomes unavailable Stop strong operations for that partition; other independent partitions can continue.
Disk synchronization stalls Bound queued requests and bytes, reject excess work and retain committed outcomes.
Replica falls behind Restore a snapshot and replay the remaining log before treating it as caught up.
Accidental deletion reaches every replica Recover from separately retained and tested backup history.

Apply tenant quotas in bytes as well as operations: a 1 MiB write costs much more than an average 1 KiB write. Reserve bandwidth for foreground requests and throttle compaction, catch-up and migration when they compete for the same disks. Unlimited buffering converts a transient stall into memory exhaustion and extreme latency.

Monitor committed-operation tail latency, conflicts, unavailable partitions, leader changes, synchronization latency, compaction backlog and hot-key concentration. Distinguish attempts from proposals and successful mutations to see retry storms clearly. Test leader loss at commit boundaries, old-leader reads, deleted-key recreation and migration during retries. A successful backup job is insufficient until a restore produces a validated state.

12Check the design against its requirements

Use the numbered requirements to check the final design. FR refers to the functional list; NFR refers to the non-functional list. Performance rows specify tests still required, not achieved benchmark results.

Requirement Design mechanism Verification and remaining limit
FR1–2; NFR3: exact-key semantics One ordered state-machine action checks the expected version and installs the value or tombstone; strong reads establish current authority. Race the two version-7 cart updates and isolate an old leader. One replacement wins and the old leader cannot return stale data as current.
FR3; NFR5: recover retries Atomically persist value/version, request fingerprint, result and applied position. Lose a reply and retry after leader replacement or migration within one hour. The same outcome returns; older uncertainty follows the documented recovery boundary.
FR4; NFR4: tolerate the chosen failure Consensus commit/election and supported membership changes preserve history and revoke old ownership. Lose one zone, then lose a majority. Verify preserved acknowledged data in the first case and refusal of unsafe strong operations in the second.
NFR1–2,5: usable capacity Partition independent keys, reserve maintenance capacity and enforce tenant/byte limits. Load-test the stated mix during compaction and recovery against 20 ms p95. Test a hot key separately; extra partitions cannot parallelize its ordered mutations.

13Rapid revision

Remember: Save the new value and request result together, so a retry can return the result of the original write.

Topic Mechanism and consequence
API scope Exact-key operations; no implicit search or cross-key transactions.
Conditional update Compare the expected version and replace the value as one command in the replicated log’s order.
Safe retry Save the request fingerprint and result with the value; reuse the request ID after a timeout.
Durability Use the replication protocol’s commit rule so acknowledged writes survive the stated failures.
Strong read Verify that the server is still leader and has applied the committed commands required for the read.
Scale Hash different keys into replica groups; competing updates to one key still run in order.
Delete Save a deletion marker in log order; keep it until recovery can no longer restore an older value.
Migration Copy values, versions, deletion markers and retry results, then safely transfer write ownership.
Operations Reserve capacity for maintenance, replicas and failures; backups help recover mistakes copied to every replica.

Close with the cart-42 conflict and lost-response example. The deliberate cost is minority-side unavailability. If asked to accept writes in disconnected regions, discuss a different merge or conflict contract before changing the topology; the original immediate conditional-update guarantee cannot simply remain on both isolated sides.

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

How is a version different from a request ID?

Reveal a model answer

The version identifies state to replace; the request ID identifies one attempted change. Their combination supports conflict detection and retry recovery.

What the answer must demonstrate: Distinguish expected state from attempted operation and reject conflicting ID reuse.

Applied · Question 2

Two clients replace version 7. Why can only one win?

Reveal a model answer

The agreed command order applies one check-and-update first, advancing the version. The next command checks the new current state and records conflict.

What the answer must demonstrate: Place version comparison and mutation inside the agreed serialized application step.

Applied · Question 3

A write commits but its reply is lost. What survives?

Reveal a model answer

The value, version and saved result remain recoverable together. Retrying the same scoped request returns that result without another change.

What the answer must demonstrate: Persist value and original outcome together across commit and response loss.

Applied · Question 4

Why might a process labeled leader be unable to serve a strong GET?

Reveal a model answer

It may be isolated while another group has elected a successor. The leader must confirm it still leads and apply the committed commands required for the read.

What the answer must demonstrate: Require current read authority and sufficient applied state after leadership change.

Foundation · Question 5

Does R+W>N fully specify linearizability?

Reveal a model answer

No. Set overlap alone does not define ordering, election safety, failed writes, application or read authority. Use a complete replication protocol.

What the answer must demonstrate: Explain why quorum overlap needs a complete ordering and election protocol.

Foundation · Question 6

Why store a deletion marker?

Reveal a model answer

Older values may remain in disk files or replicas. The marker prevents their return until safe reclamation rules allow removal.

What the answer must demonstrate: Prevent resurrection and old-version reuse across deletion and recreation.

Follow-up · Question 7

Why move retry results with values?

Reveal a model answer

A client may retry through the new owner after losing an old response. Without the saved result, the new owner cannot preserve that retry contract.

What the answer must demonstrate: Move retry history, versions and tombstones with the key’s value.

Follow-up · Question 8

Will adding partitions accelerate one heavily updated key?

Reveal a model answer

Not its ordered conditional decisions. Partitioning adds parallelism across different keys, not concurrent winners against the same state.

What the answer must demonstrate: Recognize the single-key serialization boundary despite more partitions.

Blank-page exercise · 45 minutes

Build the answer yourself

Design a strongly consistent key-value store, then lose the leader after cart-42 version 8 commits but before the caller receives success.

  • Agree numbered functional and non-functional requirements, including per-key consistency, value limits and the one-zone durability boundary. Then explain value, version and request identity using two concurrent edits.
  • Trace a committed mutation and a current-authority read through a replica group.
  • Account for storage maintenance, hot keys, migration and failure without overstating availability.
  • Trace a timeout and a concurrent request using the actual durable records.
  • Review the final architecture against every numbered requirement, including measured bottlenecks, targets still needing validation and remaining failure limits.

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 storeCart-42 changed from version 7 to 8, but the reply was lost. What must survive together?Recall first, then reveal

The value or deletion marker, version, saved request result and applied-log position. Together they let recovery return the original result without changing the cart twice.

A lost reply must not repeat the write: save state and result together.

Return to lesson
Design a distributed key-value storeWhat permits a strong read?Recall first, then reveal

The leader confirms it still leads and has applied the committed commands required for this read.

Reachable is not current.

Return to lesson
Design a distributed key-value storeWhich work can more partitions spread across machines?Recall first, then reveal

Operations on different keys. Competing updates to one key must still follow one order.

Many keys, many owners.

Return to lesson

Final revision

Summary and interview notes

Store and retrieve exact keys, check versions before changing values, and recover the same result after retries. Partition independent keys and replicate each partition while preserving durable writes and current regional reads.

Remember these points

  • Agree what later reads must see before choosing replication.
  • Commit values and retry outcomes together.
  • Validate read authority after leadership changes.
  • Spread independent keys across partitions; one hot key stays ordered, and a minority cannot serve strong operations.

Interview tips

  • Agree on what reads must see after a completed write.
  • Race two version-7 updates, lose the winning reply, then isolate the old leader.

Important qualifications

  • Traffic and latency figures are interview assumptions, not claims about a named company's deployment.

Continue after the core interview

Explore the advanced version

The advanced lesson keeps the full detailed design. Use these sections when you want to examine the stronger requirements and failure cases.

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.