← back

Raft and 2PC are essentially the same

Last week I got done with implementing Raft. This week I was studying distribution transactions and concurrency control mechanisms. Particularly, yesterday I was trying to understanding the mechanics of how 2PC (Two Phase Commit) works in depth.

And while doing so, I felt something click in my head that I absolutely have to share here. That Raft and 2PC are essentially the same, only that they serve different purposes.

If you think about 2PC, at the very core, it simply orchestrates the distribution of work but does it in a way such that you get the guarantee that either the entirety of the work is completed or it isn't. Think: you had two pieces of work A and B and you wanted to guarantee that either both commit or neither does.

Now if you think about Raft, it is a consensus algorithm used to achieve agreement about the same value across a bunch of nodes in a cluster, thereby providing fault tolerance to the cluster itself.

The sharp insight here is:

2PC makes two different participants do different pieces of work. But Raft makes every participant do the exact same piece of work.

This insight leads us to two interesting points of discussion in this writeup:


Let's compare 2PC and Raft side-by-side

2PCRaft
SetupThere's a transaction coordinator (TC) and a bunch of Participants (P).There's a Leader (L) and a bunch of Followers (F).
First RoundThe first round (or phase as it's called) is the "PREPARE" phase wherein TC asks the Ps to store a particular txn data.The first round is an AppendEntries RPC containing the new log entries which the L sends to all Fs asking them to store it in their log.
Second RoundThe second round or phase is the "COMMIT" phase in which TC asks Ps to execute the txn they had stored in the previous round.The second round is again the L sending an AppendEntries RPC but with the leaderCommit field incremented, indicating to the Fs that they may now apply the log entries to their state.
Second Round (important)The TC blocks indefinitely until it receives a thumbs up from all Ps. Which means an OK from m nodes in an m node setup.The cluster can make progress as long as the L receives an OK from a quorum (f+1 nodes in a 2f+1 node setup).
PurposeAt the end of it, we basically committed to different pieces of work atomically across distributed nodes. Hence, 2PC is also called as an atomic commit algorithm.At the end of it, we achieved consensus in a distributed system such that the same value is available on multiple machines leading us to unlock fault tolerance for our cluster.

How does Raft exploit the sharp insight to unlock superpowers?

Look at round two for 2PC in the table above. And read it carefully. It blocks indefinitely until all Ps respond back to the TC. Let’s make the difference even easier to understand.

For Raft to make progress, it needs the confirmation of only f+1 nodes in a 2f+1 node setup. For 2PC to make progress, it needs confirmation from 2f+1 nodes in a 2f+1 node setup.

And that happens for a good reason. Here's why: assume a fast P responded with an OK (meaning the txn committed on the P and the P released its lock to make further progress) but a slow P did not respond at all. In that case, the TC cannot ask the fast P to roll back or undo and it also cannot ask the slow P to roll back or undo because fast P has already committed. So it must sit there, retrying indefinitely until the slow P also gives back an OK. What if the slow P isn't slow and that it actually crashed? What if the slow P's network broke down? In distributed systems, it's hard to distinguish such cases.

But in the case of Raft, that's not needed. Since every single node (F) is performing the same piece of work, we can tolerate the failures of some Fs purely because they can catch up later on. In 2PC, every node (P) is doing a different piece of work. How is a P supposed to catch up later on if it was the only node aware of its work?

And so I rest my case.


BONUS: How can we get rid of the indefinite blocking in 2PC?

You can. And it's a very simple yet elegant solution. Put all Ps and TC in their own Raft groups. So even if multiple Ps (or the TC) fail, their Raft groups can perform a leader election and bring up the new P (or the TC). And in this way, the blocking does not remain indefinite as a new node is bound to take the place of the failed one.

The only caveat is that if multiple machines within the Raft group itself fail, 2PC becomes indefinitely blocking again.

            Transaction Coordinator
            +---------------------+
            |     Raft group      |
            |                     |
            |       Leader        |
            |       /    \        |
            |      v      v       |
            |  Replica  Replica   |
            +----------+----------+
                       |
              2PC across Raft groups
                       |
        +--------------+--------------+
        |                             |
        v                             v
Participant 1                     Participant 2
+---------------------+           +---------------------+
|     Raft group      |           |     Raft group      |
|                     |           |                     |
|       Leader        |           |       Leader        |
|       /    \        |           |       /    \        |
|      v      v       |           |      v      v       |
|  Replica  Replica   |           |  Replica  Replica   |
+---------------------+           +---------------------+

This is exactly what Spanner does for their distributed database. If you found this solution interesting, do give that paper a read.