From c8ea7a25b0e879d7cd48affccaef5014e5c9b914 Mon Sep 17 00:00:00 2001 From: Michael Han Date: Wed, 25 Sep 2019 17:45:46 -0700 Subject: [PATCH 1/4] ZOOKEEPER-3560: Add response cache to serve get children (2) requests. ZOOKEEPER-3180 introduces response cache but it only covers getData requests. This commit is to extend the response cache based on the infrastructure set up by ZOOKEEPER-3180 to so the response of get children requests can also be served out of cache. Some design decisions: * Only OpCode.getChildren2 is supported, as OpCode.getChildren does not have associated stats and current cache infra relies on stats to invalidate cache. * The children list is stored in a separate response cache object so it does not pollute the existing data cache that's serving getData requests, and this separation also allows potential separate tuning of each cache based on workload characteristics. * As a result of cache object separation, new server metrics is added to measure cache hit / miss for get children requests, that's separated from get data requests. Similar as ZOOKEEPER-3180, the get children response cache is enabled by default, with a default cache size of 400, and can be disabled (together with get data response cache.). --- .../apache/zookeeper/server/DumbWatcher.java | 3 +- .../server/FinalRequestProcessor.java | 30 +++++--- .../zookeeper/server/NIOServerCnxn.java | 27 +------ .../zookeeper/server/NettyServerCnxn.java | 5 +- .../apache/zookeeper/server/ServerCnxn.java | 57 ++++++++++++-- .../zookeeper/server/ServerMetrics.java | 9 ++- .../zookeeper/server/ZooKeeperServer.java | 7 ++ .../zookeeper/server/MockServerCnxn.java | 3 +- .../zookeeper/test/ResponseCacheTest.java | 77 +++++++++++++++++-- 9 files changed, 167 insertions(+), 51 deletions(-) diff --git a/zookeeper-server/src/main/java/org/apache/zookeeper/server/DumbWatcher.java b/zookeeper-server/src/main/java/org/apache/zookeeper/server/DumbWatcher.java index 63da494962a..9b9672a850b 100644 --- a/zookeeper-server/src/main/java/org/apache/zookeeper/server/DumbWatcher.java +++ b/zookeeper-server/src/main/java/org/apache/zookeeper/server/DumbWatcher.java @@ -61,7 +61,8 @@ public void close(DisconnectReason reason) { } @Override - public void sendResponse(ReplyHeader h, Record r, String tag, String cacheKey, Stat stat) throws IOException { + public void sendResponse(ReplyHeader h, Record r, String tag, + String cacheKey, Stat stat, int opCode) throws IOException { } @Override 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 dcb3d26570c..ccc959aee71 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 @@ -552,19 +552,29 @@ public void processRequest(Request request) { updateStats(request, lastOp, lastZxid); try { - if (request.type == OpCode.getData && path != null && rsp != null) { - // Serialized read responses could be cached by the connection object. - // Cache entries are identified by their path and last modified zxid, - // so these values are passed along with the response. - GetDataResponse getDataResponse = (GetDataResponse) rsp; + if (path == null || rsp == null) { + cnxn.sendResponse(hdr, rsp, "response"); + } else { + int opCode = request.type; Stat stat = null; - if (getDataResponse.getStat() != null) { - stat = getDataResponse.getStat(); + switch (opCode) { + case OpCode.getData : { + GetDataResponse getDataResponse = (GetDataResponse) rsp; + stat = getDataResponse.getStat(); + cnxn.sendResponse(hdr, rsp, "response", path, stat, opCode); + break; + } + case OpCode.getChildren2 : { + GetChildren2Response getChildren2Response = (GetChildren2Response) rsp; + stat = getChildren2Response.getStat(); + cnxn.sendResponse(hdr, rsp, "response", path, stat, opCode); + break; + } + default: + cnxn.sendResponse(hdr, rsp, "response"); } - cnxn.sendResponse(hdr, rsp, "response", path, stat); - } else { - cnxn.sendResponse(hdr, rsp, "response"); } + if (request.type == OpCode.closeSession) { cnxn.sendCloseSession(); } diff --git a/zookeeper-server/src/main/java/org/apache/zookeeper/server/NIOServerCnxn.java b/zookeeper-server/src/main/java/org/apache/zookeeper/server/NIOServerCnxn.java index fc69a522043..5aa427f7e9b 100644 --- a/zookeeper-server/src/main/java/org/apache/zookeeper/server/NIOServerCnxn.java +++ b/zookeeper-server/src/main/java/org/apache/zookeeper/server/NIOServerCnxn.java @@ -663,31 +663,10 @@ public static void closeSock(SocketChannel sock) { private static final ByteBuffer packetSentinel = ByteBuffer.allocate(0); - /** - * Serializes a ZooKeeper response and enqueues it for sending. - * - * Serializes client response parts and enqueues them into outgoing queue. - * - * If both cache key and last modified zxid are provided, the serialized - * response is caсhed under the provided key, the last modified zxid is - * stored along with the value. A cache entry is invalidated if the - * provided last modified zxid is more recent than the stored one. - * - * Attention: this function is not thread safe, due to caching not being - * thread safe. - * - * @param h reply header - * @param r reply payload, can be null - * @param tag Jute serialization tag, can be null - * @param cacheKey key for caching the serialized payload. a null value - * prvents caching - * @param stat stat information for the the reply payload, used - * for cache invalidation. a value of 0 prevents caching. - */ @Override - public void sendResponse(ReplyHeader h, Record r, String tag, String cacheKey, Stat stat) { + public void sendResponse(ReplyHeader h, Record r, String tag, String cacheKey, Stat stat, int opCode) { try { - sendBuffer(serialize(h, r, tag, cacheKey, stat)); + sendBuffer(serialize(h, r, tag, cacheKey, stat, opCode)); decrOutstandingAndCheckThrottle(h); } catch (Exception e) { LOG.warn("Unexpected exception. Destruction averted.", e); @@ -712,7 +691,7 @@ public void process(WatchedEvent event) { // Convert WatchedEvent to a type that can be sent over the wire WatcherEvent e = event.getWrapper(); - sendResponse(h, e, "notification", null, null); + sendResponse(h, e, "notification", null, null, -1); } /* diff --git a/zookeeper-server/src/main/java/org/apache/zookeeper/server/NettyServerCnxn.java b/zookeeper-server/src/main/java/org/apache/zookeeper/server/NettyServerCnxn.java index ac031902b1d..f092072eae1 100644 --- a/zookeeper-server/src/main/java/org/apache/zookeeper/server/NettyServerCnxn.java +++ b/zookeeper-server/src/main/java/org/apache/zookeeper/server/NettyServerCnxn.java @@ -173,13 +173,14 @@ public void process(WatchedEvent event) { } @Override - public void sendResponse(ReplyHeader h, Record r, String tag, String cacheKey, Stat stat) throws IOException { + public void sendResponse(ReplyHeader h, Record r, String tag, + String cacheKey, Stat stat, int opCode) throws IOException { // cacheKey and stat are used in caching, which is not // implemented here. Implementation example can be found in NIOServerCnxn. if (closingChannel || !channel.isOpen()) { return; } - sendBuffer(serialize(h, r, tag, cacheKey, stat)); + sendBuffer(serialize(h, r, tag, cacheKey, stat, opCode)); decrOutstandingAndCheckThrottle(h); } diff --git a/zookeeper-server/src/main/java/org/apache/zookeeper/server/ServerCnxn.java b/zookeeper-server/src/main/java/org/apache/zookeeper/server/ServerCnxn.java index 4d95c8156c8..7b95feddadc 100644 --- a/zookeeper-server/src/main/java/org/apache/zookeeper/server/ServerCnxn.java +++ b/zookeeper-server/src/main/java/org/apache/zookeeper/server/ServerCnxn.java @@ -36,9 +36,11 @@ import java.util.concurrent.atomic.AtomicLong; import org.apache.jute.BinaryOutputArchive; import org.apache.jute.Record; +import org.apache.zookeeper.metrics.Counter; import org.apache.zookeeper.Quotas; import org.apache.zookeeper.WatchedEvent; import org.apache.zookeeper.Watcher; +import org.apache.zookeeper.ZooDefs.OpCode; import org.apache.zookeeper.data.Id; import org.apache.zookeeper.data.Stat; import org.apache.zookeeper.proto.ReplyHeader; @@ -160,10 +162,34 @@ public void decrOutstandingAndCheckThrottle(ReplyHeader h) { public abstract void close(DisconnectReason reason); - public abstract void sendResponse(ReplyHeader h, Record r, String tag, String cacheKey, Stat stat) throws IOException; + /** + * Serializes a ZooKeeper response and enqueues it for sending. + * + * Serializes client response parts and enqueues them into outgoing queue. + * + * If both cache key and last modified zxid are provided, the serialized + * response is caсhed under the provided key, the last modified zxid is + * stored along with the value. A cache entry is invalidated if the + * provided last modified zxid is more recent than the stored one. + * + * Attention: this function is not thread safe, due to caching not being + * thread safe. + * + * @param h reply header + * @param r reply payload, can be null + * @param tag Jute serialization tag, can be null + * @param cacheKey Key for caching the serialized payload. A null value prevents caching. + * @param stat Stat information for the the reply payload, used for cache invalidation. + * A value of 0 prevents caching. + * @param opCode The op code appertains to the corresponding request of the response, + * used to decide which cache (e.g. read response cache, + * list of children response cache, ...) object to look up to when applicable. + */ + public abstract void sendResponse(ReplyHeader h, Record r, String tag, + String cacheKey, Stat stat, int opCode) throws IOException; public void sendResponse(ReplyHeader h, Record r, String tag) throws IOException { - sendResponse(h, r, tag, null, null); + sendResponse(h, r, tag, null, null, -1); } protected byte[] serializeRecord(Record record) throws IOException { @@ -173,11 +199,30 @@ protected byte[] serializeRecord(Record record) throws IOException { return baos.toByteArray(); } - protected ByteBuffer[] serialize(ReplyHeader h, Record r, String tag, String cacheKey, Stat stat) throws IOException { + protected ByteBuffer[] serialize(ReplyHeader h, Record r, String tag, + String cacheKey, Stat stat, int opCode) throws IOException { byte[] header = serializeRecord(h); byte[] data = null; if (r != null) { - ResponseCache cache = zkServer.getReadResponseCache(); + ResponseCache cache = null; + Counter cacheHit = null, cacheMiss = null; + switch (opCode) { + case OpCode.getData : { + cache = zkServer.getReadResponseCache(); + cacheHit = ServerMetrics.getMetrics().RESPONSE_PACKET_CACHE_HITS; + cacheMiss = ServerMetrics.getMetrics().RESPONSE_PACKET_CACHE_MISSING; + break; + } + case OpCode.getChildren2 : { + cache = zkServer.getGetChildrenResponseCache(); + cacheHit = ServerMetrics.getMetrics().RESPONSE_PACKET_GET_CHILDREN_CACHE_HITS; + cacheMiss = ServerMetrics.getMetrics().RESPONSE_PACKET_GET_CHILDREN_CACHE_MISSING; + break; + } + default: + // op codes where response cache is not supported. + } + if (cache != null && stat != null && cacheKey != null && !cacheKey.endsWith(Quotas.statNode)) { // Use cache to get serialized data. // @@ -188,9 +233,9 @@ protected ByteBuffer[] serialize(ReplyHeader h, Record r, String tag, String cac // Cache miss, serialize the response and put it in cache. data = serializeRecord(r); cache.put(cacheKey, data, stat); - ServerMetrics.getMetrics().RESPONSE_PACKET_CACHE_MISSING.add(1); + cacheMiss.add(1); } else { - ServerMetrics.getMetrics().RESPONSE_PACKET_CACHE_HITS.add(1); + cacheHit.add(1); } } else { data = serializeRecord(r); 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 1f9855ca260..10bf444a8ee 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 @@ -95,7 +95,6 @@ private ServerMetrics(MetricsProvider metricsProvider) { NODE_CHANGED_WATCHER = metricsContext.getSummary("node_changed_watch_count", DetailLevel.BASIC); NODE_CHILDREN_WATCHER = metricsContext.getSummary("node_children_watch_count", DetailLevel.BASIC); - /* * Number of dead watchers in DeadWatcherListener */ @@ -106,6 +105,8 @@ private ServerMetrics(MetricsProvider metricsProvider) { RESPONSE_PACKET_CACHE_HITS = metricsContext.getCounter("response_packet_cache_hits"); RESPONSE_PACKET_CACHE_MISSING = metricsContext.getCounter("response_packet_cache_misses"); + RESPONSE_PACKET_GET_CHILDREN_CACHE_HITS = metricsContext.getCounter("response_packet_get_children_cache_hits"); + RESPONSE_PACKET_GET_CHILDREN_CACHE_MISSING = metricsContext.getCounter("response_packet_get_children_cache_misses"); ENSEMBLE_AUTH_SUCCESS = metricsContext.getCounter("ensemble_auth_success"); @@ -338,8 +339,14 @@ private ServerMetrics(MetricsProvider metricsProvider) { public final Counter DEAD_WATCHERS_QUEUED; public final Counter DEAD_WATCHERS_CLEARED; public final Summary DEAD_WATCHERS_CLEANER_LATENCY; + + /* + * Response cache hit and miss metrics. + */ public final Counter RESPONSE_PACKET_CACHE_HITS; public final Counter RESPONSE_PACKET_CACHE_MISSING; + public final Counter RESPONSE_PACKET_GET_CHILDREN_CACHE_HITS; + public final Counter RESPONSE_PACKET_GET_CHILDREN_CACHE_MISSING; /** * Learner handler quorum packet metrics. 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 3d6c3751bb3..9c22b89b30c 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 @@ -163,6 +163,7 @@ public static void setCloseSessionTxnEnabled(boolean enabled) { private FileTxnSnapLog txnLogFactory = null; private ZKDatabase zkDb; private ResponseCache readResponseCache; + private ResponseCache getChildrenResponseCache; private final AtomicLong hzxid = new AtomicLong(0); public static final Exception ok = new Exception("No prob"); protected RequestProcessor firstProcessor; @@ -308,6 +309,8 @@ public ZooKeeperServer(FileTxnSnapLog txnLogFactory, int tickTime, int minSessio readResponseCache = new ResponseCache(); + getChildrenResponseCache = new ResponseCache(); + connThrottle = new BlueThrottle(); this.initialConfig = initialConfig; @@ -1756,6 +1759,10 @@ public ResponseCache getReadResponseCache() { return isResponseCachingEnabled ? readResponseCache : null; } + public ResponseCache getGetChildrenResponseCache() { + return isResponseCachingEnabled ? getChildrenResponseCache : null; + } + protected void registerMetrics() { MetricsContext rootContext = ServerMetrics.getMetrics().getMetricsProvider().getRootContext(); diff --git a/zookeeper-server/src/test/java/org/apache/zookeeper/server/MockServerCnxn.java b/zookeeper-server/src/test/java/org/apache/zookeeper/server/MockServerCnxn.java index 8314545406e..a82eab3764d 100644 --- a/zookeeper-server/src/test/java/org/apache/zookeeper/server/MockServerCnxn.java +++ b/zookeeper-server/src/test/java/org/apache/zookeeper/server/MockServerCnxn.java @@ -46,7 +46,8 @@ public void close(DisconnectReason reason) { } @Override - public void sendResponse(ReplyHeader h, Record r, String tag, String cacheKey, Stat stat) throws IOException { + public void sendResponse(ReplyHeader h, Record r, String tag, + String cacheKey, Stat stat, int opCode) throws IOException { } @Override diff --git a/zookeeper-server/src/test/java/org/apache/zookeeper/test/ResponseCacheTest.java b/zookeeper-server/src/test/java/org/apache/zookeeper/test/ResponseCacheTest.java index 9831e875bbf..874d7cb85d9 100644 --- a/zookeeper-server/src/test/java/org/apache/zookeeper/test/ResponseCacheTest.java +++ b/zookeeper-server/src/test/java/org/apache/zookeeper/test/ResponseCacheTest.java @@ -21,6 +21,9 @@ import static org.junit.Assert.assertArrayEquals; import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertNotSame; +import static org.junit.Assert.fail; + +import java.util.List; import java.util.Map; import org.apache.zookeeper.CreateMode; import org.apache.zookeeper.ZooDefs; @@ -48,11 +51,12 @@ public void testResponseCache() throws Exception { } } - private void checkCacheStatus(long expectedHits, long expectedMisses) { + private void checkCacheStatus(long expectedHits, long expectedMisses, + String cacheHitMetricsName, String cacheMissMetricsName) { Map metrics = MetricsUtils.currentServerMetrics(); - assertEquals(expectedHits, metrics.get("response_packet_cache_hits")); - assertEquals(expectedMisses, metrics.get("response_packet_cache_misses")); + assertEquals(expectedHits, metrics.get(cacheHitMetricsName)); + assertEquals(expectedMisses, metrics.get(cacheMissMetricsName)); } public void performCacheTest(ZooKeeper zk, String path, boolean useCache) throws Exception { @@ -78,7 +82,8 @@ public void performCacheTest(ZooKeeper zk, String path, boolean useCache) throws expectedMisses += 1; expectedHits += reads - 1; } - checkCacheStatus(expectedHits, expectedMisses); + checkCacheStatus(expectedHits, expectedMisses, "response_packet_cache_hits", + "response_packet_cache_misses"); writeData = "test2".getBytes(); writeStat = zk.setData(path, writeData, -1); @@ -91,7 +96,8 @@ public void performCacheTest(ZooKeeper zk, String path, boolean useCache) throws expectedMisses += 1; expectedHits += reads - 1; } - checkCacheStatus(expectedHits, expectedMisses); + checkCacheStatus(expectedHits, expectedMisses, "response_packet_cache_hits", + "response_packet_cache_misses"); // Create a child beneath the tested node. This won't change the data of // the tested node, but will change it's pzxid. The next read of the tested @@ -104,7 +110,66 @@ public void performCacheTest(ZooKeeper zk, String path, boolean useCache) throws } assertArrayEquals(writeData, readData); assertNotSame(writeStat, readStat); - checkCacheStatus(expectedHits, expectedMisses); + checkCacheStatus(expectedHits, expectedMisses, "response_packet_cache_hits", + "response_packet_cache_misses"); + + ServerMetrics.getMetrics().resetAll(); + expectedHits = 0; + expectedMisses = 0; + createPath(path + "/a", zk); + createPath(path + "/a/b", zk); + createPath(path + "/a/c", zk); + createPath(path + "/a/b/d", zk); + createPath(path + "/a/b/e", zk); + createPath(path + "/a/b/e/f", zk); + createPath(path + "/a/b/e/g", zk); + createPath(path + "/a/b/e/h", zk); + + checkPath(path + "/a", zk, 2); + checkPath(path + "/a/b", zk,2); + checkPath(path + "/a/c", zk,0); + checkPath(path + "/a/b/d", zk, 0); + checkPath(path + "/a/b/e", zk, 3); + checkPath(path + "/a/b/e/h", zk, 0); + + if (useCache) { + expectedMisses += 6; + } + + checkCacheStatus(expectedHits, expectedMisses, "response_packet_get_children_cache_hits", + "response_packet_get_children_cache_misses"); + + checkPath(path + "/a", zk, 2); + checkPath(path + "/a/b", zk,2); + checkPath(path + "/a/c", zk,0); + + if (useCache) { + expectedHits += 3; + } + + checkCacheStatus(expectedHits, expectedMisses, "response_packet_get_children_cache_hits", + "response_packet_get_children_cache_misses"); + } + + private void createPath(String path, ZooKeeper zk) throws Exception { + zk.create(path, "".getBytes(), ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT, null); + } + + private void checkPath(String path, ZooKeeper zk, int expectedNumberOfChildren) throws Exception { + Stat stat = zk.exists(path, false); + + List c1 = zk.getChildren(path, false); + List c2 = zk.getChildren(path, false, stat); + + if (!c1.equals(c2)) { + fail("children lists from getChildren()/getChildren2() do not match"); + } + + assertEquals(c1.size(), expectedNumberOfChildren); + + if (!stat.equals(stat)) { + fail("stats from exists()/getChildren2() do not match"); + } } } From e1edd1d0b4647859ce07e6f22d9a41c6eb3ee09e Mon Sep 17 00:00:00 2001 From: Michael Han Date: Thu, 26 Sep 2019 09:32:59 -0700 Subject: [PATCH 2/4] Fix checkstyle violations. --- .../java/org/apache/zookeeper/server/ServerCnxn.java | 2 +- .../org/apache/zookeeper/test/ResponseCacheTest.java | 9 ++++----- 2 files changed, 5 insertions(+), 6 deletions(-) diff --git a/zookeeper-server/src/main/java/org/apache/zookeeper/server/ServerCnxn.java b/zookeeper-server/src/main/java/org/apache/zookeeper/server/ServerCnxn.java index 7b95feddadc..2ba02f8d606 100644 --- a/zookeeper-server/src/main/java/org/apache/zookeeper/server/ServerCnxn.java +++ b/zookeeper-server/src/main/java/org/apache/zookeeper/server/ServerCnxn.java @@ -36,13 +36,13 @@ import java.util.concurrent.atomic.AtomicLong; import org.apache.jute.BinaryOutputArchive; import org.apache.jute.Record; -import org.apache.zookeeper.metrics.Counter; import org.apache.zookeeper.Quotas; import org.apache.zookeeper.WatchedEvent; import org.apache.zookeeper.Watcher; import org.apache.zookeeper.ZooDefs.OpCode; import org.apache.zookeeper.data.Id; import org.apache.zookeeper.data.Stat; +import org.apache.zookeeper.metrics.Counter; import org.apache.zookeeper.proto.ReplyHeader; import org.apache.zookeeper.proto.RequestHeader; import org.slf4j.Logger; diff --git a/zookeeper-server/src/test/java/org/apache/zookeeper/test/ResponseCacheTest.java b/zookeeper-server/src/test/java/org/apache/zookeeper/test/ResponseCacheTest.java index 874d7cb85d9..5d7bf5dbc3e 100644 --- a/zookeeper-server/src/test/java/org/apache/zookeeper/test/ResponseCacheTest.java +++ b/zookeeper-server/src/test/java/org/apache/zookeeper/test/ResponseCacheTest.java @@ -22,7 +22,6 @@ import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertNotSame; import static org.junit.Assert.fail; - import java.util.List; import java.util.Map; import org.apache.zookeeper.CreateMode; @@ -126,8 +125,8 @@ public void performCacheTest(ZooKeeper zk, String path, boolean useCache) throws createPath(path + "/a/b/e/h", zk); checkPath(path + "/a", zk, 2); - checkPath(path + "/a/b", zk,2); - checkPath(path + "/a/c", zk,0); + checkPath(path + "/a/b", zk, 2); + checkPath(path + "/a/c", zk, 0); checkPath(path + "/a/b/d", zk, 0); checkPath(path + "/a/b/e", zk, 3); checkPath(path + "/a/b/e/h", zk, 0); @@ -140,8 +139,8 @@ public void performCacheTest(ZooKeeper zk, String path, boolean useCache) throws "response_packet_get_children_cache_misses"); checkPath(path + "/a", zk, 2); - checkPath(path + "/a/b", zk,2); - checkPath(path + "/a/c", zk,0); + checkPath(path + "/a/b", zk, 2); + checkPath(path + "/a/c", zk, 0); if (useCache) { expectedHits += 3; From 1026426c3123660b8d4ec5bf1ef980363cdf51d1 Mon Sep 17 00:00:00 2001 From: Michael Han Date: Tue, 1 Oct 2019 18:49:15 -0700 Subject: [PATCH 3/4] Address review comments. 1. Add new property to allow tuning read and get children cache separately. 2. Update documentation. 3. Add back a comment that was removed accidentally (with slight tuning). --- .../main/resources/markdown/zookeeperAdmin.md | 9 ++++++ .../server/FinalRequestProcessor.java | 3 ++ .../zookeeper/server/ResponseCache.java | 28 +++++++++++-------- .../zookeeper/server/ZooKeeperServer.java | 11 ++++++-- .../zookeeper/test/ResponseCacheTest.java | 24 +++++++++++++++- 5 files changed, 60 insertions(+), 15 deletions(-) diff --git a/zookeeper-docs/src/main/resources/markdown/zookeeperAdmin.md b/zookeeper-docs/src/main/resources/markdown/zookeeperAdmin.md index 327ac56b789..c2c49292046 100644 --- a/zookeeper-docs/src/main/resources/markdown/zookeeperAdmin.md +++ b/zookeeper-docs/src/main/resources/markdown/zookeeperAdmin.md @@ -719,6 +719,15 @@ property, when available, is noted below. by default with a value of 400, set to 0 or a negative integer to turn the feature off. +* *maxGetChildrenResponseCacheSize* : + (Java system property: **zookeeper.maxGetChildrenResponseCacheSize**) + Similar to **maxResponseCacheSize**, but applies to get children + requests. The metrics **response_packet_get_children_cache_hits** + and **response_packet_get_children_cache_misses** can be used to tune + this value to a given workload. The feature is turned on + by default with a value of 400, set to 0 or a negative + integer to turn the feature off. + * *autopurge.snapRetainCount* : (No Java system property) **New in 3.4.0:** 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 ccc959aee71..383a7d577ae 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 @@ -557,6 +557,9 @@ public void processRequest(Request request) { } else { int opCode = request.type; Stat stat = null; + // Serialized read and get children responses could be cached by the connection + // object. Cache entries are identified by their path and last modified zxid, + // so these values are passed along with the response. switch (opCode) { case OpCode.getData : { GetDataResponse getDataResponse = (GetDataResponse) rsp; diff --git a/zookeeper-server/src/main/java/org/apache/zookeeper/server/ResponseCache.java b/zookeeper-server/src/main/java/org/apache/zookeeper/server/ResponseCache.java index 7e728b58a81..f3d546ec884 100644 --- a/zookeeper-server/src/main/java/org/apache/zookeeper/server/ResponseCache.java +++ b/zookeeper-server/src/main/java/org/apache/zookeeper/server/ResponseCache.java @@ -22,23 +22,31 @@ import java.util.LinkedHashMap; import java.util.Map; import org.apache.zookeeper.data.Stat; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; @SuppressWarnings("serial") public class ResponseCache { + private static final Logger LOG = LoggerFactory.getLogger(ResponseCache.class); // Magic number chosen to be "big enough but not too big" - private static final int DEFAULT_RESPONSE_CACHE_SIZE = 400; - + public static final int DEFAULT_RESPONSE_CACHE_SIZE = 400; + private final int cacheSize; private static class Entry { - public Stat stat; public byte[] data; - } - private Map cache = Collections.synchronizedMap(new LRUCache(getResponseCacheSize())); + private Map cache; - public ResponseCache() { + public ResponseCache(int cacheSize) { + this.cacheSize = cacheSize; + cache = Collections.synchronizedMap(new LRUCache<>(cacheSize)); + LOG.info("Response cache size is initialized with value {}.", cacheSize); + } + + public int getCacheSize() { + return cacheSize; } public void put(String path, byte[] data, Stat stat) { @@ -62,12 +70,8 @@ public byte[] get(String key, Stat stat) { } } - private static int getResponseCacheSize() { - return Integer.getInteger("zookeeper.maxResponseCacheSize", DEFAULT_RESPONSE_CACHE_SIZE); - } - - public static boolean isEnabled() { - return getResponseCacheSize() > 0; + public boolean isEnabled() { + return cacheSize > 0; } private static class LRUCache extends LinkedHashMap { 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 9c22b89b30c..1b23027a6ca 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 @@ -218,6 +218,9 @@ protected enum State { public static final int DEFAULT_STARTING_BUFFER_SIZE = 1024; public static final int intBufferStartingSizeBytes; + public static final String GET_DATA_RESPONSE_CACHE_SIZE = "zookeeper.maxResponseCacheSize"; + public static final String GET_CHILDREN_RESPONSE_CACHE_SIZE = "zookeeper.maxGetChildrenResponseCacheSize"; + static { long configuredFlushDelay = Long.getLong(FLUSH_DELAY, 0); setFlushDelay(configuredFlushDelay); @@ -307,9 +310,13 @@ public ZooKeeperServer(FileTxnSnapLog txnLogFactory, int tickTime, int minSessio listener = new ZooKeeperServerListenerImpl(this); - readResponseCache = new ResponseCache(); + readResponseCache = new ResponseCache(Integer.getInteger( + GET_DATA_RESPONSE_CACHE_SIZE, + ResponseCache.DEFAULT_RESPONSE_CACHE_SIZE)); - getChildrenResponseCache = new ResponseCache(); + getChildrenResponseCache = new ResponseCache(Integer.getInteger( + GET_CHILDREN_RESPONSE_CACHE_SIZE, + ResponseCache.DEFAULT_RESPONSE_CACHE_SIZE)); connThrottle = new BlueThrottle(); diff --git a/zookeeper-server/src/test/java/org/apache/zookeeper/test/ResponseCacheTest.java b/zookeeper-server/src/test/java/org/apache/zookeeper/test/ResponseCacheTest.java index 5d7bf5dbc3e..0b27ff59c10 100644 --- a/zookeeper-server/src/test/java/org/apache/zookeeper/test/ResponseCacheTest.java +++ b/zookeeper-server/src/test/java/org/apache/zookeeper/test/ResponseCacheTest.java @@ -30,6 +30,9 @@ import org.apache.zookeeper.data.Stat; import org.apache.zookeeper.metrics.MetricsUtils; import org.apache.zookeeper.server.ServerMetrics; +import org.apache.zookeeper.server.ZooKeeperServer; +import org.junit.After; +import org.junit.Before; import org.junit.Test; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -38,6 +41,19 @@ public class ResponseCacheTest extends ClientBase { protected static final Logger LOG = LoggerFactory.getLogger(ResponseCacheTest.class); + @Before + public void setup() throws Exception { + System.setProperty(ZooKeeperServer.GET_DATA_RESPONSE_CACHE_SIZE, "32"); + System.setProperty(ZooKeeperServer.GET_CHILDREN_RESPONSE_CACHE_SIZE, "64"); + super.setUp(); + } + + @After + public void tearDown() throws Exception { + System.clearProperty(ZooKeeperServer.GET_DATA_RESPONSE_CACHE_SIZE); + System.clearProperty(ZooKeeperServer.GET_CHILDREN_RESPONSE_CACHE_SIZE); + } + @Test public void testResponseCache() throws Exception { ZooKeeper zk = createClient(); @@ -67,9 +83,15 @@ public void performCacheTest(ZooKeeper zk, String path, boolean useCache) throws long expectedHits = 0; long expectedMisses = 0; - serverFactory.getZooKeeperServer().setResponseCachingEnabled(useCache); + ZooKeeperServer zks = serverFactory.getZooKeeperServer(); + zks.setResponseCachingEnabled(useCache); LOG.info("caching: {}", useCache); + if (useCache) { + assertEquals(zks.getReadResponseCache().getCacheSize(), 32); + assertEquals(zks.getGetChildrenResponseCache().getCacheSize(), 64); + } + byte[] writeData = "test1".getBytes(); zk.create(path, writeData, ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT, writeStat); for (int i = 0; i < reads; ++i) { From b8fd2329784be71b8782fff4753dddde949e9ca7 Mon Sep 17 00:00:00 2001 From: Michael Han Date: Wed, 16 Oct 2019 15:31:11 -0700 Subject: [PATCH 4/4] Address more code review comments. --- .../src/main/resources/markdown/zookeeperAdmin.md | 1 + .../java/org/apache/zookeeper/server/NIOServerCnxn.java | 6 +++++- .../java/org/apache/zookeeper/server/ResponseCache.java | 2 +- 3 files changed, 7 insertions(+), 2 deletions(-) diff --git a/zookeeper-docs/src/main/resources/markdown/zookeeperAdmin.md b/zookeeper-docs/src/main/resources/markdown/zookeeperAdmin.md index c2c49292046..897f9f793f7 100644 --- a/zookeeper-docs/src/main/resources/markdown/zookeeperAdmin.md +++ b/zookeeper-docs/src/main/resources/markdown/zookeeperAdmin.md @@ -721,6 +721,7 @@ property, when available, is noted below. * *maxGetChildrenResponseCacheSize* : (Java system property: **zookeeper.maxGetChildrenResponseCacheSize**) + **New in 3.6.0:** Similar to **maxResponseCacheSize**, but applies to get children requests. The metrics **response_packet_get_children_cache_hits** and **response_packet_get_children_cache_misses** can be used to tune diff --git a/zookeeper-server/src/main/java/org/apache/zookeeper/server/NIOServerCnxn.java b/zookeeper-server/src/main/java/org/apache/zookeeper/server/NIOServerCnxn.java index 5aa427f7e9b..7d4ccebfc17 100644 --- a/zookeeper-server/src/main/java/org/apache/zookeeper/server/NIOServerCnxn.java +++ b/zookeeper-server/src/main/java/org/apache/zookeeper/server/NIOServerCnxn.java @@ -35,6 +35,7 @@ import org.apache.jute.BinaryInputArchive; import org.apache.jute.Record; import org.apache.zookeeper.WatchedEvent; +import org.apache.zookeeper.ZooDefs; import org.apache.zookeeper.data.Id; import org.apache.zookeeper.data.Stat; import org.apache.zookeeper.proto.ReplyHeader; @@ -691,7 +692,10 @@ public void process(WatchedEvent event) { // Convert WatchedEvent to a type that can be sent over the wire WatcherEvent e = event.getWrapper(); - sendResponse(h, e, "notification", null, null, -1); + // The last parameter OpCode here is used to select the response cache. + // Passing OpCode.error (with a value of -1) means we don't care, as we don't need + // response cache on delivering watcher events. + sendResponse(h, e, "notification", null, null, ZooDefs.OpCode.error); } /* diff --git a/zookeeper-server/src/main/java/org/apache/zookeeper/server/ResponseCache.java b/zookeeper-server/src/main/java/org/apache/zookeeper/server/ResponseCache.java index f3d546ec884..4a76a0f89e4 100644 --- a/zookeeper-server/src/main/java/org/apache/zookeeper/server/ResponseCache.java +++ b/zookeeper-server/src/main/java/org/apache/zookeeper/server/ResponseCache.java @@ -37,7 +37,7 @@ private static class Entry { public byte[] data; } - private Map cache; + private final Map cache; public ResponseCache(int cacheSize) { this.cacheSize = cacheSize;