PrepZone Logo
PrepZone

Key-Value Store

Design a distributed KV store with partitioning, replication and tunable consistency.

Read these first

Design a KV store like Dynamo: PUT key → value, GET key → value, at billions of keys and millions of ops/sec with node failures handled gracefully.

Requirements

Functional requirements

  • put(key, value) — store or overwrite a value.
  • get(key) — retrieve value or not-found.
  • delete(key) — remove a key.
  • Support values from bytes to megabytes.

Non-functional requirements

  • Availability: 99.99% — tolerate node and rack failures.
  • Latency: p99 GET under 10 ms; p99 PUT under 20 ms.
  • Scale: Billions of keys; horizontal scale by adding nodes.
  • Durability: Survive node crash without data loss (replication factor 3).
  • Consistency: Tunable — strong for critical keys, eventual for cache-like keys.

Estimation

  • 1B keys × 1 KB average = 1 TB raw data × replication factor 3 = 3 TB cluster.
  • 1M ops/sec peak → ~1M GET/s + 100K PUT/s — each node handles ~50K ops with 20 nodes.
Hash ring0 … 2^32-1
Node Akeys 0–25%
Node Bkeys 25–50%
Node Ckeys 50–75%
Node Dkeys 75–100%

New node E joins → only ~1/N keys move, not all keys.

Adding a node only remaps keys adjacent to it on the ring.

Data partitioning

Consistent hashing in practice

  • Hash key onto ring; assign to next node clockwise.
  • Virtual nodes (100–200 per physical node) for even load distribution.
  • Adding a node moves only ~1/N keys — not full reshuffle.
  • Hot keys: replicate to multiple nodes or add random suffix sharding.
Java
def get_responsible_nodes(key: str, rf: int = 3) -> list[str]:
    hash_val = hash_key(key)
    ring_position = ring.bisect(hash_val)
    nodes = []
    for i in range(rf):
        nodes.append(ring[(ring_position + i) % len(ring)])
    return nodes

Replication and consistency

AspectQuorum read/write (R + W > N)Leader-follower (strong consistency)
WriteW nodes must ack (e.g., W=2 of N=3)Leader acks; async replicate to followers
ReadR nodes must respond (e.g., R=2); merge by version vectorRead from leader or linearizable follower
ConsistencyTunable — R+W>N gives strongStrong by default; followers may lag
FailureTolerates N-W node failures on writeLeader election needed on leader failure
  • Write

    Quorum read/write (R + W > N)W nodes must ack (e.g., W=2 of N=3)
    Leader-follower (strong consistency)Leader acks; async replicate to followers
  • Read

    Quorum read/write (R + W > N)R nodes must respond (e.g., R=2); merge by version vector
    Leader-follower (strong consistency)Read from leader or linearizable follower
  • Consistency

    Quorum read/write (R + W > N)Tunable — R+W>N gives strong
    Leader-follower (strong consistency)Strong by default; followers may lag
  • Failure

    Quorum read/write (R + W > N)Tolerates N-W node failures on write
    Leader-follower (strong consistency)Leader election needed on leader failure

Architecture

StreamHub production architecture (AWS)

HTTPSstaticmissAPICLIENT
Mobile / WebStreamHub cli…
NETWORK
Route 53GeoDNS routing
NETWORK
CloudFrontCDN + WAF edge
NETWORK
AWS ALBTLS terminati…
NETWORK
API GatewayJWT · rate li…
STORAGE
Amazon S3media origin
COMPUTE
Amazon EKSAPI · auth · …
DATABASE
ElastiCachesessions · ho…
DATABASE
RDS Postgresprimary + rep…
INTEGRATION
Amazon MSKdomain events
ANALYTICS
OpenSearchstream discov…
OPS
CloudWatchmetrics · X-R…
End-to-end path from user to data — reference this when placing any new service.

Node internals

  • Coordinator: Receives client request; hashes key; routes to responsible nodes.
  • Storage engine: LSM-tree (LevelDB/RocksDB) for write-heavy or B-tree for read-heavy.
  • Gossip protocol: Nodes discover membership and failure without central registry.
  • Hinted handoff: If target node is down, temporary node stores writes until recovery.
  • Anti-entropy: Merkle tree comparison between replicas to repair divergence.
Java
# Tunable consistency (pseudo-client config)
consistency:
  write_quorum: 2    # W
  read_quorum: 2     # R
  replication_factor: 3  # N
  # R + W > N → strong consistency (2 + 2 > 3)

Handling failures

FailureDetectionRecovery
Node crashGossip heartbeat timeoutPromote replica; hinted handoff replays writes
Network partitionQuorum can't be reachedChoose CP (reject writes) or AP (accept, repair later)
Data divergencePeriodic Merkle tree scanAnti-entropy repair syncs missing keys
Adding nodeAdmin command or auto-scaleConsistent hash remaps ~1/N keys to new node
  • Node crash

    DetectionGossip heartbeat timeout
    RecoveryPromote replica; hinted handoff replays writes
  • Network partition

    DetectionQuorum can't be reached
    RecoveryChoose CP (reject writes) or AP (accept, repair later)
  • Data divergence

    DetectionPeriodic Merkle tree scan
    RecoveryAnti-entropy repair syncs missing keys
  • Adding node

    DetectionAdmin command or auto-scale
    RecoveryConsistent hash remaps ~1/N keys to new node

Quick recall

Everything you need if you only revisit this box.

  • Partition keys via consistent hashing with virtual nodes for even load.
  • Replication factor 3; quorum W + R > N gives strong consistency.
  • Gossip protocol for membership; hinted handoff for temporary node failures.
  • Anti-entropy (Merkle trees) repairs replica divergence in background.
  • Hot keys need replication beyond consistent hashing — random suffix or local cache.
  • LSM-tree storage engine for write-heavy; B-tree for read-heavy workloads.

Test yourself

Answer these before moving on — recall is what makes it stick.