SkillByAIOpen interactive version →

Lesson 16 / 25

Design a Distributed Key-Value Store

Partitioning, replication, quorum, conflicts.

The Dynamo-style design

Requirements: put(key, value) and get(key), horizontal scale, high availability, tunable consistency. Partitioning: consistent hashing places keys on a ring, with virtual nodes to spread load and ease rebalancing. Replication: each key is stored on N nodes (the next N distinct nodes on the ring). Quorum: a write waits for W acknowledgements and a read queries R replicas; if R + W > N, read and write sets overlap, so a read will contact at least one replica that has the latest acknowledged write (subject to failure-handling details such as sloppy quorums). Conflicts: concurrent writes can diverge; resolve with last-write-wins timestamps (simple, may drop writes) or version vectors that detect siblings for the client or a merge function to resolve. Add anti-entropy (Merkle trees), hinted handoff and gossip-based membership.

Storage, search and money

These problems test distributed data fundamentals and correctness under failure.

Figure 6.1 — Key-value ring, autocomplete trie and double-entry ledger.

Quorum and conflict sketch

Illustrative configuration.

N = 3 replicas, W = 2, R = 2   ->  R + W = 4 > N  (overlapping quorums)

put(k, v):  coordinator = hash(k) on ring
            send to 3 replicas, return OK after 2 acks
get(k):     ask replicas, wait for 2 responses
            if versions differ -> return newest / siblings, repair stale replica

Version vectors
  A writes   -> {A:1}
  B writes concurrently from {A:1} -> {A:1, B:1}
  C writes concurrently from {A:1} -> {A:1, C:1}
  neither descends from the other -> conflict, keep both siblings

Tuning: W=1 fast writes, weaker; R=1 fast reads, weaker; W=N strong writes, less available

Relate choices to CAP

Say which behaviour you want during a partition: keep accepting writes (availability, conflicts later) or reject them (consistency). Quorum settings make the trade-off concrete.

Quick check: With N = 3, which setting gives overlapping read and write quorums?

  • W = 1 and R = 2
  • W = 1 and R = 1
  • W = 2 and R = 2
  • W = 0 and R = 3
Answer

W = 2 and R = 2 — R + W must exceed N.