From 6616250ce19d7177123ac8d4eda0a6ac83e81e02 Mon Sep 17 00:00:00 2001 From: Vladimir Ivic Date: Mon, 9 Dec 2019 19:22:41 -0800 Subject: [PATCH] ZOOKEEPER-3594: Ability to skip proposing requests with error transactions --- .../org/apache/zookeeper/server/ExitCode.java | 5 +- .../server/PrepRequestProcessor.java | 14 ++ .../org/apache/zookeeper/server/Request.java | 58 ++++++ .../zookeeper/server/ServerMetrics.java | 14 ++ .../zookeeper/server/ZooKeeperServer.java | 8 +- .../zookeeper/server/quorum/Follower.java | 12 ++ .../quorum/FollowerZooKeeperServer.java | 15 ++ .../zookeeper/server/quorum/Leader.java | 14 +- .../server/quorum/LeaderHandler.java | 85 +++++++++ .../server/quorum/LeaderRequestProcessor.java | 2 +- .../server/quorum/LeaderZooKeeperServer.java | 41 ++++- .../server/quorum/LearnerHandler.java | 60 ++++++- .../zookeeper/server/quorum/Observer.java | 16 +- .../server/quorum/ObserverMaster.java | 6 + .../quorum/ObserverZooKeeperServer.java | 10 ++ .../server/quorum/QuorumZooKeeperServer.java | 4 + .../server/quorum/SkipRequestHandler.java | 38 ++++ .../server/quorum/SkipRequestProcessor.java | 66 +++++++ .../server/quorum/SkippedRequestQueue.java | 78 ++++++++ .../server/PrepRequestProcessorTest.java | 45 +++++ .../quorum/QuorumPeerSkipErrorsTest.java | 169 ++++++++++++++++++ .../server/quorum/SkipRequestHandlerTest.java | 110 ++++++++++++ .../quorum/SkipRequestProcessorTest.java | 103 +++++++++++ .../zookeeper/test/MultiOperationTest.java | 48 +++++ .../zookeeper/test/ObserverMasterTest.java | 90 +++++++++- 25 files changed, 1100 insertions(+), 11 deletions(-) create mode 100644 zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/LeaderHandler.java create mode 100644 zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/SkipRequestHandler.java create mode 100644 zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/SkipRequestProcessor.java create mode 100644 zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/SkippedRequestQueue.java create mode 100644 zookeeper-server/src/test/java/org/apache/zookeeper/server/quorum/QuorumPeerSkipErrorsTest.java create mode 100644 zookeeper-server/src/test/java/org/apache/zookeeper/server/quorum/SkipRequestHandlerTest.java create mode 100644 zookeeper-server/src/test/java/org/apache/zookeeper/server/quorum/SkipRequestProcessorTest.java diff --git a/zookeeper-server/src/main/java/org/apache/zookeeper/server/ExitCode.java b/zookeeper-server/src/main/java/org/apache/zookeeper/server/ExitCode.java index 67af2c8df8f..8c379db2ab4 100644 --- a/zookeeper-server/src/main/java/org/apache/zookeeper/server/ExitCode.java +++ b/zookeeper-server/src/main/java/org/apache/zookeeper/server/ExitCode.java @@ -48,7 +48,10 @@ public enum ExitCode { QUORUM_PACKET_ERROR(13), /** Unable to bind to the quorum (election) port after multiple retry */ - UNABLE_TO_BIND_QUORUM_PORT(14); + UNABLE_TO_BIND_QUORUM_PORT(14), + + /** Used to detect if skip txn ran onto the ordering issue */ + QUORUM_PEER_SKIP_OUT_OF_ORDER(15); private final int value; diff --git a/zookeeper-server/src/main/java/org/apache/zookeeper/server/PrepRequestProcessor.java b/zookeeper-server/src/main/java/org/apache/zookeeper/server/PrepRequestProcessor.java index 70d989a34ac..c9801688160 100644 --- a/zookeeper-server/src/main/java/org/apache/zookeeper/server/PrepRequestProcessor.java +++ b/zookeeper-server/src/main/java/org/apache/zookeeper/server/PrepRequestProcessor.java @@ -844,8 +844,22 @@ protected void pRequest(Request request) throws RequestProcessorException { request.getHdr().setType(OpCode.error); request.setTxn(new ErrorTxn(Code.MARSHALLINGERROR.intValue())); } + } finally { + // When skipping error requests is enabled then we want to use next zxid only when we + // have a valid request that passed all the validations in this request processor + // otherwise we will skip the current request + if (LeaderZooKeeperServer.isSkipTxnEnabled() && request.hasError() && request.canSkip()) { + request.setSkipped(); + } else { + // When the skip feature is disabled, we will issue a new zxid for every write + // request that goes through this request processor and passes the validations + if (request.getHdr() != null) { + zks.incrementZxid(); + } + } } request.zxid = zks.getZxid(); + request.getHdr().setZxid(zks.getZxid()); ServerMetrics.getMetrics().PREP_PROCESS_TIME.add(Time.currentElapsedTime() - request.prepStartTime); nextProcessor.processRequest(request); } diff --git a/zookeeper-server/src/main/java/org/apache/zookeeper/server/Request.java b/zookeeper-server/src/main/java/org/apache/zookeeper/server/Request.java index 63cc30b7eb6..76685311aa8 100644 --- a/zookeeper-server/src/main/java/org/apache/zookeeper/server/Request.java +++ b/zookeeper-server/src/main/java/org/apache/zookeeper/server/Request.java @@ -27,6 +27,8 @@ import org.apache.zookeeper.data.Id; import org.apache.zookeeper.metrics.Summary; import org.apache.zookeeper.metrics.SummarySet; +import org.apache.zookeeper.server.quorum.LearnerHandler; +import org.apache.zookeeper.server.quorum.SkipRequestHandler; import org.apache.zookeeper.server.quorum.flexible.QuorumVerifier; import org.apache.zookeeper.server.util.AuthUtil; import org.apache.zookeeper.txn.TxnHeader; @@ -48,6 +50,8 @@ public class Request { // associated session timeout. Disabled by default. private static volatile boolean staleLatencyCheck = Boolean.parseBoolean(System.getProperty("zookeeper.request_stale_latency_check", "false")); + private SkipRequestHandler skipRequestHandler; + public Request(ServerCnxn cnxn, long sessionId, int xid, int type, ByteBuffer bb, List authInfo) { this.cnxn = cnxn; this.sessionId = sessionId; @@ -87,6 +91,8 @@ public Request(long sessionId, int xid, int type, TxnHeader hdr, Record txn, lon public final List authInfo; + private boolean isSkipped = false; + public final long createTime = Time.currentElapsedTime(); public long prepQueueStartTime = -1; @@ -105,6 +111,14 @@ public Request(long sessionId, int xid, int type, TxnHeader hdr, Record txn, lon public QuorumVerifier qv = null; + public boolean isFromLeader() { + return owner == ServerCnxn.me; + } + + public boolean isFromFollower() { + return owner instanceof LearnerHandler; + } + /** * If this is a create or close request for a local-only session. */ @@ -462,4 +476,48 @@ public String getUsers() { } return users.toString(); } + + /** + * This will check if the request's handler supports skip feature, + * the information that is passed from Learner to LearnerMaster using protocol version info + * + * See org.apache.zookeeper.server.quorum.QuorumZooKeeperServer.PROTOCOL_VERSION + * + */ + public boolean canSkip() { + return skipRequestHandler != null; + } + + /** + * Will mark this request for request skipping + */ + public void setSkipped() { + isSkipped = true; + } + + /** + * Once we mark request as skipped this method will be returning true + */ + public boolean isSkipped() { + return isSkipped; + } + + /** + * Each request has access to its handler so know who is this requests owner + * and interact with it when we need to + */ + public void setSkipRequestHandler(SkipRequestHandler skipRequestHandler) { + this.skipRequestHandler = skipRequestHandler; + } + + public SkipRequestHandler getSkipRequestHandler() { + return skipRequestHandler; + } + + /** + * Check the result of the PrepRequestProcessor validations for write requests. + */ + public boolean hasError() { + return hdr != null && e != null; + } } diff --git a/zookeeper-server/src/main/java/org/apache/zookeeper/server/ServerMetrics.java b/zookeeper-server/src/main/java/org/apache/zookeeper/server/ServerMetrics.java index 5cd4abe3b9e..5afe989fca8 100644 --- a/zookeeper-server/src/main/java/org/apache/zookeeper/server/ServerMetrics.java +++ b/zookeeper-server/src/main/java/org/apache/zookeeper/server/ServerMetrics.java @@ -151,6 +151,9 @@ private ServerMetrics(MetricsProvider metricsProvider) { WRITES_QUEUED_IN_COMMIT_PROCESSOR = metricsContext.getSummary("write_commit_proc_req_queued", DetailLevel.BASIC); COMMITS_QUEUED_IN_COMMIT_PROCESSOR = metricsContext.getSummary("commit_commit_proc_req_queued", DetailLevel.BASIC); COMMITS_QUEUED = metricsContext.getCounter("request_commit_queued"); + QUORUM_PEER_SKIP_SENT = metricsContext.getCounter("quorum_peer_skip_sent"); + QUORUM_PEER_SKIP_OUT_OF_ORDER = metricsContext.getCounter("quorum_peer_skip_out_of_order"); + QUORUM_PEER_SKIP_QUEUE_SIZE = metricsContext.getCounter("quorum_peer_skip_queue_size"); READS_ISSUED_IN_COMMIT_PROC = metricsContext.getSummary("read_commit_proc_issued", DetailLevel.BASIC); WRITES_ISSUED_IN_COMMIT_PROC = metricsContext.getSummary("write_commit_proc_issued", DetailLevel.BASIC); @@ -183,6 +186,7 @@ private ServerMetrics(MetricsProvider metricsProvider) { */ OM_PROPOSAL_PROCESS_TIME = metricsContext.getSummary("om_proposal_process_time_ms", DetailLevel.ADVANCED); OM_COMMIT_PROCESS_TIME = metricsContext.getSummary("om_commit_process_time_ms", DetailLevel.ADVANCED); + OM_SKIP_PROCESS_TIME = metricsContext.getSummary("om_skip_process_time_ms", DetailLevel.ADVANCED); /** * Time spent by the final processor. This is tracked in the commit processor. @@ -193,8 +197,10 @@ private ServerMetrics(MetricsProvider metricsProvider) { PROPOSAL_LATENCY = metricsContext.getSummary("proposal_latency", DetailLevel.ADVANCED); PROPOSAL_ACK_CREATION_LATENCY = metricsContext.getSummary("proposal_ack_creation_latency", DetailLevel.ADVANCED); COMMIT_PROPAGATION_LATENCY = metricsContext.getSummary("commit_propagation_latency", DetailLevel.ADVANCED); + SKIP_PROPAGATION_LATENCY = metricsContext.getSummary("skip_propagation_latency", DetailLevel.ADVANCED); LEARNER_PROPOSAL_RECEIVED_COUNT = metricsContext.getCounter("learner_proposal_received_count"); LEARNER_COMMIT_RECEIVED_COUNT = metricsContext.getCounter("learner_commit_received_count"); + LEARNER_SKIP_RECEIVED_COUNT = metricsContext.getCounter("learner_skip_received_count"); /** * Learner handler quorum packet metrics. @@ -218,6 +224,7 @@ private ServerMetrics(MetricsProvider metricsProvider) { QUORUM_ACK_LATENCY = metricsContext.getSummary("quorum_ack_latency", DetailLevel.ADVANCED); ACK_LATENCY = metricsContext.getSummarySet("ack_latency", DetailLevel.ADVANCED); PROPOSAL_COUNT = metricsContext.getCounter("proposal_count"); + SKIP_COUNT = metricsContext.getCounter("skip_count"); QUIT_LEADING_DUE_TO_DISLOYAL_VOTER = metricsContext.getCounter("quit_leading_due_to_disloyal_voter"); STALE_REQUESTS = metricsContext.getCounter("stale_requests"); @@ -304,8 +311,10 @@ private ServerMetrics(MetricsProvider metricsProvider) { public final Summary PROPOSAL_LATENCY; public final Summary PROPOSAL_ACK_CREATION_LATENCY; public final Summary COMMIT_PROPAGATION_LATENCY; + public final Summary SKIP_PROPAGATION_LATENCY; public final Counter LEARNER_PROPOSAL_RECEIVED_COUNT; public final Counter LEARNER_COMMIT_RECEIVED_COUNT; + public final Counter LEARNER_SKIP_RECEIVED_COUNT; public final Summary STARTUP_TXNS_LOADED; public final Summary STARTUP_TXNS_LOAD_TIME; @@ -323,6 +332,7 @@ private ServerMetrics(MetricsProvider metricsProvider) { public final Summary QUORUM_ACK_LATENCY; public final SummarySet ACK_LATENCY; public final Counter PROPOSAL_COUNT; + public final Counter SKIP_COUNT; public final Counter QUIT_LEADING_DUE_TO_DISLOYAL_VOTER; /** @@ -376,6 +386,9 @@ private ServerMetrics(MetricsProvider metricsProvider) { public final Summary WRITES_QUEUED_IN_COMMIT_PROCESSOR; public final Summary COMMITS_QUEUED_IN_COMMIT_PROCESSOR; public final Counter COMMITS_QUEUED; + public final Counter QUORUM_PEER_SKIP_SENT; + public final Counter QUORUM_PEER_SKIP_OUT_OF_ORDER; + public final Counter QUORUM_PEER_SKIP_QUEUE_SIZE; public final Summary READS_ISSUED_IN_COMMIT_PROC; public final Summary WRITES_ISSUED_IN_COMMIT_PROC; @@ -408,6 +421,7 @@ private ServerMetrics(MetricsProvider metricsProvider) { */ public final Summary OM_PROPOSAL_PROCESS_TIME; public final Summary OM_COMMIT_PROCESS_TIME; + public final Summary OM_SKIP_PROCESS_TIME; /** * Time spent by the final processor. This is tracked in the commit processor. diff --git a/zookeeper-server/src/main/java/org/apache/zookeeper/server/ZooKeeperServer.java b/zookeeper-server/src/main/java/org/apache/zookeeper/server/ZooKeeperServer.java index 1c9dda77256..96308976bb1 100644 --- a/zookeeper-server/src/main/java/org/apache/zookeeper/server/ZooKeeperServer.java +++ b/zookeeper-server/src/main/java/org/apache/zookeeper/server/ZooKeeperServer.java @@ -561,7 +561,11 @@ public SessionTracker getSessionTracker() { } long getNextZxid() { - return hzxid.incrementAndGet(); + return hzxid.get() + 1; + } + + void incrementZxid() { + hzxid.incrementAndGet(); } public void setZxid(long zxid) { @@ -1716,7 +1720,7 @@ public ProcessTxnResult processTxn(Request request) { } // do not add non quorum packets to the queue. - if (quorumRequest) { + if (quorumRequest && !request.isSkipped()) { getZKDatabase().addCommittedProposal(request); } return rc; diff --git a/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/Follower.java b/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/Follower.java index 4d80d93358f..ca632fa980c 100644 --- a/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/Follower.java +++ b/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/Follower.java @@ -195,6 +195,18 @@ protected void processPacket(QuorumPacket qp) throws Exception { ServerMetrics.getMetrics().OM_PROPOSAL_PROCESS_TIME.add(Time.currentElapsedTime() - startTime); } break; + case Leader.SKIP: + ServerMetrics.getMetrics().LEARNER_SKIP_RECEIVED_COUNT.add(1); + fzk.skip(qp.getData()); + if (om != null) { + final long startTime = Time.currentElapsedTime(); + // We need to create a new QuorumPacket object that will be queued for sending + QuorumPacket observerPacket = new QuorumPacket( + qp.getType(), qp.getZxid(), qp.getData(), qp.getAuthinfo()); + om.proposalSkipped(observerPacket); + ServerMetrics.getMetrics().OM_SKIP_PROCESS_TIME.add(Time.currentElapsedTime() - startTime); + } + break; case Leader.COMMIT: ServerMetrics.getMetrics().LEARNER_COMMIT_RECEIVED_COUNT.add(1); fzk.commit(qp.getZxid()); 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 12e552f9d7d..7e3c968818d 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 @@ -31,8 +31,10 @@ 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.server.util.SerializeUtils; import org.apache.zookeeper.txn.TxnHeader; import org.apache.zookeeper.util.ServiceUtils; import org.slf4j.Logger; @@ -108,6 +110,19 @@ public void commit(long zxid) { commitProcessor.commit(request); } + + public void skip(byte[] data) { + try { + TxnHeader hdr = new TxnHeader(); + Record txn = SerializeUtils.deserializeTxn(data, hdr); + Request request = new Request(hdr.getClientId(), hdr.getCxid(), hdr.getType(), hdr, txn, 0); + request.logLatency(ServerMetrics.getMetrics().SKIP_PROPAGATION_LATENCY); + commitProcessor.commit(request); + } catch (IOException e) { + LOG.error("Could not deserialize SKIP request"); + } + } + public synchronized void sync() { if (pendingSyncs.size() == 0) { LOG.warn("Not expecting a sync."); 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 f0eca518b9b..991a9ad78f8 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 @@ -423,6 +423,13 @@ Optional createServerSocket(InetSocketAddress address, boolean por */ static final int INFORMANDACTIVATE = 19; + /** + * This message type tells the Learner that the request has + * an error transaction in it which means that it wasn't preceded + * with a PROPOSAL and that it can be committed immediately. + */ + final static int SKIP = 101; + final ConcurrentMap outstandingProposals = new ConcurrentHashMap(); private final ConcurrentLinkedQueue toBeApplied = new ConcurrentLinkedQueue(); @@ -955,7 +962,7 @@ public synchronized boolean tryToCommit(Proposal p, long zxid, SocketAddress fol commit(zxid); inform(p); } - zk.commitProcessor.commit(p.request); + zk.getLeaderHandler().processCommit(p.request); if (pendingSyncs.containsKey(zxid)) { for (LearnerSyncRequest r : pendingSyncs.remove(zxid)) { sendSync(r); @@ -1084,6 +1091,9 @@ public void processRequest(Request request) throws RequestProcessorException { // request.zxid here because that is set on read requests to equal // the zxid of the last write op. if (request.getHdr() != null) { + if (request.isSkipped()) { + return; + } long zxid = request.getHdr().getZxid(); Iterator iter = leader.toBeApplied.iterator(); if (iter.hasNext()) { @@ -1613,6 +1623,8 @@ public static String getPacketType(int packetType) { return "PROPOSAL"; case ACK: return "ACK"; + case SKIP: + return "SKIP"; case COMMIT: return "COMMIT"; case COMMITANDACTIVATE: diff --git a/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/LeaderHandler.java b/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/LeaderHandler.java new file mode 100644 index 00000000000..341948fc28e --- /dev/null +++ b/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/LeaderHandler.java @@ -0,0 +1,85 @@ +/** + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + *

+ * http://www.apache.org/licenses/LICENSE-2.0 + *

+ * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.zookeeper.server.quorum; + +import org.apache.zookeeper.server.Request; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import java.util.List; + +/** + * This class is used to process commits on the Leader side. + * + * When -Dzookeeper.quorum.skipTxnEnabled=true: + * After commit is processed, LeaderHandler will commit + * the request and try to flush a queue of skipped transactions + * that occurred right before that commit. + */ +public class LeaderHandler implements SkipRequestHandler { + private static final Logger LOG = LoggerFactory.getLogger(LeaderHandler.class); + + private final CommitProcessor commitProcessor; + private final SkippedRequestQueue skippedRequestQueue; + + public LeaderHandler(CommitProcessor commitProcessor) { + this.commitProcessor = commitProcessor; + this.skippedRequestQueue = new SkippedRequestQueue(); + } + + /** + * This method must be called from the RP chain thread + * right after PrepRequestProcessor. Once it's known that + * the request is invalid we can skip sending PROPOSE and go directly to + * COMMIT but we need to wait to do that in order with the other requests. + * + * The request will be immediately sent back to the learner if there is no pending commits + * otherwise the request will be put inside the handler's queue to wait for pending commit + * to complete before is sent so we make sure we send all requests in order. + */ + @Override + public void skipProposing(Request request) { + List requests = skippedRequestQueue.addAndGet(request); + processSkippedRequests(requests); + } + + /** + * Leader specific commitProcessing logic + */ + public void processCommit(Request request) { + commitProcessor.commit(request); + List requests = skippedRequestQueue.setLastCommittedAndGet(request.zxid); + processSkippedRequests(requests); + } + + /** + * This method gets called from the getSkippedRequests() method + * For each request inside the queue we will run this method. + * + * Note this method is differently implemented + * in LeaderHandler and LearnerHandler. + * + * LeaderHandler commits directly in it's own commit processor. + * LearnerHandler must send SKIP message to the Learners so they can + * commit the request in their own commit processor. + */ + @Override + public void processSkip(Request request) { + commitProcessor.commit(request); + } +} \ No newline at end of file diff --git a/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/LeaderRequestProcessor.java b/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/LeaderRequestProcessor.java index e78eaab2105..142b2bfab41 100644 --- a/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/LeaderRequestProcessor.java +++ b/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/LeaderRequestProcessor.java @@ -70,7 +70,7 @@ public void processRequest(Request request) throws RequestProcessorException { if (upgradeRequest != null) { nextProcessor.processRequest(upgradeRequest); } - + request.setSkipRequestHandler(lzks.getLeaderHandler()); nextProcessor.processRequest(request); } diff --git a/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/LeaderZooKeeperServer.java b/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/LeaderZooKeeperServer.java index 27de4519177..02fdca874d4 100644 --- a/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/LeaderZooKeeperServer.java +++ b/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/LeaderZooKeeperServer.java @@ -44,12 +44,41 @@ */ public class LeaderZooKeeperServer extends QuorumZooKeeperServer { + public static final String SKIP_TXN_ENABLED = "zookeeper.quorum.skipTxnEnabled"; + private ContainerManager containerManager; // guarded by sync CommitProcessor commitProcessor; PrepRequestProcessor prepRequestProcessor; + SkipRequestProcessor skipRequestProcessor; + + LeaderHandler leaderHandler; + + /** + * Setting this to true will skip the PROPOSE, ACK, COMMIT process for all write ops that + * didn't pass validations inside the PrepRequestProcessor which will result + * in a smaller log and less network overhead. + * + * Java property zookeeper.quorum.skipTxnEnabled + */ + private static boolean skipTxnEnabled; + + static { + skipTxnEnabled = Boolean.getBoolean(SKIP_TXN_ENABLED); + LOG.info("{} = {}", SKIP_TXN_ENABLED, skipTxnEnabled); + } + + public static boolean isSkipTxnEnabled() { + return skipTxnEnabled; + } + + public static void setSkipTxnEnabled(boolean isEnabled) { + skipTxnEnabled = isEnabled; + } + + /** * @throws IOException */ @@ -61,6 +90,10 @@ public Leader getLeader() { return self.leader; } + public LeaderHandler getLeaderHandler() { + return leaderHandler; + } + @Override protected void setupRequestProcessors() { RequestProcessor finalProcessor = new FinalRequestProcessor(this); @@ -69,10 +102,16 @@ protected void setupRequestProcessors() { commitProcessor.start(); ProposalRequestProcessor proposalProcessor = new ProposalRequestProcessor(this, commitProcessor); proposalProcessor.initialize(); - prepRequestProcessor = new PrepRequestProcessor(this, proposalProcessor); + if (LeaderZooKeeperServer.isSkipTxnEnabled()) { + skipRequestProcessor = new SkipRequestProcessor(proposalProcessor, commitProcessor); + prepRequestProcessor = new PrepRequestProcessor(this, skipRequestProcessor); + } else { + prepRequestProcessor = new PrepRequestProcessor(this, proposalProcessor); + } prepRequestProcessor.start(); firstProcessor = new LeaderRequestProcessor(this, prepRequestProcessor); + leaderHandler = new LeaderHandler(commitProcessor); setupContainerManager(); } diff --git a/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/LearnerHandler.java b/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/LearnerHandler.java index e7d1d40144b..1397e66c58a 100644 --- a/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/LearnerHandler.java +++ b/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/LearnerHandler.java @@ -28,6 +28,7 @@ import java.util.Date; import java.util.Iterator; import java.util.LinkedHashMap; +import java.util.List; import java.util.Map; import java.util.Objects; import java.util.Queue; @@ -61,7 +62,7 @@ * learner. All communication with a learner is handled by this * class. */ -public class LearnerHandler extends ZooKeeperThread { +public class LearnerHandler extends ZooKeeperThread implements SkipRequestHandler { private static final Logger LOG = LoggerFactory.getLogger(LearnerHandler.class); @@ -553,6 +554,10 @@ public void run() { } peerLastZxid = ss.getLastZxid(); + + ableToSkipTxn = QuorumZooKeeperServer.isSkipTxnAvailable(learnerProtocolVersion); + LOG.info("is learner able to process skip txn packets: {}", ableToSkipTxn); + // Take any necessary action if we need to send TRUNC or DIFF // startForwarding() will be called in all cases boolean needSnap = syncFollower(peerLastZxid, learnerMaster); @@ -605,6 +610,7 @@ public void run() { queuedPackets.add(newLeaderQP); } bufferedOutput.flush(); + skippedRequestQueue.setLastCommitted(newLeaderZxid); // Start thread that blast packets in the queue to learner startSendingPackets(); @@ -703,6 +709,7 @@ public void run() { si = new Request(null, sessionId, cxid, type, bb, qp.getAuthinfo()); } si.setOwner(this); + si.setSkipRequestHandler(this); learnerMaster.submitLearnerRequest(si); requestsReceived.incrementAndGet(); break; @@ -1100,6 +1107,18 @@ private void queueOpPacket(int type, long zxid) { void queuePacket(QuorumPacket p) { queuedPackets.add(p); + + int packetType = p.getType(); + if (packetType == Leader.COMMIT || packetType == Leader.INFORM) { + // We can skip calling this method when zookeeper.quorum.skipTxnEnabled=false + // but the reason why we are not doing that is to avoid "skip out of order" issue + // since skiptTxn feature can be enabled while we are processing requests with the JMX. + // + // If the JMX part gets removed then this should be called only + // when zookeeper.quorum.skipTxnEnabled=true. + processCommit(p.getZxid()); + } + // Add a MarkerQuorumPacket at regular intervals. if (shouldSendMarkerPacketForLogging() && packetCounter.getAndIncrement() % markerPacketInterval == 0) { queuedPackets.add(new MarkerQuorumPacket(System.nanoTime())); @@ -1159,4 +1178,43 @@ public void setFirstPacket(boolean value) { needOpPacket = value; } + private final SkippedRequestQueue skippedRequestQueue = new SkippedRequestQueue(); + + /** + * This method must be called from the RP chain thread + * right after PrepRequestProcessor. Once it's known that + * the request is invalid we can skip sending PROPOSE and go directly to + * COMMIT but we need to wait to do that in order with the other requests. + * + * The request will be immediately sent back to the learner if there is no pending commits + * otherwise the request will be put inside the handler's queue to wait for pending commit + * to complete before is sent so we make sure we send all requests in order. + */ + public void skipProposing(Request request) { + List requests = skippedRequestQueue.addAndGet(request); + processSkippedRequests(requests); + } + + /** + * Learner specific commitProcessing logic + */ + public void processCommit(long zxid) { + List requests = skippedRequestQueue.setLastCommittedAndGet(zxid); + processSkippedRequests(requests); + } + + /** + * This may be called from the RP chain and Leader threads. + * Note that LeaderHandler and LearnerHandler process skipped requests differently. + * + * LeaderHandler modifies the CommitProcessor directly while + * LearnerHandler needs to send the SKIP packet over the network. + */ + @Override + public void processSkip(Request request) { + byte[] data = SerializeUtils.serializeRequest(request); + QuorumPacket packet = new QuorumPacket(Leader.SKIP, 0, data, null); + queuedPackets.add(packet); + queuedPacketsSize.addAndGet(packetSize(packet)); + } } diff --git a/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/Observer.java b/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/Observer.java index 02683e4c59b..351a96db721 100644 --- a/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/Observer.java +++ b/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/Observer.java @@ -187,11 +187,21 @@ protected void processPacket(QuorumPacket qp) throws Exception { case Leader.SYNC: ((ObserverZooKeeperServer) zk).sync(); break; - case Leader.INFORM: - ServerMetrics.getMetrics().LEARNER_COMMIT_RECEIVED_COUNT.add(1); + case Leader.SKIP: TxnHeader hdr = new TxnHeader(); + ServerMetrics.getMetrics().LEARNER_SKIP_RECEIVED_COUNT.add(1); Record txn = SerializeUtils.deserializeTxn(qp.getData(), hdr); - Request request = new Request(hdr.getClientId(), hdr.getCxid(), hdr.getType(), hdr, txn, 0); + Request request = new Request( + hdr.getClientId(), hdr.getCxid(), hdr.getType(), hdr, txn, hdr.getZxid() + ); + request.logLatency(ServerMetrics.getMetrics().SKIP_PROPAGATION_LATENCY); + ((ObserverZooKeeperServer) zk).skipRequest(request); + break; + case Leader.INFORM: + ServerMetrics.getMetrics().LEARNER_COMMIT_RECEIVED_COUNT.add(1); + hdr = new TxnHeader(); + txn = SerializeUtils.deserializeTxn(qp.getData(), hdr); + request = new Request(hdr.getClientId(), hdr.getCxid(), hdr.getType(), hdr, txn, 0); request.logLatency(ServerMetrics.getMetrics().COMMIT_PROPAGATION_LATENCY); ObserverZooKeeperServer obs = (ObserverZooKeeperServer) zk; obs.commitRequest(request); diff --git a/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/ObserverMaster.java b/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/ObserverMaster.java index 54a22c2ecbd..f5fe14a2569 100644 --- a/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/ObserverMaster.java +++ b/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/ObserverMaster.java @@ -412,6 +412,12 @@ synchronized void proposalCommitted(long zxid) { sendPacket(pkt); } + synchronized void proposalSkipped(QuorumPacket qp) { + for (LearnerHandler lh : activeObservers) { + lh.queuePacket(qp); + } + } + synchronized void informAndActivate(long zxid, long suggestedLeaderId) { QuorumPacket pkt = removeProposedPacket(zxid); if (pkt == null) { diff --git a/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/ObserverZooKeeperServer.java b/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/ObserverZooKeeperServer.java index a41a9187743..48e4fe35d37 100644 --- a/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/ObserverZooKeeperServer.java +++ b/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/ObserverZooKeeperServer.java @@ -80,6 +80,16 @@ public void commitRequest(Request request) { commitProcessor.commit(request); } + /** + * When we have a write request in the commit processor that ends up being skipped + * we want to process it in the order so it doesn't block all other requests + * + * @param request + */ + public void skipRequest(Request request) { + commitProcessor.commit(request); + } + /** * Set up the request processors for an Observer: * firstProcesor->commitProcessor->finalProcessor diff --git a/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/QuorumZooKeeperServer.java b/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/QuorumZooKeeperServer.java index c6cd93b4d18..e62e7795ec8 100644 --- a/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/QuorumZooKeeperServer.java +++ b/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/QuorumZooKeeperServer.java @@ -215,4 +215,8 @@ public void dumpMonitorValues(BiConsumer response) { response.accept("peer_state", self.getDetailedPeerState()); } + public static boolean isSkipTxnAvailable(int learnerProtocolVersion) { + return true; + } + } diff --git a/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/SkipRequestHandler.java b/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/SkipRequestHandler.java new file mode 100644 index 00000000000..084e4951be5 --- /dev/null +++ b/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/SkipRequestHandler.java @@ -0,0 +1,38 @@ +package org.apache.zookeeper.server.quorum; + +import org.apache.zookeeper.server.Request; + +import java.util.List; + +public interface SkipRequestHandler { + /** + * This method must be called from the RP chain thread + * right after PrepRequestProcessor. Once it's known that + * the request is invalid we can skip sending PROPOSE and go directly to + * COMMIT but we need to wait to do that in order with the other requests. + * + * The request will be immediately sent back to the learner if there is no pending commits + * otherwise the request will be put inside the handler's queue to wait for pending commit + * to complete before is sent so we make sure we send all requests in order. + */ + void skipProposing(Request request); + + /** + * This may be called from the RP chain and Leader threads. + * Note that LeaderHandler and LearnerHandler process skipped requests differently. + * + * LeaderHandler modifies the CommitProcessor directly while + * LearnerHandler needs to send the SKIP packet over the network. + */ + void processSkip(Request request); + + + /** + * Default method to abstract skip request processing + */ + default void processSkippedRequests(List requests) { + for (Request request : requests) { + processSkip(request); + } + } +} \ No newline at end of file diff --git a/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/SkipRequestProcessor.java b/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/SkipRequestProcessor.java new file mode 100644 index 00000000000..b94f4895677 --- /dev/null +++ b/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/SkipRequestProcessor.java @@ -0,0 +1,66 @@ +/** + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + *

+ * http://www.apache.org/licenses/LICENSE-2.0 + *

+ * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.zookeeper.server.quorum; + +import org.apache.zookeeper.server.Request; +import org.apache.zookeeper.server.RequestProcessor; +import org.apache.zookeeper.server.ServerMetrics; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +public class SkipRequestProcessor implements RequestProcessor { + + private static final Logger LOG = LoggerFactory.getLogger(SkipRequestProcessor.class); + private final RequestProcessor nextProcessor; + private final RequestProcessor commitProcessor; + + public SkipRequestProcessor(RequestProcessor nextProcessor, RequestProcessor commitProcessor) { + this.nextProcessor = nextProcessor; + this.commitProcessor = commitProcessor; + } + + @Override + public void processRequest(Request request) throws RequestProcessorException { + if (request.isSkipped()) { + /** + * In case when the request is coming from a client that is connected to the Leader, + * we need to submit the request to the Leader's CommitProcessor so we can SKIP it in order + * and return the status back to the client. + * + * We don't need to do this when the request comes from a Follower because the request will + * already be submitted to their own CommitProcessor and waiting to hear COMMIT or SKIP from + * the Leader so it could reply the request status back to the client. + */ + if (request.isFromLeader()) { + commitProcessor.processRequest(request); + } + + SkipRequestHandler handler = request.getSkipRequestHandler(); + handler.skipProposing(request); + ServerMetrics.getMetrics().SKIP_COUNT.add(1); + } else { + nextProcessor.processRequest(request); + } + } + + @Override + public void shutdown() { + LOG.info("Shutting down"); + nextProcessor.shutdown(); + } +} \ No newline at end of file diff --git a/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/SkippedRequestQueue.java b/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/SkippedRequestQueue.java new file mode 100644 index 00000000000..74a85b801a2 --- /dev/null +++ b/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/SkippedRequestQueue.java @@ -0,0 +1,78 @@ +package org.apache.zookeeper.server.quorum; + +import org.apache.zookeeper.server.ExitCode; +import org.apache.zookeeper.server.Request; +import org.apache.zookeeper.server.ServerMetrics; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import java.util.ArrayList; +import java.util.List; +import java.util.concurrent.LinkedBlockingQueue; +import java.util.concurrent.atomic.AtomicLong; + +public class SkippedRequestQueue { + private static final Logger LOG = LoggerFactory.getLogger(SkippedRequestQueue.class); + + private final AtomicLong lastCommitted = new AtomicLong(-1); + private final LinkedBlockingQueue waitingRequests = new LinkedBlockingQueue<>(); + + /** + * Setter method for lastCommitted zxid + */ + public void setLastCommitted(long zxid) { + lastCommitted.set(zxid); + } + + public List setLastCommittedAndGet(long zxid) { + synchronized (this) { + setLastCommitted(zxid); + return getSkippedRequests(); + } + } + + /** + * Adds skipped request to the end of the queue that way we ensure the order. + * Returns all the skipped requests that are ready to be sent back to the origin, + * meaning if there was a COMMIT that they were waiting for + */ + public List addAndGet(Request request) { + ServerMetrics.getMetrics().QUORUM_PEER_SKIP_QUEUE_SIZE.add(1); + synchronized (this) { + waitingRequests.add(request); + return getSkippedRequests(); + } + } + + /** + * This method flushed the error request queue. This method gets called from two different + * threads thorough skipProposing() and updateLastZxidAndGetSkippedRequests(). + * + * Leader thread calls this method every time the leader processes a COMMIT in order + * to try to flush all the pending error requests that we short circuited. + * + * The RP chain thread will try to call this when short-circuiting requests + * + */ + public List getSkippedRequests() { + List requestsToSkip = new ArrayList<>(); + long lastZxid = lastCommitted.get(); + + while (!waitingRequests.isEmpty()) { + Request request = waitingRequests.peek(); + if (lastZxid > request.zxid) { + ServerMetrics.getMetrics().QUORUM_PEER_SKIP_OUT_OF_ORDER.add(1); + LOG.error("Skip request is out of order!!!"); + System.exit(ExitCode.QUORUM_PEER_SKIP_OUT_OF_ORDER.getValue()); + } + if (lastZxid < request.zxid) { + break; + } + waitingRequests.poll(); + requestsToSkip.add(request); + ServerMetrics.getMetrics().QUORUM_PEER_SKIP_SENT.add(1); + } + + return requestsToSkip; + } +} \ No newline at end of file diff --git a/zookeeper-server/src/test/java/org/apache/zookeeper/server/PrepRequestProcessorTest.java b/zookeeper-server/src/test/java/org/apache/zookeeper/server/PrepRequestProcessorTest.java index 264601de076..5e7c3c0d190 100644 --- a/zookeeper-server/src/test/java/org/apache/zookeeper/server/PrepRequestProcessorTest.java +++ b/zookeeper-server/src/test/java/org/apache/zookeeper/server/PrepRequestProcessorTest.java @@ -22,6 +22,9 @@ import static org.junit.Assert.assertNotNull; import static org.junit.Assert.assertNull; import static org.junit.Assert.assertTrue; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.when; + import java.io.ByteArrayOutputStream; import java.io.File; import java.io.IOException; @@ -48,6 +51,8 @@ import org.apache.zookeeper.proto.RequestHeader; import org.apache.zookeeper.proto.SetDataRequest; import org.apache.zookeeper.server.ZooKeeperServer.ChangeRecord; +import org.apache.zookeeper.server.quorum.LeaderZooKeeperServer; +import org.apache.zookeeper.server.quorum.LearnerHandler; import org.apache.zookeeper.test.ClientBase; import org.apache.zookeeper.txn.ErrorTxn; import org.junit.After; @@ -241,6 +246,46 @@ public void testInvalidPath() throws Exception { assertEquals(outcome.getException().code(), KeeperException.Code.BADARGUMENTS); } + @Test + public void testRequestSkipDoesntUseZxidWhenIsOn() throws Exception { + LeaderZooKeeperServer.setSkipTxnEnabled(true); + + pLatch = new CountDownLatch(1); + processor = new PrepRequestProcessor(zks, new MyRequestProcessor()); + long lastCommittedZxid = zks.getZxid(); + + SetDataRequest record = new SetDataRequest("", new byte[0], -1); + Request req = createRequest(record, OpCode.setData); + LearnerHandler handler = mock(LearnerHandler.class); + req.setSkipRequestHandler(handler); + + processor.pRequest(req); + pLatch.await(); + + assertEquals(outcome.zxid, lastCommittedZxid); + assertEquals(outcome.getHdr().getType(), OpCode.error); + assertEquals(outcome.getException().code(), KeeperException.Code.BADARGUMENTS); + } + + @Test + public void testRequestSkipDoesUseZxidWhenIsOff() throws Exception { + LeaderZooKeeperServer.setSkipTxnEnabled(false); + + pLatch = new CountDownLatch(1); + processor = new PrepRequestProcessor(zks, new MyRequestProcessor()); + long lastCommittedZxid = zks.getZxid(); + + SetDataRequest record = new SetDataRequest("", new byte[0], -1); + Request req = createRequest(record, OpCode.setData); + + processor.pRequest(req); + pLatch.await(); + + assertEquals(outcome.getHdr().getZxid(), lastCommittedZxid + 1); + assertEquals(outcome.getHdr().getType(), OpCode.error); + assertEquals(outcome.getException().code(), KeeperException.Code.BADARGUMENTS); + } + private class MyRequestProcessor implements RequestProcessor { @Override diff --git a/zookeeper-server/src/test/java/org/apache/zookeeper/server/quorum/QuorumPeerSkipErrorsTest.java b/zookeeper-server/src/test/java/org/apache/zookeeper/server/quorum/QuorumPeerSkipErrorsTest.java new file mode 100644 index 00000000000..a6f8038b474 --- /dev/null +++ b/zookeeper-server/src/test/java/org/apache/zookeeper/server/quorum/QuorumPeerSkipErrorsTest.java @@ -0,0 +1,169 @@ +package org.apache.zookeeper.server.quorum; + +import org.apache.zookeeper.CreateMode; +import org.apache.zookeeper.KeeperException; +import org.apache.zookeeper.ZooDefs; +import org.apache.zookeeper.ZooKeeper; +import org.junit.Test; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertTrue; + +public class QuorumPeerSkipErrorsTest extends QuorumPeerTestBase { + @Test + public void testRequestSkipConsistencyAfterLeaderElection() throws Exception { + LeaderZooKeeperServer.setSkipTxnEnabled(true); + + numServers = 3; + servers = LaunchServers(numServers); + int leader = servers.findLeader(); + + // make sure there is a leader + assertTrue("There should be a leader", leader >= 0); + + int nonleader = (leader + 1) % numServers; + byte[] input = new byte[1]; + input[0] = 1; + + // this will be an existing node to be used to create error txns + servers.zk[leader].create("/client", "client-data".getBytes(), ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT); + Thread.sleep(500); + + // Shutdown every one else but the leader + for (int i = 0; i < numServers; i++) { + if (i != leader) { + servers.mt[i].shutdown(); + } + } + Thread.sleep(500); + + // shut the leader down + servers.mt[leader].shutdown(); + System.gc(); + + waitForAll(servers.zk, ZooKeeper.States.CONNECTING); + + // Start everyone but the leader + for (int i = 0; i < numServers; i++) { + if (i != leader) { + servers.mt[i].start(); + } + } + + // wait to connect to one of these + waitForOne(servers.zk[nonleader], ZooKeeper.States.CONNECTED); + + // start the old leader + servers.mt[leader].start(); + waitForOne(servers.zk[leader], ZooKeeper.States.CONNECTED); + + String mainNodePath = "/client"; + String mainNodeData = "test-data"; + + // start with an error + try { + servers.zk[leader].create(mainNodePath, mainNodeData.getBytes(), ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT); + } catch (KeeperException e) { + assertEquals("KeeperErrorCode = NodeExists for /client", e.getMessage()); + } + + // one client per server to simulate some normal load + // - create a new node that doest exist + // - try to create a node that already exists -> error transaction + Thread[] clients = new Thread[numServers]; + for (int cid = 0; cid < numServers; cid++) { + final int id = cid; + clients[cid] = new Thread(() -> { + for (int i = 0; i < 10000; i++) { + String threadNodePath = "/client-" + id + "-" + i; + String threadNodeData = "test-data"; + + try { + // valid transaction + servers.zk[id].create(threadNodePath, threadNodeData.getBytes(), ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT); + // valid read + servers.zk[id].getData(threadNodePath, null, null); + // invalid transaction + servers.zk[id].create(mainNodePath, mainNodeData.getBytes(), ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT); + } catch (InterruptedException | KeeperException e) { + } + } + }); + } + + // start the client threads + for (Thread client : clients) { + client.start(); + } + + // let the client threads finish work + for (Thread client : clients) { + client.join(); + } + + leader = servers.findLeader(); + Leader leaderPeer = servers.mt[leader].main.quorumPeer.leader; + assertEquals(leaderPeer.lastCommitted, leaderPeer.lastProposed); + assertEquals(0, leaderPeer.zk.commitProcessor.pendingRequests.size()); + assertEquals(0, leaderPeer.zk.commitProcessor.committedRequests.size()); + assertEquals(0, leaderPeer.zk.commitProcessor.queuedRequests.size()); + assertEquals(0, leaderPeer.zk.commitProcessor.queuedWriteRequests.size()); + assertEquals(0, leaderPeer.zk.commitProcessor.numRequestsProcessing.get()); + } + + @Test + public void testRequestSkipConsistency() throws Exception { + LeaderZooKeeperServer.setSkipTxnEnabled(true); + + numServers = 3; + servers = LaunchServers(numServers); + int leader = servers.findLeader(); + + // shared node for both client 1 and client 2 + servers.zk[leader].create("/shared", "shared-data".getBytes(), ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT); + assertEquals(new String(servers.zk[leader].getData("/shared", null, null)), "shared-data"); + + // znode owned by client 1 + servers.zk[leader].create("/client-1", "client-1-data".getBytes(), ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT); + assertEquals(new String(servers.zk[leader].getData("/client-1", null, null)), "client-1-data"); + + // znode owned by client 2 + servers.zk[leader].create("/client-2", "client-2-data".getBytes(), ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT); + assertEquals(new String(servers.zk[leader].getData("/client-2", null, null)), "client-2-data"); + + Thread[] clients = new Thread[numServers]; + for (int cid = 0; cid < numServers; cid++) { + final int id = cid; + clients[cid] = new Thread(() -> { + for (int i = 0; i < 10000; i++) { + try { + String threadNodePath = "/client-" + id + "-" + i; + String threadNodeData = "data-" + id + "-" + i; + + servers.zk[id].setData("/shared", threadNodeData.getBytes(), -1); + servers.zk[id].create(threadNodePath, threadNodeData.getBytes(), ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT); + servers.zk[id].create(threadNodePath, threadNodeData.getBytes(), ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT); + } catch (InterruptedException | KeeperException e) { + } + } + }); + } + + for (Thread client : clients) { + client.start(); + } + + for (Thread client : clients) { + client.join(); + } + + leader = servers.findLeader(); + Leader leaderPeer = servers.mt[leader].main.quorumPeer.leader; + assertEquals(leaderPeer.lastCommitted, leaderPeer.lastProposed); + assertEquals(0, leaderPeer.zk.commitProcessor.pendingRequests.size()); + assertEquals(0, leaderPeer.zk.commitProcessor.committedRequests.size()); + assertEquals(0, leaderPeer.zk.commitProcessor.queuedRequests.size()); + assertEquals(0, leaderPeer.zk.commitProcessor.queuedWriteRequests.size()); + assertEquals(0, leaderPeer.zk.commitProcessor.numRequestsProcessing.get()); + } +} \ No newline at end of file diff --git a/zookeeper-server/src/test/java/org/apache/zookeeper/server/quorum/SkipRequestHandlerTest.java b/zookeeper-server/src/test/java/org/apache/zookeeper/server/quorum/SkipRequestHandlerTest.java new file mode 100644 index 00000000000..78b0c3f490e --- /dev/null +++ b/zookeeper-server/src/test/java/org/apache/zookeeper/server/quorum/SkipRequestHandlerTest.java @@ -0,0 +1,110 @@ +package org.apache.zookeeper.server.quorum; + +import org.apache.zookeeper.KeeperException; +import org.apache.zookeeper.ZooDefs; +import org.apache.zookeeper.server.Request; +import org.apache.zookeeper.txn.TxnHeader; +import org.junit.Test; + +import java.util.List; + +import static org.junit.Assert.assertFalse; +import static org.junit.Assert.assertTrue; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.spy; +import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; + +public class SkipRequestHandlerTest { + private static class TestSkipRequestHandler implements SkipRequestHandler { + boolean isSkipProcessed = false; + SkippedRequestQueue skippedRequestQueue = spy(new SkippedRequestQueue()); + + @Override + public void skipProposing(Request request) { + List requestsToSkip = skippedRequestQueue.addAndGet(request); + for (Request requestToSkip : requestsToSkip) { + processSkip(requestToSkip); + } + } + + void processCommit(Request request) { + List requestsToSkip = skippedRequestQueue.setLastCommittedAndGet(request.zxid); + for (Request requestToSkip : requestsToSkip) { + processSkip(requestToSkip); + } + } + + @Override + public void processSkip(Request request) { + isSkipProcessed = true; + } + } + + @Test + public void testSkipProposingAndSendingRequestWithoutWaiting() throws Exception { + TestSkipRequestHandler handler = spy(new TestSkipRequestHandler()); + handler.skippedRequestQueue.setLastCommitted(1); + + Request request = getSampleErrorRequest(); + request.zxid = 1; + handler.skipProposing(request); + + verify(handler.skippedRequestQueue, times(1)).getSkippedRequests(); + assertTrue(handler.isSkipProcessed); + } + + @Test + public void testSkipProposingAndWaitingToSendRequest() throws Exception { + TestSkipRequestHandler handler = spy(new TestSkipRequestHandler()); + handler.skippedRequestQueue.setLastCommitted(1); + + Request request = getSampleErrorRequest(); + request.zxid = 2; + handler.skipProposing(request); + + verify(handler.skippedRequestQueue, times(1)).getSkippedRequests(); + assertFalse(handler.isSkipProcessed); + } + + @Test + public void testProcessCommit() throws Exception { + TestSkipRequestHandler handler = spy(new TestSkipRequestHandler()); + handler.skippedRequestQueue.setLastCommitted(1); + + Request errRequest = getSampleErrorRequest(); + errRequest.zxid = 2; + handler.skipProposing(errRequest); + + verify(handler.skippedRequestQueue, times(1)).getSkippedRequests(); + assertFalse(handler.isSkipProcessed); + + Request okRequest = getSampleRequest(); + okRequest.zxid = 2; + handler.processCommit(okRequest); + + assertTrue(handler.isSkipProcessed); + } + + private Request getSampleErrorRequest() { + Request request = spy(new Request(null, 0, 0, ZooDefs.OpCode.delete, null, null)); + + request.setException(new KeeperException.NoNodeException()); + request.setHdr(new TxnHeader()); + request.getHdr().setType(ZooDefs.OpCode.error); + request.setSkipRequestHandler(mock(LearnerHandler.class)); + + return request; + } + + private Request getSampleRequest() { + Request request = spy(new Request(null, 0, 0, ZooDefs.OpCode.delete, null, null)); + + request.setException(new KeeperException.NoNodeException()); + request.setHdr(new TxnHeader()); + request.getHdr().setType(ZooDefs.OpCode.create); + request.setSkipRequestHandler(mock(LearnerHandler.class)); + + return request; + } +} \ No newline at end of file diff --git a/zookeeper-server/src/test/java/org/apache/zookeeper/server/quorum/SkipRequestProcessorTest.java b/zookeeper-server/src/test/java/org/apache/zookeeper/server/quorum/SkipRequestProcessorTest.java new file mode 100644 index 00000000000..385887df7a1 --- /dev/null +++ b/zookeeper-server/src/test/java/org/apache/zookeeper/server/quorum/SkipRequestProcessorTest.java @@ -0,0 +1,103 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.zookeeper.server.quorum; + +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.spy; +import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +import org.apache.zookeeper.KeeperException; +import org.apache.zookeeper.ZooDefs.OpCode; +import org.apache.zookeeper.server.Request; +import org.apache.zookeeper.server.ServerCnxn; +import org.apache.zookeeper.txn.TxnHeader; +import org.junit.Test; + +public class SkipRequestProcessorTest { + + private SkipRequestProcessor processor; + + @Test + public void testProcessingValidRequests() throws Exception { + ProposalRequestProcessor proposalProcessor = mock(ProposalRequestProcessor.class); + CommitProcessor commitProcessor = mock(CommitProcessor.class); + + Request request = spy(new Request(null, 0, 0, OpCode.create, null, null)); + request.setHdr(new TxnHeader()); + + processor = spy(new SkipRequestProcessor(proposalProcessor, commitProcessor)); + processor.processRequest(request); + + verify(request, times(1)).isSkipped(); + verify(proposalProcessor, times(1)).processRequest(request); + verify(commitProcessor, times(0)).processRequest(request); + } + + @Test + public void testSkippingInvalidRequestsFromFollower() throws Exception { + ProposalRequestProcessor proposalProcessor = mock(ProposalRequestProcessor.class); + CommitProcessor commitProcessor = mock(CommitProcessor.class); + + Request request = spy(new Request(null, 0, 0, OpCode.delete, null, null)); + LearnerHandler handler = mock(LearnerHandler.class); + + request.setException(new KeeperException.NoNodeException()); + request.setHdr(new TxnHeader()); + request.getHdr().setType(OpCode.error); + request.setSkipRequestHandler(handler); + request.setSkipped(); + + processor = spy(new SkipRequestProcessor(proposalProcessor, commitProcessor)); + processor.processRequest(request); + + verify(request, times(1)).isSkipped(); + verify(handler, times(1)).skipProposing(request); + verify(proposalProcessor, times(0)).processRequest(request); + } + + @Test + public void testSkippingInvalidRequestsFromLeader() throws Exception { + LeaderZooKeeperServer server = mock(LeaderZooKeeperServer.class); + Leader leader = mock(Leader.class); + leader.lastCommitted = 100; + when(server.getLeader()).thenReturn(leader); + + ProposalRequestProcessor proposalProcessor = mock(ProposalRequestProcessor.class); + CommitProcessor commitProcessor = mock(CommitProcessor.class); + + Request request = spy(new Request(null, 0, 0, OpCode.error, null, null)); + LeaderHandler handler = mock(LeaderHandler.class); + + request.setException(new KeeperException.NoNodeException()); + request.setHdr(new TxnHeader()); + request.getHdr().setType(OpCode.error); + request.setSkipRequestHandler(handler); + request.setOwner(ServerCnxn.me); + request.setSkipped(); + + processor = spy(new SkipRequestProcessor(proposalProcessor, commitProcessor)); + processor.processRequest(request); + + verify(request, times(1)).isSkipped(); + verify(handler, times(1)).skipProposing(request); + verify(proposalProcessor, times(0)).processRequest(request); + verify(commitProcessor, times(1)).processRequest(request); + } +} \ No newline at end of file diff --git a/zookeeper-server/src/test/java/org/apache/zookeeper/test/MultiOperationTest.java b/zookeeper-server/src/test/java/org/apache/zookeeper/test/MultiOperationTest.java index 6e23bd62572..efaa5972714 100644 --- a/zookeeper-server/src/test/java/org/apache/zookeeper/test/MultiOperationTest.java +++ b/zookeeper-server/src/test/java/org/apache/zookeeper/test/MultiOperationTest.java @@ -59,6 +59,7 @@ import org.apache.zookeeper.data.Id; import org.apache.zookeeper.data.Stat; import org.apache.zookeeper.server.SyncRequestProcessor; +import org.apache.zookeeper.server.quorum.LeaderZooKeeperServer; import org.junit.Before; import org.junit.Test; import org.junit.runner.RunWith; @@ -908,6 +909,53 @@ public void testMixedReadAndTransaction() throws Exception { } } + @Test + public void testMultiOpWithSkipTxnEnabled() throws Exception { + LeaderZooKeeperServer.setSkipTxnEnabled(true); + + try { + multi(zk, Arrays.asList( + Op.create("/multi1", new byte[0], Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT), + Op.delete("/multi1", 1), + Op.delete("/multi1", 1) + )); + fail("delete /multi should have failed"); + } catch (KeeperException e) { + /* PASS */ + assertEquals("KeeperErrorCode = BadVersion", e.getMessage()); + } + + try { + multi(zk, Arrays.asList( + Op.create("/multi2", new byte[0], Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT), + Op.create("/multi2", new byte[0], Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT), + Op.create("/multi2", new byte[0], Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT) + )); + fail("create /multi should have failed"); + } catch (KeeperException e) { + /* PASS */ + assertEquals("KeeperErrorCode = NodeExists", e.getMessage()); + } + } + + @Test + public void testMultiOpResultWithSkipTxnEnabled() throws Exception { + LeaderZooKeeperServer.setSkipTxnEnabled(true); + + List expectedResultCodes = new ArrayList(); + expectedResultCodes.add(KeeperException.Code.OK.intValue()); + expectedResultCodes.add(KeeperException.Code.NODEEXISTS.intValue()); + expectedResultCodes.add(KeeperException.Code.RUNTIMEINCONSISTENCY.intValue()); + + List opList = Arrays.asList( + Op.create("/create_test", new byte[0], Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT), + Op.create("/create_test", new byte[0], Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT), + Op.delete("/create_test", -1) + ); + String expectedErr = KeeperException.Code.NODEEXISTS.name(); + multiHavingErrors(zk, opList, expectedResultCodes, expectedErr); + } + private static class HasTriggeredWatcher implements Watcher { private final CountDownLatch triggered = new CountDownLatch(1); diff --git a/zookeeper-server/src/test/java/org/apache/zookeeper/test/ObserverMasterTest.java b/zookeeper-server/src/test/java/org/apache/zookeeper/test/ObserverMasterTest.java index ee54a7ae57b..ac8428ca0a5 100644 --- a/zookeeper-server/src/test/java/org/apache/zookeeper/test/ObserverMasterTest.java +++ b/zookeeper-server/src/test/java/org/apache/zookeeper/test/ObserverMasterTest.java @@ -32,6 +32,7 @@ import java.util.Collections; import java.util.List; import java.util.Map; +import java.util.Random; import java.util.Set; import java.util.concurrent.CountDownLatch; import java.util.concurrent.LinkedBlockingQueue; @@ -63,6 +64,7 @@ import org.apache.zookeeper.server.admin.Commands; import org.apache.zookeeper.server.quorum.DelayRequestProcessor; import org.apache.zookeeper.server.quorum.FollowerZooKeeperServer; +import org.apache.zookeeper.server.quorum.LeaderZooKeeperServer; import org.apache.zookeeper.server.quorum.QuorumPeerConfig; import org.apache.zookeeper.server.quorum.QuorumPeerTestBase; import org.apache.zookeeper.server.util.PortForwarder; @@ -575,6 +577,71 @@ public void testAdminCommands() throws IOException, MBeanException, InstanceNotF JMXEnv.tearDown(); } + @Test + public void testInOrderCommitsWithSkipTxn() throws Exception { + LeaderZooKeeperServer.setSkipTxnEnabled(true); + setUp(-1); + + zk = new ZooKeeper("127.0.0.1:" + CLIENT_PORT_QP1, ClientBase.CONNECTION_TIMEOUT, null); + for (int i = 0; i < 10; i++) { + zk.create("/bulk" + + i, ("Initial data of some size").getBytes(), Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT); + } + zk.close(); + + q3.start(); + assertTrue( + "waiting for observer to be up", + ClientBase.waitForServerUp("127.0.0.1:" + CLIENT_PORT_OBS, CONNECTION_TIMEOUT)); + + latch = new CountDownLatch(1); + zk = new ZooKeeper("127.0.0.1:" + CLIENT_PORT_QP1, ClientBase.CONNECTION_TIMEOUT, this); + latch.await(); + assertEquals(zk.getState(), States.CONNECTED); + + zk.create("/init", "first".getBytes(), Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT); + final long zxid = q1.getQuorumPeer().getLastLoggedZxid(); + + // wait for change to propagate + waitFor("Timeout waiting for observer sync", new WaitForCondition() { + public boolean evaluate() { + return zxid == q3.getQuorumPeer().getLastLoggedZxid(); + } + }, 30); + + ZooKeeper obsZk = new ZooKeeper("127.0.0.1:" + CLIENT_PORT_OBS, ClientBase.CONNECTION_TIMEOUT, this); + int followerPort = q1.getQuorumPeer().leader == null ? CLIENT_PORT_QP1 : CLIENT_PORT_QP2; + ZooKeeper fZk = new ZooKeeper("127.0.0.1:" + followerPort, ClientBase.CONNECTION_TIMEOUT, this); + final int numTransactions = 10001; + CountDownLatch gate = new CountDownLatch(1); + CountDownLatch oAsyncLatch = new CountDownLatch(numTransactions); + Thread oAsyncWriteThread = new Thread(new AsyncWriter(obsZk, numTransactions, true, oAsyncLatch, "/obs", gate, 20)); + CountDownLatch fAsyncLatch = new CountDownLatch(numTransactions); + Thread fAsyncWriteThread = new Thread(new AsyncWriter(fZk, numTransactions, true, fAsyncLatch, "/follower", gate, 20)); + + LOG.info("ASYNC WRITES"); + oAsyncWriteThread.start(); + fAsyncWriteThread.start(); + gate.countDown(); + + oAsyncLatch.await(); + fAsyncLatch.await(); + + oAsyncWriteThread.join(ClientBase.CONNECTION_TIMEOUT); + if (oAsyncWriteThread.isAlive()) { + LOG.error("asyncWriteThread is still alive"); + } + fAsyncWriteThread.join(ClientBase.CONNECTION_TIMEOUT); + if (fAsyncWriteThread.isAlive()) { + LOG.error("asyncWriteThread is still alive"); + } + + obsZk.close(); + fZk.close(); + + shutdown(); + } + private String createServerString(String type, long serverId, int clientPort) { return "server." + serverId + "=127.0.0.1:" + PortAssignment.unique() + ":" + PortAssignment.unique() + ":" + type + ";" + clientPort; } @@ -698,6 +765,7 @@ class AsyncWriter implements Runnable { private final CountDownLatch writerLatch; private final String root; private final CountDownLatch gate; + private int errorRate; AsyncWriter(ZooKeeper client, int numTransactions, boolean issueSync, CountDownLatch writerLatch, String root, CountDownLatch gate) { this.client = client; @@ -706,10 +774,17 @@ class AsyncWriter implements Runnable { this.writerLatch = writerLatch; this.root = root; this.gate = gate; + this.errorRate = 0; + } + + AsyncWriter(ZooKeeper client, int numTransactions, boolean issueSync, CountDownLatch writerLatch, String root, CountDownLatch gate, int errorRate) { + this(client, numTransactions, issueSync, writerLatch, root, gate); + this.errorRate = errorRate; } @Override public void run() { + Random rnd = new Random(); if (gate != null) { try { gate.await(); @@ -721,7 +796,7 @@ public void run() { for (int i = 0; i < numTransactions; i++) { final boolean pleaseLog = i % 100 == 0; client.create(root - + i, "inner thread".getBytes(), Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT, new AsyncCallback.StringCallback() { + + i, "inner thread".getBytes(), Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT, new AsyncCallback.StringCallback() { @Override public void processResult(int rc, String path, Object ctx, String name) { writerLatch.countDown(); @@ -730,6 +805,19 @@ public void processResult(int rc, String path, Object ctx, String name) { } } }, null); + if (rnd.nextInt(100) < errorRate) { + // will try to produce an error by issuing the same create op + client.create(root + + i, "inner thread".getBytes(), Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT, new AsyncCallback.StringCallback() { + @Override + public void processResult(int rc, String path, Object ctx, String name) { + writerLatch.countDown(); + if (pleaseLog) { + LOG.info("creating error request on {}", path); + } + } + }, null); + } if (pleaseLog) { LOG.info("async wrote {}{}", root, i); if (issueSync) {