High Level Design

Designing a Key-Value Store

End-to-end design of a distributed key-value store — CAP theorem, consistent hashing, replication, conflict resolution with vector clocks, failure handling, and the detailed read/write path.

August 10, 2026

A key-value store is a non-relational database that uses simple key-value pairs to store data. Keys are unique; values can be strings, numbers, JSON, or binary. Key-value stores are optimised for high performance, scalability, and flexibility — ideal for caching, session management, and real-time analytics.


Step 1: Requirements and Design Scope#

  • Support basic operations: put(key, value) and get(key).
  • Key-value pair size is small (a few KBs).
  • Ability to store large data sets (a few TBs total).
  • High availability — responds quickly even during failures.
  • High scalability — can be scaled to support large data sets.
  • Automatic scaling — scales up/down based on load.
  • Tunable consistency — choose between strong and eventual per use case.
  • Low latency — fast reads and writes.

Step 2: High Level Design#

Single Server Key-Value Store#

  • Store key-value pairs in a hash table, keeping everything in memory.
  • Fast reads and writes but memory is expensive and limited.
  • Optimisations: data compression; evict cold data to disk.
  • Even with optimisations, a single server cannot handle large data sets.

Distributed Key-Value Store#

To handle large data sets, distribute data across multiple servers.

CAP Theorem#

A distributed system can guarantee at most two of:

PropertyMeaning
Consistency (C)All clients see the same data at the same time
Availability (A)Every request gets a response (may be stale)
Partition tolerance (P)System continues operating despite network partitions

Since network partitions are unavoidable in distributed systems, you choose between CP (consistency over availability) or AP (availability over consistency) depending on your use case.


System Components and Techniques#

1. Data Partitioning#

For large applications, data is spread across multiple servers.

Techniques:

MethodProsCons
Range-based — keys A–M on server 1, N–Z on server 2SimpleHot spots if certain ranges are popular
Hash modulo — hash(key) % NEven distributionLarge reshuffles when N changes
Consistent Hashing — ring-basedMinimal reshuffling; supports heterogeneous capacityMore complex

Consistent hashing is the standard choice:

  • Automatic scaling — minimal data movement when nodes join/leave.
  • Heterogeneity — servers with double the capacity get double the virtual nodes.

2. Data Replication#

  • Each key is replicated to the next N clockwise nodes on the hash ring (N is configurable).
  • If a node goes down, replicas on other nodes serve the data.
  • With virtual nodes, the first N distinct physical servers clockwise are used.
  • Replicas placed across multiple data centres for fault isolation.

3. Consistency Model#

Quorum parameters:

N → number of replicas
W → replicas that must acknowledge a write for it to succeed
R → replicas that must respond for a read to succeed
ConfigurationResult
W + R > NStrong consistency (at least one replica overlap between read and write sets)
W + R ≤ NEventual consistency (faster ops, short lag before all replicas sync)
R = 1, W = NOptimised for fast reads
W = 1, R = NOptimised for fast writes

Key-value stores typically default to eventual consistency — concurrent writes are allowed and clients reconcile versions on read.

4. Inconsistency Resolution#

In eventually consistent systems, replicas may hold conflicting versions due to concurrent writes, network partitions, or delayed replication. Systems track versions instead of overwriting blindly.

TechniqueHow it worksProsCons
Simple Version NumbersIncrement counter per update; higher counter winsSimpleLoses concurrent writes
LWW TimestampsLatest timestamp winsEasy and efficientRequires clock sync; may lose valid updates
Vector ClocksEach version carries (node, counter) pairs; tracks causal historyDetects concurrent updatesComplex; more storage overhead
CRDTsSpecial data structures with automatic merge rulesNo conflicts; strong eventual consistencyLimited to specific types (counters, sets, maps)

5. Handling Failures#

A. Failure Detection#

1. All-to-All Heartbeat Every node sends heartbeats to all others. If a heartbeat isn't received within a timeout, the node is marked failed.

heartbeat mechanism

2. Gossip Protocol (preferred at scale)

  • Each node maintains a membership list with member IDs and heartbeat counters.
  • Nodes periodically increment their own counter and send it to a random set of peers.
  • Peers propagate the info further.
  • If a heartbeat counter hasn't increased in a predefined period → node marked offline.

gossip protocol

B. Handling Temporary Failures#

  • Strict quorum: reads/writes block until required replicas are available.
  • Sloppy quorum: system picks the first W healthy servers for writes, first R for reads — continues despite partial outages.

Hinted Handoff — writes intended for an unavailable node are stored on a different node with a "hint". Once the original node recovers, hints are forwarded to it.

Read Repair — if a node detects stale data during a read compared to other replicas, it updates itself.

C. Handling Permanent Failures#

Anti-Entropy — a background process that periodically compares replicas and syncs differences using Merkle trees.

A Merkle tree is a hash tree where each non-leaf node is labelled with the hash of its children. Allows efficient and secure verification of large data structures.

  • If two replicas share the same root hash → data is identical.
  • If root hashes differ → check children recursively until the differing branches are found.
  • Only mismatched ranges need syncing — minimises network traffic.

D. Handling Data Centre Outages#

  • Replicas are placed across multiple data centres.
  • Data centres connect via high-speed networks for low-latency cross-region access.
  • If one data centre goes down, clients read from replicas elsewhere.

Step 3: Detailed Design#

Architecture Diagram#

  • Clients communicate via simple get(key) / put(key, value) APIs.
  • A coordinator node acts as proxy between client and the store.
  • All nodes are placed on a consistent hash ring.
  • The system is fully decentralised — every node has the same responsibilities.
  • Data is replicated at multiple nodes; no single point of failure.

architecture diagram

Each node carries the same responsibilities:

node responsibilities

Consistent Hashing Flow#

  1. Hash the key — hash("user123") → 856.
  2. Map nodes to the ring — Node A → 100, Node B → 500, Node C → 900.
  3. Locate the key's node — move clockwise; hash = 856 → closest node clockwise is C (900).
  4. Replication (optional) — store copies on the next N nodes clockwise.
  5. Node joins/leaves — only nearby keys need redistribution.

Write Path#

  1. Client sends put(key, value) to any coordinator node.
  2. Coordinator uses consistent hashing to find responsible node(s).
  3. Write is persisted to a commit log (append-only, ensures durability on crash).
  4. Data is saved to the memory cache.
  5. When cache is full or hits a threshold, data is flushed to SSTable on disk.

SSTable (Sorted-String Table) — a sorted list of <key, value> pairs stored on disk.

write path

Read Path#

  1. Read request arrives at the responsible node; it first checks the memory cache.

    read path step 1

  2. On a cache miss, the node consults the bloom filter to determine which SSTable files might contain the key.

    Bloom filter — a space-efficient probabilistic structure that can produce false positives but never false negatives.

  3. The bloom filter narrows down candidate SSTables; the SSTable index locates the exact offset.

  4. The value is retrieved from the SSTable on disk.

    read path step 2


Summary#

GoalTechnique
Store big dataConsistent hashing — spread load across servers
High availability readsData replication + multi-data centre setup
Highly available writesVersioning + conflict resolution with vector clocks
Dataset partitioningConsistent hashing
Incremental scalabilityConsistent hashing
Heterogeneous capacityVirtual nodes in consistent hashing
Tunable consistencyQuorum consensus (N, W, R)
Temporary failure handlingSloppy quorum + hinted handoff
Permanent failure handlingMerkle tree + anti-entropy
Data centre outage handlingCross-data centre replication