Paxos
martinfowler.com
martinfowler.com
It's a beautiful algorithm, but Raft has the advantage of being a lot simpler.
+1 - I really admire him too. I wonder, what makes people like him and Edsger Dijkstra so charismatic in this sense?
I usually enjoy Lamport's speeches, and I have said before that his book Specifying Systems should be mandatory reading for nearly anyone interested in CS, due to how approachable it is, and yet how deep it goes.
I wrote a few posts about how to understand paxos, and this is my personal favorite – https://blog.the-pans.com/understanding-paxos/.
Lamport did comment on it (via email) saying that understanding Paxos is less important than proving it works (via TLA+).
I suppose that's not wrong; if I know for a fact that something is going to behave how I expect it too, I don't generally care how it works. I can treat the algorithm as a black box and just go from there.
That said, I dont' completely agree. If I am implementing Paxos in NewLangOfTheMonth, it's likely that a straight one-to-one port of another version will not fully work due to tiny differences in the language semantics. If I don't understand Paxos, then it's unlikely I'll be able to correct any mistakes in the implementation.
Unless you're one of the rare group who gets paid to implement it in x language/library/technology, it's more of a mental exercise than anything, something nice to go "ahhh cool!" when you understand it, then move on.
However, there is a difficulty in expressing exactly what RMW-Paxos does, because it may end up "modifying" a value that was never "committed".
For example, assume that I have three acceptors storing a replicated counter. The proposer executes phase-1, increments the value received from the acceptor with the highest ballot, and sends the value to all acceptors in phase-2. Consider now an execution where all acceptors initially hold 0. The proposer manages to store a new value 1 in one acceptor, and then crashes. Thus, 1 was not committed and a reader may conclude that the consensus value is still 0. Now a new proposer arrives, receives (1, 0, 0) from the three acceptors, chooses 1 as the newest value, and writes (2, 2, 2) to all acceptors. Thus "2" is now committed, and the commit of the "2" retroactively commits the "1". Thus, the condition for being committed is no longer "there exists a phase-2 quorum ...", but it is more complicated and depends on future history.
I have verified with TLA+ that this algorithm does indeed implement an atomic read-modify-write operation, but only under a complicated notion of being "committed" that boils down to either having a phase-2 quorum, or one of the successor operations having a phase-2 quorum.
Did you ever figure out a simpler way to say this?
A note worthy point is that, Paxos actually expects proposers to drop their value and carry-forward someone-else's value to complete the consensus. This is typically not an intended behavior from client's point of view, but proposers role is to form the consensus, not push their value.
Read phase is basically understanding if there is any value that could've been chosen when F nodes are unavailable.
Determine phase is to choose to carry-forward any value (if-any) from the Read phase or use the new value from the client.
Write phase is to actually ask majority of the acceptors to save the value selected in the Determine phase.
However, the proposer can also modify the value proposed by somebody else, and it turns out that this is a perfectly valid algorithm for implementing a distributed read-modify-write. This variant does more than just consensus---it basically behaves like a fault-tolerant memory with an atomic update operation (think compare-and-swap or load-linked/store-conditional). In his blog post, GP also mentions that paxos is a read-modify-write transaction, but perhaps I read too much into what he is saying. Either way, this RMW variant is correct, I have verified it exhaustively with TLA+, and it is used in a few production systems that I have implemented.
The difficulty is now to say exactly what this RMW variant does. I was hoping that GP (or anybody else) may have some insights.
I first learned about this as a project called gryadka and I had trouble believing that extending Paxos from consensus (write once) to a full linearizable register worked until model checking a TLA+ spec for it that I wrote.
[1]: https://martinfowler.com/articles/patterns-of-distributed-sy...
edit: found it
>-You think that’s going to be around in 6 months?
bahaha.
> It may be difficult if you can’t trust anything and the entire concept of happiness is a lie designed by unseen overlords of endless deceptive power.
https://vadosware.io/post/paxosmon-gotta-concensus-them-all/
It’s a bit tricky to explain precisely in an HN comment but basically, something is monotonic if, once you learn that something is true, nothing can disprove it. For example, let’s say you know of a bunch of nodes on a graph, and you’re slowly learning about edges between the nodes. Once you see a cycle, no amount of new edges will ever “disprove” the cycle. So this means that, if what you care about is cycles, then a CRDT is fine. However, let’s say you’re trying to do disturbed garbage collection of nodes that have no edges leading to them. (And as before you’re slowly learning about new edges.) You can see a node has no edges leading to it that you know of, but if you think a node of garbage-collectible you might always be proved wrong by learning about a new edge. So the problem of distributed garbage collection can’t be solved using CRDTs.
There’s a paper on this that you’ll find if you google for the CALM theorem.
In 2016 some academics at Stony Brook University did, along with a machine-checked TLAPS proof, updated in 2019:
https://arxiv.org/abs/1606.01387
¹Paxos consensus refines specs for consensus and voting. See the spec and accompanying material here: https://lamport.azurewebsites.net/tla/paxos-algorithm.html
Otherwise I think it's a pretty good description of Paxos. I like that it walks through a good number of failure scenarios and mentions interesting corner cases such as a lack of an explicit commit in the original papers and that reading also needs to run the full Paxos algorithm.
Paxos also has drawbacks like being slow or lacking liveness guarantees. Some alternatives are discussed in https://aws.amazon.com/builders-library/leader-election-in-d...
The way to do this is to asynconously write to all nodes in real-time, here you can try it out: http://root.rupy.se
Go to http://root.rupy.se/node?make&info and press Make and you'll se the replication in real-time! (although it's too fast so you'll only see it as being delivered in one go, if you wireshark it you'll see multiple HTTP chunks...)
The source is also here: https://github.com/tinspin/rupy in Root.java
It feels kinda strange to see an author other than Martin Fowler blogging on martinfowler.com....
https://martinfowler.com/articles/patterns-of-distributed-sy... - other content by the same author under title Patterns of Distributed Systems.
I used to hate people that preached his ways, but I recognized they're often pretty good. People just take them too much as laws and not opinions.
People have to remember that he's just human like the rest of us.