Raft: Understandable Distributed Consensus
thesecretlivesofdata.com
thesecretlivesofdata.com
While the candidate in the smaller split receives votes from a majority of the split, there is no true majority, so no leader. The cluster is configured with the total number of nodes.
What could happen is that an already elected leader continues to think it's the leader for a while while the rest fo the cluster elects a new leader. The split leader will however fail to commit its log, and throw them away once it rejoins.
Another important detail that's missing is that node only votes once pere term, and only for a node that has an equal or higher term than itself. It will never vote twice or vote for an outdated node.
Changing the configuration is in fact handled in a special way at the end of the raft paper in a way that avoids split-brain.
[edit] Oh, the 2-node split was in fact already the leader, so it does exactly what I described. Dur...
It seems to me that homogeneity amongst nodes in a distributed consensus mechanism is of utmost importance; once you implement a hierarchical power structure dishonesty becomes difficult to deal with... especially when "voting" for a "leader" is involved in mapping the node landscape.
edit:spelling
> VR uses a leaderbased approach with many similarities to Raft.
> However, Raft has less mechanism that VR or ZooKeeper because it minimizes the functionality in non-leaders. For example, log entries in Raft flow in only one direction: outward from the leader in AppendEntries RPCs. In VR log entries flow in both directions (leaders can receive log entries during the election process); this results in additional mechanism and complexity
> Raft has fewer message types than any other algorithm for consensus-based log replication that we are aware of. For example, we counted the message types VR and ZooKeeper use for basic consensus and membership changes (excluding log compaction and client interaction, as these are nearly independent of the algorithms). VR and ZooKeeper each define 10 different message types, while Raft has only 4 message types (two RPC requests and their responses).
Honestly, what I think happened is this: They first explained paxos to the poor students, then asked questions. In a later session explained Raft and asked questions. Can't it be the students started processing the problem of distributed consensus between the sessions so they got a better grasp of the topic? This would mean paxos helped them understand raft better. Anyway I'm nitpicking etc etc.
> Each student watched one video, took the corresponding quiz, watched the second video, and took the second quiz. About half of the participants did the Paxos portion first and the other half did the Raft portion first in order to account for both individual differences in performance and experience gained from the first portion of the study. We compared participants’ scores on each quiz to deter- mine whether participants showed a better understanding of Raft.
The reason I'm so picky about this claim is that before you know it, you have mythical pseudo-statistical claims like "some programmers are more than 10 times as good as others" that will live a life of their own. CS has way too many of those.
If you want to dive even deeper "FLP" and "Leslie Lamport" should also open up a can of interesting worms
Is this something the full algorithm handles differently to the way the diagrams would indicate?
1. Client sends log entry to leader
2. Leader appends log entry, forwards it to followers
3. Majority of followers confirm
4. Leader commits the log entry
5. Leader confirms the commit to the client
6. Followers commit on the next heartbeat
What happens if the leader goes away between 5 and 6? To my eyes, it looks like the followers will time out, elect a new leader, and have to roll back the last log entry.
Roughly, a newly elected leader will have all committed entries (guaranteed by the "election restriction", 5.4.2) but it does not know precisely which are committed. The new leader will commit a no-op log entry (section 8) and after it has received replies from a majority of the cluster it will know which entries have already been committed.
ie. P(client queuing up too many requests and crashing because db is too slow waiting for disc) > P(5 distributed db nodes with the transaction in memory crashing simultaneously before it was written to disc)
EDIT: and it's not P(5...), it's P(Leader...) in the case I'm worrying about.
Why can't I read at a normal pace instead of being interrupted all the time and having to wait while the next sentence is shown?
[edit] This behavior would be very suitable if I was making a presentation to an audience with this content - but it's quite contrary to what's needed for the audience to view the content themselves at their own pacing.
It would be even nicer if the "Continue" button had a permanent position (and if I could use enter/space/pagedown/... instead of mouse), but I didn't notice that it was hidden between animations. I guess I am slower than you are. :)
I am not sure if the concept is valid (some other comments have issue with that), but it was well presented. Good job, OP - keep it up!
The presentation originally moved forward on its own but as you get into later sections there's a lot going on visually so it's easy to miss key points in the explanation. There's also an issue that the presentation is attached to the wall clock in later sections so you need the visual in a certain state before moving forward on the explanation.
In hindsight, I agree that there are better ways to present this information. Part of this project is to help me learn how to best communicate complicated topics like distributed consensus in the most effective way. Most existing resources are 20 page PhD dissertations which are not very accessible to beginners so I'm just figuring this out as I go.
I'm changing the format in future presentations. I'm working on infrastructure now to allow D3.js visualizations to be embedded into Medium blog posts. That'll give the best of both worlds -- read at your own pace and interact with mini visualizations as you go along.
I'm always looking for CS topics to visualize and help explain better. If you (or anyone) has any suggestions, please let me know. I'm @benbjohnson on Twitter.
[1] According to http://raftconsensus.github.io.
Contrary some of the folks here, I found the presentation very cool. But that maybe because I'm a slow learner.
https://github.com/benbjohnson/playback.js
looks interesting.
To give an example, say I have n machines in datacenter A, and n*.99 in datacenter B. datacenter A gets destroyed, permanently. Does datacenter B now reject all (EDIT: where reject = not commit) requests until a human comes along to tell it that datacenter A isn't coming back?
Of CAP, you are now choosing CP with Raft. So yes, the system is unavailable until an external agent fixes it. In other words, the system needs to have a majority of nodes online to be "available".
* needs to have the same data as the other nodes
* needs a round of Raft to notify its presence to other nodes
So you can only add new nodes (automatically) when you have a 'live' system.
majority = ceil((2n + 1)/2) : so by getting the number of available nodes in the partition, nodes can figure out if they are in the majority or minority cluster.
See section 6 in the paper for details of its implementation.
Is this used in production somewhere already? Would love to hear more of the details about use cases and deployment.
It's used in etcd, consul, serf and probably more.
Considering it basically relies on random chance (I.E. who receives the message first) to elect a master, has basically no real way of resolving a conflict in election (I.E. if two nodes receive the same amount of votes, we do a re-election ad infinitum) and does not address the situation of two nodes having conflicting sets of data (for instance from network partition).
Considering all that, this protocol doesn't seem very interesting (from a use-case point of view).
The conflicting data under partition is covered I think - the set that doesn't meet quorum won't accept writes and will see it has a lower election term than the other partition when it re-joins.
EDIT:
From a use vase point of view, it simplifies the construction of CP systems. This has directly led to etcd and consul, which would be many times more complex had their authors had to implement paxos.
Both etcd and consul are still young software, but if you take a look at the 'Call me Maybe' series of blog posts it's pretty apparent that there's a massive deficiency in current systems handling of network partitions.
http://aphyr.com/posts/281-call-me-maybe-carly-rae-jepsen-an...
I believe it does address this. Each log entry is either committed or not; an entry can only be committed if it has been replicated to a majority of nodes. Any node that lacks a committed entry cannot be elected master because of the election rules: a node will not vote for another node less complete than itself. Since a committed entry has been replicated to a majority, a node lacking that entry cannot receive a majority of the votes. (Thus the committed log entries will always be the same on all nodes (though some may be behind, and may only have a subset), which is the purpose of the protocol.)
> Considering it basically relies on random chance (I.E. who receives the message first) to elect a master, has basically no real way of resolving a conflict in election
This is mostly true. The PDF slides I link to below recommend that the election timeout be much greater than the broadcast time, the idea being that things should work out in the long run.
Highly recommend the PDF slides here, as they explain it better than I can: https://ramcloud.stanford.edu/~ongaro/userstudy/ — there's also a YouTube talk here: https://www.youtube.com/watch?v=YbZ3zDzDnrw
> Considering all that, this protocol doesn't seem very interesting (from a use-case point of view).
I'd love to hear of alternatives.
What happens if (for instance) a 4 node cluster splits into 2 node clusters (I.E. a network fault between two data centers)- does each cluster choose a leader? how are is "majority" calculated? is the raft protocol unable to handle half of it's nodes being taken down? What happens if two clusters break off, both choose a leader (if it's possible), both gets new writes and then both clusters come back together?
> I'd love to hear of alternatives.
I no of no protocols per se, but for implementations of a master-slave protocol, there's mongo's replica-set algorithm (one notable change is that each node can have a priority).
There are also master-master implementations (such as cassandra's) that require no election, and serve IMO more interesting use-cases.
A Raft cluster must have an odd number of nodes.
> how are is "majority" calculated?
ceil(nodes/2).
> is the raft protocol unable to handle half of it's nodes being taken down? What happens if two clusters break off, both choose a leader (if it's possible), both gets new writes and then both clusters come back together?
They cannot each choose a leader, see above.
what about a 7 to 3 / 3 / 1 split?
Theorem. With 2n + 1 nodes, there can not be two separate majorities after a net split.
Proof. By way of contradiction, assume there are two separate majorities. Each separate majority would contain at least ceil((2n + 1)/2) = n + 1 nodes. This implies that there are in total at least 2(n + 1) = 2n + 2 nodes in the system, contradiction.
Why must a Raft cluster have an odd number of nodes?
> > how are is "majority" calculated?
> ceil(nodes/2).
A majority is defined as having greater than half the votes. I.e., you need nodes / 2 + ((nodes + 1) % 2) votes, or more simply votes > nodes / 2. Even in an even-sized cluster that can only hold true for one node, and not cause splits.
No majority is possible here.
> how are is "majority" calculated?
The definition of majority is greater than half the set (that is the meaning of the word). If you have four members, 3 is the lowest such number that is greater than 4 / 2.
> If I've understood the presentation correctly, Raft is a
> master-slave protocol that determines how to choose a
> master.
You haven't understood the presentation correctly.