Cluster State & Leader Election: How the Minister of Magic is Chosen
Imagine the Wizarding World without a Minister of Magic. Who would maintain the official records of who belongs in which Hogwarts house? Who would enforce rules on what spells are forbidden? Who would coordinate the aurors to detect dark magic outbreaks?
Without a central authority, the ministry would dissolve into chaos. A cluster of OpenSearch nodes faces the exact same challenge. It requires a leader—the cluster-manager—to govern, issue decrees, and maintain the single source of truth: the Cluster State.
In this edition, we explore how OpenSearch elects its leader and propagates state changes, mapping the complex term-based consensus protocol to the inner workings of the Ministry of Magic.
The Ministry Mapping
- Cluster State: The Ministry's Official Records containing index metadata and shard locations.
- Leader Election: The democratic process of choosing the Minister of Magic.
- Leader (Master Node): The active Minister who dictates updates to the records.
- Followers (Data Nodes): Ministry employees who record and follow the Minister's decrees.
Interactive Magic Consensus Diagram
Consensus Cycle in Print: Node A starts an election term, requests votes from Followers (B & C). Once the votes return and a quorum is achieved, A is crowned leader (Minister) and issues state heartbeat decrees.
The Rajya Patra of OpenSearch: Cluster State
The Cluster State is an immutable snapshot of everything the cluster knows about itself. It contains lists of active nodes, details about index settings, the routing table of which shards reside on which nodes, and the current election term.
In our analogy, it is the Ministry’s official registry. Every time a new index is created or a node falls offline, the record must be updated. However, to prevent conflicting records, only the elected leader can modify this registry.
// From ClusterState.java
public class ClusterState {
private final long version; // Edition number
private final String stateUUID; // Royal seal
private final RoutingTable routingTable;
private final DiscoveryNodes nodes;
private final Metadata metadata; // Laws & Treaties
}
The Election Protocol: Step by Step
OpenSearch uses a term-based consensus protocol similar to Raft. An election is triggered when the active leader goes offline, prompting nodes to run for office.
Phase 1: Pre-Vote (Seeking Allies)
Before a candidate announces their candidacy, they send messengers (probes) to other nodes asking, “Do you agree the leader is gone?” This prevents partitioned nodes from disrupting the cluster by unnecessarily incrementing the term.
Phase 2: Start Election (Optimistic Claim)
If a majority agrees, the candidate increments the term (Yuga) and formally declares: “I am running for Minister!” The candidate votes for himself first and persists this decision to disk to survive unexpected crashes.
// Coordinator.java --- startElection()
private void startElection() {
synchronized (mutex) {
if (mode == Mode.CANDIDATE) {
final long newTerm = getCurrentTerm() + 1;
persistedState.setCurrentTerm(newTerm);
persistedState.setVoteFor(getLocalNode());
// Broadcast candidacy
final StartJoinRequest request = new StartJoinRequest(
getLocalNode(), newTerm
);
sendToAllNodes(request);
}
}
}
Phase 3: Casting Ballots
Each voting-eligible node can only vote once per term. Upon receiving the candidacy request, if the candidate’s term is higher than the node’s current term, the node pledges its vote and responds with a Join request.
Phase 4: Coronation (Rajyabhishek)
Once the candidate receives votes from a quorum (majority) of voting nodes, they become the leader. The trumpets sound: a new leader is crowned, and they immediately publish a fresh Cluster State declaring their reign.
The Quorum Math
A quorum ensures that at most one leader can be elected in a term, preventing the split-brain scenario. The formula is:
Quorum = (N / 2) + 1
For a 5-node cluster, a candidate needs at least 3 votes. Since any two majorities must overlap, it is mathematically impossible for two candidates to gather 3 votes each simultaneously.
Decree Propagation: Two-Phase Commit
When the Minister issues a new law (Cluster State change), they use a two-phase commit protocol to publish it:
- Publish Phase: The leader sends the new state (or the diff, to save network bandwidth) to all nodes. The nodes write the update but do not apply it yet; they send back an acknowledgment.
- Commit Phase: The leader waits for a quorum of acknowledgments. Once met, the leader sends a commit message. Only then do the nodes apply the changes and update their active memory.
If the leader falls mid-publish, the uncommitted state is safely discarded, avoiding partial updates and data corruption.
Failure Detection: The Aurors
To maintain stability, the cluster runs continuous bidirectional health checks:
- LeaderChecker: Followers ping the leader periodically. If the leader fails to respond after 3 attempts, the followers declare “the king is dead!” and trigger a new election.
- FollowersChecker: The leader pings all followers. If a follower stops responding, the leader removes them from the cluster state, updating the registry to reflect their absence.
Classifieds: The Quorum Balancer Game
A distributed consensus network is forming. Help elect a leader by adjusting the sliders. Identify the minimum number of nodes required for a majority (Quorum) for your selected cluster size.