computer science//distributed systems//distributed consensus
Distributed consensus is the problem of making a set of replicas agree exactly on one value, or on the order of a log of values, despite crashes and lost messages, and it is the mechanism that lets three or five servers act as one database that survives losing any one of them. The value might be who owns order 4512, which configuration is current, or which server leads; once it is agreed it is never taken back.
Distributed consensus is the problem of making a set of replicas agree exactly on one value, or on the order of a log of values, despite crashes and lost messages, and it is the mechanism that lets three or five servers act as one database that survives losing any one of them. The value might be who owns order 4512, which configuration is current, or which server leads; once it is agreed it is never taken back.
The core idea is the majority quorum: a value counts as chosen once more than half of the replicas have accepted it. Any two majorities share at least one replica, so two different values can never both be chosen, and a cluster of 2f+12f+12f+1 replicas keeps working with fff of them down (three tolerate one, five tolerate two). Paxos, due to Leslie Lamport, was the first practical algorithm built on it: it never commits two values and makes progress whenever the network behaves well enough for its timeouts to settle. Raft reorganises the same guarantees around an elected leader and is the one most new systems use.
Control consensus averages approximately among neighbours; Raft votes exactly by majority. They share a name and little else.
The consensus protocol of a drone fleet moves continuous values toward each other, tolerates loss and bounded delay, and its small transient error is harmless. Distributed consensus decides a discrete value or an order, tolerates the crash of a minority, and its failure is two truths at once, such as an order delivered twice. Raft to average positions is overkill; averaging to decide who owns an order buys a customer two deliveries.
Two theorems set the limits. FLP impossibility means no algorithm can guarantee progress in a fully asynchronous network, so practical ones keep safety always and progress only with timeouts; the CAP theorem means the minority side of a network partition stops answering.
Consensus is the engine of state machine replication: agree on each log entry in turn and every replica computes the same state.
The failure model sets the size. Crash faults need 2f+12f+12f+1 replicas; nodes that may lie need 3f+13f+13f+1 and different algorithms (Byzantine fault tolerance). Blockchains solve that harder version among parties who do not trust each other (blockchain consensus, and its block-signing quorum).
Every decision costs at least one round trip to a majority, so writes are slower than on a single server, and a cluster spread over distant sites pays that latency on every commit. The consensus layer is for small, critical state (ownership, configuration, leases); bulk data such as telemetry flows without it.