Distributed SystemUnit 513 min read
Distributed Commit & Transaction Management: 2PC, ACID, Replication
Unit 5 of Distributed System: Explores how distributed systems ensure atomicity across multiple nodes via two-phase commit (2PC), transaction models (ACID), replication strategies, and failure recovery, with real-world ties to banking, e-commerce, and cloud services.
TAKEAWAYS:
- Distributed commit protocols like 2PC enforce atomicity in transactions spanning multiple nodes, using a prepare-and-commit handshake to avoid partial updates.
- ACID properties (Atomicity, Consistency, Isolation, Durability) define reliable transaction processing, while BASE (for eventual consistency) is used in scalable systems like WhatsApp’s message delivery.
- Replication strategies (primary-backup, multi-leader) balance consistency and availability, with quorum-based systems (e.g., NEPSE’s stock trading) ensuring data integrity.
- Failure recovery in distributed systems relies on checkpoints, logs, and recovery blocks to roll back or forward transactions after crashes.
- Conflict resolution in replicated systems uses vector clocks, causal ordering, or timestamp-based methods to merge updates without data loss.
- Practical examples: eSewa’s transaction finality (2PC), Daraz’s inventory updates (multi-leader replication), and NTC’s billing systems (checkpoint recovery).
1. Introduction to Distributed Transactions
Distributed transactions occur when a single logical operation spans multiple nodes (e.g., debiting an account in Bank A and crediting Bank B). Without coordination, partial updates can corrupt data (e.g., a flight booking where only the passenger’s seat is reserved but the airline’s inventory isn’t updated). Distributed commit protocols ensure all nodes agree on the outcome (commit or abort) before finalizing.
Key Challenge: How to guarantee atomicity (all-or-nothing) when nodes may fail or communicate asynchronously?
2. Two-Phase Commit (2PC) Protocol
The gold standard for distributed commit, 2PC divides the process into two phases:
Phase 1: Prepare Phase
- Coordinator (e.g., a central server) sends a prepare request to all participants (e.g., banks, inventory systems).
- Each participant votes "yes" (ready to commit) or "no" (abort) after checking local constraints (e.g., sufficient funds).
- Participants lock resources (e.g., account balances) to prevent conflicts.
Phase 2: Commit Phase
- If all participants vote "yes", the coordinator sends a commit request. All participants finalize the transaction.
- If any participant votes "no", the coordinator sends an abort request. All participants roll back changes.
Mermaid Sequence Diagram:
sequenceDiagram
participant Coordinator
participant BankA
participant BankB
participant Inventory
Coordinator->>BankA: Prepare(Transaction T)
Coordinator->>BankB: Prepare(Transaction T)
Coordinator->>Inventory: Prepare(Transaction T)
BankA-->>Coordinator: Vote Yes (locks account)
BankB-->>Coordinator: Vote Yes (locks account)
Inventory-->>Coordinator: Vote Yes (reserves seat)
Coordinator->>BankA: Commit
Coordinator->>BankB: Commit
Coordinator->>Inventory: CommitWhy 2PC?
- Atomicity: Ensures all nodes commit or abort together.
- Durability: Uses logs (e.g., transaction journals) to recover after crashes.
- Consistency: Prevents partial updates (e.g., a bank transfer where only one account is updated).
Limitations:
- Blocking: If the coordinator fails during Phase 2, participants may be stuck in a "prepared" state (requiring a recovery manager).
- Performance: High latency due to synchronous voting.
- Single Point of Failure: Coordinator’s crash halts the system.
3. ACID Properties in Distributed Transactions
ACID properties define reliable transaction processing in distributed systems:
| Property | Definition | Example |
|---|---|---|
| Atomicity | Transaction is treated as a single unit (all-or-nothing). | eSewa’s payment: either debits your account and credits the merchant, or neither. |
| Consistency | Transaction moves the system from one valid state to another. | A bank’s total money supply remains unchanged after a transfer. |
| Isolation | Concurrent transactions do not interfere (e.g., no dirty reads). | Two users booking the same flight seat simultaneously are blocked. |
| Durability | Once committed, changes survive system failures (e.g., via logs). | NEPSE’s stock trades persist even after a server reboot. |
Conflict Example: Without isolation, dirty reads can occur:
- Transaction A reads an uncommitted update from Transaction B.
- If B rolls back, A’s data becomes invalid.
4. Replication Strategies for Distributed Transactions
Replication improves availability and fault tolerance but introduces consistency trade-offs. Common strategies:
A. Primary-Backup Replication
- Primary node handles all writes; backups replicate changes asynchronously.
- Use Case: Traditional databases (e.g., Ncell’s customer database).
- Pros: Simple, low latency for reads.
- Cons: Primary failure halts writes; backups may fall behind (stale reads).
B. Multi-Leader Replication
- Multiple nodes accept writes; conflicts resolved via last-write-wins or application logic.
- Use Case: Daraz’s inventory system (multiple warehouses update stock in real time).
- Pros: High availability; works offline.
- Cons: Risk of split-brain (inconsistent data if leaders disagree).
C. Quorum-Based Replication
- Writes require majority of nodes to agree (e.g., 3/5 nodes for a commit).
- Use Case: NEPSE’s trading system (ensures no partial order execution).
- Pros: Strong consistency; fault tolerance.
- Cons: Higher latency; stricter write rules.
Mermaid Comparison Table:
5. Handling Failures in Distributed Transactions
Distributed systems must recover from:
- Node crashes (e.g., a bank server going down mid-transaction).
- Network partitions (e.g., NTC’s regional offices losing connectivity).
- Disk failures (e.g., a log file corrupting).
A. Checkpointing
- Periodically snapshot the system state (e.g., every 5 minutes).
- On recovery, replay logs since the last checkpoint.
- Example: A bank’s nightly backup of all transactions.
B. Transaction Logs
- Append-only logs record all changes (e.g.,
T1: Debit Account X by Rs. 1000). - On restart, logs are replayed to recover the last committed state.
- Example: WhatsApp’s message delivery logs ensure no message is lost.
C. Recovery Blocks
- Atomic actions wrapped in a recovery block (e.g.,
BEGIN TRANSACTION→COMMIT/ROLLBACK). - If a step fails, the entire transaction is aborted.
- Example: Pathao’s ride booking system rolls back if payment fails.
Mermaid State Diagram for Recovery:
stateDiagram-v2
[*] --> Active
Active --> Checkpoint: Periodically save state
Active --> Log: Append transaction to log
Active --> Crash: Node fails
Crash --> Recovery: Restart from checkpoint + replay logs
Recovery --> Active
Active --> Abort: Failure detected
Abort --> Rollback: Undo transaction
Rollback --> Active6. Conflict Resolution in Replicated Systems
When multiple nodes update the same data concurrently, conflicts arise. Solutions:
A. Vector Clocks
- Track happened-before relationships between events.
- Example: Two users editing a Google Doc simultaneously; vector clocks ensure changes are merged without loss.
B. Causal Ordering
- Preserves causal dependencies (e.g., if A sends a message to B, B’s response must follow A’s action).
- Example: WhatsApp’s message delivery ensures replies come after the original message.
B. Timestamp-Based
- Assign logical clocks to events (e.g., Lamport timestamps).
- Example: NEPSE’s trading system orders transactions by timestamp to avoid race conditions.
Mermaid Vector Clock Example:
graph TD
A["User 1: Update Price"] -->|"Clock=1"| B["Node 1: Receives Update"]
A -->|"Clock=1"| C["Node 2: Receives Update"]
B -->|"Clock=2"| D["Node 1: Applies Update"]
C -->|"Clock=3"| E["Node 2: Applies Update"]
D -->|"Clock=4"| F["Conflict Detected: Merge Changes"]7. Practical Example: eSewa’s Transaction Finality
Scenario: You pay Rs. 1000 to a merchant via eSewa. The transaction involves:
- Debit from your account (Bank A).
- Credit to the merchant’s account (Bank B).
- Update eSewa’s transaction log.
How 2PC Ensures Atomicity:
- eSewa (Coordinator) sends prepare to Bank A and Bank B.
- Both banks lock your account and the merchant’s account.
- Banks vote "yes" (sufficient funds).
- eSewa sends commit; both banks finalize the transfer.
- If Bank A had insufficient funds, it would vote "no", and eSewa would abort, returning your money.
Failure Scenario:
- If eSewa crashes during Phase 2, banks are stuck in a "prepared" state.
- On recovery, eSewa checks logs and forces a commit (or aborts if timeout expires).
8. BASE vs. ACID: Trade-offs
| Property | ACID (Traditional) | BASE (Eventual Consistency) |
|---|---|---|
| Consistency | Strong (immediate) | Eventual (tolerates temporary inconsistency) |
| Availability | Lower (blocks during conflicts) | Higher (allows offline writes) |
| Performance | Higher latency (synchronous) | Lower latency (asynchronous) |
| Use Case | Banking (eSewa), Stock Trading (NEPSE) | Social Media (WhatsApp), E-commerce (Daraz) |
Example: WhatsApp uses BASE for message delivery:
- Your message is delivered eventually (not immediately).
- If offline, it syncs later (no blocking).
9. Exam Tip: How This Unit is Tested
Definitions:
- Expect questions on 2PC, ACID, replication strategies, and conflict resolution.
- Example: "Define distributed commit and explain why 2PC uses a two-phase handshake."
Diagrams:
- Always draw sequence diagrams for 2PC or state diagrams for recovery.
- Example: Show the prepare-commit flow with votes and rollback paths.
Real-World Applications:
- Link concepts to eSewa (2PC), Daraz (multi-leader), or NEPSE (quorum-based).
- Example: "How does NEPSE’s stock trading system use quorum replication to prevent partial order execution?"
Comparison Tables:
- Compare ACID vs. BASE, primary-backup vs. multi-leader, or vector clocks vs. timestamps.
- Example:
Method Pros Cons 2PC Atomicity Blocking, coordinator risk BASE Scalability Eventual inconsistency
Failure Scenarios:
- Describe how a system recovers from node crashes or network partitions.
- Example: "If the coordinator in 2PC fails during Phase 2, how does the system recover?"
Worked Example:
- Solve a transaction trace (e.g., "Given a log of prepare/commit votes, what happens if Bank B votes 'no'?").
- Example:
Coordinator -> BankA: Prepare(T1) BankA -> Coordinator: Yes (locks Rs. 5000) Coordinator -> BankB: Prepare(T1) BankB -> Coordinator: No (insufficient funds) Coordinator -> BankA: Abort(T1)
In the Real World
eSewa’s Payment Finality:
- Uses 2PC to ensure your money is either debited and credited or neither.
- Why it matters: Prevents fraud where a merchant gets paid but your account isn’t updated.
Daraz’s Inventory Replication:
- Employs multi-leader replication across warehouses to update stock in real time.
- Why it matters: If a product is sold in Kathmandu, the inventory in Pokhara updates instantly (no overselling).
NEPSE’s Stock Trading:
- Relies on quorum-based replication (e.g., 3/5 brokers must agree) to finalize trades.
- Why it matters: Ensures no partial order execution (e.g., a stock isn’t bought but the seller’s account isn’t credited).
Visual Summary
Based on the TU BCA syllabus for Distributed System (CACS352), unit 5.
Discussion
Loading…