In Search of an Understandable Consensus Algorithm (2014) [pdf]
raft.github.io
raft.github.io
Paxos: https://www.youtube-nocookie.com/embed/JEpsBg0AO6o
Raft: https://www.youtube-nocookie.com/embed/YbZ3zDzDnrw
Then there's this presentation-animation, too, which I find to be pretty nifty: http://thesecretlivesofdata.com/raft/
And yeah wow, Raft is dead simple.
It is concise and presents a clear view of what makes code easier to work with. I've found its concepts useful when trying to think of what to say during code reviews.
[1] https://arxiv.org/abs/2004.05074 [2] https://news.ycombinator.com/item?id=22994420
Is it perfect? No, as another comment here points out, Raft might not be as fast as some other more complicated Paxos variants but that's okay for most use cases and is much easier to reason about.
To give an example of how powerful this is, I started implementing high availability for an opensource instant search engine I've been working on (https://github.com/typesense/typesense) and was able to get a basic clustering solution working in just a few days.
Pre-raft this would have taken a few weeks or would have required an external coordinator like Zookeeper. This way, Raft has unlocked so many opportunities.
First of all remember that Paxos is a family of protocols for solving consensus. When doing research it's useful to reduce a problem into smaller and smaller parts. The "standard" Paxos algorithm is a very simple consensus algorithm which can only decide a single value once. It's not practical at all, but provides a good framework for thinking about consensus.
When this article proposes "Raft vs Paxos" they are actually comparing Raft against a standard way of configuring Paxos with a leader (MultiPaxos). Note that MultiPaxos allows a lot of nuances in the implementations while still being called "MultiPaxos". MultiPaxos is not a spec you implement; it's a set of ideas.
Raft on the other hand is a concrete protocol with well-defined, specified behavior. In fact, Raft is essentially an implementation of MultiPaxos[1]. This is a very good thing! Paxos provides a framework for thinking about consensus, while Raft puts some of these ideas into a concrete specification which is easy to implement. And it is a good point that we should make the knowledge in the field of consensus available for a wider audience. Yay, Raft is good!
And here comes the problem: A lot of people have read the Raft paper and made the conclusion that "Raft is the best way of solving consensus". Raft is (relatively) easy to implement and get started with and gives you a very simple model to program for (a log of commands), but it's far from a panacea.
The most important thing to know about Raft is that it's not performant (every command has to be sent to a single leader which becomes a bottleneck) nor scalable (every command needs to be processed by all nodes). Etcd supports "1000s of writes" and recommends up to "7 nodes".
This doesn't mean that Raft is bad; it's just a trade-off you need to be aware of. Simplicity vs performance. If you're integrating Raft into your stack and aim for scalability/performance you must always be very weary of when you use it. You should minimize writes at all costs. Unfortunately many developers gets the impression that you can just plug Raft into an existing system and suddenly have a performant and scalable distributed system.
A good example is CockrouchDB: They're using plain Raft for writes, but uses "leader leases" for scaling reads. Suddenly things become a lot more complicated (for instance see this issue about how leader leases are implemented in the Go library for Raft: https://github.com/hashicorp/raft/issues/108). I'm sorry, but you're going to have to get your hands dirty if you want something that's both fast and correct.
The end result is that you have two choices: (1) You can use a library which provides a simple model (a log of commands), but doesn't scale well or (2) you can use a more complicated consensus algorithm and then deal with all of the Hard Problems™ that comes with it. If you're going for the second option, you might as well take advantage of all of the research discovered in the last few years (see https://vadosware.io/post/paxosmon-gotta-concensus-them-all/)
It should also be noted that even though the consensus algorithm doesn't scale, it doesn't mean your system can't scale. Scalog (https://www.usenix.org/system/files/nsdi20-paper-ding.pdf) is an example of a system which uses a consensus algorithm in a constant way (i.e. regardless of load of the system). Once again: Focus on how you can avoiding using a consensus algorithm due to the way your system works.
TL;DR: Hard problems are still hard. Don't think that Raft is a magical piece of software which solves all of your problems.
[1]: There are some nuances between Raft and "standard" MultiPaxos as mentioned in https://arxiv.org/abs/2004.05074. I would still consider Raft to be in the same class as MultiPaxos (compared to other solutions of consensus).
https://github.com/rqlite/rqlite
Your point about "the impression that you can just plug Raft into an existing system and suddenly have a performant and scalable distributed system" is an excellent one. One the biggest misconceptions I come across is that rqlite is distributed for performance when it's actually all about reliability. In fact performance is significantly reduced versus just writing to SQLite itself.
> First of all remember that Paxos is a family of protocols for solving consensus... Raft on the other hand is a concrete protocol with well-defined, specified behavior. In fact, Raft is essentially an implementation of MultiPaxos... You have two choices: (1) You can use a library which provides a simple model (a log of commands), but doesn't scale well or (2) You can use a more complicated consensus algorithm and then deal with all of the Hard Problems™ that comes with it.
At AWS [0], everyone (I spoke to) who worked on distributed consensus had this exact same opinion, so you're not at all off the mark.
> A good example is CockroachDB: They're using plain Raft for writes, but uses "leader leases" for scaling reads.
The Chubby paper by Google [1] goes in to excruciating details of running a production Paxos system.
> Focus on how you can avoiding using a consensus algorithm due to the way your system works.
Amazon SQS may be one such example: I'd presume, it scales by avoiding consensus, in a way, simply maintaining multiple copies [2] and by placing guard-rails around delivery [3][4], ingestion, and duration of storage [5].
[0] https://aws.amazon.com/builders-library/leader-election-in-d...
[1] https://blog.acolyer.org/2015/02/13/the-chubby-lock-service-...
[2] https://docs.aws.amazon.com/AWSSimpleQueueService/latest/SQS...
[3] https://docs.aws.amazon.com/AWSSimpleQueueService/latest/SQS...
We generalize every day because everyone can't be expected to learn everything, so we abstract away, define constraints and best practices and we select our library-fied tools and apply them. We do this with operating systems, file systems, encryption, and also with distributed systems. The basics of distributed systems and their invariants and trade offs are not that hard if you give yourself time to study them properly, but when you want high performance and scale it does get Challenging.
It is important that any abstractions we make, any generalizations, constraint definitions and best practices that we expect people to read, accept and adhere to are as correct as possible without breaking that abstraction. So thank you for your comment. Please note that mine is not a pro-Paxos comment either, just one appreciating good information being spread so that people can make good choices and trade offs.
Distributed systems are hard. To paraphrase: In the data flow, no one can hear your canary msg scream.
https://github.com/etcd-io/etcd/tree/master/raft#raft-librar...
Users:
- cockroachdb A Scalable, Survivable, Strongly-Consistent SQL Database
- dgraph A Scalable, Distributed, Low Latency, High Throughput Graph Database
- etcd A distributed reliable key-value store
- tikv A Distributed transactional key value database powered by Rust and Raft
- swarmkit A toolkit for orchestrating distributed systems at any scale.
https://github.com/etcd-io/etcd/tree/master/raft#notable-use...
Go docs
I hadn't heard about LCR before so thanks for the pointer. It seems to be first described in "The complexity of reliable distributed storage"[2]. I found one reference to it in the Ring-Paxos paper[3] and they explain the disadvantages as follows:
> Of special interest is LCR, which arranges processes along a ring and uses vector clocks for message ordering. This protocol has slightly better throughput than Ring Paxos but exhibits a higher latency, which increases linearly with the number of processes in the ring, and requires perfect failure detection: erroneously suspecting a process to have crashed is not tolerated. Perfect failure detection implies strong synchrony assumptions about processing and message transmission times.
I'll definitely read the paper and try to understand how LCR works!
[1]: https://dl.acm.org/doi/abs/10.1145/3269981
Ring Paxos uses similar ideas (ring overlay for the acceptors), but requires IP-multicast (to minimize bandwidth for broadcasting values to the acceptors and learners), and hence isn't practical in most public cloud environments. (The variant of Ring Paxos that doesn't use IP-multicast doesn't seem to me to have any significant advantages over LCR, while being less symmetric and more complicated.) I do think that Ring Paxos is probably the best approach I've seen to improving the throughput of Paxos while preserving its latency advantages (but it still has issues with tail latency and transient faults because the acceptors communicate via a ring rather than quorums).
Very interesting! I'll definitely have to check this out. Thanks for the references/links.