CACS352 Distributed System

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:

Prepare PhaseLocks acquiredCommit PhaseFinal decision sentRecovery PhaseCheckpointing & logs
Phases of 2PC with key actions in each phase.

Phase 1: Prepare Phase

  1. Coordinator (e.g., a central server) sends a prepare request to all participants (e.g., banks, inventory systems).
  2. Each participant votes "yes" (ready to commit) or "no" (abort) after checking local constraints (e.g., sufficient funds).
  3. 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: Commit

Why 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:

[object Object][object Object][object Object][object Object][object Object][object Object]Primary NodeBackup Node 1Backup Node 2Client
Primary-backup replication workflow for synchronous writes.

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 --> Active

6. 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:

  1. Debit from your account (Bank A).
  2. Credit to the merchant’s account (Bank B).
  3. Update eSewa’s transaction log.

How 2PC Ensures Atomicity:

  1. eSewa (Coordinator) sends prepare to Bank A and Bank B.
  2. Both banks lock your account and the merchant’s account.
  3. Banks vote "yes" (sufficient funds).
  4. eSewa sends commit; both banks finalize the transfer.
  5. 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

  1. Definitions:

    • Expect questions on 2PC, ACID, replication strategies, and conflict resolution.
    • Example: "Define distributed commit and explain why 2PC uses a two-phase handshake."
  2. Diagrams:

    • Always draw sequence diagrams for 2PC or state diagrams for recovery.
    • Example: Show the prepare-commit flow with votes and rollback paths.
  3. 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?"
  4. 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
  5. 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?"
  6. 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

  1. 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.
  2. 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).
  3. 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…