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.
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.
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)
Sends PREPARE to all participants
YES ✅
Write to WAL
Lock rows
Hold state
YES ✅
Write to WAL
Lock rows
Hold state
YES ✅
Write to WAL
Lock rows
Hold state
✅ Phase 2a: COMMIT (All said YES)
Sends COMMIT to all participants
COMMITTED
Apply to data
Release locks
ACK
COMMITTED
Apply to data
Release locks
ACK
COMMITTED
Apply to data
Release locks
ACK
💥 Phase 2b: ABORT (Any participant said NO)
Release locks
Release locks
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?
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 for | Small number of participants, low latency network, strong consistency required |
| Used in | RDBMS distributed transactions, JTA (Java), XA protocol, MSDTC (.NET), PostgreSQL postgres_fdw |
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.
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
Order
(DB1)
Inventory
(DB2)
Payment
(DB3)
Confirm
(DB4)
COMPLETE
💥 Failure Path — Payment Fails → Compensating Transactions Execute in Reverse
Order ✅
Inventory ✅
Payment ❌
FAILED!
Inventory
(undo step 2)
Order
(undo step 1)
ROLLED BACK
🎭 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.
🎯 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.
🚗 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.
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 = 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 ✅
(Balanced — Most Common)
Read needs 2/3 nodes
W+R=4 > N=3 ✅
Strong consistency
Tolerates 1 failure
(Fast Reads)
Read needs just 1
W+R=4 > N=3 ✅
Strong consistency
Slow writes, fast reads
(Eventual Consistency)
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
⚡ Quorum in Raft/etcd — Leader Election and Log Replication
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.
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.
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.
A increments own counter
[1, 0, 0]
A increments own counter
Sends vector with message
[2, 0, 0]
Merge: max(2,2)=2, max(0,2)=2
Increment A: 3
[3, 2, 0]
msg+[2,0,0]
msg+[2,2,0]
B increments own counter
[0, 1, 0]
Merge: max(0,2)=2, max(1,0)=1
Increment B: 2
[2, 2, 0]
B increments own counter
Sends vector with reply
[2, 3, 0]
C knows nothing about A or B
[0, 0, 1]
Neither [0,0,1] < [3,2,0]
Nor [3,2,0] < [0,0,1]
→ They are CONCURRENT.
CONFLICT detected!
📦 Amazon Dynamo — Vector Clocks in Production
📦 How Amazon Dynamo Uses Vector Clocks for Shopping Cart Conflicts
⏱️ Lamport Timestamps vs Vector Clocks
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).
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
📐 Section 7: Core Design Principles from Consistency Patterns
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.
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.
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.
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
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.
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.
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).
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.
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.
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
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
Post a Comment