Skip to main content

Distributed Systems Consistency: 2PC, Saga, Quorum & Vector Clocks Explained

Calculating read time…

You transfer ₹10,000 from your bank account to a friend's account. Your bank deducts ₹10,000. A network failure occurs before crediting your friend. The money has left your account but never arrived. Where did it go?

This is the fundamental consistency challenge of distributed systems. In a single database, a transaction is atomic — all or nothing. In a distributed system with ten databases across five data centres, making all of them agree on the same state at the same time is one of the hardest problems in all of computer science.

The four patterns in this post represent four different answers to that challenge, each making different trade-offs between consistency, availability, and performance. Understanding them deeply is what separates senior distributed systems engineers from everyone else.



💡 The CAP Theorem Context

Before diving in, one essential backdrop: the CAP Theorem states that a distributed system can guarantee at most two of these three properties simultaneously:

Consistency — every read sees the most recent write
Availability — every request gets a (non-error) response
Partition tolerance — the system works despite network partitions

Since network partitions are inevitable in real systems, you always choose between CP (consistent but may be unavailable during partitions) and AP (always available but may serve stale data).

🤝 2PC → CP (consistent, blocks on partition)
📖 Saga → AP (available, eventually consistent)
🗳️ Quorum → tunable (dial between CP and AP)
🕐 Vector Clocks → conflict detection tool for AP systems

🤝 Section 1: Two-Phase Commit (2PC) — The Classic Distributed Transaction

Two-Phase Commit is the original solution to distributed transactions. When you need all participants to either ALL commit or ALL abort — guaranteed, with no partial results — 2PC is the foundational protocol.

💡 The Wedding Coordinator Analogy

Imagine you're planning a wedding that requires a venue, catering, band, and florist — all committed on the same date or the wedding cannot happen.

Phase 1 (Prepare): The coordinator calls all vendors: "Can you guarantee availability on June 15th? Lock it down, don't take other bookings — but don't confirm yet." Each vendor either says YES (locked in, waiting) or NO (can't do it).

Phase 2 (Commit/Abort): If all vendors said YES → coordinator calls everyone: "Confirmed! It's on." If ANY vendor said NO → coordinator calls everyone: "Called off. Release your holds."

The critical insight: during Phase 1, all vendors are holding resources while waiting. If the coordinator dies after Phase 1 — everyone is stuck holding their bookings forever. That's 2PC's fundamental weakness.

🤝 Two-Phase Commit — Complete Protocol Flow

📋 Phase 1: PREPARE (Voting)

🎯 COORDINATOR
Sends PREPARE to all participants
⬇️    ⬇️    ⬇️
🗄️ DB-1
YES ✅
Write to WAL
Lock rows
Hold state
🗄️ DB-2
YES ✅
Write to WAL
Lock rows
Hold state
🗄️ DB-3
YES ✅
Write to WAL
Lock rows
Hold state

✅ Phase 2a: COMMIT (All said YES)

🎯 COORDINATOR
Sends COMMIT to all participants
⬇️    ⬇️    ⬇️
🗄️ DB-1
COMMITTED
Apply to data
Release locks
ACK
🗄️ DB-2
COMMITTED
Apply to data
Release locks
ACK
🗄️ DB-3
COMMITTED
Apply to data
Release locks
ACK

💥 Phase 2b: ABORT (Any participant said NO)

🎯 Coordinator sends ROLLBACK to ALL participants
→
DB-1: ROLLBACK ↩️
Release locks
DB-2: ROLLBACK ↩️
Release locks
DB-3: ROLLBACK ↩️
Release locks

🛡️ Result: No partial commits. All-or-nothing guaranteed. Data remains consistent across all participants.

☠️ The Fatal Weakness: Coordinator Failure

☠️ What Happens When the Coordinator Dies After Phase 1?

✅ All three participants replied YES. All have locked their rows.
💥 COORDINATOR CRASHES! Before sending Phase 2.
🗄️ DB-1: Locked, waiting for commit/abort... ⏳
🗄️ DB-2: Locked, waiting for commit/abort... ⏳
🗄️ DB-3: Locked, waiting for commit/abort... ⏳
💀 ALL PARTICIPANTS ARE BLOCKED. Holding locks. Cannot proceed. Cannot roll back safely. System is stuck until coordinator recovers.

This is why 2PC is called a blocking protocol. Participants hold locks indefinitely while waiting for the coordinator's Phase 2 message. Any transaction requiring those locked rows is also blocked. In high-traffic systems, this can cause cascading delays. This blocking is 2PC's fundamental and unsolvable limitation.

Property Two-Phase Commit (2PC)
Consistency✅ Strong — all or nothing, guaranteed
Availability❌ Blocks if coordinator or any participant fails
Latency⚠️ 2 network round-trips minimum
Single point of failure❌ Coordinator is a SPOF
Locks held❌ Entire duration of protocol (potentially minutes)
Best forSmall number of participants, low latency network, strong consistency required
Used inRDBMS distributed transactions, JTA (Java), XA protocol, MSDTC (.NET), PostgreSQL postgres_fdw
✅ Three-Phase Commit (3PC) — The Attempted Fix

3PC adds a third phase (CanCommit → PreCommit → Commit) that adds a "non-blocking" property — participants can make progress even if the coordinator fails during PreCommit. But: 3PC assumes no network partitions. Under network partition, 3PC can still reach inconsistent states. This is why 3PC is rarely used in practice — and why modern systems prefer Paxos/Raft consensus instead.

📖 Section 2: The Saga Pattern — Distributed Transactions Without Locks

The Saga pattern was originally proposed in 1987 for long-running database transactions, then re-discovered by the microservices community as the primary alternative to 2PC. Instead of holding locks across services for the duration of a transaction, a Saga breaks the transaction into a sequence of smaller local transactions — each with a compensating transaction that can undo its work if something fails later.

💡 The Travel Booking Analogy

Book a full holiday: flight + hotel + car rental. You can't hold an airline reservation, hotel room, and car atomically — each is a separate company, separate system.

Instead: book flight → on success, book hotel → on success, book car → done!

If the car rental fails: cancel hotel booking (compensate) → then cancel flight (compensate). No distributed lock. No coordinator. Each step has a local transaction + a cancellation procedure.

Saga = sequence of local transactions + compensating transactions for rollback.

📖 Saga: Order Processing Example — Happy Path & Failure Compensation

✅ Happy Path — All Steps Succeed

🛒 Create
Order
(DB1)
✅→
📦 Reserve
Inventory
(DB2)
✅→
💳 Charge
Payment
(DB3)
✅→
📧 Send
Confirm
(DB4)
→
🎉 Saga
COMPLETE

💥 Failure Path — Payment Fails → Compensating Transactions Execute in Reverse

🛒 Create
Order ✅
→
📦 Reserve
Inventory ✅
→
💳 Charge
Payment ❌
FAILED!
↩️ Compensating transactions execute in REVERSE order:
↩️ Release
Inventory
(undo step 2)
←
↩️ Cancel
Order
(undo step 1)
←
🔄 Saga
ROLLED BACK
⚠️ Important: Between steps, other processes could observe intermediate states (order created, inventory reserved, payment not yet charged). This is eventual consistency — the system reaches a consistent state eventually, not immediately. Design your system to handle this transient inconsistency gracefully.

🎭 Two Saga Implementation Styles

🕺 Choreography-Based Saga

No central coordinator. Each service listens for events and knows what to do next. Services publish events. Other services react to those events. Like a dance — everyone knows their role, no choreographer during the performance.

📨 OrderService → publishes "order.created"
📦 InventoryService ← listens → reserves → publishes "inventory.reserved"
💳 PaymentService ← listens → charges → publishes "payment.completed"
📧 NotifService ← listens → sends email
✅ Simple — no central component
✅ Services fully decoupled
❌ Hard to track overall state
❌ Complex compensation logic distributed everywhere

🎯 Orchestration-Based Saga

A central Saga Orchestrator tells each participant what to do and when. The orchestrator tracks state and decides next steps based on responses. Like a film director — one person coordinates the whole production.

🎯 SAGA ORCHESTRATOR
1. Tell InventoryService: "Reserve items"
↓ success? → 2. Tell PaymentService: "Charge card"
↓ failure? → Tell InventoryService: "Release items"
↓ success? → 3. Tell NotifService: "Send email"
✅ Easy to track saga state
✅ Compensation logic centralised
❌ Orchestrator is a central component (not SPOF, but complex)
✅ Used by Uber Cadence, Netflix Conductor
✅ Who Uses Saga in Production?

🚗 Uber: Cadence workflow engine runs orchestration-based Sagas for trip creation, payment processing, and driver assignment — each step with defined compensations.
🎬 Netflix: Conductor (open-sourced orchestration engine) manages Sagas for content ingestion pipelines and subscriber billing workflows.
🛒 Amazon: Choreography-based Sagas via SNS/SQS for order fulfilment — warehouse, payment, and delivery services react to events independently.

Best for: Microservices architecture, long-running business processes, cases where eventual consistency is acceptable.

🗳️ Section 3: Quorum — Majority Rules in Distributed Systems

Quorum is the most elegant consistency mechanism in distributed systems. Instead of requiring ALL nodes to agree (which fails if any node is down), quorum requires a majority — enough nodes that any two quorums must overlap, guaranteeing shared knowledge.

💡 The Jury Analogy

A jury of 12 people doesn't require ALL 12 to agree — just a majority (9 out of 12, in many jurisdictions). Why? Because even if 3 jurors are absent or unreachable, the remaining 9 can still reach a verdict.

And here's the mathematical magic: if you need 7 of 12 to write a verdict (quorum for writes) and 7 of 12 to read it (quorum for reads), then any write quorum and any read quorum MUST share at least two members in common — the verdict readers always include at least one person who participated in writing it.

This overlap is what guarantees consistency. You always read from at least one node that saw the write.

🗳️ The Quorum Formula — The Foundation of Distributed Consistency

W + R > N
N = Total number of replicas
W = Write quorum (min nodes that must ACK write)
R = Read quorum (min nodes that must respond to read)
When W + R > N → every read overlaps with every write → STRONG CONSISTENCY ✅
🏆 N=3, W=2, R=2
(Balanced — Most Common)
✅
✅
❓
Write needs 2/3 nodes
Read needs 2/3 nodes
W+R=4 > N=3 ✅
Strong consistency
Tolerates 1 failure
✍️ N=3, W=3, R=1
(Fast Reads)
W
W
W
Write needs ALL 3
Read needs just 1
W+R=4 > N=3 ✅
Strong consistency
Slow writes, fast reads
⚡ N=3, W=1, R=1
(Eventual Consistency)
W
?
?
Write needs just 1
Read needs just 1
W+R=2 NOT > N=3 ❌
Eventual consistency only
Fastest performance

🗄️ Quorum in Apache Cassandra — Real Production Example

🗄️ Cassandra Consistency Levels — Quorum Tunable Per Query

ONE — Read/write succeeds if 1 replica responds. W+R=2, N=3. Fastest. Eventually consistent.
QUORUM — Read/write needs majority (⌊N/2⌋+1 = 2 of 3). W+R=4>3. Strong consistency. Most popular.
LOCAL_QUORUM — Quorum within the local data centre only. Best for multi-DC deployments wanting speed.
ALL — All N replicas must respond. Strongest consistency. Highest latency. Fails if ANY replica is down.
💡 Practical tip: Use QUORUM for reads and writes in most cases. Use ONE for non-critical high-frequency writes (analytics events, counters). Use ALL only when you absolutely cannot tolerate stale reads (never in practice, as it breaks availability).

⚡ Quorum in Raft/etcd — Leader Election and Log Replication

Raft uses quorum (majority) for two critical operations:
1. Leader Election
A candidate becomes leader only if it receives votes from a majority (N/2+1) of nodes. In a 5-node cluster: needs 3 votes. This guarantees at most ONE leader at any time — two candidates cannot both achieve majority simultaneously.
2. Log Replication
A log entry is "committed" (safe to apply) only when the leader has replicated it to a majority of nodes. Even if the leader crashes, the committed entry will be present in any node that can win the next election (because quorums overlap).

🕐 Section 4: Vector Clocks — Tracking Causality in Distributed Systems

In a distributed system, events happen on different machines simultaneously. There is no global clock everyone agrees on. How do we know which event happened first? Did two events happen concurrently? Did one cause the other? Vector Clocks answer these questions precisely.

💡 The Detective Story Analogy

Three detectives — Alice, Bob, Carol — are investigating across different cities. Each keeps a notebook tracking how many updates they've made and how many they've heard from each colleague.

When they share notes, they also share these counters: "I've updated my notes 3 times. Last I heard from Bob, he had made 2 updates. Last I heard from Carol, 1 update."

If two detectives have conflicting conclusions, the counters tell us: did one see the other's notes before reaching their conclusion (one caused the other) or did both write independently without knowing what the other concluded (concurrent)?

This is exactly what Vector Clocks track — causality between distributed events.

🕐 Vector Clock — Step-by-Step Mechanics (3 Nodes: A, B, C)

Each node maintains a vector [A_count, B_count, C_count]. On each local event: increment own counter. On message receive: merge vectors (take max per entry), then increment own counter.

🖥️ Node A
Event 1: Local write
A increments own counter
[1, 0, 0]
Event 3: Sends msg to B
A increments own counter
Sends vector with message
[2, 0, 0]
Event 6: Receives msg from B
Merge: max(2,2)=2, max(0,2)=2
Increment A: 3
[3, 2, 0]
→
msg+[2,0,0]
←
msg+[2,2,0]
🖥️ Node B
Event 2: Local write
B increments own counter
[0, 1, 0]
Event 4: Receives from A
Merge: max(0,2)=2, max(1,0)=1
Increment B: 2
[2, 2, 0]
Event 5: Sends reply to A
B increments own counter
Sends vector with reply
[2, 3, 0]
🖥️ Node C (concurrent)
Event 2b: Local write (concurrent with everything!)
C knows nothing about A or B
[0, 0, 1]
CONCURRENT with A and B events!
Neither [0,0,1] < [3,2,0]
Nor [3,2,0] < [0,0,1]
→ They are CONCURRENT.
CONFLICT detected!
🔍 How to Compare Two Vector Clocks:
VC(A) "happens-before" VC(B) if: every entry of VC(A) ≤ corresponding entry of VC(B), AND at least one entry strictly less. Notation: VC(A) → VC(B).
CONCURRENT: Neither VC(A) → VC(B) nor VC(B) → VC(A). The events happened independently. This signals a conflict that must be resolved by the application!

📦 Amazon Dynamo — Vector Clocks in Production

📦 How Amazon Dynamo Uses Vector Clocks for Shopping Cart Conflicts

Scenario: Two mobile devices edit a shopping cart offline simultaneously. Device A adds "milk". Device B adds "eggs". Network reconnects.
Without Vector Clocks: Last-write-wins. One update is silently discarded. Customer loses either the milk or the eggs. Silent data loss. 😱
With Vector Clocks: Dynamo stores both versions with their vector clocks. When user reads the cart, the system detects the conflict ([milk]) and ([eggs]) are concurrent. It returns BOTH versions to the application. The application (or user) merges them. Result: cart contains [milk, eggs]. Both additions preserved! ✅
Semantic reconciliation: The application knows how to merge "add to cart" operations — union of both. Different operations need different merge strategies. Vector clocks detect the conflict; your code resolves it.

⏱️ Lamport Timestamps vs Vector Clocks

⏱️ Lamport Timestamps (Simple)

Single counter per node. On send: max(local, received)+1. Orders events: if A → B then L(A) < L(B). BUT: if L(A) < L(B), does NOT guarantee A → B (could be concurrent). Cannot detect concurrency. Used for total ordering of events (e.g., distributed logging).

🕐 Vector Clocks (Complete)

Counter per node, for ALL nodes. Captures full causality. A → B iff VC(A) ≤ VC(B). Can detect concurrent events. If VC(A) and VC(B) are incomparable → CONCURRENT → CONFLICT. More expressive but more data per message (size grows with number of nodes). Used for conflict detection (DynamoDB, Riak).


🗺️ Section 6: Everything Together — Consistency Patterns in Production

⚖️ Consistency Patterns — Where Each One Lives in a Distributed System

🤝 Two-Phase Commit — RDBMS Distributed Transactions
📊 PostgreSQL shard 1
↔️ 2PC Coordinator ↔️
📊 PostgreSQL shard 2
Use for: banking transactions, inventory adjustments, small-cluster atomic operations
📖 Saga Pattern — Microservices Business Transactions
🎯 Saga Orchestrator
→
📦 Inventory Svc
→
💳 Payment Svc
→
📧 Notif Svc
Use for: order flows, travel booking, subscription management, any multi-service business process
🗳️ Quorum — Replicated Databases and Consensus
✅ Replica 1 (ACK)
✅ Replica 2 (ACK)
❓ Replica 3 (slow)
W=2: quorum achieved! ✅
Use for: Cassandra, DynamoDB, etcd/Raft, Zookeeper — tune per-query for consistency vs availability trade-off
🕐 Vector Clocks — Conflict Detection in AP Systems
Node A: [3,2,0]
⚡ CONCURRENT
Node C: [0,0,1]
→ conflict detected → merge both versions
Use for: Amazon DynamoDB, Riak, CouchDB — detect and resolve concurrent writes in AP distributed systems

📐 Section 7: Core Design Principles from Consistency Patterns

🤝 Strong Consistency is Expensive — Pay Only When You Must

2PC requires multiple round-trips, holds locks, and blocks on failures. Only use it when partial writes are truly unacceptable (financial transactions, reservation systems). For most microservice operations, eventual consistency + Saga is both sufficient and far more resilient. Don't default to strong consistency — default to understanding your actual consistency requirements.

📖 Design Compensations Before Writing Forward Steps

When designing Sagas, always define the compensating transaction before implementing the forward step. If you can't define a meaningful compensation, you may need a different consistency model. Some operations cannot be compensated (you can't "unsend" an email, only send a correction). Compensatable operations are a fundamental design constraint of the Saga pattern.

🗳️ Tune Quorum Per Operation, Not Per System

In Cassandra, you can use QUORUM for reads but ONE for writes on the same table. Or QUORUM writes and ONE reads (if you accept stale reads). The optimal choice depends on the operation: user authentication needs strong reads, analytics events need fast writes. Per-operation tuning is one of quorum's greatest strengths — use it.

🕐 Detect Conflicts Explicitly — Never Silently Discard Writes

Last-write-wins (LWW) is the most common cause of silent data loss in distributed systems. When two concurrent writes exist and your system just picks the "latest" one, real user data is silently destroyed. Vector clocks (or similar mechanisms like CRDTs) make conflicts visible so they can be resolved intentionally — merge, alert, or reject. Silent data loss is always worse than an explicit conflict.


🎓 Section 8: Cheat Sheet

❓ "Explain Two-Phase Commit and its limitations"

Answer: 2PC is a distributed transaction protocol with two phases. Phase 1 (Prepare): coordinator asks all participants if they can commit — each votes YES/NO and locks resources. Phase 2 (Commit/Abort): if all YES → send COMMIT to all; if any NO → send ROLLBACK to all. Guarantees atomicity across distributed nodes. Key limitations: (1) Blocking — if coordinator fails after Phase 1, all participants hold locks indefinitely. (2) Coordinator SPOF. (3) Two network round-trips add latency. Best for small clusters, low latency networks.

❓ "How does the Saga pattern differ from 2PC?"

Answer: Saga replaces one distributed transaction with a sequence of local transactions, each with a compensating transaction for rollback. No distributed locks, no coordinator SPOF. On failure: compensations execute in reverse order. Two styles: choreography (event-driven, decentralised) and orchestration (central coordinator). Key difference from 2PC: Saga provides eventual consistency (intermediate states visible), while 2PC provides strong consistency (atomic across all participants). Saga is preferred for microservices.

❓ "Explain Quorum reads and writes. What is W+R>N?"

Answer: In a system with N replicas, W = write quorum (minimum nodes to ACK write), R = read quorum (minimum nodes to respond to read). When W+R > N, every read overlaps with every write quorum — reads always see the latest write → strong consistency. Example: N=3, W=2, R=2 → W+R=4 > 3 → strong consistency, tolerates 1 node failure. Used in: Cassandra (QUORUM level), DynamoDB (ConsistentRead), Raft consensus (majority). Trade-off: higher quorum = more consistency, less availability (fail if nodes unavailable).

❓ "What are vector clocks and what problem do they solve?"

Answer: Vector clocks track causality between events in distributed systems. Each node maintains a vector of counters, one per node. On local event: increment own counter. On send: increment own counter, attach full vector. On receive: merge (take max), increment own. Comparing two vectors: VC(A) < VC(B) means A happened before B. If incomparable → concurrent. Used by Amazon Dynamo to detect concurrent shopping cart writes — stores both versions, lets application merge. Solves silent data loss from last-write-wins. Detects conflicts explicitly.

❓ "How does Raft achieve consensus? What role does quorum play?"

Answer: Raft is a consensus algorithm that uses quorum for two operations. Leader Election: a candidate becomes leader only with votes from majority (N/2+1 nodes), preventing split-brain (two leaders simultaneously). Log Replication: leader commits an entry only after majority of nodes have written it to their WAL. This guarantees that any new leader (winner of next election) will have all committed entries — because the leader-election and log-replication quorums always overlap. etcd (Kubernetes) and CockroachDB use Raft.

❓ "Design a payment system that guarantees money is never lost"

Answer: Use a Saga pattern with an idempotent orchestrator. Steps: (1) Create payment record (PENDING state, idempotency key). (2) Call payment gateway (with idempotency key). (3) Update order status. Each step has a compensation: refund, mark failed, etc. Add outbox pattern: write intent to DB table atomically with business data, then publish to message queue. Consumers are idempotent. Use QUORUM reads for balance checks to ensure latest balance is always seen. Monitor DLQ for failed compensations.


🎉 Final Summary

🤝 Two-Phase Commit: Prepare (vote) → Commit/Abort. All-or-nothing across distributed nodes. Strong consistency. Blocks on coordinator failure. Use for small clusters, RDBMS distributed transactions. XA protocol standard.
📖 Saga Pattern: Sequence of local transactions + compensating rollbacks. Choreography (event-driven) or Orchestration (central coordinator). Eventual consistency. No locks. Best for microservices. Used by Uber Cadence, Netflix Conductor.
🗳️ Quorum: W+R > N guarantees strong consistency (read/write quorums overlap). Tune per-query: QUORUM for consistency, ONE for availability. Powers Cassandra, DynamoDB, Raft/etcd leader election and log replication.
🕐 Vector Clocks: Each node tracks counters for all nodes. Merge on receive (take max). Compare to determine causality: before, after, or concurrent. Detects conflicts — never silently loses writes. Used by Amazon Dynamo, Riak.
⚖️ CAP Context: 2PC = CP (consistent, may block). Saga = AP (available, eventually consistent). Quorum = tunable between CP and AP. Vector Clocks = conflict detection tool for AP systems.
🔑 Compensation is not rollback: In Saga, a compensation is a new forward action that logically undoes a previous step. It may involve new events, notifications, refunds — not a database-level rollback. Design compensations explicitly.
🌐 Never silent last-write-wins: When concurrent writes occur, detect them with vector clocks or CRDTs, then resolve with domain logic. Silent data loss is always worse than an explicit conflict that can be resolved.
✅ The Most Important Insight from Distributed Consistency:

These four patterns exist because there is no single correct answer to consistency. The right choice depends entirely on your consistency requirements, failure tolerance, latency budget, and the nature of the operations you're coordinating.

The best distributed systems engineers don't memorise patterns — they understand the trade-offs so deeply that they can derive the right pattern from first principles for any new problem they encounter.

That's the level this post aims to take you to. 🎯


Happy Learning! Keep Building! 🔥

Comments