Skip to content

Key value store

System Design Task: Highly Available, Strongly Consistent Distributed Key-Value Store

Problem Statement

Design a distributed key-value store that is both highly available and strongly consistent — the storage layer that dozens of internal teams will build on for configuration, metadata, session state, ledgers, inventory, and idempotency keys.

The store must survive node failures, rack failures, availability-zone loss, and cross-region network partitions while still returning correct answers. It must scale horizontally to hundreds of terabytes and millions of operations per second, and it must be operable by a small team.

The headline requirement contains a deliberate tension. CAP says you cannot have both total availability and linearizability during a partition. Part of the task is to resolve that tension explicitly rather than hand-wave it: state which side you give up, on which operations, for how long, and show the arithmetic that gets you to the availability target anyway.

Assume this is a platform — once teams depend on it, the API and its guarantees are effectively frozen for a decade.


Functional Requirements

Your system must support:

  1. Core Key-Value Operations

  2. Get(key) — read the current value

  3. Put(key, value) — write or overwrite
  4. Delete(key)
  5. Keys: arbitrary byte strings up to 4 KB
  6. Values: arbitrary bytes up to 1 MB

  7. Conditional / Atomic Operations

  8. CompareAndSwap(key, expected_version, new_value)

  9. PutIfAbsent(key, value) — for locks, leader election, idempotency keys
  10. Atomic counters (Increment(key, delta))
  11. Every mutation returns a version / revision the client can reason about

  12. Range and Batch Access

  13. Scan(start_key, end_key, limit) in sorted key order

  14. BatchGet(keys[]) — up to 1000 keys per call
  15. BatchWrite(mutations[]) — atomic within one partition
  16. Define precisely what atomicity you offer across partitions, if any

  17. Watches / Change Notification

  18. Clients can subscribe to a key or key prefix and receive ordered change events

  19. Events must be gap-free and resumable from a revision after a disconnect

  20. Expiry and Lifecycle

  21. Per-key TTL

  22. Leases: a set of keys tied to a client heartbeat, deleted when it stops

  23. Tenancy and Isolation

  24. Namespaces per team, with independent quotas

  25. Per-tenant rate limiting so one team cannot starve another

  26. Read Consistency Levels (client-selectable per call)

  27. linearizable — reads see all completed writes

  28. bounded_staleness(t) — cheaper, may lag by at most t
  29. eventual — cheapest, any replica
  30. Explain what each costs and which failure modes each survives

Non-Functional Requirements

  1. Scale

  2. 100 billion keys, ~50 TB logical data (pre-replication)

  3. Sustained: 1M writes/sec, 5M reads/sec
  4. Peak: 2M writes/sec, 10M reads/sec
  5. Single keys up to 1 MB; single namespaces up to 10 TB

  6. Latency (single-region client, same-region data)

  7. Put P99 ≤ 10 ms, P99.9 ≤ 50 ms

  8. Linearizable Get P99 ≤ 5 ms
  9. Bounded-staleness Get P99 ≤ 2 ms
  10. Scan of 1000 keys P99 ≤ 50 ms

  11. Availability

  12. ≥ 99.99% for reads and writes (≈ 52 minutes/year)

  13. Survive with zero downtime: single node loss, single rack loss, single AZ loss
  14. Survive with degraded but defined behavior: full region loss, cross-region partition
  15. State the availability of each consistency level separately — they are not the same number

  16. Consistency

  17. Default: linearizable for single-key operations

  18. Session guarantees (read-your-writes, monotonic reads) must hold even when a client is bounced between replicas
  19. No lost updates, no resurrected deletes, no split-brain double-writes

  20. Durability

  21. An acknowledged write must survive the immediate loss of any one node and any one AZ

  22. Target ≤ 1 durability incident per 10⁹ writes
  23. Point-in-time restore to any second within the last 7 days

  24. Operability

  25. Rolling upgrades with no write downtime

  26. Add or drain a node without operator-authored rebalancing
  27. Every guarantee must be testable, and you should say how

What You Should Deliver

  1. Requirement clarification & assumptions

  2. Define "consistent" precisely — linearizability, sequential, causal, read-your-writes — and say which you are promising for which call

  3. State the failure model: crash-stop or Byzantine, fail-fast disks, clock assumptions

  4. The CAP decision, made explicit

  5. Which side you sacrifice during a partition, and for which operations

  6. PACELC: what you trade for latency when there is no partition
  7. Why the alternative was rejected

  8. High-level architecture

  9. Components, request path, control plane vs data plane

  10. Where the routing decision is made (client, proxy, or server redirect)

  11. Partitioning

  12. Hash vs range partitioning, and the consequences of each for Scan

  13. Partition sizing, split and merge policy
  14. How partitions are placed across nodes, racks, AZs, and regions

  15. Replication and consensus

  16. Replication protocol, replication factor, and quorum sizing

  17. Membership changes without losing quorum
  18. How a leader is elected and how long failover takes — with numbers

  19. Write path

  20. Every hop from client call to acknowledgement

  21. Where the durability point is (WAL fsync? quorum ack? both?)
  22. Batching, pipelining, and their effect on P99

  23. Read path

  24. How a linearizable read is served without paying a full consensus round for every read — and what that optimization assumes

  25. How bounded-staleness reads are made safe
  26. Follower and cross-region reads

  27. Storage engine

  28. LSM vs B-tree for this workload, with justification

  29. Key encoding, MVCC, tombstones, and garbage collection
  30. Compaction strategy and its effect on tail latency and space amplification

  31. Failure detection and recovery

  32. How a dead node is distinguished from a slow one

  33. Re-replication policy and how you avoid a metastable failure cascade
  34. What happens to in-flight writes during a failover

  35. Multi-region

    • Replica topology, and the latency arithmetic of a cross-region commit
    • Data pinning / locality
    • Behavior during a region partition, per consistency level
  36. Hot keys and hot partitions

    • Detection and mitigation
    • Why a single hot key is a fundamentally different problem than a hot partition
  37. Capacity estimates

    • Node count, disk, memory, network — and the calculation, not just the answer
    • Partition count and its per-node overhead
  38. Client library design

    • Routing cache and staleness handling
    • Retries, backoff, idempotency, and how you avoid a retry storm
  39. Anti-entropy, backup, restore

    • Detecting and repairing silent divergence
    • Backup mechanics and restore time for 50 TB
  40. Observability and SLOs

    • The handful of metrics that actually tell you the store is healthy
    • How you would prove linearizability holds in production
  41. Trade-offs

    • What you deliberately did not build, and what breaks if a team needs it

Variants (Design These Too)

These are separate design tasks, not footnotes. Each one changes the answer.

Variant A — The AP Store. The same API, but availability wins during a partition: every replica accepts writes at all times. Design the conflict model — last-writer-wins, version vectors, or CRDTs — and show which of the functional requirements above (CAS, watches, leases, atomic counters) you can still honor, which become approximate, and which you must remove from the API entirely.

Variant B — Global Linearizability. Replicas span three regions and the store must be linearizable globally. Do the latency arithmetic for a cross-region commit. Then design the escape hatches (region-pinned partitions, leader placement, follower reads with closed timestamps) that make it usable, and state the write latency a client in the wrong region must accept.

Variant C — Small Scale. The same guarantees, 3 nodes, 100 GB, 10k ops/sec, one AZ, one part-time operator. What collapses out of the design? Justify every component you keep, and name the exact scale threshold at which each dropped component must come back.

Variant D — The Migration. An existing team runs 4 TB on a sharded MySQL setup with application-level routing and no cross-shard transactions. Design the cutover to your store with zero downtime and a working rollback at every step.


Stretch Problems

  1. Multi-key transactions. Extend to serializable transactions across partitions. Which protocol, what it costs, and how you keep it from becoming the default path teams reach for.
  2. The clock question. Your fast-read optimization probably assumes bounded clock skew. Design the version that assumes nothing about clocks and quantify what it costs. Then design what happens when a machine's clock jumps 40 seconds.
  3. Quorum loss. Two of three replicas of one partition are permanently gone and the data is not in a backup. Design the unsafe-recovery procedure, its guardrails, and what you tell the affected tenant.
  4. Metastable failure. A cache-miss storm after a failover pushes the cluster into a state where load stays above capacity even after the trigger is removed. Design the load-shedding that prevents it.

Expectations

  • Do the arithmetic. Quorum sizes, failover budgets, replica counts per node, and cross-region RTTs should appear as numbers, not adjectives.
  • Name concrete mechanisms — Raft, leader leases, ReadIndex, closed timestamps, hinted handoff, Merkle anti-entropy, LSM compaction — and say what each one buys.
  • Be precise about guarantees. "Strongly consistent" without a definition is the single most common failure in this design.
  • Show the failure walkthrough. For each failure class, state what a client in flight observes.
  • Prefer a boring design that a small team can operate over an elegant one that needs its own on-call rotation.
  • Assume this system will be maintained for a decade by people who did not write it.

Interview Kit

Read first: Solution · Consensus §5–§7, §9.4 CAS fencing · LSM trees · Replication

Curveballs. The interviewer changes one thing mid-design. The hint in italics is what a strong answer reaches for:

  1. A team uses your store for a distributed lock and double-processes payments during a GC pause. Whose bug is it, and what API do you add? (Leases plus fencing tokens / CAS, §9.4.)
  2. Your Raft nodes acknowledge before fsync "for speed". What can be lost, and when? (Jepsen NATS 2025: a coordinated power loss loses acknowledged writes.)
  3. One tenant's key is read 500K times per second. (Hot-key detection per tenant, read leases or follower reads, client caching.)
  4. Losing a whole region must not lose acknowledged writes. What does that cost in write latency?

Must answer (security, privacy, operations):

  • Tenant isolation enforced server-side (key prefixing, clipped scans), and per-tenant encryption keys
  • Backups that can't bring back a crypto-shredded tenant

Phase it (MVP → Growth → Scale): MVP: single Raft group (etcd-like) for less than 8 GB. Growth: multi-Raft ranges with splits. Scale: multi-region placement, follower reads, tenant quotas.

Score yourself with the rubric.