CACS352 Distributed System

Distributed SystemUnit 712 min read

Election Algorithms & Leader Selection: Bully, Ring, Chord, and Practical Scenarios

Unit 7 of Distributed System: Explores how distributed systems elect leaders (coordinators) when nodes fail or join, covering algorithms like Bully, Ring, and Chord, their protocols, trade-offs, and real-world applications in peer-to-peer networks, cloud services, and Nepal’s eSewa/Khalti payment systems.

Why Election Algorithms?

In distributed systems, nodes must agree on a leader (coordinator) for tasks like:

  • Resource allocation (e.g., assigning tasks in a cloud cluster).
  • Conflict resolution (e.g., settling disputes in blockchain).
  • Fault recovery (e.g., restarting a failed service in eSewa’s payment network).

Without a leader, the system risks splits (partitioned consensus) or chaos (no agreement). Election algorithms solve this by:

  1. Detecting failures (e.g., a node crashes).
  2. Triggering elections (e.g., when the current leader fails).
  3. Selecting a new leader (e.g., via majority vote or ring-based rotation).

1. Definitions and Core Concepts

Key Terms

  • Leader (Coordinator): A node elected to manage critical tasks (e.g., commit transactions in distributed databases).
  • Election Trigger: An event that starts a new election (e.g., leader crash, network partition).
  • Quorum: A minimum number of nodes required to elect a leader (e.g., 2f+1 in Byzantine fault tolerance).
  • Logical Time: A way to order events across nodes (e.g., Lamport timestamps for causality).

Why Logical Time Matters

Elections rely on consistent event ordering. Without it:

  • Node A might declare itself leader while Node B is still processing a previous election.
  • Example: In eSewa’s payment system, if two nodes elect leaders simultaneously, transactions could conflict.

2. Physical vs. Logical Clocks

Event A (t=3):Node A sends messageEvent B (t=4):Node B receives messagEvent C (t=5):Node B repliesEvent D (t=6):Node A receives reply
Lamport clock progression with message exchange (t_A = max(t_A, t_B) + 1)

Physical Clock

  • Definition: Real-world time (e.g., wall clock, GPS time).
  • Problem: Nodes’ clocks drift due to hardware differences (e.g., Node 1’s clock runs 10% faster than Node 2).
  • Solution: Clock synchronization protocols (e.g., NTP, PTP) align clocks across nodes.

Logical Clock

  • Definition: A virtual timestamp assigned to events to preserve causality (if A → B, then timestamp(A) < timestamp(B)).
  • Types:
    • Lamport Clock: Incremented on each event (send/receive).
    • Vector Clock: Tracks causality across multiple nodes (used in distributed databases like Cassandra).

Example: Lamport Clock in Action

Consider two nodes, A and B:

  1. A sends a message to B at time t_A = 3.
  2. B receives it and increments its clock: t_B = max(t_B, t_A) + 1 = 4.
  3. B sends a reply at t_B = 5.
  4. A receives the reply and updates: t_A = max(t_A, t_B) + 1 = 6.

Visualization:

sequenceDiagram
    participant A
    participant B
    A->>B: t_A=3 (send)
    B-->>A: t_B=4 (receive)
    B->>A: t_B=5 (reply)
    A-->>B: t_A=6 (receive)

3. Election Algorithms

Pros: Fast, no ring structureCons: Requires node IDs, single point of failureBully AlgorithmPros: Distributed, no ID requirementCons: Slower, higher message overheadRing AlgorithmPros: Scalable, hybrid approachCons: Complex finger tablesChord AlgorithmElection Algorithms

A. Bully Algorithm

How it works:

  1. A node declares itself leader if no response is received from higher-priority nodes within a timeout.
  2. Higher-priority nodes (e.g., Node 3 > Node 2) respond with their status.
  3. If a higher-priority node is alive, it becomes the leader; otherwise, the declaring node wins.

When to use:

  • Systems with hierarchical priorities (e.g., cloud data centers where servers have fixed IDs).
  • Not suitable for dynamic peer-to-peer networks (e.g., BitTorrent).

Worked Example: Bully Election

Nodes: 1, 2, 3, 4 (priority order).

  1. Node 2 crashes. Node 1 declares itself leader.
  2. Node 1 sends "Who is the leader?" to Nodes 3 and 4.
  3. Node 3 responds: "I am the leader."
  4. Node 1 steps down; Node 3 becomes leader.
sequenceDiagram
  participant N1
  participant N3
  participant N4
  N1->>N3: "Who is leader?" (timeout)
  N1->>N4: "Who is leader?"
  N3-->>N1: "I am leader"
  N1-->>N3: "Step down"
  note right of N1: Node 2 crashed
  note right of N3: Highest ID wins

Advantages:

  • Simple to implement.
  • Works well in structured environments.

Disadvantages:

  • Single point of failure: If the highest-priority node fails, the system may split.
  • No fault tolerance: No recovery if multiple nodes fail simultaneously.

B. Ring Algorithm

How it works:

  1. A node detects a failure (e.g., no heartbeat from leader).
  2. It sends an election message clockwise around the ring.
  3. The node with the highest ID (or lowest IP) wins and becomes leader.

When to use:

  • Peer-to-peer networks (e.g., Chord, Pastry).
  • Systems with dynamic membership (nodes join/leave frequently).

Worked Example: Ring Election

Nodes: A, B, C, D (ordered clockwise).

  1. Node A detects leader failure.
  2. A sends election message to B → B → C → D → A.
  3. D has the highest ID, wins, and becomes leader.
ElectionElectionElectionElectionABCD
Clockwise ring topology with election propagation (D wins as highest ID)

Advantages:

  • Decentralized: No single point of failure.
  • Scalable: Works for large networks (e.g., 10,000+ nodes).

Disadvantages:

  • Slower: Election can take O(n) time (where n = nodes).
  • Complexity: Requires consistent ring topology.

C. Chord Algorithm (Hybrid Approach)

How it works:

  • Combines ring-based election with distributed hash tables (DHT).
  • Nodes are placed on a circular ID space (e.g., 0–2^160).
  • Each node maintains a finger table to locate successors quickly.

When to use:

  • Large-scale peer-to-peer systems (e.g., BitTorrent, IPFS).
  • Low-latency lookups (e.g., finding a file in a distributed database).

Example: Chord Lookup

  • Key: 12345 (hex).
  • Node IDs: 0x1000, 0x2000, 0x3000, 0x4000.
  • Successor: 0x2000 (next node in the ring).
  • Finger Table Entry: 0x1000 → 0x2000 (responsible for keys 0x1000–0x2000).
SuccessorSuccessorSuccessorSuccessorSuccessor0x00000x10000x20000x30000x4000
Chord ring with successor links and finger table entry (0x1000 → 0x2000)

Advantages:

  • O(log n) lookup time (vs. O(n) in ring algorithm).
  • Self-healing: Nodes join/leave dynamically.

Disadvantages:

  • Complex implementation (finger tables, hashing).
  • Not suitable for small systems (overhead).

4. Comparison Table: Election Algorithms

Algorithm Best For Leader Selection Fault Tolerance Time Complexity Example Use Case
Bully Structured environments Highest-priority node Low O(1) Cloud data centers
Ring Peer-to-peer networks Clockwise highest ID Medium O(n) BitTorrent, Chord
Chord Large-scale DHTs Hybrid ring + hashing High O(log n) IPFS, Distributed databases

5. Real-World Applications

In the Real World

  1. eSewa/Khalti Payment Systems

    • Idea: Bully Algorithm for electing a primary node to validate transactions.
    • How: If the primary node fails, a higher-priority backup (e.g., a server in Kathmandu) takes over to prevent downtime.
    • Worked Example: During the Dashain festival (high transaction volume), if Node A (primary) crashes, Node B (backup) detects the failure and triggers a Bully election. Node C (highest priority) becomes primary, ensuring no payment is lost.
  2. Daraz Order Processing

    • Idea: Ring Algorithm for assigning order queues.
    • How: Orders are routed to the next available node in a circular fashion. If a node fails, the next node in the ring takes over.
    • Worked Example: During a sale, if Node 1 (handling orders) fails, Node 2 detects the failure and becomes the leader, ensuring no order is delayed.
  3. NTC/Ncell Network Failover

    • Idea: Chord-like DHT for routing calls.
    • How: Mobile towers act as nodes in a distributed system. If one tower fails, the network reassigns calls using a precomputed successor list (like Chord’s finger table).
    • Worked Example: During a landslide in Sindhupalchok, a tower fails. The network uses Chord’s lookup to redirect calls to the next available tower, minimizing call drops.

6. Exam Tip

How to Score Full Marks

  1. Diagrams are Mandatory

    • Always draw sequence diagrams for Bully/Ring elections.
    • Use ring topology graphs for Ring/Chord.
    • Example:
  2. Compare Algorithms

    • Highlight trade-offs (e.g., Bully is fast but centralized; Ring is decentralized but slow).
    • Use the comparison table above to structure your answer.
  3. Real-World Tie-Ins

    • Link algorithms to Nepalese examples (eSewa, Daraz, NTC).
    • Example answer:

      "In eSewa’s payment system, the Bully algorithm ensures that if the primary server fails, a backup server with a higher priority takes over, preventing transaction failures during peak hours like Dashain."

  4. Clock Synchronization

    • If asked about physical/logical clocks, explain:
      • Physical clocks (NTP) align real-time.
      • Logical clocks (Lamport) ensure causality.
    • Example:

      "Without Lamport clocks, Node A might declare itself leader while Node B is still processing a previous election, leading to conflicts in distributed transactions."

  5. Avoid Common Mistakes

    • Don’t confuse Bully with Ring: Bully uses priority; Ring uses ID order.
    • Don’t forget fault tolerance: Always mention how an algorithm handles node failures.
    • Don’t skip examples: Exams love worked traces (e.g., "Show how Node 3 becomes leader in a 4-node Bully election").

Sample Exam Answer (Full Marks)

Question: Explain the Bully algorithm with a suitable diagram.

Answer: The Bully algorithm is a leader election protocol used in distributed systems to select a coordinator when a node fails. It works as follows:

  1. Election Trigger: A node declares itself leader if it detects a failure (e.g., no heartbeat from current leader).
  2. Priority Check: The declaring node sends "Who is the leader?" messages to higher-priority nodes.
  3. Response Handling:
    • If a higher-priority node responds, it becomes the leader.
    • If no response is received within a timeout, the declaring node wins.

Diagram:

sequenceDiagram
  participant N1
  participant N2
  participant N3
  N1->>N2: "Who is leader?" (N2 > N1)
  N1->>N3: "Who is leader?" (N3 > N1)
  alt N3 responds
    N3-->>N1: "I am leader"
    N1-->>N3: "Step down"
  else Timeout
    N1-->>N1: "I win"
  end

Advantages:

  • Simple and efficient for structured systems.
  • Ensures the highest-priority node becomes leader.

Disadvantages:

  • Single point of failure: If the highest-priority node fails, the system may split.
  • No fault tolerance: Requires manual recovery for multiple failures.

Real-World Example: In eSewa’s payment system, the Bully algorithm ensures that if the primary server (Node 3) fails, Node 2 (next highest priority) detects the failure and triggers an election. If Node 2 is also down, Node 1 becomes the leader, preventing transaction delays during peak hours like Dashain.


Final Checklist for Exam

  • Draw sequence diagrams for Bully/Ring elections.
  • Compare Bully vs. Ring vs. Chord in a table.
  • Explain logical clocks with Lamport timestamps.
  • Tie to Nepalese examples (eSewa, Daraz, NTC).
  • Mention fault tolerance and trade-offs.

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

Discussion

Loading…