Glossary · System coordination, integration and orchestration
Distributed consensus
Also known as: consensus algorithm, consensus protocol
German: Verteilter Konsens
In distributed systems, distributed consensus is the problem, and the family of algorithms such as Paxos and Raft that solve it, of getting several nodes to agree on a single value or sequence of decisions even when some nodes fail or messages are delayed.
- System integration
In one sentence
Distributed consensus lets several nodes agree on one value or sequence of decisions even when some nodes fail or messages are delayed.
Example
Three edge servers use a Raft-based configuration store; a new line configuration becomes valid only when at least two of the three have accepted it.
How it applies
- Architecture: Consensus underlies many high-availability components such as configuration stores, cluster managers and replicated databases. It typically needs a majority (quorum) of nodes, so three or five nodes are common.
- Engineering: Consensus-based systems stay consistent under failures but may become unavailable when no majority can communicate. Place nodes so that a single network or power fault cannot remove the majority.
- Documentation: Operations documentation should explain quorum requirements, how many nodes may be down for maintenance at the same time and how to recover when quorum is lost. These procedures are rarely used and must be clear when they are.
Distributed consensus vs. leader election
Leader election chooses one node to coordinate; it is one use of consensus. Consensus more generally ensures agreement on any sequence of decisions, for example the order of writes to a replicated log.