Design a Key-Value Store

⭐ Interview Importance: HIGH
⏱️ Revision Time: 5 min

Concept

The Problem: Design a distributed key-value store (like DynamoDB or Cassandra) that can store trillions of records, scale infinitely across thousands of servers, and guarantee high availability.

This question is the definitive test of your deep understanding of the CAP Theorem, Consistent Hashing, and Vector Clocks.

1. Requirements

Functional:

  • put(key, value)
  • get(key)

Non-Functional:

  • Highly Available (If 5 servers crash, the system stays up).
  • Highly Scalable (You can add 10 new servers without rewriting data).
  • Low Latency (Millisecond response times).

2. Partitioning (Consistent Hashing)

You cannot store 10 Petabytes of data on one server. You must partition it.
If you use simple modulo hashing (hash(key) % N), adding a single new server changes the N value, requiring you to physically move 99% of your data to new servers, which will crash the system.

The Solution: Consistent Hashing

  1. Imagine a conceptual ring (a circle) going from 0 to 360 degrees.
  2. Hash your Server IPs and place the servers on the ring (e.g., Node A at 10°, Node B at 100°).
  3. When put("user123", "data") is called, hash the key "user123". It lands at 50°.
  4. Walk clockwise around the ring from 50° until you hit a server. You hit Node B. Store the data there.
    Benefit: If you add Node C at 60°, you only have to move a tiny fraction of data from Node B to Node C. The rest of the cluster is untouched.

3. Replication & High Availability

If Node B crashes, "user123" is lost forever. We must replicate data.
When we walk clockwise and hit Node B, we store a copy there. But we don’t stop. We continue walking clockwise and store identical copies on the next two physical servers (e.g., Node C and Node D). We now have a Replication Factor of 3.

4. The CAP Theorem & Quorums

In a distributed system, network partitions (the ‘P’ in CAP) are unavoidable. A router will break. When it does, you must choose between Consistency (C) and Availability (A).
Key-Value stores like DynamoDB/Cassandra choose Availability (AP). They will accept your put request even if some servers are offline, at the cost of Eventual Consistency.

We manage this using Quorum Math (N, W, R):

  • N = Replication Factor (e.g., 3 copies).
  • W = Write Quorum. How many nodes must acknowledge a put before returning “Success” to the user? (e.g., 2).
  • R = Read Quorum. How many nodes must we query during a get to find the most recent data? (e.g., 2).

The Golden Rule of Consistency: If W + R > N, you are mathematically guaranteed to read the most recent data (Strong Consistency). If you set W=1 and R=1 to make the database lightning fast, 1+1 is less than 3, so you will experience Stale Reads (Eventual Consistency).

5. Conflict Resolution (Vector Clocks)

Because we prioritize Availability (W=1), two different servers might accept conflicting writes for the exact same key during a network partition.
Server A thinks name="Alice". Server B thinks name="Bob".

How do we resolve this when the network heals?
We cannot use physical timestamps (Clock Skew makes them unreliable). We must use Vector Clocks.
A Vector Clock is an array of logical counters attached to the data: [ServerA: 2, ServerB: 1].
By comparing the arrays mathematically, the database can definitively determine if one update is a direct ancestor of the other. If the updates are entirely divergent, the database cannot resolve it safely. It will return both “Alice” and “Bob” to the client application and force the client to write custom business logic to merge the conflict.

Interview Questions

Q: A server dies permanently. Its replacement boots up totally empty. How does the new server quickly figure out which millions of records it is missing without asking other nodes to send their entire 1TB datasets over the network?
A: Using Merkle Trees. A Merkle Tree is a cryptographic hash tree of the data. The new server asks a neighbor for its Merkle Tree. By comparing the root hashes, and then walking down the differing branches, the servers can instantly pinpoint the exact 2 Megabytes of data that are missing in O(log⁡N)O(\log N) time, and only transmit those specific bytes over the network to repair the node.

Q: What is a “Gossip Protocol” and why is it necessary in a cluster of 1,000 servers?
A: If you have 1,000 servers, you cannot have one “Master” node tracking the health of everyone; it becomes a bottleneck and single point of failure.
Instead, Cassandra uses a Gossip Protocol. Every second, Server A picks a random server (Server D) and sends it a tiny message: “Here is what I know about the health of the cluster.” Server D merges that info with its own, and then gossips with Server G. Like a real-life rumor, knowledge of a crashed server mathematically propagates through a 1,000-node cluster in less than 2 seconds, completely decentralized.