Paxos Made Moderately Complex
paxos.systems
paxos.systems
For a proposal i with value x:
1. Ask everyone not to participate in proposals older than i, get a majority to agree and send back any previously accepted values.
2. Ask everyone to accept proposal i with the newest previously accepted value or x if none, get a majority to accept.
- On failure, restart with i > max(previous proposals)
That is the Paxos algorithm.
Its fault-tolerance comes from the fact that you can tolerate losing the minority (1 out of 3 / 2 out of 5). There's no magic. If the majority fails, Paxos blocks. If two proposers keep blocking each other's proposals by going through round 1, Paxos can take forever.
The goal of Paxos is to agree on 1 value. Once consensus is reached (majority agrees), the value can never be changed, since no one can complete round 1 without getting back the newest accepted value from at least one node. However, you can run Paxos again and again to agree on a sequence of values (e.g. transitions in a state machine). This is Multi-Paxos.
The Paxos algorithm is optimal, but Multi-Paxos allows for a lot of optimizations and there is no clearly defined implementation. A pretty typical optimization is to elect a leader (through Paxos) which makes the proposals. This avoids duelling proposers and ensures you have a node that is always aware of the latest accepted value. Making this optimization requires you to make decisions that fall outside the scope of Paxos, such as how long a node remains leader, what happens if the leader fails, and what happens if two nodes think they are the leader. These optimizations can easily cause you to lose the fault-tolerance and consistency guarantees that Paxos gives you within the scope of reaching consensus on a single value.
It's important to remember that Paxos itself is easy to understand and a lot of the problems in Multi-Paxos can be solved by running Paxos. You need to familiarize yourself with it as though it was a programming construct.
I highly recommend that anyone interesting in learning Paxos study this pseudo code from MIT's distributed systems class; it was the single thing that helped me understand Paxos most:
Implementing the protocol in an industrial setting brings up lots of ... interesting design and engineering issues for sure!
Has anyone used both? What are the pros and cons?
Fundamentally, both protocols (multi-decree, multi-paxos, and Raft are equivalent (in performance, and capability)). There are other variants of Paxos (ePaxos) that have advantages to Paxos in some cases, especially in the WAN.
Script: nodes {1, 2, 3*} where 2 and 3 are partitioned, and 3 is the current leader.
Node 2 fires its leader election timer, broadcasts RequestVotes to 1 and 3, only 1 gets it, but now 3 is ignored for the term of its leadership and 2 is the only node capable of quorum. In a decent implementation 3 will be nack'd the next time it sends an AppendEntries to 1, and 1 will pass along the current term. 3 isn't hearing from the leader so broadcasts RequestVotes (either failing once due to failure to hear the current term from 1, or jumping straight into the next term, which actually increases the ratio of cluster livelock). Leadership bounces back and forth rapidly, making your cluster worthless.
The current hotness is spec paxos, which gets the positive trade-offs of both cheap paxos and fast paxos, if you're able to actually implement it in your DC (has networking assumptions).
Besides that, I genuinely wonder how often does this happen in the real world? I see a typical case where parts of your cluster is behind a NAT. Otherwise, what could cause an asymmetric partition?
That's not such an uncommon situation. E.g. quite a few large to mid-sized companies with servers in three locations does not have the in-house skillset to either get BGP set up and/or set up VPNs between the locations and a means to update routes automatically on outages. Some larger shops will have "metro LAN" type setups with separate ports for dedicated connections between racks in different data centres, and won't even have capacity enough on their public ports to handle a failover of traffic normally going on the private connection between two of the data centres over the public port - the expectation is often that the data centre operator will handle redundancy... (until they don't...)
It's not a terribly hard thing to fix, but it's also something it seems few people think about until they reach much larger size.
Sounds interesting, is "spec paxos" like Spanner in requiring special hardware? If you could share some links or references to more information about it, that would be great, thanks.
If you're interested in Raft, take a look at Kontiki (https://github.com/NicolasT/kontiki). Could use some maintenance though...
https://github.com/scalien/scaliendb/tree/master/src/Framewo...
If you just want to agree on a leader, there's PaxosLease:
https://github.com/scalien/scaliendb/tree/master/src/Framewo...
2. We got foobared by potential investors, we focused on one group but they eventually walked away, and by that time we ran out of money. The other side of this was long enterprise deal lifecycles. Trying to convince somebody who is big enough to pay for DB software to use your alpha thing is 6-12 months, that's the same timescale you run out of money.
3. We did have some bugs, so we were claiming consistency/reliability, but it was an alpha product with bugs. (Of course it was, it was just us, 2 guys writing it, we didn't even have money to buy testing infrastructure!)
In the end we quit after ~3.5 years, after we exhaused all options, at great personal (but not monetary) cost, eg. I got divorced shortly after. Fortunately we were able to exit gracefully from all business affairs.
Consensus can be used for SMR, atomic broadcast, leader election. One can argue all of these replicate a sort of state machine. However, that dilutes the meaning of state machine replication -- replication of state with durability and performance in mind
Jokes aside, there are two ways out: huge simulators or writing very very simple code that matches 1-to-1 the steps of the proof. And Lamport's writing style for proofs [1] lends itself to simple matching implementations.
[1] "How to Write a 21st Century Proof" http://research.microsoft.com/en-us/um/people/lamport/pubs/p...
Last time on HN: https://news.ycombinator.com/item?id=8806835