CACS352 Distributed System

Distributed SystemUnit 69 min read

Fault Tolerance & Recovery: Mechanisms, Models & Real-World Resilience

Unit 6 of Distributed System: Explores how distributed systems detect, isolate, and recover from faults (hardware/software failures, network partitions) using redundancy, checkpointing, and consensus protocols—with worked examples from banking (NEPSE), e-commerce (Daraz), and telecom (NTC).

TAKEAWAYS:

  • Fault tolerance relies on redundancy (backup nodes, replicated data) and graceful degradation (systems keep running despite failures).
  • Recovery mechanisms include checkpointing, logging, and atomic commit protocols (e.g., 2PC) to restore consistency after failures.
  • Failure classification (crash, omission, timing, Byzantine) dictates which recovery strategy works (e.g., NTC’s redundant fiber routes handle link failures).
  • Consensus algorithms (Paxos, Raft) ensure agreement among nodes even if some fail (used in Daraz’s order processing).
  • State recovery (rollback, forward recovery) depends on whether the system can undo or redo failed operations.
  • Trade-offs: Higher redundancy improves fault tolerance but increases cost and latency (e.g., NEPSE’s replicated trading servers vs. slower response times).

1. Definitions and Failure Models

Fault tolerance is the ability of a distributed system to continue operating correctly despite failures. Failures can be classified into four categories:

PhysicalBitsData LinkFramesNetworkPacketsTransportSegmentsApplicationData
OSI model layers showing where different failure types (e.g., link failure at Data Link) occur in a distributed system.
stateDiagram-v2
    [*] --> Crash
    Crash: Node stops responding
    [*] --> Omission
    Omission: Message lost or not delivered
    [*] --> Timing
    Timing: Message arrives too late
    [*] --> Byzantine
    Byzantine: Node behaves arbitrarily (malicious or buggy)

Key terms:

  • Fault: An error or flaw in a component (e.g., a server crashing).
  • Failure: The observable effect of a fault (e.g., a service becoming unavailable).
  • Recovery: Restoring the system to a correct state after a failure.

2. Fault Tolerance Mechanisms

Systems use redundancy and diversity to tolerate faults. Common mechanisms include:

stateDiagram-v2
    [*] --> Crash
    Crash --> Node_Redundancy
    Node_Redundancy --> Primary_Server
    Node_Redundancy --> Backup_Server
    Backup_Server --> Takes_Over
    Takes_Over --> [*]
    Crash --> Data_Redundancy
    Data_Redundancy --> Replicated_Data
    Replicated_Data --> Sync_Updates
    Sync_Updates --> [*]
    Crash --> Link_Redundancy
    Link_Redundancy --> Fiber_Ring
    Fiber_Ring --> Traffic_Reroute
    Traffic_Reroute --> [*]
How redundancy mechanisms handle different failure types in distributed systems.

A. Redundancy

  • Node redundancy: Extra servers handle workload if primary fails (e.g., NEPSE’s replicated trading nodes).
  • Data redundancy: Replicating data across nodes (e.g., Khalti’s payment gateways).
  • Link redundancy: Multiple network paths (e.g., NTC’s fiber rings).

Visual: NTC’s Fiber Network Redundancy


B. Checkpointing and Logging

  • Checkpointing: Periodically saving the system state to disk (e.g., bank transactions before processing).
  • Logging: Recording operations for replay (e.g., Daraz’s order logs for recovery).

Example: Bank Transaction Recovery

  1. A bank saves a checkpoint after processing 100 transactions.
  2. If a failure occurs, the bank rolls back to the checkpoint and replays logs from there.

C. Atomic Commit Protocols

Ensures all nodes agree on a global state change (e.g., transferring money between accounts). Two-Phase Commit (2PC):

  1. Prepare phase: Coordinator asks all nodes if they can commit.
  2. Commit phase: If all say "yes," all commit; else, all abort.
sequenceDiagram
    participant BankA
    participant BankB
    participant Coordinator
    Coordinator->>BankA: Prepare (Transfer $100)
    Coordinator->>BankB: Prepare (Receive $100)
    BankA-->>Coordinator: ACK
    BankB-->>Coordinator: ACK
    Coordinator->>BankA: Commit
    Coordinator->>BankB: Commit

Disadvantage: Coordinator becomes a single point of failure.


3. Recovery Mechanisms

Recovery depends on whether the system can undo (rollback) or redo (forward recovery) failed operations.

Checkpoint 1Bank state saved:A=₹10,000, B=₹5,000Transaction 101Transfer ₹5,000(A→B) initiatedPrepare PhaseCoordinator asks A& B to prepareCommit PhaseA & B confirm →Commit signals sentFailureLink break(A→Coordinator) → RollRecoveryA reverts toCheckpoint 1, B unchan

A. Rollback Recovery

  • Reverts to a previous checkpoint (e.g., undoing a failed Daraz order).
  • Requires undo logs (e.g., "Cancel order #12345").

B. Forward Recovery

  • Replays logs from the last checkpoint (e.g., reprocessing a failed payment in eSewa).
  • Requires redo logs (e.g., "Process payment again").

Comparison Table:

Mechanism Use Case Pros Cons
Rollback Undo failed transactions Fast recovery Data loss if checkpoint old
Forward Reprocess from checkpoint No data loss Slower if many retries needed
2PC Atomic commits (banks) Strong consistency Coordinator bottleneck
Checkpointing Periodic state saves Low overhead Stale data if failure occurs

4. Handling Different Failures

Failure Type Example Scenario Recovery Strategy
Crash Server reboot (Ncell’s base station) Restart from checkpoint
Omission Message lost (WhatsApp delivery) Retransmission + acknowledgments
Timing Late message (YouTube buffer) Timeout + retry
Byzantine Malicious node (hacked Daraz server) Consensus (Paxos/Raft)

5. Real-World Examples

A. NEPSE’s Replicated Trading System

  • Idea: Uses state machine replication to ensure all nodes agree on stock trades.
  • How: If one node fails, others continue trading without interruption.

B. Pathao’s Ride Dispatching

  • Idea: Uses checkpointing to save driver assignments before failures.
  • How: If the system crashes, it replays logs to reassign rides.

C. NTC’s Redundant Fiber Routes

  • Idea: Link redundancy ensures calls remain connected even if one fiber breaks.
  • How: Traffic automatically reroutes via backup paths.

6. Worked Example: Distributed Commit in a Bank

Scenario: Transferring ₹5,000 from Account A to Account B. Steps:

  1. Bank saves a checkpoint: A = ₹10,000, B = ₹5,000.
  2. Coordinator asks A and B to prepare the transfer.
  3. Both A and B confirm they can deduct/credit.
  4. Coordinator sends commit signals.
  5. Failure occurs: Link between A and coordinator breaks.
  6. Recovery:
    • A rolls back to checkpoint (A = ₹10,000).
    • B remains unchanged (B = ₹5,000).
    • Coordinator restarts 2PC from checkpoint.
sequenceDiagram
    participant BankA
    participant BankB
    participant Coordinator
    BankA->>Coordinator: Prepare (Transfer ₹5,000)
    BankB->>Coordinator: Prepare (Receive ₹5,000)
    Coordinator->>BankA: ACK
    Coordinator->>BankB: ACK
    Coordinator->>BankA: Commit
    alt Link Failure
    BankA->>Coordinator: Timeout
    Coordinator->>BankA: Abort
    BankA->>BankA: Rollback to Checkpoint
    end

In the Real World

  1. Daraz’s Order Processing:

    • Uses checkpointing to save unfulfilled orders before crashes.
    • If the system fails, it replays logs to resume processing from the last checkpoint.
    • Idea: Forward recovery ensures no orders are lost.
  2. Ncell’s Call Routing:

    • Employs redundant base stations to handle cell tower failures.
    • If one tower crashes, calls reroute via backup towers.
    • Idea: Link redundancy maintains service continuity.
  3. NEPSE’s Trading System:

    • Implements state machine replication to prevent inconsistent trades.
    • If a node fails, others continue trading and sync later.
    • Idea: Consensus ensures all nodes agree on the market state.

Exam Tip

  • Focus on 2PC: Explain its two phases and why it’s used for atomic commits (e.g., banking).
  • Compare recovery methods: Know when to use rollback vs. forward recovery (e.g., Daraz vs. banks).
  • Link to real systems: Mention NEPSE, NTC, or Daraz in answers—examiners love context!
  • Diagrams: Always draw sequence diagrams for 2PC or state diagrams for failure types.
  • Trade-offs: Discuss pros/cons of redundancy (e.g., cost vs. reliability in NTC’s fiber network).

Key Formula to Remember: For 2PC success rate: Where:

  • = Probability all nodes respond to prepare.
  • = Probability all nodes commit after prepare.

Final Visual: Fault Tolerance Layers


Based on the TU BCA syllabus for Distributed System (CACS352), unit 6.

Discussion

Loading…