1 of 16

XI INTERNATIONAL CONFERENCE

“INFORMATION TECHNOLOGY AND IMPLEMENTATION” (IT&I-2024)

Stateful cluster leader failover models and methods based on Replica State Discovery Protocol

Serhii Toliupa, Maksym Kotov, Serhii Buchyk, Juliy Boiko, Serhii Shtanenko

2 of 16

2

Introduction

High availability is a cornerstone of fault tolerance in production clusters. The following work delves into the novel methods and models to achieve rapid cluster leader failover based on the Replica State Discovery Protocol (RSDP).

​

  • Firstly, RSDP is described and evaluated as a quintessential method of achieving consensus within homogenous multiagent distributed systems.
  • Secondly, a new state reducer is developed that allows to perform a synchronized leader election process. Its mathematical model and code implementation written in JavaScript are provided and comply with an established extension interface described within the confines of RSDP’s State Manager component.
  • Lastly, this work delves into the practical implications of the mentioned state reducer in the context of the stateful cluster leader failover. Three different approaches and models based on the novel consensus algorithm to mitigate spontaneous critical events are modeled and assessed. Based on failure probability, failover duration, and communication overhead mathematical models, the said approaches were compared, and recommendations for their application were provided.

3 of 16

3

Local Area Network Simulation based on AMQP

The fundamental idea of AMQP is the separation of client and server, which are called producer and consumer. Instead of direct communication, the message goes through the broker and its queue, thus allowing for alleviating direct dependency between clients and servers. Additionally, AMQP describes the achievement of fundamental communication properties such as resilience, durability, congestion control, security, and the list goes on .

Figure 1: Advanced Message Queuing Protocol’s interactions.

4 of 16

4

Local Area Network Simulation based on AMQP

SLAN operation basis relies on the two main capabilities described within the context of AMQP: “Fanout” and “Direct” exchanges. These abstract interfaces allow to either send a message to all binded queues (e.g., broadcast the message) or send the message to a single binded queue by the routing key (direct communication routing).

Figure 2: Local Area Network Simulation based on AMQP.

5 of 16

5

RSDP phases and consensus process

To provide capabilities of cluster-wide state operations, RSDP describes its lifecycle in a few distinct phases. These phases include “DEBATES”, “SHARE”, and “CLOSE”. Each of the mentioned phases is respectively responsible for the introduction of the new node, state sharing, and final state derivation. The last phase is responsible for handling the shutdown lifecycle event.

Since each distinct phase somehow interferes with the cluster state stored on each replica individually, every instance has a concurrency control mechanism based on the “interPhaseMutex”. That mutex prevents multiple simultaneous cluster events from interfering with each other by restricting stored state access to a single active phase.

Figure 3: Replica State Discover Protocol phases.

6 of 16

6

RSDP formal definition and properties

 

7 of 16

7

Leader Election Reducer for the RSDP

 

8 of 16

8

LER Based on “Latest Source of Truth”

 

9 of 16

9

LER Based on “Popular Vote Decision”

 

The second is the “Popular Vote Decision”. This approach counts every incoming proposed state rather than relying on the latest and could be described as follows:

10 of 16

10

LER Based on “Ranked Vote with Exponential Weights”

 

 

11 of 16

11

Stateful Cluster Failover Models

A stateful cluster in the context of this work is a distributed system where each node has its own subset of the system state. The subset might be either a unique unknown portion for every other cluster member or, as a more common case, a subset of another node’s state. Stateful clusters often establish a leader-follower model to achieve high consistency.

​

While the entire cluster follows a single leader, it becomes a single point of failure. The principles of fault tolerance in that context require establishing a failover mechanism as a contingency. During its passive phase called “monitoring”, the active cluster nodes should be probed and tested to detect any issues promptly. As soon as the critical event on the leader node is detected, the mechanism switches to the active phase of achieving consensus. The entire network has to agree upon a new leader of a cluster to continue its operation.

12 of 16

12

Self-regulated mutual health evaluation

We will first evaluate a model based on a single logical plane. Each node in such a cluster is responsible for operational execution, monitoring, and governance. In such topology, every node must have a communication link with every other in the system to successfully achieve consensus and monitor other instances.

The cluster could be preconfigured to initiate health probes in a specified interval but with different initial timestamp shifts based on a node position in the network. This allows to efficiently utilize the repetitive status probes and decrease failure detection time. Additionally, since every node conducts the monitoring, the network can tolerate to up to N – 1 failed nodes, where N is a node set cardinality.

In that regard, LER provides all the necessary data needed to establish successful monitoring and election in an automated way. The capabilities of LER already provided data for the dynamic node discovery. Hence, every node has a list of the cluster members that they must probe. LER automatically adjusts the states of nodes in a cluster as soon as some subset leaves, but the new leader election cycle could also be triggered in case the nodes suffered critical events.

Figure 4: Self-regulated mutual health evaluation.

13 of 16

13

Centralized observer health monitoring

The topology with a centralized observer includes a set of stateful worker nodes in a cluster and a single coordinating machine that manages the entire network. This topology introduces the division between operational and control planes and thus fosters the single responsibility principle in the system.

This model imposes smaller communication overhead and reduces coordination complexity but suffers from a single point of failure.

Figure 5: Centralized observer health monitoring.

14 of 16

14

Distributed observer health assessment

The following discussed topology is also comprised of two distinct execution plains. In such networks, a consensus algorithm is used between observers themselves and decides on the responsible node that must coordinate the workers plane.

This model provides a balance between communication overhead, complexity, separation of concerns and fault tolerance.

Figure 6: Distributed observer health assessment.

15 of 16

15

Conclusions

Mathematical models of the RSDP and multiple variants of the leader election reducer were developed and described to further simplify integration of the said consensus-achieving protocol into the modern coordination systems. The provided model thoroughly describes properties and processes that constitute the protocol itself, allowing for an effective decision-making process.

Given the models and graphs for the three different cluster failover models and approaches, it is fair to assume that none of those could be called objectively superior in every plain of comparison.

    • Method involving self-regulated mutual health evaluation is most suitable when high additional infrastructure incurrence and the critical event detection time are the most influential metrics towards the successful operation of the system. Though it is worth noting that the communication overhead grows exponentially with the number of cluster members.
    • The second method, based on centralized observer health monitoring, is an appropriate solution only in cases where higher infrastructure and communication overhead costs are a primary decision factor. Since the entire stability of the system depends on a single centralized external observer, the very same observer becomes a single point of failure, which could lead to a disaster when high availability is a hard requirement.
    • Lastly, a method based on distributed observer health assessment serves as a trade-off between high availability, infrastructure cost incurrence, communication overhead, and failover delay by offloading the decision-making process to the parallel distributed layer of coordination. Such an approach is mostly suitable for current cloud infrastructure demands due to its flexibility and clearly established separation of concerns.

16 of 16

16

Thank you for attention!