Elective Distributed Networking

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:

  1. Consistency: How do we ensure all replicas agree on the same data?
  2. Availability: How do we keep the system running if some nodes fail?
  3. 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
    end

Worked Example: Ncell’s Billing System Ncell replicates billing records across three data centers (Kathmandu, Pokhara, Biratnagar). If the primary DC fails:

  1. A backup DC takes over within <2 seconds.
  2. Customers’ call credits and usage logs remain consistent.
  3. 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, where N = 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: ACK

Worked 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 order

Worked Example: WhatsApp’s End-to-End Encryption WhatsApp uses state machine replication for its Signal Protocol:

  1. Your message is encrypted and sent to a WhatsApp server cluster.
  2. All servers in the cluster execute the same encryption/decryption steps (deterministic).
  3. 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: Restored

Worked Example: NTC’s Network Switches NTC’s core routers use checkpointing to survive power outages:

  1. Every 10 minutes, the router saves its routing table and active connections.
  2. 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
    end

Worked Example: eSewa’s Transaction Logs When you pay a bill via eSewa:

  1. eSewa’s system logs the transaction before deducting money.
  2. If the system crashes mid-transaction, eSewa replays logs to ensure your payment is either fully completed or rolled back.

In the Real World

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

  1. 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.
  2. 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").
  3. 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?").
  4. 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").
  5. 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.

Based on the TU BSc CSIT syllabus for Distributed Networking, unit 6.

Discussion

Loading…