Repository navigation
Replies: 5 comments 4 replies
|
Between the previous log index and opaque config version approaches, I personally much prefer the opaque conf versions. It's roughly equivalent to the previous log index approach but strictly more general. Callers could use the log index of the previous membership change as a version and it would work the same. |
|
One potential answer might be, just have the application solve this issue by implementing one of the above ideas in the application. I've thought about that, but still can't figure out how to do it correctly. Basically all of the above approaches boil down to the following two steps:
If I'm understanding correctly, those two steps MUST be done atomically, otherwise there's a race where the change is invalidated before it's appended. A solution to that could be for the application to handle all membership changes sequentially, maybe on a single task/thread. However, that just eliminates the race from happening on a single node, not across nodes. In-between steps (1) and (2), the node could lose leadership, the change becomes invalidates, then the node regains leadership and appends it to the log. I think the only way for (1) and (2) to be done atomically is form within |
|
Both of the problems you raise, unique node IDs and stale membership changes, come from one root cause: a node rejoining the cluster under an ID it already used, after its state was reverted. Unique node IDs. Raft assumes a node never loses state it has already acknowledged. Wiping a node's disk and letting it rejoin under the same ID is a state reversion dressed up as a membership change, and after such a reversion the protocol guarantees nothing: the split brain you describe is one symptom, silent loss of committed logs is another symptom. No code change is needed. I am adding the rule that a node ID must never be reused to the dynamic membership guide and to the FAQ. Stale membership changes. Your step-8 sequence, where a delayed duplicate re-adds a node that then wipes itself, is real, and only openraft itself can prevent that sequence, since validating a membership change and appending it must be one atomic step. Of your four proposals I would take Previous Log Index. Your objection to Previous Log Index, that a caller cannot easily obtain the log index of the previous membership change, does not apply to openraft: Previous Log Index only removes the stranded cluster of your step 8, which is an availability failure. Previous Log Index does not make ID reuse safe: a node that rejoins under a used ID with reverted state still breaks consensus, because a precondition orders membership changes but does not restore the log that the node erased. Allocating node IDs. To obtain an ID that no node has ever held, write a blank entry to the cluster and use the log index returned for that blank entry. The Raft log is the one counter every node already agrees on, and no committed log index is handed out twice. A small project of mine allocates node IDs this way: ezraft If openraft enforces that a node ID never rejoins the membership once it has left, that enforcement also covers your step 8: the delayed add for |
|
It sounds like the previous log index approach is the right way forward. Do you plan on making the change? If not I'd be happy to submit a PR myself. |
Here's an interesting alternative idea that could work after the previous Log Index check is implemented. Applications can use the log index of the previous membership change as both the node ID and as the validation. That would guarantee unique node IDs and avoid writing blank entries. |
Uh oh!
There was an error while loading. Please reload this page.
Hi, my team and I have been building a system using openraft (I've opened some discussions and issues before). Recently we've been implementing dynamic membership and I think there are some improvements that
openraftcan make around this feature. I'm opening this discussion to see what people think.Unique Node IDs
This one is maybe obvious in hindsight, but it is not valid to reuse node IDs. I.e. callers should avoid doing something like
n1.n1.n1.unless the
n1in step (3) has the exact same log contents as then1from step (1). Otherwise it's as ifn1had reverted its log which can lead to split brain situations. For example here's a split brain scenario:{n3}.{n1, n2, n3},n1is the leader.n1removesn3from group making membership{n1, n2}.n3confirms that it has been removed and wipes its disk.n1adds a newn3to group making membership{n1, n2, n3}.n3, starting with an empty disk, replays log up to the point where membership was{n3}.{n1, n2}and{n3}. Each can start committing writes leading to a split brain.I don't actually think any code changes are needed here, but it would be helpful to explicitly call this out somewhere in the documentation. It's maybe implicit in the FAQ when it says not to wipe the data of a node, but it would be helpful if this was explicitly stated in the dynamic membership docs
Rejecting Stale Membership Changes
As it turns out, it's actually a bit difficult to avoid re-adding a removed node. Imagine that I'm a human or computer operator and I want to change the membership of my cluster. Then consider the following series of events:
n, initialize the on-disk state, and start anopenraftengine.n. This message times out, maybe it's lost in the network.n. This message arrives successfully.n.n.n.nto stop itsopenraftengine and delete its on-disk state.nback into the group membership.nis still running, it still has log entries, so nothing blocks it from being promoted from a learner to a voter.n.nstops itsopenraftengine and deletes its on-disk state.We've now entered a really bad scenario.
nis part of the group membership, but it is dead. Its log is deleted and it will not respond to any heartbeats. If the group membership only had two members, the group can no longer achieve a quorum. For larger group membership numbers, you can repeat the above multiple times to also get into a stuck position.How can I prevent this situation? I spent a lot of time trying to figure this out, and I don't think there's a simple solution with the current API. So here are some possible updates to the API to help with this. I'm curious if you agree that this is a problem, if the current API makes this hard/impossible to solve, what you think about the following proposals, and if you have any other ideas?
Previous Log Index
This approach adds an optional field to all membership change requests of the log index of the previous membership change. Callers would include this field when they send a membership change request.
openraftthen rejects the request if log index is not accurate, i.e. there has been a new membership change since the index provided in the field.A drawback of this approach is that it may be difficult for callers to actually get the value of the log index of the previous membership change.
This is the approach that Hashicorp Raft takes:
Opaque Configuration Version
This approach is similar to the "Previous Log Index" approach, but instead of adding a log index to requests, an opaque version number is added to every request. Each membership change is appended to the log with this opaque version number.
openraftthen rejects the request if the version is less than or equal to the previously appended version number.A drawback of this approach is that callers are responsible for generating monotonically increasing version numbers somehow.
This is roughly the approach that TiKV takes:
Removed ID Set
This approach involves storing a set of removed IDs somewhere in the log. Every time a remove node is appended to the log, it updates this set.
openraftthen rejects add node requests if it's a member of the removed set.There are a couple of drawbacks with this approach. First of all the set strictly grows and I don't think it's ever safe to clean it up. Also, callers may have a legitimate reason for re-adding a node, as long as the log hasn't reverted. This prevents them from doing that. For those reasons I don't think it's a good fit for
openraft, but I wanted to mention it anyway.This is the approach
etcdtakes:Deterministic Node IDs
This is an interesting approach, I don't think it's a good fit for
openraft, but it's interesting enough to mention.This approach stores a
next_node_idfield somewhere in the log. Add node requests DO NOT contain any node ID. Whenopenraftreceives an add node request, it assigns it the node ID ofnext_node_idand incrementsnext_node_idin the same log entry. That way all nodes can deterministically derive the node ID from the log.The drawbacks are that this would be a really big breaking change and it would prevent callers from supplying their own node IDs, which they may want to do for some reason.
This is the approach CockroachDB takes:
All reactions