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:
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
- A bank saves a checkpoint after processing 100 transactions.
- 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):
- Prepare phase: Coordinator asks all nodes if they can commit.
- 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: CommitDisadvantage: 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.
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:
- Bank saves a checkpoint: A = ₹10,000, B = ₹5,000.
- Coordinator asks A and B to prepare the transfer.
- Both A and B confirm they can deduct/credit.
- Coordinator sends commit signals.
- Failure occurs: Link between A and coordinator breaks.
- 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
endIn the Real World
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.
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.
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…