In search of a simple consensus algorithm
rystsov.info
rystsov.info
> When it is applied to the algorithms, it means that an algorithm with the shortest implementation is simpler.
That is misleading. Kolmogorov complexity is the length of the shortest program (in a pre-defined language) that produces a given object. So, if the shortest program that produces Algorithm A is smaller than the shortest program that produces Algorithm B, then the Algorithm A is less Kolmogorov-complex ("simpler") than Algorithm A.
This does not mean you can take two existing implementations (in C, say) and compare the implementation length and declare one is "simpler," unless you are claiming that both implementations are as short as possible. Since Kolmogorov complexity is not computable, that seems like a tall order.
Maybe they are right that Single-Decree Paxos is simpler (either in the sense of Kolmogorov complexity or in some other sense, who knows), but invoking Kolmogorov complexity here seems totally unwarranted -- it doesn't add anything substantive.
The problem isn't just the incomputable quality of Kolmogorov complexity but that fact that Kolmogorov complexity applies only to finite strings or things that can be meaningfully mapped to them. Especially, Kolmogorov doesn't apply directly to abstract algorithms or programs with multiple implementations.
Additionally, if you limit your language to a recursive language, you can compute complexity (which is no longer Komnogorov) directly. Simply begin enumerating all programs in order of length and then run them checking the output to see if it is that string. While this is by no means efficient, it works for recursive functions since they must halt. For Komogorov complexity, there isn't really a notion of inputs, merely that some particular string should be produced as output.
For recursively enumerable (i.e. Turing complete) languages, the former method will not work because a program might run forever.
Edit: but if someone gives you the primitive recursive Kolmogorov complexity of a program, you can check it by running all shorter programs on all inputs until you found a counterexample for each of them. So it is semi-decidable.
Edit to edit: This would even work for the general Turing-machine definition of Kolmogorov complexity.
This may be true, but:
If you use primitive recursive functions (instead of turing complete) because of practicality, with the same reasoning you can cap the inputs at some insanely large number. Then, these functions still "can compute anything you'd ever want to run on large inputs".
In that setting, equivalence is decidable, because you can simply run both functions over the finite set of all possible inputs.
However, this really gets amortized in most workloads if the leader changes only rarely. Additionally, in an environment with a good network connection between nodes (a few ms), you can set the timeout to be much less than a few seconds (could be less than a second actually), this way you have shorter unavailability.
There's another point he touches, about the unavailability of the whole cluster when the leader is down. Really, this isn't something dependent on the protocols, but on the applications. If you have one paxos/raft group per replica, you actually only get a small unavailability. Additionally, even consistent reads do not need a living master to be possible.
It's worth reading the spanner paper to get more insight into high availability/consistency achieved with a paxos based db implementation ( https://static.googleusercontent.com/media/research.google.c... ).
EDIT: And in my opionion, calling it a SPOF is missing the point a little bit.
It's wrong. If you're fine with stale reads then you don't a living master, but if you want to have a guarantee that the read value is up-to-date then the living master is necessary.
Etcd has a bug related to this - https://github.com/coreos/etcd/issues/741
(Not talking about etcd, check out the spanner paper)
In case there are off the database communications then this level of consistency isn't enough.
It isn't a small unavailability, it's an unavailability of the whole replica. If you don't have a lot of data then it's the whole cluster :)
Of cause, we can introduce something like virtual replicas and eventually end up with a replica per key which is almost a Single Decree Paxos but with an overhead on log compaction and snapshotting.
Btw, thanks for answering my comments actively, I really appreciate it!
https://arxiv.org/pdf/1608.06696.pdf
The gist of it is that there are realistically 2 quorums to be had: One for read/write to a cluster, and one for leader election. Assuming leader election is a relatively uncommon event (as it is in Paxos), you can require more quorum nodes for leader election while allowing less quorum nodes for quorum read/write. This means you get strong consistency with better performance.
Distributed systems are one of my favorite things. I literally spent part of an evening with a beer reading a small handful of papers like this one. If you want an absolute treasure trove of similar works, visit the blog of Adrian Coyler:
He has a good writeup on the Flexible Paxos paper as well:
https://blog.acolyer.org/2016/09/27/flexible-paxos-quorum-in...
Just one? :)
Consensus algorithms are _fascinating_! And +1 for https://blog.acolyer.org/ , I'll often read his blog over lunch.
Another great paper is https://research.google.com/archive/paxos_made_live.html which talks about what it's like to realistically run Paxos. There's a ton of corner cases, engineering realities, and optimizations to be done to make Paxos work well in practice.
Another good resource is the Raft GitHub page[1] which links to the paper, has an interactive visualization, and a plethora of talks by various people.
Raft is the backbone of opensource, I'd be curious to hear from any Googlers in the know whether there's deficiencies in Raft that lead to continued use of Paxos, or if it's experience (and already battle-tested code) with Paxos that leads them to continue deploying Paxos-backed systems.
Chubby paper: 2006 http://dl.acm.org/citation.cfm?id=1298487
Raft paper: 2014 http://dl.acm.org/citation.cfm?id=2643666
If you've got something battletested, there isn't a lot of value to rip it all up if it works within the given business requirements. Also they figured out how to implement Multi-Paxos, which is known for being difficult.
Distributed systems learning ftw
If I squint at the scheme proposed here, it looks like ePaxos taken to the extreme of a single key/value per paxos group, so conflict tracking becomes trivial. It has demonstrated good performance in the case of uncontended writes; I'd be curious to see how it behaves under contention, or when a downed node rejoins the cluster.
Yeah, that was my understanding as well but without having read the paper in question (as of yet) I'm entirely short on the particular technical details of the trade-offs.
Gryadka does 1 roundtrip to write a value and its performance (4720 rps, 1.68ms latency) is very similar to Etcd (5227 rps, 1.55ms).
If you can't afford stale reads then you should ask Etcd to wait for confirmation from the majority of followers before acknowledging a read with ?quorum=true. It makes Etcd do 1 round trip for reads (just like Gryadka).
The same is applicable to other products (you can't relay on time in distributed systems, unless you're Google, consequently you can't relay on read leases)
So there is no performance penalty.
You can elect a leader for a set period of time (say, 10 seconds) and serve strong reads for a lesser period of time (5 seconds) if you have reasonable assumptions of how good your local oscillators work and avoid jumps. If you don't trust your local clock to any level of accuracy, why do you trust your local CPU?
Cockroach DB does this correctly. https://github.com/cockroachdb/cockroach/blob/master/docs/de...
An instance may be running in virtual environment where time freezes are possible. A human may make an error and rollback time to 1970.
It's impossible to eliminate all these factors so yes I don't trust time but I trust CPU.
Maybe CockroachDB is doing it correctly but the terrible default settings make this optimisation negligible because when the leader dies the system hangs for 12 seconds.
Hibari does have a master orchestrator, similar to GFS master server. But it's only needed in reconfiguration events.
In Raft, it's possible for multiple nodes to prevent leader election progress by being overly aggressive when requesting votes. It is also possible for a follower to knock out a perfectly healthy leader by being too quick time out the leader and start a new term.
Both of these limitations stem from the simplicity of Raft's leader election algorithm. To compensate, most Raft implementations I've seen have more conservative follower timeouts that extend the time to detect leader failure and elect a new one.
It's possible for a more optimized algorithm to get sub-second latencies for detecting and re-electing a leader, even in a latent (e.g. geo-distributed) environment. In other words, well within the commit window for the replica set based on network hop latencies.
Also, while the latency for individual writes in single-decree paxos can be closer to strong leader protocols, it is non-trivial to achieve the same level of throughput that is possible in Raft et al when writing to an ordered log, as you cannot start a paxos instance for a new log entry until all prior instances have been resolved. Raft can just add new values in the next append call (or just spam out appends for new messages w/o waiting for the replies for previous ones).
IME, I'd say both single-decree paxos and raft are probably equivalent in terms of understandability, but raft is a better base on which to build a fast high-throughput consensus protocol.
It's wrong. In Grydka (Single-decree Paxos) all the keys are independent so it's possible to update them at the same time without blocking.
Grydka's throughput is comparable to Etcd on the same type of machines (4720 vs 5227 rps) and I never optimized for it (my goal was to fit 500 lines) so it's also wrong that it's "non-trivial to achieve the same level of throughput" - I did it by accident.
So I don't understand why Raft is a better base to build a fast high-throughput consensus protocol.
1. There are tasks which don't require atomic multi-key updates
2. Atomic multi-key updates can be implemented on the client side (see RAMP, Percolator transactions or the Saga pattern)
3. Once the data overgrow the size of one machine you need to shard the log and at this time you're in the same situation
My main point was that the deficiencies of strong-leader-based consensus protocols are overstated, and despite a (minor IMO) level of additional starting complexity, a raft-like protocol is going to be quite a bit simpler than a paxos-based protocol of equivalent capability.
1. I am technically pretty strong but I have no idea what this paper is about
2. So many people know this is about that it shot up to #1 on HN
Can someone give a pointer (a link or two) to the lay, interested audience here about what the field IS. Just a sort of intro guide to someone who knows about programming and math, but has never heard the term paxos?
I am curious, and I am sure many others are as well.
edit: spelling
"[Paxos] provides a new way of implementing the state-machine approach to the design of distributed systems."
Original Paxos paper (written with somewhat whimsical style):
https://www.microsoft.com/en-us/research/publication/part-ti...
Follow up paper to explain Paxos more simply:
https://www.microsoft.com/en-us/research/publication/paxos-m...
See also: https://en.wikipedia.org/wiki/Paxos_(computer_science)
The basic idea is how to coordinate multiple independent agents over an unreliable network. For example, multiple servers trying to manage a shared database. Lots of HN people know about this because it's a key building block of reliable and scalable cloud computing.
Raft is a consensus algorithm that is touted as being easy to understand. It is well specified, compared to paxos, which leaves many implementation details up to the creator. Raft, however, is fairly explicitly specified on page 4 of this paper: https://raft.github.io/raft.pdf
Consensus algorithms allow us to build reliable castles out of constantly failing servers (sand). Systems like zookeeper, etcd, chubby, consul, and others use these algorithms to achieve high availability AND strong consistency (linearizability, the strongest possible) despite up to a majority of the cluster failing.
This is a great introductory course about it in my opinion: https://www.coursera.org/learn/cloud-computing
Then, it's just reading papers like, which you mentioned, "The part time parliment" (paxos)
That is a key insight. I often wonder if people who implemented Raft and discarded Paxos right off the bat knew this? Also I think "Paxos Made Live" scared everyone away from Paxos for a long time.
But what is often missing is that Google implemented a distributed log system. Paxos doesn't do a distributed log by default and just deals with reaching consensus on a value. In practice there is often a need for a log, but not always. If a distributed log is not need Paxos becomes less scary.
One advantage not mentioned in the article, of leader based consensus algorithms is the ability to more easily implement read leases for faster reads.
Read leases can allow for fresh reads without having to run them through the quorum protocol, by trading off availability (due to leader failure) and also correctness in certain edge cases (since read leases will depend on ability for individual machines to measure time delta with reasonable accuracy, which may not be true on some weird VM scenarios).
Here's my blog post on the issue: http://hack.systems/2017/03/13/pocdb/
The entire implementation is 1100 lines of code including comments.
This blog post we're discussing agrees: "I planned to finish it in a couple of days, but the whole endeavor lasted a couple of months, the first version has consistency issues, so I had to mock the network and to introduce fault injections to catch the bugs. [...] Single Decree Paxos seems simpler than Raft and Multi-Paxos but remains complex enough to spend months of weekends in chasing consistency."
Your approach of using many small Paxos instances is very interesting, though! I'd love to see some real world performance comparisons.
My main projects using Paxos are Replicant[1] and Consus[2], both of which have consumed significantly more time to get Paxos correct.
[1] https://github.com/rescrv/replicant [2] http://consus.io/
It's worth reading his description of the paper to understand why: http://lamport.azurewebsites.net/pubs/pubs.html#lamport-paxo...
In short, not everyone has the same sense of humor.
Paxos made simple is much easier in comparison.
One thing I think is worth touching on is the pipelining ability of multi-paxos that is missing when you have a single register with the Synod protocol. For key-value operations, it's not a problem. For a true replicated state machine, this can hinder performance.
In Replicant[1] I added pipelining to ensure Paxos was unlikely to be a bottleneck to replicated state machines. For complex state machines, the state machine becomes CPU bound.
What do you mean by pipelining?
With single-decree Paxos you can write a storage which provides the following API:
function changeQuery(key, change, query) {
var value=this.db.get(key); value = change(value);
this.db.commit(key, value);
return query(value);
}
By providing different implementations of change & query you can achieve different beviour including CAS.
Folk from Elastifile demonstrated it in their Bizur paper - https://arxiv.org/abs/1702.04242v1