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

DecisionUpsideDownside
Higher W and RFresher reads, stronger guarantees.Higher latency (wait for slower replicas) and failure when too few replicas respond.
Lower W and RFaster and available with more replicas down.Stale reads become possible.

Real World

SystemHow it's used
Apache CassandraLets each query set its consistency level, from ONE to QUORUM to ALL.
Amazon DynamoUsed 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?
  1. W = 1, R = 1
  2. W = 2, R = 2
  3. W = 1, R = 2
  4. 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?
  1. 0
  2. 1
  3. 2
  4. 3

Three of five must respond, so up to two can fail.

What does read repair do?
  1. Rebuilds the ring
  2. Writes the newest value back to stale replicas found during a read
  3. Deletes old keys
  4. Elects a leader

The read compares versions and fixes replicas that were behind.

Why can a sloppy quorum return stale reads?
  1. It uses LWW
  2. Writes may land on nodes outside the key's home replicas
  3. It requires all replicas
  4. It disables versions

A read of the home replicas may not overlap nodes that temporarily accepted the write.

Cassandra LOCAL_QUORUM is used to…
  1. Wait for every region
  2. Avoid cross-region latency by using a majority in the local data centre
  3. Skip replication
  4. 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 lesson

Two-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 lesson

Sagas

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 lesson

Clocks & 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 lesson

Practice 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.