From caeda310094185c22bc469f587e12016429183bf Mon Sep 17 00:00:00 2001 From: maoling Date: Fri, 20 Nov 2020 11:25:18 +0800 Subject: [PATCH 1/2] ZOOKEEPER-3600 support the complete linearizable read and multiply read consistency level --- .../src/main/resources/zookeeper.jute | 8 + .../org/apache/zookeeper/AsyncCallback.java | 21 ++ .../java/org/apache/zookeeper/ClientCnxn.java | 11 + .../java/org/apache/zookeeper/ZooDefs.java | 6 + .../java/org/apache/zookeeper/ZooKeeper.java | 86 ++++++ .../zookeeper/client/ReadConsistencyMode.java | 55 ++++ .../org/apache/zookeeper/client/ZNode.java | 66 +++++ .../server/FinalRequestProcessor.java | 52 +++- .../server/PrepRequestProcessor.java | 2 + .../org/apache/zookeeper/server/Request.java | 12 +- .../zookeeper/server/RequestThrottler.java | 2 +- .../zookeeper/server/ZooKeeperServer.java | 3 +- .../server/quorum/CommitProcessor.java | 25 +- .../zookeeper/server/quorum/Follower.java | 6 +- .../quorum/FollowerRequestProcessor.java | 9 +- .../quorum/FollowerZooKeeperServer.java | 3 + .../zookeeper/server/quorum/Leader.java | 248 +++++++++++++++++- .../zookeeper/server/quorum/Learner.java | 6 + .../server/quorum/LearnerHandler.java | 24 +- .../server/quorum/LearnerMaster.java | 2 + .../server/quorum/LearnerSyncRequest.java | 78 ++++++ .../server/quorum/ObserverMaster.java | 7 + .../quorum/ObserverRequestProcessor.java | 3 + .../quorum/ProposalRequestProcessor.java | 12 +- .../quorum/ReadOnlyRequestProcessor.java | 2 + .../quorum/SendAckRequestProcessor.java | 9 +- .../util/RequestPathMetricsCollector.java | 6 + .../zookeeper/test/StandaloneReadTest.java | 155 +++++++++++ 28 files changed, 887 insertions(+), 32 deletions(-) create mode 100644 zookeeper-server/src/main/java/org/apache/zookeeper/client/ReadConsistencyMode.java create mode 100644 zookeeper-server/src/main/java/org/apache/zookeeper/client/ZNode.java create mode 100644 zookeeper-server/src/test/java/org/apache/zookeeper/test/StandaloneReadTest.java diff --git a/zookeeper-jute/src/main/resources/zookeeper.jute b/zookeeper-jute/src/main/resources/zookeeper.jute index 796ea396755..f64fb6e274b 100644 --- a/zookeeper-jute/src/main/resources/zookeeper.jute +++ b/zookeeper-jute/src/main/resources/zookeeper.jute @@ -223,6 +223,14 @@ module org.apache.zookeeper.proto { buffer data; org.apache.zookeeper.data.Stat stat; } + class SyncedReadResponse { + buffer data; + org.apache.zookeeper.data.Stat stat; + } + class LinearizableReadResponse { + buffer data; + org.apache.zookeeper.data.Stat stat; + } class GetChildrenResponse { vector children; } diff --git a/zookeeper-server/src/main/java/org/apache/zookeeper/AsyncCallback.java b/zookeeper-server/src/main/java/org/apache/zookeeper/AsyncCallback.java index 513238cfb59..1d391be348c 100644 --- a/zookeeper-server/src/main/java/org/apache/zookeeper/AsyncCallback.java +++ b/zookeeper-server/src/main/java/org/apache/zookeeper/AsyncCallback.java @@ -19,7 +19,9 @@ package org.apache.zookeeper; import java.util.List; +import java.util.concurrent.CompletableFuture; import org.apache.yetus.audience.InterfaceAudience; +import org.apache.zookeeper.client.ZNode; import org.apache.zookeeper.data.ACL; import org.apache.zookeeper.data.Stat; @@ -349,3 +351,22 @@ interface EphemeralsCallback extends AsyncCallback { } } + +//TODO +class DataFutureCallback implements AsyncCallback.DataCallback { + + private final CompletableFuture future; + + public DataFutureCallback(CompletableFuture future) { + this.future = future; + } + + @Override + public void processResult(int rc, String path, Object ctx, byte[] data, Stat stat) { + if (rc != KeeperException.Code.OK.intValue()) { + future.completeExceptionally(KeeperException.create(KeeperException.Code.get(rc)).fillInStackTrace()); + } else { + future.complete(new ZNode(path, data, stat)); + } + } +} diff --git a/zookeeper-server/src/main/java/org/apache/zookeeper/ClientCnxn.java b/zookeeper-server/src/main/java/org/apache/zookeeper/ClientCnxn.java index ecd2a0a9d09..83da81d5844 100644 --- a/zookeeper-server/src/main/java/org/apache/zookeeper/ClientCnxn.java +++ b/zookeeper-server/src/main/java/org/apache/zookeeper/ClientCnxn.java @@ -53,6 +53,7 @@ import org.apache.zookeeper.AsyncCallback.ChildrenCallback; import org.apache.zookeeper.AsyncCallback.Create2Callback; import org.apache.zookeeper.AsyncCallback.DataCallback; +import org.apache.zookeeper.DataFutureCallback; import org.apache.zookeeper.AsyncCallback.EphemeralsCallback; import org.apache.zookeeper.AsyncCallback.MultiCallback; import org.apache.zookeeper.AsyncCallback.StatCallback; @@ -82,12 +83,14 @@ import org.apache.zookeeper.proto.GetDataResponse; import org.apache.zookeeper.proto.GetEphemeralsResponse; import org.apache.zookeeper.proto.GetSASLRequest; +import org.apache.zookeeper.proto.LinearizableReadResponse; import org.apache.zookeeper.proto.ReplyHeader; import org.apache.zookeeper.proto.RequestHeader; import org.apache.zookeeper.proto.SetACLResponse; import org.apache.zookeeper.proto.SetDataResponse; import org.apache.zookeeper.proto.SetWatches; import org.apache.zookeeper.proto.SetWatches2; +import org.apache.zookeeper.proto.SyncedReadResponse; import org.apache.zookeeper.proto.WatcherEvent; import org.apache.zookeeper.server.ByteBufferInputStream; import org.apache.zookeeper.server.ZooKeeperThread; @@ -723,6 +726,14 @@ private void processEvent(Object event) { } else { cb.processResult(rc, p.ctx, null); } + } else if (p.response instanceof SyncedReadResponse) { + DataFutureCallback cb = (DataFutureCallback) p.cb; + SyncedReadResponse rsp = (SyncedReadResponse) p.response; + cb.processResult(rc, clientPath, p.ctx, rsp.getData(), rsp.getStat()); + } else if (p.response instanceof LinearizableReadResponse) { + DataFutureCallback cb = (DataFutureCallback) p.cb; + LinearizableReadResponse rsp = (LinearizableReadResponse) p.response; + cb.processResult(rc, clientPath, p.ctx, rsp.getData(), rsp.getStat()); } else if (p.cb instanceof VoidCallback) { VoidCallback cb = (VoidCallback) p.cb; cb.processResult(rc, clientPath, p.ctx); diff --git a/zookeeper-server/src/main/java/org/apache/zookeeper/ZooDefs.java b/zookeeper-server/src/main/java/org/apache/zookeeper/ZooDefs.java index 9cf70787668..7eeca7625fe 100644 --- a/zookeeper-server/src/main/java/org/apache/zookeeper/ZooDefs.java +++ b/zookeeper-server/src/main/java/org/apache/zookeeper/ZooDefs.java @@ -55,6 +55,9 @@ public interface OpCode { int sync = 9; + // sync + int syncedRead = 10; + int ping = 11; int getChildren2 = 12; @@ -95,6 +98,9 @@ public interface OpCode { int whoAmI = 107; + // linearable + int linearizableRead = 108; + int createSession = -10; int closeSession = -11; diff --git a/zookeeper-server/src/main/java/org/apache/zookeeper/ZooKeeper.java b/zookeeper-server/src/main/java/org/apache/zookeeper/ZooKeeper.java index 9fba7a5ebd4..22e5b4e7191 100644 --- a/zookeeper-server/src/main/java/org/apache/zookeeper/ZooKeeper.java +++ b/zookeeper-server/src/main/java/org/apache/zookeeper/ZooKeeper.java @@ -29,6 +29,7 @@ import java.util.List; import java.util.Map; import java.util.Set; +import java.util.concurrent.CompletableFuture; import org.apache.jute.Record; import org.apache.yetus.audience.InterfaceAudience; import org.apache.zookeeper.AsyncCallback.ACLCallback; @@ -44,8 +45,10 @@ import org.apache.zookeeper.Watcher.WatcherType; import org.apache.zookeeper.client.ConnectStringParser; import org.apache.zookeeper.client.HostProvider; +import org.apache.zookeeper.client.ReadConsistencyMode; import org.apache.zookeeper.client.StaticHostProvider; import org.apache.zookeeper.client.ZKClientConfig; +import org.apache.zookeeper.client.ZNode; import org.apache.zookeeper.client.ZooKeeperSaslClient; import org.apache.zookeeper.common.PathUtils; import org.apache.zookeeper.data.ACL; @@ -72,6 +75,7 @@ import org.apache.zookeeper.proto.GetDataResponse; import org.apache.zookeeper.proto.GetEphemeralsRequest; import org.apache.zookeeper.proto.GetEphemeralsResponse; +import org.apache.zookeeper.proto.LinearizableReadResponse; import org.apache.zookeeper.proto.RemoveWatchesRequest; import org.apache.zookeeper.proto.ReplyHeader; import org.apache.zookeeper.proto.RequestHeader; @@ -81,6 +85,7 @@ import org.apache.zookeeper.proto.SetDataResponse; import org.apache.zookeeper.proto.SyncRequest; import org.apache.zookeeper.proto.SyncResponse; +import org.apache.zookeeper.proto.SyncedReadResponse; import org.apache.zookeeper.proto.WhoAmIResponse; import org.apache.zookeeper.server.DataTree; import org.apache.zookeeper.server.EphemeralType; @@ -2026,6 +2031,87 @@ public void getData(final String path, Watcher watcher, DataCallback cb, Object cnxn.queuePacket(h, new ReplyHeader(), request, response, cb, clientPath, serverPath, ctx, wcb); } + /** + * TODO + * @param path + * @param watcher + * @param readMode + * @return + * @throws KeeperException + * @throws InterruptedException + * @since + */ + public CompletableFuture getData(ReadConsistencyMode readMode, final String path, Watcher watcher) throws KeeperException, InterruptedException { + CompletableFuture future = new CompletableFuture<>(); + + switch (readMode) { + case SEQUENTIAL_READ: + Stat stat = new Stat(); + byte[] data = getData(path, watcher, stat); + ZNode n = new ZNode(path, data, stat); + future.complete(n); + return future; + case ORDERED_SEQUENTIAL_READ: + return syncedRead(path, watcher); + case LINEARIZABLE_READ: + return linearizableRead(path, watcher); + default: + throw new IllegalArgumentException("Invalid Read Consistency Mode: " + readMode.getName()); + } + } + + private CompletableFuture syncedRead(final String path, Watcher watcher) { + + final String clientPath = path; + PathUtils.validatePath(clientPath); + + // the watch contains the un-chroot path + WatchRegistration wcb = null; + if (watcher != null) { + wcb = new DataWatchRegistration(watcher, clientPath); + } + + final String serverPath = prependChroot(clientPath); + + RequestHeader h = new RequestHeader(); + h.setType(ZooDefs.OpCode.syncedRead); + GetDataRequest request = new GetDataRequest(); + SyncedReadResponse response = new SyncedReadResponse(); + request.setPath(serverPath); + request.setWatch(watcher != null); + CompletableFuture future = new CompletableFuture<>(); + DataFutureCallback callback = new DataFutureCallback(future); + cnxn.queuePacket(h, new ReplyHeader(), request, response, callback, clientPath, serverPath, null, wcb); + + return future; + } + + private CompletableFuture linearizableRead(final String path, Watcher watcher) { + + final String clientPath = path; + PathUtils.validatePath(clientPath); + + // the watch contains the un-chroot path + WatchRegistration wcb = null; + if (watcher != null) { + wcb = new DataWatchRegistration(watcher, clientPath); + } + + final String serverPath = prependChroot(clientPath); + + RequestHeader h = new RequestHeader(); + h.setType(ZooDefs.OpCode.linearizableRead); + GetDataRequest request = new GetDataRequest(); + LinearizableReadResponse response = new LinearizableReadResponse(); + request.setPath(serverPath); + request.setWatch(watcher != null); + CompletableFuture future = new CompletableFuture<>(); + DataFutureCallback callback = new DataFutureCallback(future); + cnxn.queuePacket(h, new ReplyHeader(), request, response, callback, clientPath, serverPath, null, wcb); + + return future; + } + /** * The asynchronous version of getData. * diff --git a/zookeeper-server/src/main/java/org/apache/zookeeper/client/ReadConsistencyMode.java b/zookeeper-server/src/main/java/org/apache/zookeeper/client/ReadConsistencyMode.java new file mode 100644 index 00000000000..f230eacab40 --- /dev/null +++ b/zookeeper-server/src/main/java/org/apache/zookeeper/client/ReadConsistencyMode.java @@ -0,0 +1,55 @@ +/* + * 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.client; + +import org.apache.yetus.audience.InterfaceAudience; + +/** + * TODO + * @since 3.7.0 + */ +@InterfaceAudience.Public +public enum ReadConsistencyMode { + SEQUENTIAL_READ("Sequential Read"), + ORDERED_SEQUENTIAL_READ("Ordered Sequential Read"), + LINEARIZABLE_READ("Linearizable Read"), + LEASE_READ("Lease Read"), + DUMMY_READ("Dummy Read Used For Testing"); + + public static final ReadConsistencyMode DEFAULT_READ_MODE = SEQUENTIAL_READ; + + private String name; + + ReadConsistencyMode(String name) { + this.name = name; + } + + public String getName() { + return name; + } + + public static ReadConsistencyMode fromString(String name) { + for (ReadConsistencyMode c : values()) { + if (c.getName().compareToIgnoreCase(name) == 0) { + return c; + } + } + return DEFAULT_READ_MODE; + } +} \ No newline at end of file diff --git a/zookeeper-server/src/main/java/org/apache/zookeeper/client/ZNode.java b/zookeeper-server/src/main/java/org/apache/zookeeper/client/ZNode.java new file mode 100644 index 00000000000..71e1bcf959b --- /dev/null +++ b/zookeeper-server/src/main/java/org/apache/zookeeper/client/ZNode.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.client; + +import java.util.Arrays; +import edu.umd.cs.findbugs.annotations.SuppressFBWarnings; +import org.apache.yetus.audience.InterfaceAudience; +import org.apache.zookeeper.data.Stat; + +/** + * TODO + * @since 3.7.0 + * + */ +@InterfaceAudience.Public +@SuppressFBWarnings(value = {"EI_EXPOSE_REP", "EI_EXPOSE_REP2"}, justification = "The field:data is mutable object, " + + "but a clone of it may have performance issue") +public class ZNode { + + private final String path; + private final byte[] data; + private final Stat stat; + + public ZNode(String path, byte[] data, Stat stat) { + this.path = path; + this.data = data;//.clone() + this.stat = stat; + } + + public String getPath() { + return path; + } + + public byte[] getData() { + return data;//.clone() + } + + public Stat getStat() { + return stat; + } + + @Override + public String toString() { + return "ZNode{" + + "path='" + path + '\'' + + ", data=" + Arrays.toString(data) + + ", stat=" + stat + + '}'; + } +} \ No newline at end of file diff --git a/zookeeper-server/src/main/java/org/apache/zookeeper/server/FinalRequestProcessor.java b/zookeeper-server/src/main/java/org/apache/zookeeper/server/FinalRequestProcessor.java index 6db245e060d..3c31b748573 100644 --- a/zookeeper-server/src/main/java/org/apache/zookeeper/server/FinalRequestProcessor.java +++ b/zookeeper-server/src/main/java/org/apache/zookeeper/server/FinalRequestProcessor.java @@ -69,6 +69,7 @@ import org.apache.zookeeper.proto.GetDataResponse; import org.apache.zookeeper.proto.GetEphemeralsRequest; import org.apache.zookeeper.proto.GetEphemeralsResponse; +import org.apache.zookeeper.proto.LinearizableReadResponse; import org.apache.zookeeper.proto.RemoveWatchesRequest; import org.apache.zookeeper.proto.ReplyHeader; import org.apache.zookeeper.proto.SetACLResponse; @@ -77,6 +78,7 @@ import org.apache.zookeeper.proto.SetWatches2; import org.apache.zookeeper.proto.SyncRequest; import org.apache.zookeeper.proto.SyncResponse; +import org.apache.zookeeper.proto.SyncedReadResponse; import org.apache.zookeeper.proto.WhoAmIResponse; import org.apache.zookeeper.server.DataTree.ProcessTxnResult; import org.apache.zookeeper.server.quorum.QuorumZooKeeperServer; @@ -283,7 +285,7 @@ public void processRequest(Request request) { subResult = new GetChildrenResult(((GetChildrenResponse) rec).getChildren()); break; case OpCode.getData: - rec = handleGetDataRequest(readOp.toRequestRecord(), cnxn, request.authInfo); + rec = handleReadRequest(readOp.toRequestRecord(), cnxn, request.authInfo, request.type); GetDataResponse gdr = (GetDataResponse) rec; subResult = new GetDataResult(gdr.getData(), gdr.getStat()); break; @@ -380,10 +382,30 @@ public void processRequest(Request request) { GetDataRequest getDataRequest = new GetDataRequest(); ByteBufferInputStream.byteBuffer2Record(request.request, getDataRequest); path = getDataRequest.getPath(); - rsp = handleGetDataRequest(getDataRequest, cnxn, request.authInfo); + rsp = handleReadRequest(getDataRequest, cnxn, request.authInfo, request.type); requestPathMetricsCollector.registerRequest(request.type, path); break; } + case OpCode.syncedRead: { + lastOp = "SYRD"; + GetDataRequest syncedReadRequest = new GetDataRequest(); + //System.out.println("fuck_FinalRequestProcessor_receive_OpCode.syncedRead request.request" + request.request); + ByteBufferInputStream.byteBuffer2Record(request.request, syncedReadRequest); + path = syncedReadRequest.getPath(); + rsp = handleReadRequest(syncedReadRequest, cnxn, request.authInfo, request.type); + requestPathMetricsCollector.registerRequest(request.type, syncedReadRequest.getPath()); + break; + } + case OpCode.linearizableRead: { + lastOp = "LIRD"; + GetDataRequest linearizableReadRequest = new GetDataRequest(); + //System.out.println("fuck_FinalRequestProcessor_receive_OpCode.linearizableRead request.request" + request.request); + ByteBufferInputStream.byteBuffer2Record(request.request, linearizableReadRequest); + path = linearizableReadRequest.getPath(); + rsp = handleReadRequest(linearizableReadRequest, cnxn, request.authInfo, request.type); + requestPathMetricsCollector.registerRequest(request.type, linearizableReadRequest.getPath()); + break; + } case OpCode.setWatches: { lastOp = "SETW"; SetWatches setWatches = new SetWatches(); @@ -655,17 +677,32 @@ private Record handleGetChildrenRequest(Record request, ServerCnxn cnxn, List authInfo) throws KeeperException, IOException { - GetDataRequest getDataRequest = (GetDataRequest) request; - String path = getDataRequest.getPath(); + private Record handleReadRequest(Record request, ServerCnxn cnxn, List authInfo, int requestType) throws KeeperException, IOException { + GetDataRequest readRequest = (GetDataRequest) request; + + String path = readRequest.getPath(); DataNode n = zks.getZKDatabase().getNode(path); if (n == null) { throw new KeeperException.NoNodeException(); } zks.checkACL(cnxn, zks.getZKDatabase().aclForNode(n), ZooDefs.Perms.READ, authInfo, path, null); Stat stat = new Stat(); - byte[] b = zks.getZKDatabase().getData(path, stat, getDataRequest.getWatch() ? cnxn : null); - return new GetDataResponse(b, stat); + byte[] b = zks.getZKDatabase().getData(path, stat, readRequest.getWatch() ? cnxn : null); + + switch (requestType) { + case OpCode.getData: { + return new GetDataResponse(b, stat); + } + case OpCode.syncedRead: { + return new SyncedReadResponse(b, stat); + } + case OpCode.linearizableRead: { + return new LinearizableReadResponse(b, stat); + } + default: + throw new IllegalArgumentException("Invalid request type " + requestType); + } + //return null; } private boolean closeSession(ServerCnxnFactory serverCnxnFactory, long sessionId) { @@ -692,5 +729,4 @@ private void updateStats(Request request, String lastOp, long lastZxid) { zks.serverStats().updateLatency(request, currentTime); request.cnxn.updateStatsForResponse(request.cxid, lastZxid, lastOp, request.createTime, currentTime); } - } 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 6eb7b96ec53..5039813e5ac 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 @@ -902,6 +902,8 @@ private void pRequestHelper(Request request) throws RequestProcessorException { //All the rest don't need to create a Txn - just verify session case OpCode.sync: + case OpCode.syncedRead: + case OpCode.linearizableRead: case OpCode.exists: case OpCode.getData: case OpCode.getACL: 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 a68203b207c..4e19e48197c 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 @@ -269,6 +269,8 @@ static boolean isValid(int type) { case OpCode.setWatches: case OpCode.setWatches2: case OpCode.sync: + case OpCode.syncedRead: + case OpCode.linearizableRead: case OpCode.checkWatches: case OpCode.removeWatches: case OpCode.addWatch: @@ -362,14 +364,20 @@ public static String op2String(int op) { return "auth"; case OpCode.setWatches: return "setWatches"; - case OpCode.setWatches2: - return "setWatches2"; case OpCode.sasl: return "sasl"; case OpCode.getEphemerals: return "getEphemerals"; case OpCode.getAllChildrenNumber: return "getAllChildrenNumber"; + case OpCode.setWatches2: + return "setWatches2"; + case OpCode.addWatch: + return "addWatch"; + case OpCode.syncedRead: + return "syncedRead"; + case OpCode.linearizableRead: + return "linearizableRead"; case OpCode.createSession: return "createSession"; case OpCode.closeSession: diff --git a/zookeeper-server/src/main/java/org/apache/zookeeper/server/RequestThrottler.java b/zookeeper-server/src/main/java/org/apache/zookeeper/server/RequestThrottler.java index 32863d92a61..3c5d805c36c 100644 --- a/zookeeper-server/src/main/java/org/apache/zookeeper/server/RequestThrottler.java +++ b/zookeeper-server/src/main/java/org/apache/zookeeper/server/RequestThrottler.java @@ -18,8 +18,8 @@ package org.apache.zookeeper.server; -import edu.umd.cs.findbugs.annotations.SuppressFBWarnings; import java.util.concurrent.LinkedBlockingQueue; +import edu.umd.cs.findbugs.annotations.SuppressFBWarnings; import org.apache.zookeeper.common.Time; import org.apache.zookeeper.util.ServiceUtils; import org.slf4j.Logger; 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 7ef687ed740..8505e5a0b38 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 @@ -18,7 +18,6 @@ package org.apache.zookeeper.server; -import edu.umd.cs.findbugs.annotations.SuppressFBWarnings; import java.io.ByteArrayOutputStream; import java.io.File; import java.io.IOException; @@ -39,6 +38,7 @@ import java.util.concurrent.atomic.AtomicLong; import java.util.function.BiConsumer; import javax.security.sasl.SaslException; +import edu.umd.cs.findbugs.annotations.SuppressFBWarnings; import org.apache.jute.BinaryInputArchive; import org.apache.jute.BinaryOutputArchive; import org.apache.jute.Record; @@ -1767,7 +1767,6 @@ public ProcessTxnResult processTxn(Request request) { final boolean writeRequest = (hdr != null); final boolean quorumRequest = request.isQuorum(); - // return fast w/o synchronization when we get a read if (!writeRequest && !quorumRequest) { return new ProcessTxnResult(); diff --git a/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/CommitProcessor.java b/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/CommitProcessor.java index 86dce2b7f85..392ebf9f15e 100644 --- a/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/CommitProcessor.java +++ b/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/CommitProcessor.java @@ -18,7 +18,6 @@ package org.apache.zookeeper.server.quorum; -import edu.umd.cs.findbugs.annotations.SuppressFBWarnings; import java.util.ArrayDeque; import java.util.Deque; import java.util.HashMap; @@ -27,6 +26,7 @@ import java.util.Set; import java.util.concurrent.LinkedBlockingQueue; import java.util.concurrent.atomic.AtomicInteger; +import edu.umd.cs.findbugs.annotations.SuppressFBWarnings; import org.apache.zookeeper.ZooDefs.OpCode; import org.apache.zookeeper.common.Time; import org.apache.zookeeper.server.ExitCode; @@ -184,6 +184,8 @@ protected boolean needCommit(Request request) { case OpCode.check: return true; case OpCode.sync: + case OpCode.syncedRead: + case OpCode.linearizableRead: return matchSyncs; case OpCode.createSession: case OpCode.closeSession: @@ -215,6 +217,8 @@ public void run() { */ commitIsWaiting = !committedRequests.isEmpty(); requestsToProcess = queuedRequests.size(); + //System.out.println("fuck_CommitProcessor commitIsWaiting " + commitIsWaiting + ",committedRequests: " + committedRequests + //+",queuedRequests: " + queuedRequests + ", queuedWriteRequests: " + queuedWriteRequests); // Avoid sync if we have something to do if (requestsToProcess == 0 && !commitIsWaiting) { // Waiting for requests to process @@ -249,14 +253,18 @@ public void run() { && (maxReadBatchSize < 0 || readsProcessed <= maxReadBatchSize) && (request = queuedRequests.poll()) != null) { requestsToProcess--; + //System.out.println("fuck__CommitProcessor_needCommit(request):" + needCommit(request) + ",pendingRequests.containsKey(request.sessionId):"+ + //pendingRequests.containsKey(request.sessionId)); if (needCommit(request) || pendingRequests.containsKey(request.sessionId)) { // Add request to pending Deque requests = pendingRequests.computeIfAbsent(request.sessionId, sid -> new ArrayDeque<>()); requests.addLast(request); + //System.out.println("fuck_CommitProcessor commitIsWaiting got a write request request:" + request); ServerMetrics.getMetrics().REQUESTS_IN_SESSION_QUEUE.add(requests.size()); } else { readsProcessed++; numReadQueuedRequests.decrementAndGet(); + //System.out.println("fuck_CommitProcessor commitIsWaiting got a read request request(Or don't needCommit):" + request); sendToNextProcessor(request); } /* @@ -284,7 +292,7 @@ public void run() { if (!commitIsWaiting) { commitIsWaiting = !committedRequests.isEmpty(); } - + //System.out.println("fuck_CommitProcessor commitIsWaiting " + commitIsWaiting); /* * Handle commits, if any. */ @@ -323,12 +331,14 @@ public void run() { * a commit for a local write, as commits are received in order. Else * it must be a commit for a remote write. */ + //System.out.println("fuck_CommitProcessor queuedWriteRequests: " + queuedWriteRequests); if (!queuedWriteRequests.isEmpty() && queuedWriteRequests.peek().sessionId == request.sessionId && queuedWriteRequests.peek().cxid == request.cxid) { /* * Commit matches the earliest write in our write queue. */ + //System.out.println("fuck_CommitProcessor Commit matches(I love it) request:"+request); Deque sessionQueue = pendingRequests.get(request.sessionId); ServerMetrics.getMetrics().PENDING_SESSION_QUEUE_SIZE.add(pendingRequests.size()); if (sessionQueue == null || sessionQueue.isEmpty() || !needCommit(sessionQueue.peek())) { @@ -337,6 +347,9 @@ public void run() { * Either there are reads pending in this session, or we * haven't gotten to this write yet. */ + System.out.println("* Can't process this write yet.\n" + + " * Either there are reads pending in this session, or we\n" + + " * haven't gotten to this write yet."); break; } else { ServerMetrics.getMetrics().REQUESTS_IN_SESSION_QUEUE.add(sessionQueue.size()); @@ -365,6 +378,7 @@ public void run() { } // Only decrement if we take a request off the queue. numWriteQueuedRequests.decrementAndGet(); + //System.out.println("fuck_CommitProcessor queuedWriteRequests.poll() and add to queuesToDrain"); queuedWriteRequests.poll(); queuesToDrain.add(request.sessionId); } @@ -378,6 +392,7 @@ public void run() { commitsProcessed++; // Process the write inline. + //System.out.println("fuck_CommitProcessor processWrite(request) request:"+request); processWrite(request); commitIsWaiting = !committedRequests.isEmpty(); @@ -391,10 +406,13 @@ public void run() { * empty. */ readsProcessed = 0; + //System.out.println("fuck_CommitProcessor queuesToDrain:"+queuesToDrain); for (Long sessionId : queuesToDrain) { Deque sessionQueue = pendingRequests.get(sessionId); + //System.out.println("fuck_CommitProcessor sessionQueue:"+sessionQueue); int readsAfterWrite = 0; while (!stopped && !sessionQueue.isEmpty() && !needCommit(sessionQueue.peek())) { + //System.out.println("fuck_CommitProcessor queuesToDrain -> sendToNextProcessor request:"+sessionQueue.peek()); numReadQueuedRequests.decrementAndGet(); sendToNextProcessor(sessionQueue.poll()); readsAfterWrite++; @@ -474,6 +492,7 @@ private void processWrite(Request request) throws RequestProcessorException { processCommitMetrics(request, true); long timeBeforeFinalProc = Time.currentElapsedTime(); + //System.out.println("fuck_CommitProcessor_processWrite_send_to_nextProcessor"); nextProcessor.processRequest(request); ServerMetrics.getMetrics().WRITE_FINAL_PROC_TIME.add(Time.currentElapsedTime() - timeBeforeFinalProc); } @@ -604,7 +623,7 @@ public void processRequest(Request request) { if (stopped) { return; } - LOG.debug("Processing request:: {}", request); + //LOG.info("fuck_CommitProcessor#processRequest Processing request:: {}", request); request.commitProcQueueStartTime = Time.currentElapsedTime(); queuedRequests.add(request); // If the request will block, add it to the queue of blocking requests 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 971710c91ab..aa49d529acb 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 @@ -162,6 +162,9 @@ protected void processPacket(QuorumPacket qp) throws Exception { case Leader.PING: ping(qp); break; + case Leader.HEARTBEAT: + sendHeartBeat(qp); + break; case Leader.PROPOSAL: ServerMetrics.getMetrics().LEARNER_PROPOSAL_RECEIVED_COUNT.add(1); TxnLogEntry logEntry = SerializeUtils.deserializeTxn(qp.getData()); @@ -181,7 +184,7 @@ protected void processPacket(QuorumPacket qp) throws Exception { QuorumVerifier qv = self.configFromString(new String(setDataTxn.getData(), UTF_8)); self.setLastSeenQuorumVerifier(qv, true); } - + //System.out.println("fuck_I'm Follower I receive the PROPOSAL from the leader hdr:"+hdr); fzk.logRequest(hdr, txn, digest); if (hdr != null) { /* @@ -202,6 +205,7 @@ protected void processPacket(QuorumPacket qp) throws Exception { } break; case Leader.COMMIT: + //System.out.println("fuck_Follower#processPacket I'm follower I receive leader's COMMIT"); ServerMetrics.getMetrics().LEARNER_COMMIT_RECEIVED_COUNT.add(1); fzk.commit(qp.getZxid()); if (om != null) { diff --git a/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/FollowerRequestProcessor.java b/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/FollowerRequestProcessor.java index 58ca9905dfe..f92fcfb5c4b 100644 --- a/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/FollowerRequestProcessor.java +++ b/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/FollowerRequestProcessor.java @@ -83,7 +83,6 @@ public void run() { // the request to the leader so that we are ready to receive // the response maybeSendRequestToNextProcessor(request); - if (request.isThrottled()) { continue; } @@ -91,11 +90,14 @@ public void run() { // We now ship the request to the leader. As with all // other quorum operations, sync also follows this code // path, but different from others, we need to keep track - // of the sync operations this follower has pending, so we + // of the sync/?/? operations this follower has pending, so we // add it to pendingSyncs. switch (request.type) { case OpCode.sync: + case OpCode.syncedRead: + case OpCode.linearizableRead: zks.pendingSyncs.add(request); + //System.out.println("fuck_FollowerRequestProcessor follow get SYNC/syncedRead/linearizableRead request forward to leader, add to pendingSyncs queue request:"+request); zks.getFollower().request(request); break; case OpCode.create: @@ -132,6 +134,8 @@ private void maybeSendRequestToNextProcessor(Request request) throws RequestProc if (skipLearnerRequestToNextProcessor && request.isFromLearner()) { ServerMetrics.getMetrics().SKIP_LEARNER_REQUEST_TO_NEXT_PROCESSOR_COUNT.add(1); } else { + //System.out.println("fuck_FollowerRequestProcessor nextProcessor.processRequest to commit-RequestProcessor request:"+request + //+",skipLearnerRequestToNextProcessor:"+skipLearnerRequestToNextProcessor+",request.isFromLearner:"+request.isFromLearner()); nextProcessor.processRequest(request); } } @@ -141,6 +145,7 @@ public void processRequest(Request request) { } void processRequest(Request request, boolean checkForUpgrade) { + //System.out.println("fuck_FollowerRequestProcessor#processRequest Request:" + request); if (!finished) { if (checkForUpgrade) { // Before sending the request, check if the request requires a 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..84adf4ef429 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 @@ -85,6 +85,7 @@ public void logRequest(TxnHeader hdr, Record txn, TxnDigest digest) { if ((request.zxid & 0xffffffffL) != 0) { pendingTxns.add(request); } + //System.out.println("fuck_FollowerZooKeeperServer#syncProcessor.processRequest request:"+request); syncProcessor.processRequest(request); } @@ -119,8 +120,10 @@ public synchronized void sync() { Request r = pendingSyncs.remove(); if (r instanceof LearnerSyncRequest) { LearnerSyncRequest lsr = (LearnerSyncRequest) r; + //System.out.println("fuck_FollowerZooKeeperServer send SYNC????(when this is called wuwuw ObserverMaster???) LearnerSyncRequest:"+r); lsr.fh.queuePacket(new QuorumPacket(Leader.SYNC, 0, null, null)); } + //System.out.println("fuck_Follower#processPacket I'm follower I receive leader's SYNC and apply it to state space(pendingSyncs.remove())!!!#commitProcessor.commit(r) Request:" + r); commitProcessor.commit(r); } 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 2de2ceeb8b0..ddc8fd6b98f 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 @@ -3,7 +3,7 @@ * 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 + * to you under the Apache License, Version 2.0 (theLeader * "License"); you may not use this file except in compliance * with the License. You may obtain a copy of the License at * @@ -51,6 +51,8 @@ import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicLong; +import java.util.concurrent.locks.Lock; +import java.util.concurrent.locks.ReentrantLock; import java.util.stream.Collectors; import javax.security.sasl.SaslException; import org.apache.zookeeper.KeeperException; @@ -99,6 +101,57 @@ public String toString() { } + public static class ReadStateContext { + long lastProposed; + long sessionId; + int xid; + int type; + + //SyncedLearnerTracker syncedAckSet; + + public ReadStateContext(long lastProposed, long sessionId, int xid, int type) { + this.lastProposed = lastProposed; + this.sessionId = sessionId; + this.xid = xid; + this.type = type; + //syncedAckSet = new SyncedLearnerTracker(); + } + + @Override + public boolean equals(Object o) { + if (this == o) return true; + if (o == null || getClass() != o.getClass()) return false; + ReadStateContext that = (ReadStateContext) o; + return lastProposed == that.lastProposed && + sessionId == that.sessionId && + xid == that.xid && + type == that.type; + } + + @Override + public int hashCode() { + return Objects.hash(lastProposed, sessionId, xid, type); + } + + @Override + public String toString() { + return "ReadStateContext{" + + "lastProposed=" + lastProposed + + ", sessionId=" + sessionId + + ", xid=" + xid + + ", type=" + type + + '}'; + } + +// public byte[] getBytes() { +// return this.toString().getBytes(UTF_8); +// } + +// public SyncedLearnerTracker getSyncedAckSet() { +// return syncedAckSet; +// } + } + // log ack latency if zxid is a multiple of ackLoggingFrequency. If <=0, disable logging. private static final String ACK_LOGGING_FREQUENCY = "zookeeper.leader.ackLoggingFrequency"; private static int ackLoggingFrequency; @@ -215,6 +268,9 @@ public void resetObserverConnectionStats() { // Pending sync requests. Must access under 'this' lock. private final Map> pendingSyncs = new HashMap>(); + //private final Map> waitingQuorumSyncs = new ConcurrentHashMap<>(1000); + private final Map waitingQuorumSyncs = new ConcurrentHashMap<>();//new ConcurrentHashMap<>(1000); + public synchronized int getNumPendingSyncs() { return pendingSyncs.size(); } @@ -421,6 +477,9 @@ Optional createServerSocket(InetSocketAddress address, boolean por */ static final int COMMITANDACTIVATE = 9; + // TODO(Heartbeat) + static final int HEARTBEAT = 10; + /** * Similar to INFORM, only for a reconfig operation. */ @@ -451,7 +510,7 @@ public void run() { if (!stop.get() && !serverSockets.isEmpty()) { ExecutorService executor = Executors.newFixedThreadPool(serverSockets.size()); CountDownLatch latch = new CountDownLatch(serverSockets.size()); - + System.out.println("fuck_Leader#LearnerCnxAcceptor#serverSockets.size():" + serverSockets.size()); serverSockets.forEach(serverSocket -> executor.submit(new LearnerCnxAcceptorHandler(serverSocket, latch))); @@ -749,7 +808,7 @@ void lead() throws IOException, InterruptedException { syncedAckSet.addAck(f.getSid()); } } - + //System.out.println("fuck_Leader#this.isRunning():"+this.isRunning()); // check leader running status if (!this.isRunning()) { // set shutdown flag @@ -767,6 +826,7 @@ void lead() throws IOException, InterruptedException { } tickSkip = !tickSkip; } + //System.out.println("fuck_Leader#ping_all_followers"); for (LearnerHandler f : getLearners()) { f.ping(); } @@ -949,13 +1009,28 @@ public synchronized boolean tryToCommit(Proposal p, long zxid, SocketAddress fol informAndActivate(p, designatedLeader); } else { p.request.logLatency(ServerMetrics.getMetrics().QUORUM_ACK_LATENCY); + //System.out.println("fuck_Leader I receive enough acks start to send commit to followes (zxid): "+zxid); commit(zxid); inform(p); } zk.commitProcessor.commit(p.request); + + //TODO 这一段整体封装成一个方法 + //System.out.println("fuck_Leader zk.commitProcessor.commit(p.request): "+p.request +",pendingSyncs.containsKey(zxid):" + pendingSyncs.containsKey(zxid)); if (pendingSyncs.containsKey(zxid)) { + System.out.println("fuck_Leader pendingSyncs.containsKey(zxid) send SYNC to all learns(III), we have committed zxid: " + + zxid + ",pendingSyncs.size: "+ pendingSyncs.get(zxid).size()); for (LearnerSyncRequest r : pendingSyncs.remove(zxid)) { - sendSync(r); + if (r.type == OpCode.linearizableRead) { //&& r.isFromLeader() + //System.out.println("fuck_tryToCommit@ syncRequest.fh == null zk.commitProcessor.commit(syncRequest) at leader locally(linearizableRead) " + + // "r.setIsCommited(true); r.getCountDownLatch().countDown()"); + //zk.commitProcessor.commit(r); + r.setCanCommitted(true); + r.getCountDownLatch().countDown(); + } else { + //System.out.println("fuck_tryToCommit@ sendSync to other followers r.type:" + r.type); + sendSync(r); + } } } @@ -972,6 +1047,7 @@ public synchronized boolean tryToCommit(Proposal p, long zxid, SocketAddress fol */ @Override public synchronized void processAck(long sid, long zxid, SocketAddress followerAddr) { + //LOG.info("fuck_processAck_at_the_begin, ThreadName:" + Thread.currentThread().getName()+",sid:"+sid); if (!allowedToCommit) { return; // last op committed was a leader change - from now on } @@ -1039,6 +1115,108 @@ public synchronized void processAck(long sid, long zxid, SocketAddress followerA } } } + //LOG.info("fuck_processAck_at_the_end, ThreadName:" + Thread.currentThread().getName() + ",sid:"+sid); + } + + private Lock lock = new ReentrantLock(); + + @Override + public void processHeartbeat(long sid, QuorumPacket qp) { + // 获取锁 + lock.lock(); + try { + processHeartbeatInline(sid, qp); + } finally { + lock.unlock(); + } + } + + public void processHeartbeatInline(long sid, QuorumPacket qp) { + //LOG.info("fuck_processHeartbeat_at_the_begin, ThreadName:" + Thread.currentThread().getName()+",sid:"+sid); + + //public synchronized 问题问题问题问题问题问题问题 + //Deque learnerSyncRequestsQueue = waitingQuorumSyncs.get(zxid); + //方法 + ByteArrayInputStream bis = new ByteArrayInputStream(qp.getData()); + DataInputStream dis = new DataInputStream(bis); + long lastProposed = -1L; + long sessionId = -1L; + int xid = -1; + int type = -1; + try { + lastProposed = dis.readLong(); + sessionId = dis.readLong(); + xid = dis.readInt(); + type = dis.readInt(); + } catch (IOException e) { + LOG.warn("Unexpected exception", e); + } + + ReadStateContext readStateContext = new ReadStateContext(lastProposed, sessionId, xid, type); + LearnerSyncRequest learnerSyncRequest = waitingQuorumSyncs.get(readStateContext); + //System.out.println("fuck_leader_processPing zxid:" + zxid+",sid:"+sid+",ThreadName:" + Thread.currentThread().getName()); + //Iterator iterator = learnerSyncRequestsQueue.iterator(); + if (learnerSyncRequest == null) { + //System.out.println("fuck_processPing#learnerSyncRequest == null return;lastProposed:" + lastProposed+",sid:"+sid+",ThreadName:" + Thread.currentThread().getName()); + return; + } +// if (learnerSyncRequestsQueue.isEmpty()) { +// System.out.println("fuck_learnerSyncRequestsQueue.isEmpty() return;zxid:" + zxid+",sid:"+sid+",ThreadName:" + Thread.currentThread().getName()); +// waitingQuorumSyncs.remove(zxid); +// return; +// } + +// while (!learnerSyncRequestsQueue.isEmpty()) { +// +// } + //LearnerSyncRequest syncRequest = learnerSyncRequestsQueue.peek(); + learnerSyncRequest.getSyncedAckSet().addAck(sid); + +// System.out.println("fuck_Leader_syncRequest.syncedAckSet_" +// + "sid:" + sid +// + ",ThreadName:" + Thread.currentThread().getName() +// + ",syncedAckSet:" + learnerSyncRequest.syncedAckSet.ackSetsToString() +// //+ ", queue.size():"+ syncRequest.size() +// + ",I'm leader and I receive the follower's Heartbeart" + //); + + if (!learnerSyncRequest.getSyncedAckSet().hasAllQuorums()) {// || syncRequest.isQuorum + return; + } + System.out.println("fuck_Leader syncRequest.syncedAckSet: Leader gets a quorum of ping(I love it)" + learnerSyncRequest.syncedAckSet.ackSetsToString() + +",ThreadName:" + Thread.currentThread().getName() +", readStateContext:"+readStateContext); + System.out.println();System.out.println();System.out.println();System.out.println();System.out.println(); +// if(syncRequest.initialized.compareAndSet(false, true)) { +// +// } + waitingQuorumSyncs.remove(readStateContext); + //learnerSyncRequestsQueue.poll(); + //syncRequest.setQuorum(true); + + //III + long start = Time.currentWallTime(); + System.out.println("fuck_processHeartbeat_start_to_try_to_commit_lastProposed-zxid:"+lastProposed+"_canCommitted:" + learnerSyncRequest.canCommitted + + ", learnerSyncRequest.isFromLeader: " + learnerSyncRequest.isFromLeader()); + if (!learnerSyncRequest.canCommitted) { + //wait for apply + try { + learnerSyncRequest.getCountDownLatch().await(5000, TimeUnit.MILLISECONDS); + System.out.println("fuck_processHeartbeat_await_time:" + (System.currentTimeMillis() - start)+" ms"); + //learnerSyncRequest.getCountDownLatch().await(5000, TimeUnit.MILLISECONDS); + } catch (InterruptedException e) { + e.printStackTrace(); + //TODO timeout + return; + } + } + + if (learnerSyncRequest.isFromLeader()) { + System.out.println("processHeartbeat: syncRequest.fh == null zk.commitProcessor.commit(syncRequest) locally!!!"); + zk.commitProcessor.commit(learnerSyncRequest); + } else { + sendSync(learnerSyncRequest); + } + //LOG.info("fuck_processHeartbeat_at_the_end, ThreadName:" + Thread.currentThread().getName()+",sid:"+sid); } static class ToBeAppliedRequestProcessor implements RequestProcessor { @@ -1083,6 +1261,7 @@ public void processRequest(Request request) throws RequestProcessorException { if (request.getHdr() != null) { long zxid = request.getHdr().getZxid(); Iterator iter = leader.toBeApplied.iterator(); + //System.out.println("fuck_Leader.ToBeAppliedRequestProcessor.ToBeAppliedRequestProcessor: " + leader.toBeApplied); if (iter.hasNext()) { Proposal p = iter.next(); if (p.request != null && p.request.zxid == zxid) { @@ -1141,6 +1320,7 @@ public void commit(long zxid) { lastCommitted = zxid; } QuorumPacket qp = new QuorumPacket(Leader.COMMIT, zxid, null, null); + //System.out.println("fuck_Leader send COMMIT packet"); sendPacket(qp); ServerMetrics.getMetrics().COMMIT_COUNT.add(1); } @@ -1250,6 +1430,7 @@ public Proposal propose(Request request) throws XidRolloverException { lastProposed = p.packet.getZxid(); outstandingProposals.put(lastProposed, p); + //System.out.println("fuck_Leader send PROPOSAL packet lastProposed:"+lastProposed+",outstandingProposals:"+outstandingProposals); sendPacket(pp); } ServerMetrics.getMetrics().PROPOSAL_COUNT.add(1); @@ -1261,13 +1442,39 @@ public Proposal propose(Request request) throws XidRolloverException { * * @param r the request */ - public synchronized void processSync(LearnerSyncRequest r) { +// System.out.println("fuck_Leader processSync#ReadRequests outstandingProposals:" + outstandingProposals.size() +// + ", LearnerSyncRequest:"+r); if (outstandingProposals.isEmpty()) { - sendSync(r); + System.out.println("fuck_Leader outstandingProposals.isEmpty() leader sends SYNC at once" + + ",LearnerSyncRequest:" + r.toString() +",bb:" + r.request +",r.type:"+r.type); + if (r.type == OpCode.sync || r.type == OpCode.syncedRead) { + sendSync(r); + } else { + r.setCanCommitted(true); + } } else { pendingSyncs.computeIfAbsent(lastProposed, k -> new ArrayList<>()).add(r); } + + if (r.type == OpCode.linearizableRead) { + //Quorum heartbeat(整体封装成方法) + r.getSyncedAckSet().addQuorumVerifier(self.getQuorumVerifier()); + if (self.getQuorumVerifier().getVersion() < self.getLastSeenQuorumVerifier().getVersion()) { + r.getSyncedAckSet().addQuorumVerifier(self.getLastSeenQuorumVerifier()); + } + + //waitingQuorumSyncs.computeIfAbsent(lastProposed, k -> new ArrayDeque<>()).addLast(r); + //r.setLastProposed(lastProposed); + ReadStateContext readStateContext = new ReadStateContext(lastProposed, r.sessionId, r.cxid, r.type); + waitingQuorumSyncs.put(readStateContext, r); +// System.out.println("fuck_Leader#processSync:" + ",ThreadName:" + Thread.currentThread().getName() +// + ", waitingQuorumSyncs:" + waitingQuorumSyncs + ",LearnerSyncRequest:" + r +// + ", readStateContext:"+readStateContext); + + r.getSyncedAckSet().addAck(self.getId()); + sendHeartBeat(readStateContext); + } } /** @@ -1275,9 +1482,38 @@ public synchronized void processSync(LearnerSyncRequest r) { */ public void sendSync(LearnerSyncRequest r) { QuorumPacket qp = new QuorumPacket(Leader.SYNC, 0, null, null); + //System.out.println("fuck_Leader leader sends SYNC to one follower. LearnerSyncRequest:"+r); r.fh.queuePacket(qp); } + /** + * TODO + */ + public void sendHeartBeat(ReadStateContext readStateContext) { + ByteArrayOutputStream baos = new ByteArrayOutputStream(); + DataOutputStream dos = new DataOutputStream(baos); + try { + dos.writeLong(readStateContext.lastProposed); + dos.writeLong(readStateContext.sessionId); + dos.writeInt(readStateContext.xid); + dos.writeInt(readStateContext.type); + dos.close(); + } catch (IOException e) { + LOG.warn("Unexpected exception", e); + } + QuorumPacket qp = new QuorumPacket(Leader.HEARTBEAT, getLastProposed(), baos.toByteArray(), null); + //System.out.println("fuck_Leader leader sends ---> sendHeartBeat to one follower."); + + sendPacket(qp); + //r.fh.queuePacket(qp); + //参考 propose() method +// synchronized (forwardingFollowers) { +// for (LearnerHandler f : forwardingFollowers) { +// f.queuePacket(qp); +// } +// } + } + /** * lets the leader know that a follower is capable of following and is done * syncing 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 ef36b2f9000..031b2526979 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 @@ -844,6 +844,12 @@ protected void ping(QuorumPacket qp) throws IOException { writePacket(pingReply, true); } + protected void sendHeartBeat(QuorumPacket qp) throws IOException { + // TODO fuck + //System.out.println("fuck_follower#sendHeartBeat yuan-feng-bu-dong qp:" + qp); + writePacket(qp, true); + } + /** * Shutdown the Peer */ 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 b91319e0ca0..901f3122c6c 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 @@ -159,6 +159,7 @@ public synchronized void start() { } public synchronized void updateProposal(long zxid, long time) { + //System.out.println("fuck_LearnerHandler#updateProposal_zxid:"+zxid+",time:"+time); if (!started) { return; } @@ -172,6 +173,7 @@ public synchronized void updateProposal(long zxid, long time) { } public synchronized void updateAck(long zxid) { + //System.out.println("fuck_LearnerHandler#updateAck_zxid:"+zxid); if (currentZxid == zxid) { currentTime = nextTime; currentZxid = nextZxid; @@ -192,6 +194,9 @@ public synchronized boolean check(long time) { return true; } else { long msDelay = (time - currentTime) / 1000000; + if (!(msDelay < learnerMaster.syncTimeout())) { + System.out.println("fuck#LearnerHandler#SyncLimitCheck#check msDelay:"+msDelay); + } return (msDelay < learnerMaster.syncTimeout()); } } @@ -588,7 +593,7 @@ public void run() { ServerMetrics.getMetrics().DIFF_COUNT.add(1); } - LOG.debug("Sending NEWLEADER message to {}", sid); + //LOG.info("fuck Sending NEWLEADER message to {}", sid); // the version of this quorumVerifier will be set by leader.lead() in case // the leader is just being established. waitForEpochAck makes sure that readyToStart is true if // we got here, so the version was set @@ -618,7 +623,7 @@ public void run() { return; } - LOG.debug("Received NEWLEADER-ACK message from {}", sid); + //LOG.info("fuck Received NEWLEADER-ACK message from {}", sid); learnerMaster.waitForNewLeaderAck(getSid(), qp.getZxid()); @@ -644,7 +649,7 @@ public void run() { // so we need to mark when the peer can actually start // using the data // - LOG.debug("Sending UPTODATE message to {}", sid); + LOG.info("Sending UPTODATE message to {}", sid); queuedPackets.add(new QuorumPacket(Leader.UPTODATE, -1, null, null)); while (true) { @@ -674,6 +679,7 @@ public void run() { LOG.debug("Received ACK from Observer {}", this.sid); } syncLimitCheck.updateAck(qp.getZxid()); + //System.out.println("fuck_LearnerHandler leader had received ACK"); learnerMaster.processAck(this.sid, qp.getZxid(), sock.getLocalSocketAddress()); break; case Leader.PING: @@ -686,6 +692,15 @@ public void run() { learnerMaster.touch(sess, to); } break; + case Leader.HEARTBEAT: + //TODO + //System.out.println("fuck_I'm leader and I receive the follower the PUREPING"); + //fuck 模拟ping确认的网络延迟 + //Thread.sleep(300); + syncLimitCheck.updateAck(qp.getZxid()); + + learnerMaster.processHeartbeat(this.sid, qp); + break; case Leader.REVALIDATE: ServerMetrics.getMetrics().REVALIDATE_COUNT.add(1); learnerMaster.revalidateSession(qp, this); @@ -697,12 +712,13 @@ public void run() { type = bb.getInt(); bb = bb.slice(); Request si; - if (type == OpCode.sync) { + if (type == OpCode.sync || type == OpCode.syncedRead || type == OpCode.linearizableRead) { si = new LearnerSyncRequest(this, sessionId, cxid, type, bb, qp.getAuthinfo()); } else { si = new Request(null, sessionId, cxid, type, bb, qp.getAuthinfo()); } si.setOwner(this); + //System.out.println("fuck_LearnerHandler leader had received REQUEST learnerMaster.submitLearnerRequest: " + si); learnerMaster.submitLearnerRequest(si); requestsReceived.incrementAndGet(); break; diff --git a/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/LearnerMaster.java b/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/LearnerMaster.java index 9bf6032af68..aa77e382f97 100644 --- a/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/LearnerMaster.java +++ b/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/LearnerMaster.java @@ -188,6 +188,8 @@ public LearnerSyncThrottler getLearnerDiffSyncThrottler() { */ abstract void processAck(long sid, long zxid, SocketAddress localSocketAddress); + abstract void processHeartbeat(long sid, QuorumPacket qp); + /** * mark session as alive * @param sess session id diff --git a/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/LearnerSyncRequest.java b/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/LearnerSyncRequest.java index d4c83aeab7b..d7098a03dec 100644 --- a/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/LearnerSyncRequest.java +++ b/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/LearnerSyncRequest.java @@ -20,16 +20,94 @@ import java.nio.ByteBuffer; import java.util.List; +import java.util.Objects; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.atomic.AtomicBoolean; import org.apache.zookeeper.data.Id; import org.apache.zookeeper.server.Request; +import org.apache.zookeeper.server.ServerCnxn; public class LearnerSyncRequest extends Request { LearnerHandler fh; + + SyncedLearnerTracker syncedAckSet; + + //命名 + volatile boolean canCommitted = false; + + CountDownLatch countDownLatch = new CountDownLatch(1); + //long lastProposed = -1L; + + //boolean isQuorum = false; + public LearnerSyncRequest( LearnerHandler fh, long sessionId, int xid, int type, ByteBuffer bb, List authInfo) { super(null, sessionId, xid, type, bb, authInfo); this.fh = fh; + syncedAckSet = new SyncedLearnerTracker(); + //isQuorum = new AtomicBoolean(false); + } + + public LearnerSyncRequest( + ServerCnxn cnxn, LearnerHandler fh, long sessionId, int xid, int type, ByteBuffer bb, List authInfo) { + super(cnxn, sessionId, xid, type, bb, authInfo); + this.fh = fh; + syncedAckSet = new SyncedLearnerTracker(); + //isQuorum = new AtomicBoolean(false); + } + + public static LearnerSyncRequest makeLearnerSyncRequest (Request r) { + return new LearnerSyncRequest(r.cnxn, null, r.sessionId, r.cxid, r.type, r.request, r.authInfo); } + public SyncedLearnerTracker getSyncedAckSet() { + return syncedAckSet; + } + +// @Override +// public boolean equals(Object o) { +// if (this == o) return true; +// if (o == null || getClass() != o.getClass()) return false; +// LearnerSyncRequest that = (LearnerSyncRequest) o; +// return lastProposed == that.lastProposed && +// sessionId == that.sessionId && +// cxid == that.cxid && +// type == that.type; +// } +// +// @Override +// public int hashCode() { +// return Objects.hash(lastProposed, sessionId, cxid, type); +// } +// +// @Override +// public String toString() { +// return "LearnerSyncRequest{" + +// "lastProposed=" + lastProposed + +// ", sessionId=" + sessionId + +// ", cxid=" + cxid + +// ", type=" + type + +// '}'; +// } + +// public void setQuorum(boolean quorum) { +// isQuorum = quorum; +// } + +// public void setLastProposed(long lastProposed) { +// this.lastProposed = lastProposed; +// } + + public boolean isFromLeader() { + return this.fh == null; + } + + public void setCanCommitted(boolean canCommitted) { + this.canCommitted = canCommitted; + } + + public CountDownLatch getCountDownLatch() { + return countDownLatch; + } } 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 98c0d0c08bf..ca27f6f2731 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 @@ -27,6 +27,7 @@ import java.net.ServerSocket; import java.net.Socket; import java.net.SocketAddress; +import java.sql.SQLOutput; import java.util.ArrayList; import java.util.Arrays; import java.util.Collections; @@ -232,6 +233,12 @@ public void processAck(long sid, long zxid, SocketAddress localSocketAddress) { throw new RuntimeException("Observers shouldn't send ACKS ack = " + Long.toHexString(zxid)); } + @Override + public void processHeartbeat(long sid, QuorumPacket qp) { + // Nothing + System.out.println("fuck_ObserverMaster#processPing-do-Nothing"); + } + @Override public void touch(long sess, int to) { zks.getSessionTracker().touchSession(sess, to); diff --git a/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/ObserverRequestProcessor.java b/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/ObserverRequestProcessor.java index 0e7071b91b7..78d645b2d2a 100644 --- a/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/ObserverRequestProcessor.java +++ b/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/ObserverRequestProcessor.java @@ -96,6 +96,9 @@ public void run() { // add it to pendingSyncs. switch (request.type) { case OpCode.sync: + case OpCode.syncedRead: + case OpCode.linearizableRead: + //System.out.println("fuck_ObserverRequestProcessor_zks.pendingSyncs.add(request):"+ request); zks.pendingSyncs.add(request); zks.getObserver().request(request); break; diff --git a/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/ProposalRequestProcessor.java b/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/ProposalRequestProcessor.java index c1e2fe16e43..b8c1f824f16 100644 --- a/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/ProposalRequestProcessor.java +++ b/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/ProposalRequestProcessor.java @@ -18,6 +18,7 @@ package org.apache.zookeeper.server.quorum; +import org.apache.zookeeper.ZooDefs; import org.apache.zookeeper.server.Request; import org.apache.zookeeper.server.RequestProcessor; import org.apache.zookeeper.server.ServerMetrics; @@ -74,18 +75,26 @@ public void processRequest(Request request) throws RequestProcessorException { * call processRequest on the next processor. */ if (request instanceof LearnerSyncRequest) { + //System.out.println("fuck_Leader#ProposalRequestProcessor#processSync request from other followers instanceof LearnerSyncRequest bb:" + request.request); zks.getLeader().processSync((LearnerSyncRequest) request); + } else if (request.type == ZooDefs.OpCode.linearizableRead) { + //System.out.println("fuck_Leader#ProposalRequestProcessor#request from leader's linearizableRead, bb:" + request.request); + zks.getLeader().processSync(LearnerSyncRequest.makeLearnerSyncRequest(request)); } else { + //System.out.println("fuck_Leader#ProposalRequestProcessor request.getHdr():" + request.getHdr() + ",bb:" + request.request +",request:"+request); if (shouldForwardToNextProcessor(request)) { + //System.out.println("fuck_ProposalRequestProcessor shouldForwardToNextProcessor to next commit processor request:" + request); nextProcessor.processRequest(request); } if (request.getHdr() != null) { // We need to sync and get consensus on any transactions try { + //System.out.println("fuck_ProposalRequestProcessor(zks.getLeader().propose(request)) if this's a write request: "+request); zks.getLeader().propose(request); } catch (XidRolloverException e) { throw new RequestProcessorException(e.getMessage(), e); } + //System.out.println("fuck_ProposalRequestProcessor#syncProcessor.processRequest"); syncProcessor.processRequest(request); } } @@ -105,6 +114,7 @@ private boolean shouldForwardToNextProcessor(Request request) { ServerMetrics.getMetrics().REQUESTS_NOT_FORWARDED_TO_COMMIT_PROCESSOR.add(1); return false; } - return true; + //fuck + return false; } } diff --git a/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/ReadOnlyRequestProcessor.java b/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/ReadOnlyRequestProcessor.java index 3aec97f72e2..b692be962f6 100644 --- a/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/ReadOnlyRequestProcessor.java +++ b/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/ReadOnlyRequestProcessor.java @@ -75,6 +75,8 @@ public void run() { // filter read requests switch (request.type) { case OpCode.sync: + case OpCode.syncedRead: + case OpCode.linearizableRead: case OpCode.create: case OpCode.create2: case OpCode.createTTL: diff --git a/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/SendAckRequestProcessor.java b/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/SendAckRequestProcessor.java index 8218ddae4ae..bff6dcd0787 100644 --- a/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/SendAckRequestProcessor.java +++ b/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/SendAckRequestProcessor.java @@ -38,11 +38,16 @@ public class SendAckRequestProcessor implements RequestProcessor, Flushable { } public void processRequest(Request si) { - if (si.type != OpCode.sync) { + //System.out.println("fuck_SendAckRequestProcessor——before(I'm a follower) Request: " + si + ",si.type:"+si.type); + if (si.type == OpCode.sync || si.type == OpCode.syncedRead || si.type == OpCode.linearizableRead) { + //fuck + System.out.println("fuck_SendAckRequestProcessor——before(I'm a follower) Request: " + si + ", (happy-men)si.type:" + si.type); + } + if (si.type != OpCode.sync && si.type != OpCode.syncedRead && si.type != OpCode.linearizableRead) { QuorumPacket qp = new QuorumPacket(Leader.ACK, si.getHdr().getZxid(), null, null); try { si.logLatency(ServerMetrics.getMetrics().PROPOSAL_ACK_CREATION_LATENCY); - + //System.out.println("fuck_SendAckRequestProcessor(I'm a follower) Request: " + si); learner.writePacket(qp, false); } catch (IOException e) { LOG.warn("Closing connection to leader, exception during packet send", e); diff --git a/zookeeper-server/src/main/java/org/apache/zookeeper/server/util/RequestPathMetricsCollector.java b/zookeeper-server/src/main/java/org/apache/zookeeper/server/util/RequestPathMetricsCollector.java index f3ec1fcea92..509d3a6f026 100644 --- a/zookeeper-server/src/main/java/org/apache/zookeeper/server/util/RequestPathMetricsCollector.java +++ b/zookeeper-server/src/main/java/org/apache/zookeeper/server/util/RequestPathMetricsCollector.java @@ -29,11 +29,14 @@ import static org.apache.zookeeper.ZooDefs.OpCode.getChildren; import static org.apache.zookeeper.ZooDefs.OpCode.getChildren2; import static org.apache.zookeeper.ZooDefs.OpCode.getData; +import static org.apache.zookeeper.ZooDefs.OpCode.linearizableRead; import static org.apache.zookeeper.ZooDefs.OpCode.removeWatches; import static org.apache.zookeeper.ZooDefs.OpCode.setACL; import static org.apache.zookeeper.ZooDefs.OpCode.setData; import static org.apache.zookeeper.ZooDefs.OpCode.setWatches2; import static org.apache.zookeeper.ZooDefs.OpCode.sync; +import static org.apache.zookeeper.ZooDefs.OpCode.syncedRead; + import java.io.PrintWriter; import java.util.Arrays; import java.util.Collection; @@ -134,12 +137,15 @@ public RequestPathMetricsCollector(boolean accurateMode) { requestsMap.put(Request.op2String(removeWatches), new PathStatsQueue(removeWatches)); requestsMap.put(Request.op2String(setWatches2), new PathStatsQueue(setWatches2)); requestsMap.put(Request.op2String(sync), new PathStatsQueue(sync)); + requestsMap.put(Request.op2String(syncedRead), new PathStatsQueue(syncedRead)); + requestsMap.put(Request.op2String(linearizableRead), new PathStatsQueue(linearizableRead)); this.immutableRequestsMap = java.util.Collections.unmodifiableMap(requestsMap); } static boolean isWriteOp(int requestType) { switch (requestType) { case ZooDefs.OpCode.sync: + case ZooDefs.OpCode.syncedRead: case ZooDefs.OpCode.create: case ZooDefs.OpCode.create2: case ZooDefs.OpCode.createContainer: diff --git a/zookeeper-server/src/test/java/org/apache/zookeeper/test/StandaloneReadTest.java b/zookeeper-server/src/test/java/org/apache/zookeeper/test/StandaloneReadTest.java new file mode 100644 index 00000000000..488c228024e --- /dev/null +++ b/zookeeper-server/src/test/java/org/apache/zookeeper/test/StandaloneReadTest.java @@ -0,0 +1,155 @@ +/* + * 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.test; + +import static org.junit.jupiter.api.Assertions.assertThrows; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; +import org.apache.zookeeper.CreateMode; +import org.apache.zookeeper.WatchedEvent; +import org.apache.zookeeper.Watcher; +import org.apache.zookeeper.ZooDefs; +import org.apache.zookeeper.ZooKeeper; +import org.apache.zookeeper.client.ReadConsistencyMode; +import org.apache.zookeeper.client.ZNode; +import org.junit.Assert; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.Timeout; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * TODO Standalone server tests. + */ +public class StandaloneReadTest extends ClientBase { + + protected static final Logger LOG = LoggerFactory.getLogger(StandaloneReadTest.class); + + @BeforeEach + public void setup() { + //System.setProperty("zookeeper.DigestAuthenticationProvider.superDigest", "super:D/InIHSb7yEEbrWz8b9l71RjZJU="/* password is 'test'*/); + //QuorumPeerConfig.setReconfigEnabled(true); + } + + + @Test + public void testWrongReadConsistencyMode() throws Exception { + ZooKeeper zk = createClient(); + String path = "/foo"; + + assertThrows(NullPointerException.class, () -> zk.getData(null, path, null)); + assertThrows(IllegalArgumentException.class, () -> zk.getData(ReadConsistencyMode.DUMMY_READ, path, null)); + + // clean-up + zk.close(); + } + + @Test + @Timeout(value = 30) + public void testSequentialReadWithoutWatch() throws Exception { + testReadWithoutWatch(ReadConsistencyMode.SEQUENTIAL_READ); + } + + @Test + @Timeout(value = 30) + public void testSequentialReadWithWatch() throws Exception { + testReadWithWatch(ReadConsistencyMode.SEQUENTIAL_READ); + } + + @Test + @Timeout(value = 30) + public void testOrderedSequentialReadWithoutWatch() throws Exception { + testReadWithoutWatch(ReadConsistencyMode.ORDERED_SEQUENTIAL_READ); + } + + @Test + @Timeout(value = 30) + public void testOrderedSequentialReadWithWatch() throws Exception { + testReadWithWatch(ReadConsistencyMode.ORDERED_SEQUENTIAL_READ); + } + + @Test + @Timeout(value = 30) + public void testLinearizableReadWithoutWatch() throws Exception { + testReadWithoutWatch(ReadConsistencyMode.LINEARIZABLE_READ); + } + + @Test + @Timeout(value = 30) + public void testLinearizableReadWithWatch() throws Exception { + testReadWithWatch(ReadConsistencyMode.LINEARIZABLE_READ); + } + + private void testReadWithWatch(ReadConsistencyMode readMod) throws Exception { + ZooKeeper zk = createClient(); + String path = "/foo"; + String data = "bar"; + zk.create(path, data.getBytes(), ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT); + + CountDownLatch countDownLatch = new CountDownLatch(1); + MyWatcher watcher = new MyWatcher(countDownLatch); + CompletableFuture future = zk.getData(readMod, path, watcher); + ZNode zNode = future.get(); + Assert.assertEquals(path, zNode.getPath()); + Assert.assertEquals(data, new String(zNode.getData())); + Assert.assertEquals(data.length(), zNode.getStat().getDataLength()); + + zk.setData(path, "test-watch".getBytes(), -1); + if (!countDownLatch.await(10, TimeUnit.SECONDS)) { + Assert.fail("the client doesn't receive a watch event."); + } + // clean-up + zk.close(); + } + + private void testReadWithoutWatch(ReadConsistencyMode readMode) throws Exception { + + ZooKeeper zk = createClient(); + String path = "/foo"; + String data = "bar"; + zk.create(path, data.getBytes(), ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT); + + CompletableFuture future = zk.getData(readMode, path, null); + ZNode zNode = future.get(); + Assert.assertEquals(path, zNode.getPath()); + Assert.assertEquals(data, new String(zNode.getData())); + Assert.assertEquals(data.length(), zNode.getStat().getDataLength()); + + // clean-up + zk.close(); + } + + private class MyWatcher implements Watcher { + CountDownLatch countDownLatch; + + public MyWatcher(CountDownLatch countDownLatch) { + this.countDownLatch = countDownLatch; + } + + public void process(WatchedEvent event) { + if (event.getType() != Event.EventType.None) { + countDownLatch.countDown(); + } + } + + } + +} From f7f347b541f4221f1f46635c5f97c643520ed55b Mon Sep 17 00:00:00 2001 From: maoling Date: Wed, 23 Dec 2020 17:58:40 +0800 Subject: [PATCH 2/2] 12-23-xiawu-leiya-heriwancheng-haiyouyige-tla-plus-shenwukelian --- .../java/org/apache/zookeeper/server/quorum/Leader.java | 7 +++++-- .../org/apache/zookeeper/server/quorum/LearnerHandler.java | 2 +- 2 files changed, 6 insertions(+), 3 deletions(-) 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 ddc8fd6b98f..da98b15a7e1 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 @@ -578,6 +578,8 @@ private void acceptConnections() throws IOException { BufferedInputStream is = new BufferedInputStream(socket.getInputStream()); LearnerHandler fh = new LearnerHandler(socket, is, Leader.this); + System.out.println("fuck_Leader#LearnerCnxAcceptorHandler.acceptConnections start one LearnerHandler:" + fh.toString()); + fh.start(); } catch (SocketException e) { error = true; @@ -1019,7 +1021,7 @@ public synchronized boolean tryToCommit(Proposal p, long zxid, SocketAddress fol //System.out.println("fuck_Leader zk.commitProcessor.commit(p.request): "+p.request +",pendingSyncs.containsKey(zxid):" + pendingSyncs.containsKey(zxid)); if (pendingSyncs.containsKey(zxid)) { System.out.println("fuck_Leader pendingSyncs.containsKey(zxid) send SYNC to all learns(III), we have committed zxid: " - + zxid + ",pendingSyncs.size: "+ pendingSyncs.get(zxid).size()); + + zxid + ",pendingSyncs.size: "+ pendingSyncs.get(zxid).size()+ ",time:" + System.currentTimeMillis()); for (LearnerSyncRequest r : pendingSyncs.remove(zxid)) { if (r.type == OpCode.linearizableRead) { //&& r.isFromLeader() //System.out.println("fuck_tryToCommit@ syncRequest.fh == null zk.commitProcessor.commit(syncRequest) at leader locally(linearizableRead) " + @@ -1196,7 +1198,7 @@ public void processHeartbeatInline(long sid, QuorumPacket qp) { //III long start = Time.currentWallTime(); System.out.println("fuck_processHeartbeat_start_to_try_to_commit_lastProposed-zxid:"+lastProposed+"_canCommitted:" + learnerSyncRequest.canCommitted - + ", learnerSyncRequest.isFromLeader: " + learnerSyncRequest.isFromLeader()); + + ", learnerSyncRequest.isFromLeader: " + learnerSyncRequest.isFromLeader() + ",time:" + System.currentTimeMillis()); if (!learnerSyncRequest.canCommitted) { //wait for apply try { @@ -1204,6 +1206,7 @@ public void processHeartbeatInline(long sid, QuorumPacket qp) { System.out.println("fuck_processHeartbeat_await_time:" + (System.currentTimeMillis() - start)+" ms"); //learnerSyncRequest.getCountDownLatch().await(5000, TimeUnit.MILLISECONDS); } catch (InterruptedException e) { + System.out.println("fuck_processHeartbeat_await_time#InterruptedException:" + (System.currentTimeMillis() - start)+" ms"); e.printStackTrace(); //TODO timeout return; 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 901f3122c6c..2fcd3528eff 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 @@ -1088,7 +1088,7 @@ public void ping() { queuePacket(ping); } else { LOG.warn("Closing connection to peer due to transaction timeout."); - shutdown(); + //shutdown(); } }