Agreeing on Multiple Values
Extending single decree Paxos in the most transparent way possible
In the last post we covered single decree Paxos and how it enables nodes in a distributed system to agree on a single value. Multi decree Paxos allows us to chain this primitive together to build stateful distributed systems. In this post we will naively extend single decree Paxos to multi decree Paxos.
A Replicated State Machine
A deterministic state machine can be represented by a 4-tuple (S, C, s0 , δ) where S is a set of states, C is a set of commands, s0 is the initial state, and δ is a transition function S × C → S that maps a state and a command to a next state. Determinism means that δ is a pure function. Extend this idea to command sequences σ ∈ C*. Let ε represent an empty sequence, c :: σ represent the cons operation, and δ* : S × C* → S be the extension of δ to sequences, mapping a state and a command sequence to the state reached by applying each command in order. We can express this like so:
Consider a set of nodes in a network that all start with the same initial state and observe the same set of commands, which are applied in the same order. Then each node reproduces an identical final state, without needing to communicate with other nodes to stay consistent. Then the problem of replicating state across a network collapses to the problem of agreeing on a sequence of commands.
Naively Extending Single Decree Paxos
We already know that single decree Paxos allows all nodes in a system to safely agree one value. Consider how we might extend this to agreeing on a sequence of values. Suppose we have a distributed key-value database with an unspecified (but ≥ 2) number of nodes that have no leader, that each implement a replicated state machine, and let the possible commands that can be sent to this system be get k, set k v, and delete k. Each node implements the state machine and has the same initial state. The nodes keep track of the commands sent in a durable log. Such a log might look something like this:
set "leslie" "lamport"set "grace" "hopper"get "grace"set "ronald" "rivest"get "ronald"delete "ronald"get "ronald"Suppose that in our distributed system, clients send requests to some unspecified load balancer which then routes a request to exactly one node in our system. Let N1 and N2 be two nodes in the system, let r1 (carrying get "leslie") and r2 (carrying delete "leslie") be two requests that arrive at around the same time, and suppose r1 is routed to N1 while r2 is routed to N2. When N1 receives r1, it proposes get "leslie" for the next slot in its log (in this case 8). N2 does the same with delete "leslie". There are three things that can happen at this junction:
- Slot 8 might already be filled by a different node in the network by the time N1 and N2 begin executing Paxos. In this case, both N1 and N2 will discover this value in phase one of Paxos, and after phase two they will know for certain that a different value was chosen for slot 8. They will then add that command as slot 8 in their log, and will attempt to run Paxos for the next slot with the command they received in their respective requests.
- Slot 8 is empty, and N1 wins (i.e. receives a quorum of acceptances for slot 8). Then all nodes (including N2) will learn that the command in slot 8 must be
get "leslie". N1 will actually go fulfill the request, returning"lamport"to the client. N2 will retry with slot 9. It may or may not succeed. In a real system, the request would probably be bounded by a timeout, so if N2 doesn’t get a slot for its command before the timeout, it would likely return an error to the client. - Slot 8 is empty, and N2 wins. Then slot 8 for all nodes will eventually be known to be
delete "leslie". N1 would learn this fact, and retry its command for the next slot. Assume it wins. Then sincedelete "leslie"was chosen beforeget "leslie", N1 would return something isomorphic toNot Foundto the client, since the key would’ve been deleted before the request to get that key’s value would execute.
While what we have outlined above works, there are a few notable things that we can do to improve and optimize the algorithm in a real implementation, including but not limited to amortizing phase one with a leader, and log compaction. We will cover both of these and more in the next post, where we will extend our naive and unoptimized specification of multi Paxos into something more usable.
Thank you for reading.