Advanced System design concept · Distributed Systems & Data · 55 mins read
Consistency in Distributed Systems
How replicas agree on values, how nodes agree on a leader, how transactions span services, and how events are ordered without a shared clock.
Quorum Reads & Writes
Choose how many replicas must confirm a write and answer a read so that reads overlap the latest write.
Intuition
In a leaderless store, a write goes to several replicas, but some may be down or slow and miss it. A read that asks only one replica might hit one that missed the write and return an old value.
Quorums let you tune each operation between consistency, latency and availability. Cassandra, DynamoDB-style systems and many coordination services rely on them.
Mental Model
With N replicas, a write waits for W confirmations and a read asks R replicas. If W + R > N, every read set overlaps every write set in at least one replica, so the read sees the latest acknowledged write and picks it by version. Common choice: N = 3, W = 2, R = 2. Lowering W or R makes operations faster and more available but allows stale reads. Think of it like: Three notebooks: you write in at least two and read from at least two. Whichever two you read, one of them has the latest note.
Building Blocks
- N, W, R: Replicas per key, write acknowledgements required, and replicas consulted on read.
- Overlap rule: W + R > N guarantees read and write sets intersect.
- Versions: Each value carries a timestamp or version vector so the reader can pick the newest.
- Read repair: When a read sees a stale replica, it writes the newer value back to it.
Definitions
- Quorum
-
The minimum number of replicas that must respond for an operation to succeed.
- Often a majority: ⌊N/2⌋ + 1.
- Sloppy quorum
-
During failures, writes are accepted by any reachable nodes, not just the key's home replicas.
- Improves write availability.
- Breaks the overlap guarantee until hinted handoff completes.
- Tunable consistency
-
Choosing W and R per request.
- For example Cassandra's ONE, QUORUM and ALL.
Patterns
- Majority writes and reads — When reads must reflect the latest acknowledged write.
- Fast writes, careful reads — Write-heavy data that is rarely read.
- Local quorum in multi-region — When cross-region latency on every request is unacceptable.
Strategies
- Pick W and R from the requirement When: For each class of data. How: Decide whether stale reads are acceptable, then choose the smallest W and R that meet it; W + R > N for fresh reads, smaller values for speed and availability. Example: Session data uses ONE; account balances use QUORUM.
- Keep replicas converging When: Always, because quorums leave some replicas behind. How: Combine read repair for hot keys with scheduled anti-entropy repair for keys that are rarely read. Example: Cassandra runs incremental repair regularly in addition to read repair.
Where W + R > N is not enough
Overlap guarantees that a read contacts a replica holding the latest acknowledged write, but edge cases remain. With sloppy quorums, a write may land on nodes outside the key's normal replica set, so a later read of the home replicas can miss it. If a write fails on some replicas but succeeded on fewer than W, it is reported as failed yet may still be read later. Concurrent writes need versions to resolve, and last-write-wins can discard one.
So quorums give 'reads usually see the latest value', not the full linearizability a consensus-based system provides. When you need strict guarantees — unique usernames, money — use a leader-based or consensus-backed store, or add a linearizable check on top.
Tradeoffs
| Decision | Upside | Downside |
|---|---|---|
| Higher W and R | Fresher reads, stronger guarantees. | Higher latency (wait for slower replicas) and failure when too few replicas respond. |
| Lower W and R | Faster and available with more replicas down. | Stale reads become possible. |
Real World
| System | How it's used |
|---|---|
| Apache Cassandra | Lets each query set its consistency level, from ONE to QUORUM to ALL. |
| Amazon Dynamo | Used configurable N, W, R with sloppy quorums and hinted handoff for an always-writable shopping cart. |
Interview
Questions interviewers ask
- With 5 replicas, what W and R would you choose?
- Why does W + R > N give fresh reads?
- Is a quorum read linearizable?
What a strong answer covers
Do the overlap math, explain failure tolerance for chosen values and name the edge cases (sloppy quorums, concurrent writes).
Common traps
- Claiming quorums alone give linearizability.
- Forgetting that higher quorums reduce availability.
Quiz
With N = 3, which choice guarantees a read overlaps the latest write?
- W = 1, R = 1
- W = 2, R = 2
- W = 1, R = 2
- W = 2, R = 1
2 + 2 = 4 > 3, so every read set shares a replica with every write set.
N = 5, W = 3, R = 3. How many replicas can be down while reads and writes still succeed?
- 0
- 1
- 2
- 3
Three of five must respond, so up to two can fail.
What does read repair do?
- Rebuilds the ring
- Writes the newest value back to stale replicas found during a read
- Deletes old keys
- Elects a leader
The read compares versions and fixes replicas that were behind.
Why can a sloppy quorum return stale reads?
- It uses LWW
- Writes may land on nodes outside the key's home replicas
- It requires all replicas
- It disables versions
A read of the home replicas may not overlap nodes that temporarily accepted the write.
Cassandra LOCAL_QUORUM is used to…
- Wait for every region
- Avoid cross-region latency by using a majority in the local data centre
- Skip replication
- Read from one node
It keeps majority semantics within one region without waiting for remote ones.
Consensus & Leader Election
How Raft-style consensus lets a group of nodes agree on one leader and one ordered log, even when some fail.
This section is part of the full PRISM roadmap, with worked examples, trade-off tables, interview questions and a quiz.
Unlock the full lessonTwo-Phase Commit
Atomically commit a transaction across several databases, and why the protocol can block when the coordinator fails.
This section is part of the full PRISM roadmap, with worked examples, trade-off tables, interview questions and a quiz.
Unlock the full lessonSagas
Run a business process across services as a sequence of local transactions, undoing completed steps with compensations when a later step fails.
This section is part of the full PRISM roadmap, with worked examples, trade-off tables, interview questions and a quiz.
Unlock the full lessonClocks & Event Ordering
Why wall clocks cannot order events across machines, and how logical clocks, version vectors and fencing tokens do it safely.
This section is part of the full PRISM roadmap, with worked examples, trade-off tables, interview questions and a quiz.
Unlock the full lessonPractice consistency in distributed systems in PRISM
Concepts stick when you watch them fail. Build an architecture that depends on consistency in distributed systems, push traffic through it in the PRISM simulator, and see the latency and error rates change as you adjust the design.