Distributed NetworkingUnit 611 min read
Reliability & Replication in Distributed Systems
Unit 6 of Distributed Networking explores how distributed systems achieve fault tolerance through replication, consistency models, and recovery techniques—covering primary-backup, quorum-based, and state machine approaches with real-world examples from banking, cloud storage, and messaging apps.
TAKEAWAYS:
- Replication copies data/processes across nodes to survive failures, but introduces trade-offs between consistency, availability, and partition tolerance (CAP theorem).
- Primary-backup systems designate a leader node for writes, while backups replicate state, but failover must be fast to avoid downtime.
- Quorum-based protocols (e.g., Raft, Paxos) ensure consistency by requiring majority agreement for reads/writes, but increase latency and network overhead.
- State machine replication synchronizes all nodes’ states via deterministic execution of commands, guaranteeing identical outputs even if some nodes fail.
- Checkpointing and logging enable recovery by saving system state periodically and replaying logs after crashes.
- Consistency models (strong, eventual, causal) define how updates propagate, with trade-offs between performance and correctness.
Core Concepts: Why Replication?
Distributed systems cannot rely on a single node for critical operations. If one node fails, the entire system crashes unless redundancy is built in. Replication solves this by maintaining multiple copies of data or processes across nodes. However, replication introduces new challenges:
- Consistency: How do we ensure all replicas agree on the same data?
- Availability: How do we keep the system running if some nodes fail?
- Partition Tolerance: How do we handle network splits (e.g., a router failure isolating some nodes)?
These trade-offs are formalized in the CAP theorem:
mindmap
root((CAP Theorem))
CA: Consistency + Availability (but not Partition Tolerant)
CP: Consistency + Partition Tolerance (but not Available)
AP: Availability + Partition Tolerance (but not Consistent)Replication Strategies
1. Primary-Backup Replication
How it works:
- One primary node handles all write requests and replicates updates to backup nodes.
- Backups apply updates asynchronously (or synchronously) to stay in sync.
- If the primary fails, a backup is promoted to take over.
Visual: Primary-Backup Flow
sequenceDiagram
participant Client
participant Primary
participant Backup1
participant Backup2
Client->>Primary: Write Request (e.g., "Deposit $100")
Primary->>Backup1: Replicate Log
Primary->>Backup2: Replicate Log
Primary-->>Client: ACK
loop Crash Recovery
Primary) fails
Backup1->>Backup1: Promote to Primary
Backup1->>Backup2: Sync State
endWorked Example: Ncell’s Billing System Ncell replicates billing records across three data centers (Kathmandu, Pokhara, Biratnagar). If the primary DC fails:
- A backup DC takes over within <2 seconds.
- Customers’ call credits and usage logs remain consistent.
- Trade-off: Synchronous replication adds ~50ms latency to updates.
Advantages:
- Simple to implement.
- Low overhead for read-heavy workloads (reads can go to any replica).
Disadvantages:
- Single point of failure: Primary node failure causes downtime.
- Stale reads: Backups may lag behind the primary.
2. Quorum-Based Replication (Raft, Paxos)
How it works:
- Writes require a majority of nodes (
W + R > N, whereN= total nodes) to acknowledge. - Reads may also need a quorum to ensure up-to-date data.
- Example: In a 5-node cluster, a write needs 3 ACKs (quorum = 3).
Visual: Quorum Write Process
sequenceDiagram
participant Client
participant Node1
participant Node2
participant Node3
Client->>Node1: Write Request
Client->>Node2: Write Request
Client->>Node3: Write Request
Note over Client,Node1: Quorum = 3/5 nodes
Node1-->>Client: ACK
Node2-->>Client: ACK
Node3-->>Client: ACKWorked Example: Daraz’s Order Processing Daraz’s inventory system uses quorum-based replication across its Singapore, India, and Nepal warehouses:
- When you place an order, Daraz’s system writes to 3 out of 5 replicas (quorum).
- If one warehouse (e.g., Nepal) goes offline, orders still process via other warehouses.
- Trade-off: Requires 50%+1 nodes to be up for writes (e.g., 3/5).
Advantages:
- Fault tolerance: Survives up to
(N/2 - 1)failures. - No single point of failure.
Disadvantages:
- High latency: Requires network round-trips to multiple nodes.
- Complexity: Harder to implement than primary-backup.
3. State Machine Replication
How it works:
- All nodes execute the same sequence of commands in the same order.
- Commands are deterministic (same input → same output).
- Uses total order broadcast (e.g., via Raft or Paxos) to agree on command order.
Visual: State Machine Replication
sequenceDiagram
participant Client
participant NodeA
participant NodeB
participant NodeC
Client->>NodeA: Command (e.g., "Transfer $50")
NodeA->>NodeB: Broadcast Command
NodeA->>NodeC: Broadcast Command
NodeB->>NodeA: ACK (with log index)
NodeC->>NodeA: ACK (with log index)
NodeA-->>Client: ACK (after quorum)
Note over NodeA,NodeB: All nodes execute command in same orderWorked Example: WhatsApp’s End-to-End Encryption WhatsApp uses state machine replication for its Signal Protocol:
- Your message is encrypted and sent to a WhatsApp server cluster.
- All servers in the cluster execute the same encryption/decryption steps (deterministic).
- If one server fails, others continue processing the message.
- Why it matters: Ensures your messages are decrypted correctly even if some servers crash.
Advantages:
- Strong consistency: All nodes see the same state.
- Recoverable: Nodes can replay logs to restore state after a crash.
Disadvantages:
- High overhead: Requires agreement on every command.
- Not suitable for high-latency systems.
Consistency Models
Different systems tolerate different levels of inconsistency. The key models:
| Model | Definition | Example Use Case | Trade-offs |
|---|---|---|---|
| Strong | All reads return the most recent write. | Banking transactions (NMB, Global IME) | High latency, complex to implement. |
| Eventual | Replicas will eventually converge, but may temporarily diverge. | Social media (Facebook posts) | Fast writes, but stale reads possible. |
| Causal | Preserves the "happens-before" relationship between events. | Chat apps (WhatsApp, Slack) | Balances consistency and performance. |
| Sequential | Reads appear to execute in a total order. | Stock trading systems (NEPSE) | High coordination overhead. |
Recovery Techniques
Even with replication, systems crash. Here’s how they recover:
1. Checkpointing
- Periodically save the system’s state (e.g., every 5 minutes).
- Recovery: Restore the last checkpoint and replay logs since then.
Visual: Checkpointing Process
stateDiagram-v2
[*] --> Active
Active --> Checkpoint: Save State
Checkpoint --> Active: Resume
Active --> Crash: Failure
Crash --> Recover: Load Last Checkpoint
Recover --> ReplayLogs: Apply Logs
ReplayLogs --> Active: RestoredWorked Example: NTC’s Network Switches NTC’s core routers use checkpointing to survive power outages:
- Every 10 minutes, the router saves its routing table and active connections.
- If the router crashes, it reloads the last checkpoint and replays connection logs to restore service.
2. Logging and Write-Ahead Logging (WAL)
- Before modifying data, write the change to a stable log (e.g., disk).
- Recovery: Replay logs to redo uncommitted changes.
Visual: Write-Ahead Logging
sequenceDiagram
participant App
participant Log
participant Database
App->>Log: Write Log Entry (e.g., "UPDATE balance")
Log-->>App: ACK
App->>Database: Apply Update
Database-->>App: ACK
loop Crash Recovery
Database) fails
Database->>Log: Replay Logs
Log-->>Database: Redo Updates
endWorked Example: eSewa’s Transaction Logs When you pay a bill via eSewa:
- eSewa’s system logs the transaction before deducting money.
- If the system crashes mid-transaction, eSewa replays logs to ensure your payment is either fully completed or rolled back.
In the Real World
Khalti’s Payment System
- Idea Used: Primary-backup replication with quorum-based writes.
- How: When you transfer money, Khalti writes to 3 out of 5 data centers (quorum). If one DC fails, your transaction still goes through via the others.
- Real Impact: Ensures 99.99% uptime even during Nepal’s frequent power cuts.
Pathao’s Ride-Matching Algorithm
- Idea Used: State machine replication for driver availability.
- How: Pathao’s servers in Kathmandu, Pokhara, and Delhi execute the same ride-matching logic in lockstep. If one server crashes, others continue matching rides without inconsistency.
- Real Impact: Prevents duplicate ride assignments or lost orders.
NEPSE’s Stock Trading Platform
- Idea Used: Strong consistency with checkpointing.
- How: Every trade is logged and replicated across three servers. If a server fails, NEPSE restores from the last checkpoint and replays trades to keep all investors’ portfolios accurate.
- Real Impact: Avoids disputed trades or incorrect share prices.
Comparison Table: Replication Strategies
| Feature | Primary-Backup | Quorum-Based (Raft) | State Machine Replication |
|---|---|---|---|
| Consistency | Eventual (if async) | Configurable | Strong |
| Fault Tolerance | Survives backup failures | Survives (N/2 - 1) failures |
Survives any node failure |
| Latency | Low (async) / High (sync) | High (quorum waits) | Very High (total order) |
| Complexity | Low | Medium | High |
| Use Case | Web apps (Daraz), DBs | Critical systems (Ncell) | Encrypted messaging (WhatsApp) |
Exam Tip
Define Key Terms Clearly:
- Distinguish between primary-backup, quorum-based, and state machine replication.
- Explain the CAP theorem and which systems prioritize CA, CP, or AP.
Worked Examples Are Gold:
- Expect 2-3 marks for applying replication to real scenarios (e.g., "How would you design Khalti’s replication?").
- Trace a write/read operation step-by-step (e.g., "Show how a Daraz order is replicated").
Consistency Models:
- 1 mark each for naming models (strong, eventual, causal) and 2 marks for matching them to use cases (e.g., "Why does WhatsApp use causal consistency?").
Recovery Techniques:
- Checkpointing vs. logging: Know when to use each (e.g., checkpointing for long-running processes, logging for databases).
- Draw a sequence diagram for recovery (e.g., "Show how NTC restores routes after a crash").
Common Pitfalls:
- Don’t confuse quorum size with majority: Quorum =
W + R > N(e.g., for 5 nodes, quorum = 3). - State machine replication ≠ primary-backup: The former requires total order, not just leader election.
- Don’t confuse quorum size with majority: Quorum =
Based on the TU BSc CSIT syllabus for Distributed Networking, unit 6.
Discussion
Loading…