From b0608e0092b3a8630c340ccb8bc82296ff752329 Mon Sep 17 00:00:00 2001 From: Kezhu Wang Date: Sat, 2 Apr 2022 13:01:18 +0800 Subject: [PATCH 1/2] ZOOKEEPER-3023: Sync and commit diff log entries before NEWLEADER ack ZOOKEEPER-2678 could skip snapshot in diff sync, but diff txns are logged and committed after NEWLEADER ack. ZOOKEEPER-3911 moves txn logging before NEWLEADER ack, but the txn logging is asynchronous. So it is indeterminate whether diff txns have been persisted to disk or not after NEWLEADER ack. This commit try to sync and commit txn logs synchronously before ack to NEWLEADER thus provides strong guarantee that follower is in sync with leader after NEWLEADER ack received. This behavior is consistent with pre ZOOKEEPER-2678 and easy to test. --- .../apache/zookeeper/server/TxnLogEntry.java | 6 ++++ .../quorum/FollowerZooKeeperServer.java | 16 ++++++++++ .../zookeeper/server/quorum/Leader.java | 2 +- .../zookeeper/server/quorum/Learner.java | 31 ++++++++++++++++--- .../zookeeper/server/quorum/Zab1_0Test.java | 17 ++-------- 5 files changed, 51 insertions(+), 21 deletions(-) diff --git a/zookeeper-server/src/main/java/org/apache/zookeeper/server/TxnLogEntry.java b/zookeeper-server/src/main/java/org/apache/zookeeper/server/TxnLogEntry.java index 352eb81da90..d792948bead 100644 --- a/zookeeper-server/src/main/java/org/apache/zookeeper/server/TxnLogEntry.java +++ b/zookeeper-server/src/main/java/org/apache/zookeeper/server/TxnLogEntry.java @@ -47,4 +47,10 @@ public TxnHeader getHeader() { public TxnDigest getDigest() { return digest; } + + public Request toRequest() { + Request request = new Request(header.getClientId(), header.getCxid(), header.getType(), header, txn, header.getZxid()); + request.setTxnDigest(digest); + return request; + } } diff --git a/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/FollowerZooKeeperServer.java b/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/FollowerZooKeeperServer.java index 8d371ae5790..ae1bc303f37 100644 --- a/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/FollowerZooKeeperServer.java +++ b/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/FollowerZooKeeperServer.java @@ -19,6 +19,7 @@ package org.apache.zookeeper.server.quorum; import java.io.IOException; +import java.util.List; import java.util.concurrent.ConcurrentLinkedQueue; import java.util.concurrent.LinkedBlockingQueue; import javax.management.JMException; @@ -31,6 +32,7 @@ import org.apache.zookeeper.server.RequestProcessor; import org.apache.zookeeper.server.ServerMetrics; import org.apache.zookeeper.server.SyncRequestProcessor; +import org.apache.zookeeper.server.TxnLogEntry; import org.apache.zookeeper.server.ZKDatabase; import org.apache.zookeeper.server.persistence.FileTxnSnapLog; import org.apache.zookeeper.txn.TxnDigest; @@ -88,6 +90,20 @@ public void logRequest(TxnHeader hdr, Record txn, TxnDigest digest) { syncProcessor.processRequest(request); } + public void syncAndCommitInitialLogEntries(List logEntries) throws IOException { + State state = this.state; + if (state != State.INITIAL) { + String msg = String.format("illegal state %s to sync initial log entries", state); + throw new IllegalStateException(msg); + } + for (TxnLogEntry logEntry : logEntries) { + Request request = logEntry.toRequest(); + getZKDatabase().append(request); + processTxn(logEntry.getHeader(), logEntry.getTxn()); + } + getZKDatabase().commit(); + } + /** * When a COMMIT message is received, eventually this method is called, * which matches up the zxid from the COMMIT with (hopefully) the head of diff --git a/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/Leader.java b/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/Leader.java index ce8f7999c45..fc27f030e0a 100644 --- a/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/Leader.java +++ b/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/Leader.java @@ -361,7 +361,7 @@ Optional createServerSocket(InetSocketAddress address, boolean por /** * This message type is sent by the leader to indicate that the follower is - * now uptodate andt can start responding to clients. + * now uptodate and can start responding to clients. */ static final int UPTODATE = 12; diff --git a/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/Learner.java b/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/Learner.java index b6eeb758ac9..e1ec5eb287f 100644 --- a/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/Learner.java +++ b/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/Learner.java @@ -30,7 +30,9 @@ import java.net.Socket; import java.nio.ByteBuffer; import java.util.ArrayDeque; +import java.util.ArrayList; import java.util.Deque; +import java.util.List; import java.util.Map; import java.util.Map.Entry; import java.util.Set; @@ -82,6 +84,9 @@ static class PacketInFlight { Record rec; TxnDigest digest; + TxnLogEntry toLogEntry() { + return new TxnLogEntry(rec, hdr, digest); + } } QuorumPeer self; @@ -754,12 +759,28 @@ protected void syncWithLeader(long newLeaderZxid) throws Exception { sock.setSoTimeout(self.tickTime * self.syncLimit); self.setSyncMode(QuorumPeer.SyncMode.NONE); zk.startupWithoutServing(); - if (zk instanceof FollowerZooKeeperServer) { - FollowerZooKeeperServer fzk = (FollowerZooKeeperServer) zk; - for (PacketInFlight p : packetsNotCommitted) { - fzk.logRequest(p.hdr, p.rec, p.digest); + if (zk instanceof FollowerZooKeeperServer && !packetsCommitted.isEmpty()) { + List entries = new ArrayList<>(packetsCommitted.size()); + // Pop log entries from packetsNotCommitted according to packetsCommitted. + // In case of mismatch, log warning and keep packetsNotCommitted untouched. + while (!packetsCommitted.isEmpty()) { + long zxid = packetsCommitted.removeFirst(); + pif = packetsNotCommitted.peekFirst(); + if (pif == null) { + LOG.warn("Committing 0x{}, but got no proposal", Long.toHexString(zxid)); + continue; + } else if (pif.hdr.getZxid() != zxid) { + LOG.warn( + "Committing 0x{}, but next proposal is 0x{}", + Long.toHexString(zxid), + Long.toHexString(pif.hdr.getZxid())); + continue; + } + packetsNotCommitted.removeFirst(); + entries.add(pif.toLogEntry()); } - packetsNotCommitted.clear(); + FollowerZooKeeperServer fzk = (FollowerZooKeeperServer) zk; + fzk.syncAndCommitInitialLogEntries(entries); } writePacket(new QuorumPacket(Leader.ACK, newLeaderZxid, null, null), true); diff --git a/zookeeper-server/src/test/java/org/apache/zookeeper/server/quorum/Zab1_0Test.java b/zookeeper-server/src/test/java/org/apache/zookeeper/server/quorum/Zab1_0Test.java index 3bdbcd908dc..12aea6061b6 100644 --- a/zookeeper-server/src/test/java/org/apache/zookeeper/server/quorum/Zab1_0Test.java +++ b/zookeeper-server/src/test/java/org/apache/zookeeper/server/quorum/Zab1_0Test.java @@ -743,8 +743,8 @@ public void converseWithFollower(InputArchive ia, OutputArchive oa, Follower f) readPacketSkippingPing(ia, qp); assertEquals(Leader.ACKEPOCH, qp.getType()); - assertEquals(0, qp.getZxid()); - assertEquals(ZxidUtils.makeZxid(0, 0), ByteBuffer.wrap(qp.getData()).getInt()); + assertEquals(ZxidUtils.makeZxid(0, 0), qp.getZxid()); + assertEquals(0, ByteBuffer.wrap(qp.getData()).getInt()); assertEquals(1, f.self.getAcceptedEpoch()); assertEquals(0, f.self.getCurrentEpoch()); @@ -779,24 +779,11 @@ public void converseWithFollower(InputArchive ia, OutputArchive oa, Follower f) assertEquals(1, f.self.getAcceptedEpoch()); assertEquals(1, f.self.getCurrentEpoch()); - //Wait for the transactions to be written out. The thread that writes them out - // does not send anything back when it is done. - long start = System.currentTimeMillis(); - while (createSessionZxid != f.fzk.getLastProcessedZxid() - && (System.currentTimeMillis() - start) < 50) { - Thread.sleep(1); - } - assertEquals(createSessionZxid, f.fzk.getLastProcessedZxid()); // Make sure the data was recorded in the filesystem ok ZKDatabase zkDb2 = new ZKDatabase(new FileTxnSnapLog(logDir, snapDir)); - start = System.currentTimeMillis(); zkDb2.loadDataBase(); - while (zkDb2.getSessionWithTimeOuts().isEmpty() && (System.currentTimeMillis() - start) < 50) { - Thread.sleep(1); - zkDb2.loadDataBase(); - } LOG.info("zkdb2 sessions:{}", zkDb2.getSessions()); LOG.info("zkdb2 with timeouts:{}", zkDb2.getSessionWithTimeOuts()); assertNotNull(zkDb2.getSessionWithTimeOuts().get(4L)); From aa7e625ceff16b3c83729b20c290f922e5ef1057 Mon Sep 17 00:00:00 2001 From: Kezhu Wang Date: Thu, 4 May 2023 22:58:57 +0800 Subject: [PATCH 2/2] fixup! ZOOKEEPER-3023: Sync and commit diff log entries before NEWLEADER ack --- .../org/apache/zookeeper/server/quorum/Learner.java | 12 ++++++++---- 1 file changed, 8 insertions(+), 4 deletions(-) diff --git a/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/Learner.java b/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/Learner.java index e1ec5eb287f..bb29b697b7f 100644 --- a/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/Learner.java +++ b/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/Learner.java @@ -755,10 +755,6 @@ protected void syncWithLeader(long newLeaderZxid) throws Exception { //Anything after this needs to go to the transaction log, not applied directly in memory isPreZAB1_0 = false; - // ZOOKEEPER-3911: make sure sync the uncommitted logs before commit them (ACK NEWLEADER). - sock.setSoTimeout(self.tickTime * self.syncLimit); - self.setSyncMode(QuorumPeer.SyncMode.NONE); - zk.startupWithoutServing(); if (zk instanceof FollowerZooKeeperServer && !packetsCommitted.isEmpty()) { List entries = new ArrayList<>(packetsCommitted.size()); // Pop log entries from packetsNotCommitted according to packetsCommitted. @@ -783,6 +779,14 @@ protected void syncWithLeader(long newLeaderZxid) throws Exception { fzk.syncAndCommitInitialLogEntries(entries); } + // We almost complete the synchronization phase, all that's left is UPTODATE + // which is a client serving valve and interleaved with broadcast phase. + // + // We are ready for broadcast phase on our behalf now except serving client requests. + sock.setSoTimeout(self.tickTime * self.syncLimit); + self.setSyncMode(QuorumPeer.SyncMode.NONE); + zk.startupWithoutServing(); + writePacket(new QuorumPacket(Leader.ACK, newLeaderZxid, null, null), true); break; }