HA clients quorum¶
This page explains the quorum concept for HA clients in KX Sensors v3, including how quorum checks work using local etcd read requests, and walks through worked examples of quorum determination under various failure scenarios.
Overview¶
Connectivity among cooperating nodes in a computing cluster is key to maintaining data coherence within the cluster. If this connectivity is broken, various classes of problems relating to data consistency can arise. For example, inconsistencies are likely to occur when the interconnections between two or more nodes fail, but the processes on the nodes continue to operate. This might result in some nodes being able to communicate over a functioning part of the network, but not being able to communicate with other nodes in a different part of the network.
Events such as these are referred to as dual primary, split-brain, or network partition conditions. This chapter discusses how KX Sensors handles software failure or network partition situations where it is unclear which nodes are operating correctly and which are not.
The basic goal is to handle the situation where a network partition exists between two or more nodes, one of which is primary for at least one feed. Without some additional mechanism, one node on each side of the partition can potentially assume the role of primary, unbeknown to the other.
Consider a simple three-node system. The HA implementation in KX Sensors chooses a primary through a user-specified preference order and internode messages that provide status and health indicators to peer processes. If the three processes cannot communicate, one of the secondary processes may assume that the primary has failed and become acting primary. Multiple acting primaries would result in inconsistencies. Quorum prevents this and guarantees consistency.
HA clients¶
The concept of quorums applies to HA clients. An HA client is any q process that loads the HA framework to determine which process among the set responsible for a feed should act as the primary process. Examples of service classes that are HA clients in the core implementation include SDLs, MDLs, REPLs, and MQs.
Note
In the code, these are referred to as HACTL clients, not HA clients.
Quorum checks¶
KX Sensors v3 implements just a single quorum policy: every HA client continuously checks whether it is part of the quorum by issuing a read request to the node's local QETCD process. (A single QETCD process runs on every node in the cluster.) QETCD in turn issues a request to the etcd cluster. A successful read means the HA client is part of the quorum.
A quorum check may fail due to the following situations:
- failure of the local QETCD instance (that is, the instance on the same node as the HA client)
- loss of connectivity between the local QETCD instance and the etcd cluster
- failure of the etcd cluster
When quorum is confirmed by an HA client, its quorum membership is retained for a configurable period (.cfgp.haQuorumCheckFreq).
HA failover determination¶
When a process is no longer part of a quorum, it is logically prevented from performing its function. For example, a primary SDL continuously evaluates whether it is in quorum on receipt of messages. If it is not, it does not give up its primary status; however, it stops ingesting data it might receive from its configured data sources.
When a secondary HA client process determines that it should potentially become primary (based on various HA heuristics), it also considers whether it is part of the quorum. If it is, the promotion to acting primary is allowed; if it is not, no failover takes place (but the condition is checked and retried until the situation changes, either through quorum membership adjustment or changes in HA state). Primary status is only claimed by secondary processes; it is not relinquished by the current primary process.
A settling period (.cfgp.haQuorumSettle) is observed by HA clients after quorum acquisition and before any state change decision is made. This period helps ensure that heartbeats from other processes have a chance to arrive before full feed assessment occurs, in the event a network partition has recently healed.
Depending on the quorum and system state, no process may be entitled to become primary for a feed. This represents a disruption of service; however, it is required to maintain the guarantee of consistency.
A transition of primary status from one HA client to another is called a failover (or failback, depending on the direction of the transition based on the failover order). The failover order is determined by the node order as defined in the feed for the set of responsible HA clients. For example, given the feed configuration below (feed.yaml), the HA clients responsible for feed1 are kxsSDL_A1, kxsSDL_B1, and kxsSDL_C1. The order of precedence is kxsSDL_A1 > kxsSDL_B1 > kxsSDL_C1:
feed:
values:
feed1:
id: 100
sc: sdl
nodes: [ A, B, C ]
procs: [ kxsSDL_A1, kxsSDL_B1, kxsSDL_C1 ]
What's new in v3¶
The implementation of quorum is significantly streamlined in KX Sensors v3:
.cfgp.haRequireQuorumno longer exists — quorum is no longer optional; v3 always guarantees consistency.- The Data Arbiter (DA) process was previously involved in quorum-related communication. DA is no longer part of the v3 architecture.
- The majority, index, and size quorum policies no longer exist. They are replaced by a single policy based on communication with the discovery layer (etcd via QETCD).
- The use of a ping device and related settings (
.cfgp.haPingDevice,.cfgp.haPingTO,.cfgp.haPingRetries,.cfgp.haPingLat) has been removed from v3.
Minimum etcd nodes
The mechanism underlying the new quorum policy is based on communication with etcd, which in turn achieves quorum using the Raft algorithm. The Raft algorithm requires a majority of members for quorum. High availability in KX Sensors v3 requires at least three etcd nodes for this reason.
Quorum examples¶
The following illustrations show the quorum determination outcomes under various topologies and failure conditions. Nodes are named A, B, C, and D and are specified in that order for the feed in question; that is, under normal conditions, A takes the highest precedence and is the expected primary. The outcomes assume that all other HA heuristics are normal, except in cases where processes are explicitly marked as down. Some topologies include nodes E, F, and G, which host an etcd cluster separate from the SDLs.