-
Notifications
You must be signed in to change notification settings - Fork 7.3k
[ZOOKEEPER-3642] Fix potential data inconsistency due to DIFF sync after partial SNAP sync #1224
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from all commits
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -1620,11 +1620,139 @@ public void testFaultyMetricsProviderOnConfigure() throws Exception { | |
| assertTrue("complains about metrics provider MetricsProviderLifeCycleException", found); | ||
| } | ||
|
|
||
| /** | ||
| * If learner failed to do SNAP sync with leader before it's writing | ||
| * the snapshot to disk, it's possible that it might have DIFF sync | ||
| * with new leader or itself being elected as a leader. | ||
| * | ||
| * This test is trying to guarantee there is no data inconsistency for | ||
| * this case. | ||
| */ | ||
| @Test | ||
| public void testDiffSyncAfterSnap() throws Exception { | ||
| final int ENSEMBLE_SERVERS = 3; | ||
| MainThread[] mt = new MainThread[ENSEMBLE_SERVERS]; | ||
| ZooKeeper[] zk = new ZooKeeper[ENSEMBLE_SERVERS]; | ||
|
|
||
| try { | ||
| // 1. start a quorum | ||
| final int[] clientPorts = new int[ENSEMBLE_SERVERS]; | ||
| StringBuilder sb = new StringBuilder(); | ||
| String server; | ||
|
|
||
| for (int i = 0; i < ENSEMBLE_SERVERS; i++) { | ||
| clientPorts[i] = PortAssignment.unique(); | ||
| server = "server." + i + "=127.0.0.1:" + PortAssignment.unique() | ||
| + ":" + PortAssignment.unique() | ||
| + ":participant;127.0.0.1:" + clientPorts[i]; | ||
| sb.append(server + "\n"); | ||
| } | ||
| String currentQuorumCfgSection = sb.toString(); | ||
|
|
||
| // start servers | ||
| Context[] contexts = new Context[ENSEMBLE_SERVERS]; | ||
| for (int i = 0; i < ENSEMBLE_SERVERS; i++) { | ||
| final Context context = new Context(); | ||
| contexts[i] = context; | ||
| mt[i] = new MainThread(i, clientPorts[i], currentQuorumCfgSection, false) { | ||
| @Override | ||
| public TestQPMain getTestQPMain() { | ||
| return new CustomizedQPMain(context); | ||
| } | ||
| }; | ||
| mt[i].start(); | ||
| zk[i] = new ZooKeeper("127.0.0.1:" + clientPorts[i], ClientBase.CONNECTION_TIMEOUT, this); | ||
| } | ||
| waitForAll(zk, States.CONNECTED); | ||
| LOG.info("all servers started"); | ||
|
|
||
| final String nodePath = "/testDiffSyncAfterSnap"; | ||
|
|
||
| // 2. find leader and a follower | ||
| int leaderId = -1; | ||
| int followerA = -1; | ||
| for (int i = ENSEMBLE_SERVERS - 1; i >= 0; i--) { | ||
| if (mt[i].main.quorumPeer.leader != null) { | ||
| leaderId = i; | ||
| } else if (followerA == -1) { | ||
| followerA = i; | ||
| } | ||
| } | ||
|
|
||
| // 3. stop follower A | ||
| LOG.info("shutdown follower {}", followerA); | ||
| mt[followerA].shutdown(); | ||
| waitForOne(zk[followerA], States.CONNECTING); | ||
|
|
||
| // 4. issue some traffic | ||
| int index = 0; | ||
| int numOfRequests = 10; | ||
| for (int i = 0; i < numOfRequests; i++) { | ||
| zk[leaderId].create(nodePath + index++, | ||
| new byte[1], Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT); | ||
| } | ||
|
|
||
| CustomQuorumPeer leaderQuorumPeer = (CustomQuorumPeer) mt[leaderId].main.quorumPeer; | ||
|
|
||
| // 5. inject fault to cause the follower exit when received NEWLEADER | ||
| contexts[followerA].newLeaderReceivedCallback = new NewLeaderReceivedCallback() { | ||
| boolean processed = false; | ||
| @Override | ||
| public void process() throws IOException { | ||
| if (processed) { | ||
| return; | ||
| } | ||
| processed = true; | ||
| System.setProperty(LearnerHandler.FORCE_SNAP_SYNC, "false"); | ||
| throw new IOException("read timedout"); | ||
| } | ||
| }; | ||
|
|
||
| // 6. force snap sync once | ||
| LOG.info("force snapshot sync"); | ||
| System.setProperty(LearnerHandler.FORCE_SNAP_SYNC, "true"); | ||
|
|
||
| // 7. start follower A | ||
| mt[followerA].start(); | ||
| waitForOne(zk[followerA], States.CONNECTED); | ||
| LOG.info("verify the nodes are exist in memory"); | ||
| for (int i = 0; i < index; i++) { | ||
| assertNotNull(zk[followerA].exists(nodePath + i, false)); | ||
| } | ||
|
|
||
| // 8. issue another request which will be persisted on disk | ||
| zk[leaderId].create(nodePath + index++, | ||
| new byte[1], Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT); | ||
|
|
||
| // wait some time to let this get written to disk | ||
| Thread.sleep(500); | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. I think in general we advocate not using sleep (at least, not using sleep alone) in unit tests as it's not reliable and proved to be a source of flaky-ness. Is there a better approach to wait for the log being flushed to disk (preferably also verify that fact in test)?
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. That's a reasonable concern, I'll try to find ways to check that in the test here. |
||
|
|
||
| // 9. reload data from disk and make sure it's still consistent | ||
| LOG.info("restarting follower {}", followerA); | ||
| mt[followerA].shutdown(); | ||
| waitForOne(zk[followerA], States.CONNECTING); | ||
| mt[followerA].start(); | ||
| waitForOne(zk[followerA], States.CONNECTED); | ||
|
|
||
| for (int i = 0; i < index; i++) { | ||
| assertNotNull("node " + i + " should exist", zk[followerA].exists(nodePath + i, false)); | ||
| } | ||
|
|
||
| } finally { | ||
| System.clearProperty(LearnerHandler.FORCE_SNAP_SYNC); | ||
| for (int i = 0; i < ENSEMBLE_SERVERS; i++) { | ||
| mt[i].shutdown(); | ||
| zk[i].close(); | ||
| } | ||
| } | ||
| } | ||
|
|
||
| static class Context { | ||
|
|
||
| boolean quitFollowing = false; | ||
| boolean exitWhenAckNewLeader = false; | ||
| NewLeaderAckCallback newLeaderAckCallback = null; | ||
| NewLeaderReceivedCallback newLeaderReceivedCallback = null; | ||
|
|
||
| } | ||
|
|
||
|
|
@@ -1634,6 +1762,10 @@ interface NewLeaderAckCallback { | |
|
|
||
| } | ||
|
|
||
| interface NewLeaderReceivedCallback { | ||
| void process() throws IOException; | ||
| } | ||
|
|
||
| interface StartForwardingListener { | ||
|
|
||
| void start(); | ||
|
|
@@ -1707,6 +1839,14 @@ void writePacket(QuorumPacket pp, boolean flush) throws IOException { | |
| } | ||
| super.writePacket(pp, flush); | ||
| } | ||
|
|
||
| @Override | ||
| void readPacket(QuorumPacket qp) throws IOException { | ||
| super.readPacket(qp); | ||
| if (qp.getType() == Leader.NEWLEADER && context.newLeaderReceivedCallback != null) { | ||
| context.newLeaderReceivedCallback.process(); | ||
| } | ||
| } | ||
| }; | ||
| } | ||
|
|
||
|
|
||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
It's not obvious to me how we could end up in this code path - if we are here then the shutdown must have been called and that specific shut down must not be a full shutdown. Do we have a concrete example when we will hit this code path?
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
In case the follower exited before finish syncing with leader, the ZK server will not start, canShutdown will return false, but we still need to clear the db here.