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.
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:
| Property | Meaning |
|---|---|
| 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:
| Method | Pros | Cons |
|---|---|---|
| Range-based — keys A–M on server 1, N–Z on server 2 | Simple | Hot spots if certain ranges are popular |
| Hash modulo — hash(key) % N | Even distribution | Large reshuffles when N changes |
| Consistent Hashing — ring-based | Minimal reshuffling; supports heterogeneous capacity | More 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
| Configuration | Result |
|---|---|
| W + R > N | Strong consistency (at least one replica overlap between read and write sets) |
| W + R ≤ N | Eventual consistency (faster ops, short lag before all replicas sync) |
| R = 1, W = N | Optimised for fast reads |
| W = 1, R = N | Optimised 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.
| Technique | How it works | Pros | Cons |
|---|---|---|---|
| Simple Version Numbers | Increment counter per update; higher counter wins | Simple | Loses concurrent writes |
| LWW Timestamps | Latest timestamp wins | Easy and efficient | Requires clock sync; may lose valid updates |
| Vector Clocks | Each version carries (node, counter) pairs; tracks causal history | Detects concurrent updates | Complex; more storage overhead |
| CRDTs | Special data structures with automatic merge rules | No conflicts; strong eventual consistency | Limited 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.

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.

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.

Each node carries the same responsibilities:

Consistent Hashing Flow#
- Hash the key — hash("user123") → 856.
- Map nodes to the ring — Node A → 100, Node B → 500, Node C → 900.
- Locate the key's node — move clockwise; hash = 856 → closest node clockwise is C (900).
- Replication (optional) — store copies on the next N nodes clockwise.
- Node joins/leaves — only nearby keys need redistribution.
Write Path#
- Client sends put(key, value) to any coordinator node.
- Coordinator uses consistent hashing to find responsible node(s).
- Write is persisted to a commit log (append-only, ensures durability on crash).
- Data is saved to the memory cache.
- 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.

Read Path#
-
Read request arrives at the responsible node; it first checks the memory cache.

-
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.
-
The bloom filter narrows down candidate SSTables; the SSTable index locates the exact offset.
-
The value is retrieved from the SSTable on disk.

Summary#
| Goal | Technique |
|---|---|
| Store big data | Consistent hashing — spread load across servers |
| High availability reads | Data replication + multi-data centre setup |
| Highly available writes | Versioning + conflict resolution with vector clocks |
| Dataset partitioning | Consistent hashing |
| Incremental scalability | Consistent hashing |
| Heterogeneous capacity | Virtual nodes in consistent hashing |
| Tunable consistency | Quorum consensus (N, W, R) |
| Temporary failure handling | Sloppy quorum + hinted handoff |
| Permanent failure handling | Merkle tree + anti-entropy |
| Data centre outage handling | Cross-data centre replication |