Agreeing on a Single Value
How single decree Paxos provides a base for building consensus in a distributed system
Introduction
In this series, I want to build a distributed key-value store, like etcd, from first principles. We will be implementing Paxos, which is an algorithm that’s eluded me for a long time, and our language of choice will be Go. In this post, I want to put Paxos into my own words. If you aren’t already familiar with the algorithm, I strongly suggest reading Lamport’s original paper Paxos Made Simple. We are going to talk about single decree Paxos in this post, though as we continue we will implement multi decree Paxos with the full replicated state machine.
Paxos is considered a notoriously difficult algorithm to understand. It is difficult to test, and if you are not careful, you will implement it incorrectly. In fact, the first time you implement it, you will probably screw up and have some difficult debugging to do. The most important thing is to fully understand Lamport’s paper. That’s what will give you the foundation to confidently debug, trace, and test your implementation. In the sections that follow, we will diagram the ideas behind single decree Paxos and illustrate why they work.
A Distributed System with Five Processes
Consider the following non-Byzantine distributed system with five independent processes P1 through P5, running on different machines, at different speeds, at different levels of reliability, and over a network connection with inconsistent reliability.
graph LR
P1((P1))
P2((P2))
P3((P3))
P4((P4))
P5((P5))
P1 ~~~ P2 ~~~ P3 ~~~ P4 ~~~ P5P1..P5 can each propose, accept, learn, and fail
independently.Suppose that this set of processes needs to elect exactly one of them to become leader. The Paxos algorithm proceeds as follows.
Phase 1
Any of these processes is eligible to become leader. Each sends requests to all the others with a prepare request containing their proposal number (in this case, assume the proposal number is (<round>, <process-id>)). Then P1 sends prepare((1, 1)) to all processes (including itself), P2 sends prepare((1, 2)) and so on. This is where the first fork occurs. One of the following three things can happen:
- No quorum is reached (let’s say P1 and P2 respond affirmatively to P1’s prepare request, P3 and P4 respond affirmatively to P3’s prepare request, and P5 responds affirmatively to P5’s prepare request). In this scenario, messages were dropped (if none were dropped,
prepare((1, 5))would trigger an affirmative response to P5 due to (1, 5) being the highest proposal number). - A quorum can respond affirmatively to one of the processes’ prepare requests, say, P1, with a promise that it will not accept a proposal request with a proposal number lower than
(1, 1), and none of them report an already-accepted proposal. - A quorum can respond affirmatively to one of the processes’ requests with a promise that it will not accept a proposal request with a proposal number lower than
(1, 1)and at least one of them reports an already-accepted proposal.
Suppose no quorum is reached. Zoom into P1, supposing that the mechanics of the prepare inside P1 was to kick off five goroutines with a timeout. When the timeout is reached, P1 can assume that its request failed, and repeat the process, incrementing the round number to 2 and kicking off another 5 goroutines, each sending prepare((2, 1)) to the other processes (and itself). Other processes that reach their own timeouts may do the same, continuing the process until a quorum is reached by at least one process, which then allows that one to continue. Livelock is possible, however, randomizing or jittering the timeouts, or using exponential backoff causes the probability for livelock to vanish in practice, and eventually a quorum will be reached.
sequenceDiagram
participant P1
participant P2
participant P3
participant P4
participant P5
P1->>P1: prepare((1, 1))
P1->>P2: prepare((1, 1))
P1->>P3: prepare((1, 1))
P1->>P4: prepare((1, 1))
P1->>P5: prepare((1, 1))
P1-->>P1: promise, nothing accepted
P2-->>P1: promise, nothing accepted
P3-->>P1: reject, promised (1, 3)
Note over P4,P5: dropped or slow
Note over P1: 2 of 5, so retry above (1, 3)sequenceDiagram
participant P1
participant P2
participant P3
participant P4
participant P5
P1->>P1: prepare((2, 1))
P1->>P2: prepare((2, 1))
P1->>P3: prepare((2, 1))
P1->>P4: prepare((2, 1))
P1->>P5: prepare((2, 1))
P1-->>P1: promise, nothing accepted
P2-->>P1: promise, nothing accepted
P3-->>P1: promise, nothing accepted
Note over P4,P5: dropped or slow
Note over P1: quorum, all empty, so v is freesequenceDiagram
participant P1
participant P2
participant P3
participant P4
participant P5
P1->>P1: prepare((2, 1))
P1->>P2: prepare((2, 1))
P1->>P3: prepare((2, 1))
P1->>P4: prepare((2, 1))
P1->>P5: prepare((2, 1))
P1-->>P1: promise, nothing accepted
P2-->>P1: promise, accepted ((1, 3), 3)
P3-->>P1: promise, accepted ((1, 2), 2)
Note over P4,P5: dropped or slow
Note over P1: highest is (1, 3), so v is forced to be 3Phase 2
We suppose now that a quorum has responded affirmatively to P1’s prepare((2, 1)) request. This authorizes P1 to move to phase 2 of Paxos, but does not mean that P2..P5 cannot continue making prepare requests. In fact, in the absence of failures, they will continue. Regardless, now that P1 has a quorum of affirmative responses, it can now send an accept request with the proposal number (2, 1) and a value v. There are two possibilities for the value of v:
- If all responses in the quorum P1 received indicate they have not accepted any proposal, then P1 is free to send whatever value it wants to send in its accept requests (let’s say 1, since this is a leader election and 1 is P1’s process ID).
- If any of the responses in the quorum P1 received indicated that they’ve already accepted a value, then P1 must send the value associated with the highest proposal number in its accept requests.
Assuming a Go implementation again, it kicks off five goroutines to all processes, including itself. Three things can happen at this stage:
- A quorum responds affirmatively to P1’s accept requests. At this point the value is locked in.
- The process is preempted. Consider that while P1 was sending its accept requests, P3 was able to to complete phase 1 and achieve a quorum, allowing it to proceed to phase 2. Then P1’s accepts will time out as the processes that P3 has received affirmative prepare responses from will ignore them and P1 will have to retry phase 1 with a higher proposal number, essentially scrapping the progress it has made and starting from scratch (if no responses to any of its prepares indicated an already-accepted proposal) or resending the highest-numbered already-accepted value v (if any of the responses indicated an already-accepted proposal). Note that P1 (or some other process) can preempt P3’s progress if it completes phase 1 before P3 is able to complete phase 2. This is called dueling proposers and is a problem that surfaces when using Paxos where there is not yet a distinguished proposer (AKA leader). A solution to this is jittered/randomized timeouts or exponential backoff.
- Not enough responses are received, forcing P1 to retry.
sequenceDiagram
participant P1
participant P2
participant P3
participant P4
participant P5
P1->>P1: accept((2, 1), 1)
P1->>P2: accept((2, 1), 1)
P1->>P3: accept((2, 1), 1)
P1->>P4: accept((2, 1), 1)
P1->>P5: accept((2, 1), 1)
P1-->>P1: accepted
P2-->>P1: accepted
P3-->>P1: accepted
Note over P4,P5: dropped or slow
Note over P1: 3 of 5 accepted (2, 1), so value 1 is chosensequenceDiagram
participant P1
participant P2
participant P3
participant P4
participant P5
Note over P3,P5: P3 completed phase 1 at (3, 3) with these three promises
P1->>P1: accept((2, 1), 1)
P1->>P2: accept((2, 1), 1)
P1->>P3: accept((2, 1), 1)
P1->>P4: accept((2, 1), 1)
P1->>P5: accept((2, 1), 1)
P1-->>P1: accepted
P2-->>P1: accepted
Note over P3,P5: ignored, promised (3, 3)
Note over P1: 2 of 5, so retry phase 1 above (3, 3)((2, 1), 1), and the quorum P3, P4, P5 contains neither.sequenceDiagram
participant P1
participant P2
participant P3
participant P4
participant P5
P1->>P1: accept((2, 1), 1)
P1->>P2: accept((2, 1), 1)
P1->>P3: accept((2, 1), 1)
P1->>P4: accept((2, 1), 1)
P1->>P5: accept((2, 1), 1)
P1-->>P1: accepted
P2-->>P1: accepted
Note over P3,P4: dropped or slow
Note over P5: dropped or slow
Note over P1: 2 of 5 — retry phase 1 at (3, 1)Learning the Value
After the value is chosen, the processes all need a way to learn it. There are multiple approaches, but since we are only discussing single decree Paxos, we will present the most straightforward one. Assume that one of the processes has its value accepted, meaning a quorum responded affirmatively to its accept(n, v) messages. The other processes do not know a value has been chosen, so they continue to run Paxos. Zoom in on P4: it runs phase 1 until it receives a quorum of responses, and since a value has already been accepted, at least one of those responses indicates as much, forcing P4 to use that value in its accept broadcast. When it receives a quorum of accept responses, it knows that the value is locked in. Eventually every process learns the value this way.
sequenceDiagram
participant P1
participant P2
participant P3
participant P4
participant P5
Note over P4: wants to know the chosen value
P4->>P4: prepare((3, 4))
P4->>P1: prepare((3, 4))
P4->>P2: prepare((3, 4))
P4->>P3: prepare((3, 4))
P4->>P5: prepare((3, 4))
P4-->>P4: promise, nothing accepted
P1-->>P4: promise, accepted ((2, 1), 1)
P2-->>P4: promise, accepted ((2, 1), 1)
Note over P3: dropped or slow
Note over P5: dropped or slow
Note over P4: quorum reached, highest is ((2, 1), 1): propose 1
P4->>P4: accept((3, 4), 1)
P4->>P1: accept((3, 4), 1)
P4->>P2: accept((3, 4), 1)
P4->>P3: accept((3, 4), 1)
P4->>P5: accept((3, 4), 1)
P4-->>P4: accepted
P1-->>P4: accepted
P2-->>P4: accepted
Note over P3: dropped or slow
Note over P5: dropped or slow
Note over P4: majority accepted: value is 1In the next post, we will extend the ideas of single decree Paxos to multi decree Paxos and cover how a replicated state machine can work in a fault tolerant way.
Thank you for reading.