From 498a0a7ef56bde4c37e5151cb046f4fcc3782b55 Mon Sep 17 00:00:00 2001 From: Matteo Merli Date: Tue, 16 Nov 2021 16:15:07 -0800 Subject: [PATCH 1/4] Ensure cache is refreshed (and not just invalidated) after a store write (#12788) ### Motivation When we're doing a write to the store from outside the `MetadataCache`, we are immediately invalidating the cache to ensure read-after-write consistency through the cache. The only issue is that the invalidation, will not trigger a reloading of the value. Instead it is relying on the next call to `cache.get()` which will see the cache miss and it will load the new value into the cache. This means that calls `cache.getIfCached()`, which is not triggering a cache load, will keep seeing the key as missing. ### Modification Ensure we're calling refresh on the cache to get the value automatically reloaded in background and make sure the `getIfCached()` will eventually return the new value, even if there are no calls to `cache.get()`. (cherry picked from commit 2bc449933f72f28dfae24ca8a6fb022be7c55a44) --- .../pulsar/metadata/api/MetadataCache.java | 7 ++++ .../cache/impl/MetadataCacheImpl.java | 17 ++++++---- .../metadata/impl/AbstractMetadataStore.java | 2 +- .../pulsar/metadata/MetadataCacheTest.java | 34 +++++++++++++++++-- 4 files changed, 50 insertions(+), 10 deletions(-) diff --git a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/api/MetadataCache.java b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/api/MetadataCache.java index 360c092b6f1b4..1272130eb76f2 100644 --- a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/api/MetadataCache.java +++ b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/api/MetadataCache.java @@ -148,4 +148,11 @@ public interface MetadataCache { * @param path the path of the object in the metadata store */ void invalidate(String path); + + /** + * Invalidate and reload an object in the metadata cache. + * + * @param path the path of the object in the metadata store + */ + void refresh(String path); } diff --git a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/cache/impl/MetadataCacheImpl.java b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/cache/impl/MetadataCacheImpl.java index 16b419fe9ed91..0ab56bb1bc18d 100644 --- a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/cache/impl/MetadataCacheImpl.java +++ b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/cache/impl/MetadataCacheImpl.java @@ -168,8 +168,7 @@ public CompletableFuture readModifyUpdateOrCreate(String path, Function { - objCache.synchronous().invalidate(path); - objCache.synchronous().refresh(path); + refresh(path); }).thenApply(__ -> newValueObj); }), path); } @@ -198,8 +197,7 @@ public CompletableFuture readModifyUpdate(String path, Function modifyF } return store.put(path, newValue, Optional.of(expectedVersion)).thenAccept(__ -> { - objCache.synchronous().invalidate(path); - objCache.synchronous().refresh(path); + refresh(path); }).thenApply(__ -> newValueObj); }), path); } @@ -220,7 +218,7 @@ public CompletableFuture create(String path, T value) { // In addition to caching the value, we need to add a watch on the path, // so when/if it changes on any other node, we are notified and we can // update the cache - objCache.get(path).whenComplete( (stat2, ex) -> { + objCache.get(path).whenComplete((stat2, ex) -> { if (ex == null) { future.complete(null); } else { @@ -261,6 +259,12 @@ public void invalidate(String path) { objCache.synchronous().invalidate(path); } + @Override + public void refresh(String path) { + objCache.synchronous().invalidate(path); + objCache.synchronous().refresh(path); + } + @VisibleForTesting public void invalidateAll() { objCache.synchronous().invalidateAll(); @@ -275,8 +279,7 @@ public void accept(Notification t) { if (objCache.synchronous().getIfPresent(path) != null) { // Trigger background refresh of the cached item, but before make sure // to invalidate the entry so that we won't serve a stale cached version - objCache.synchronous().invalidate(path); - objCache.synchronous().refresh(path); + refresh(path); } break; diff --git a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/AbstractMetadataStore.java b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/AbstractMetadataStore.java index 40c209f4d3e87..ac7feb4747f67 100644 --- a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/AbstractMetadataStore.java +++ b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/AbstractMetadataStore.java @@ -252,7 +252,7 @@ public final CompletableFuture put(String path, byte[] data, Optional c.invalidate(path)); + metadataCaches.forEach(c -> c.refresh(path)); return stat; }); } diff --git a/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/MetadataCacheTest.java b/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/MetadataCacheTest.java index 322c9bb779bde..e768029656371 100644 --- a/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/MetadataCacheTest.java +++ b/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/MetadataCacheTest.java @@ -278,6 +278,34 @@ public void insertionDeletion(String provider, Supplier urlSupplier) thr assertEquals(objCache.get(key1).join(), Optional.empty()); } + @Test(dataProvider = "impl") + public void insertionWithInvalidation(String provider, Supplier urlSupplier) throws Exception { + @Cleanup + MetadataStore store = MetadataStoreFactory.create(urlSupplier.get(), MetadataStoreConfig.builder().build()); + MetadataCache objCache = store.getMetadataCache(MyClass.class); + + String key1 = newKey(); + + assertEquals(objCache.getIfCached(key1), Optional.empty()); + assertEquals(objCache.get(key1).join(), Optional.empty()); + + MyClass value1 = new MyClass("a", 1); + store.put(key1, ObjectMapperFactory.getThreadLocal().writeValueAsBytes(value1), Optional.of(-1L)).join(); + + Awaitility.await().untilAsserted(() -> { + assertEquals(objCache.getIfCached(key1), Optional.of(value1)); + assertEquals(objCache.get(key1).join(), Optional.of(value1)); + }); + + MyClass value2 = new MyClass("a", 2); + store.put(key1, ObjectMapperFactory.getThreadLocal().writeValueAsBytes(value2), Optional.of(0L)).join(); + + Awaitility.await().untilAsserted(() -> { + assertEquals(objCache.getIfCached(key1), Optional.of(value2)); + assertEquals(objCache.get(key1).join(), Optional.of(value2)); + }); + } + @Test(dataProvider = "impl") public void insertionOutsideCache(String provider, Supplier urlSupplier) throws Exception { @Cleanup @@ -310,8 +338,10 @@ public void insertionOutsideCacheWithGenericType(String provider, Supplier { + assertEquals(objCache.getIfCached(key1), Optional.of(v)); + assertEquals(objCache.get(key1).join(), Optional.of(v)); + }); } @Test(dataProvider = "impl") From 6ee1b0ed23d3e6ccc44b324cba2da2c2a190790c Mon Sep 17 00:00:00 2001 From: JiangHaiting Date: Wed, 1 Dec 2021 05:09:08 +0800 Subject: [PATCH 2/4] fix-12894 (#12896) Co-authored-by: Jiang Haiting (cherry picked from commit 2b939b7a8b0bf9ed2fd31353067033d049ea5350) --- .../pulsar/metadata/cache/impl/MetadataCacheImpl.java | 7 +++++-- .../java/org/apache/pulsar/metadata/MetadataCacheTest.java | 3 ++- 2 files changed, 7 insertions(+), 3 deletions(-) diff --git a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/cache/impl/MetadataCacheImpl.java b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/cache/impl/MetadataCacheImpl.java index 0ab56bb1bc18d..fccbe2ca38a6c 100644 --- a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/cache/impl/MetadataCacheImpl.java +++ b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/cache/impl/MetadataCacheImpl.java @@ -261,8 +261,11 @@ public void invalidate(String path) { @Override public void refresh(String path) { - objCache.synchronous().invalidate(path); - objCache.synchronous().refresh(path); + // Refresh object of path if only it is cached before. + if (objCache.getIfPresent(path) != null) { + objCache.synchronous().invalidate(path); + objCache.synchronous().refresh(path); + } } @VisibleForTesting diff --git a/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/MetadataCacheTest.java b/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/MetadataCacheTest.java index e768029656371..a4ba88852b0e3 100644 --- a/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/MetadataCacheTest.java +++ b/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/MetadataCacheTest.java @@ -325,13 +325,14 @@ public void insertionOutsideCache(String provider, Supplier urlSupplier) } @Test(dataProvider = "impl") - public void insertionOutsideCacheWithGenericType(String provider, Supplier urlSupplier) throws Exception { + public void updateOutsideCacheWithGenericType(String provider, Supplier urlSupplier) throws Exception { @Cleanup MetadataStore store = MetadataStoreFactory.create(urlSupplier.get(), MetadataStoreConfig.builder().build()); MetadataCache> objCache = store.getMetadataCache(new TypeReference>() { }); String key1 = newKey(); + objCache.get(key1); Map v = new TreeMap<>(); v.put("a", "1"); From c82bf93ac7b38437d7563ea82c34c63dcbbe4dad Mon Sep 17 00:00:00 2001 From: Kai Wang Date: Tue, 15 Feb 2022 18:42:13 +0800 Subject: [PATCH 3/4] Fix metadata cache inconsistency on do refresh (#14283) (cherry picked from commit 2e16b4341682083a329b63ffaae6d2df2eeec3c5) --- .../pulsar/metadata/cache/impl/MetadataCacheImpl.java | 11 ++--------- 1 file changed, 2 insertions(+), 9 deletions(-) diff --git a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/cache/impl/MetadataCacheImpl.java b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/cache/impl/MetadataCacheImpl.java index fccbe2ca38a6c..8bf5a729b7694 100644 --- a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/cache/impl/MetadataCacheImpl.java +++ b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/cache/impl/MetadataCacheImpl.java @@ -262,10 +262,7 @@ public void invalidate(String path) { @Override public void refresh(String path) { // Refresh object of path if only it is cached before. - if (objCache.getIfPresent(path) != null) { - objCache.synchronous().invalidate(path); - objCache.synchronous().refresh(path); - } + objCache.asMap().computeIfPresent(path, (oldKey, oldValue) -> readValueFromStore(path)); } @VisibleForTesting @@ -279,11 +276,7 @@ public void accept(Notification t) { switch (t.getType()) { case Created: case Modified: - if (objCache.synchronous().getIfPresent(path) != null) { - // Trigger background refresh of the cached item, but before make sure - // to invalidate the entry so that we won't serve a stale cached version - refresh(path); - } + refresh(path); break; case Deleted: From a0416b780446f3ff472b55c6c2533b537840aa20 Mon Sep 17 00:00:00 2001 From: Lari Hotari Date: Thu, 8 Sep 2022 21:21:46 +0300 Subject: [PATCH 4/4] cherry-pick/#17401 --- .../broker/MultiBrokerTestZKBaseTest.java | 61 ++++++++++ ...ltiBrokerLeaderElectionExpirationTest.java | 114 ++++++++++++++++++ .../MultiBrokerLeaderElectionTest.java | 38 +----- .../metadata/api/MetadataCacheConfig.java | 50 ++++++++ .../pulsar/metadata/api/MetadataStore.java | 57 ++++++++- .../cache/impl/MetadataCacheImpl.java | 26 ++-- .../coordination/impl/LeaderElectionImpl.java | 39 +++--- .../metadata/impl/AbstractMetadataStore.java | 15 ++- .../impl/FaultInjectionMetadataStore.java | 13 +- 9 files changed, 330 insertions(+), 83 deletions(-) create mode 100644 pulsar-broker/src/test/java/org/apache/pulsar/broker/MultiBrokerTestZKBaseTest.java create mode 100644 pulsar-broker/src/test/java/org/apache/pulsar/broker/loadbalance/MultiBrokerLeaderElectionExpirationTest.java create mode 100644 pulsar-metadata/src/main/java/org/apache/pulsar/metadata/api/MetadataCacheConfig.java diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/MultiBrokerTestZKBaseTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/MultiBrokerTestZKBaseTest.java new file mode 100644 index 0000000000000..e6cf86c05b5eb --- /dev/null +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/MultiBrokerTestZKBaseTest.java @@ -0,0 +1,61 @@ +/** + * 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.pulsar.broker; + +import lombok.extern.slf4j.Slf4j; +import org.apache.pulsar.metadata.TestZKServer; +import org.apache.pulsar.metadata.api.MetadataStoreConfig; +import org.apache.pulsar.metadata.api.MetadataStoreException; +import org.apache.pulsar.metadata.api.extended.MetadataStoreExtended; + +/** + * Multiple brokers with a real test Zookeeper server (instead of the mock server) + */ +@Slf4j +public abstract class MultiBrokerTestZKBaseTest extends MultiBrokerBaseTest { + TestZKServer testZKServer; + + @Override + protected void doInitConf() throws Exception { + super.doInitConf(); + testZKServer = new TestZKServer(); + } + + @Override + protected void onCleanup() { + super.onCleanup(); + if (testZKServer != null) { + try { + testZKServer.close(); + } catch (Exception e) { + log.error("Error in stopping ZK server", e); + } + } + } + + @Override + protected MetadataStoreExtended createLocalMetadataStore() throws MetadataStoreException { + return MetadataStoreExtended.create(testZKServer.getConnectionString(), MetadataStoreConfig.builder().build()); + } + + @Override + protected MetadataStoreExtended createConfigurationMetadataStore() throws MetadataStoreException { + return MetadataStoreExtended.create(testZKServer.getConnectionString(), MetadataStoreConfig.builder().build()); + } +} diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/loadbalance/MultiBrokerLeaderElectionExpirationTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/loadbalance/MultiBrokerLeaderElectionExpirationTest.java new file mode 100644 index 0000000000000..65962868b87c9 --- /dev/null +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/loadbalance/MultiBrokerLeaderElectionExpirationTest.java @@ -0,0 +1,114 @@ +/** + * 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.pulsar.broker.loadbalance; + +import static org.mockito.Mockito.spy; +import static org.mockito.Mockito.when; +import static org.testng.Assert.assertEquals; +import static org.testng.Assert.assertTrue; +import java.util.Optional; +import java.util.concurrent.ExecutionException; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.TimeoutException; +import lombok.extern.slf4j.Slf4j; +import org.apache.pulsar.broker.MultiBrokerTestZKBaseTest; +import org.apache.pulsar.broker.PulsarService; +import org.apache.pulsar.metadata.api.MetadataCacheConfig; +import org.apache.pulsar.metadata.api.MetadataStoreException; +import org.apache.pulsar.metadata.api.extended.MetadataStoreExtended; +import org.awaitility.Awaitility; +import org.testng.annotations.Test; + +@Slf4j +@Test(groups = "broker") +public class MultiBrokerLeaderElectionExpirationTest extends MultiBrokerTestZKBaseTest { + private static final long EXPIRE_AFTER_WRITE_MILLIS_IN_TEST = 2000L; + private static final long REFRESH_AFTER_WRITE_MILLIS_IN_TEST = 1000L; + + @Override + protected int numberOfAdditionalBrokers() { + return 9; + } + + @Test + public void shouldElectOneLeader() { + int leaders = 0; + for (PulsarService broker : getAllBrokers()) { + if (broker.getLeaderElectionService().isLeader()) { + leaders++; + } + } + assertEquals(leaders, 1); + } + + @Override + protected MetadataStoreExtended createLocalMetadataStore() throws MetadataStoreException { + return changeDefaultMetadataCacheConfig(super.createLocalMetadataStore()); + } + + @Override + protected MetadataStoreExtended createConfigurationMetadataStore() throws MetadataStoreException { + return changeDefaultMetadataCacheConfig(super.createConfigurationMetadataStore()); + } + + MetadataStoreExtended changeDefaultMetadataCacheConfig(MetadataStoreExtended metadataStore) { + MetadataStoreExtended spy = spy(metadataStore); + when(spy.getDefaultMetadataCacheConfig()).thenReturn(MetadataCacheConfig + .builder() + .refreshAfterWriteMillis(REFRESH_AFTER_WRITE_MILLIS_IN_TEST) + .expireAfterWriteMillis(EXPIRE_AFTER_WRITE_MILLIS_IN_TEST) + .build()); + return spy; + } + + @Test + public void shouldAllBrokersBeAbleToGetTheLeaderAfterExpiration() + throws ExecutionException, InterruptedException, TimeoutException { + + // if you want to see this test fail, modify the line in LeaderElectionImpl constructor for creating + // the metadata cache to not skip expirations: + // this.cache = store.getMetadataCache(clazz); + + // Given that all brokers have the leader elected + Awaitility.await().untilAsserted(() -> { + for (PulsarService broker : getAllBrokers()) { + Optional currentLeader = broker.getLeaderElectionService().getCurrentLeader(); + assertTrue(currentLeader.isPresent(), "Leader wasn't known on broker " + broker.getBrokerServiceUrl()); + } + }); + + // Wait for metadata cache entries to expire + Thread.sleep(EXPIRE_AFTER_WRITE_MILLIS_IN_TEST); + + // then leader should be known on all brokers and it should be the same leader + LeaderBroker leader = null; + for (PulsarService broker : getAllBrokers()) { + Optional currentLeader = + broker.getLeaderElectionService().readCurrentLeader().get(1, TimeUnit.SECONDS); + assertTrue(currentLeader.isPresent(), "Leader wasn't known on broker " + broker.getBrokerServiceUrl()); + if (leader != null) { + assertEquals(currentLeader.get(), leader, + "Different leader on broker " + broker.getBrokerServiceUrl()); + } else { + leader = currentLeader.get(); + } + } + } +} diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/loadbalance/MultiBrokerLeaderElectionTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/loadbalance/MultiBrokerLeaderElectionTest.java index 462b640c17511..36ed1ca6c6190 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/loadbalance/MultiBrokerLeaderElectionTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/loadbalance/MultiBrokerLeaderElectionTest.java @@ -34,55 +34,21 @@ import java.util.stream.IntStream; import lombok.Cleanup; import lombok.extern.slf4j.Slf4j; -import org.apache.pulsar.broker.MultiBrokerBaseTest; +import org.apache.pulsar.broker.MultiBrokerTestZKBaseTest; import org.apache.pulsar.broker.PulsarService; import org.apache.pulsar.client.admin.PulsarAdmin; import org.apache.pulsar.client.admin.PulsarAdminException; -import org.apache.pulsar.metadata.TestZKServer; -import org.apache.pulsar.metadata.api.MetadataStoreConfig; -import org.apache.pulsar.metadata.api.MetadataStoreException; -import org.apache.pulsar.metadata.api.extended.MetadataStoreExtended; import org.awaitility.Awaitility; import org.testng.annotations.Test; @Slf4j @Test(groups = "broker") -public class MultiBrokerLeaderElectionTest extends MultiBrokerBaseTest { +public class MultiBrokerLeaderElectionTest extends MultiBrokerTestZKBaseTest { @Override protected int numberOfAdditionalBrokers() { return 9; } - TestZKServer testZKServer; - - @Override - protected void doInitConf() throws Exception { - super.doInitConf(); - testZKServer = new TestZKServer(); - } - - @Override - protected void onCleanup() { - super.onCleanup(); - if (testZKServer != null) { - try { - testZKServer.close(); - } catch (Exception e) { - log.error("Error in stopping ZK server", e); - } - } - } - - @Override - protected MetadataStoreExtended createLocalMetadataStore() throws MetadataStoreException { - return MetadataStoreExtended.create(testZKServer.getConnectionString(), MetadataStoreConfig.builder().build()); - } - - @Override - protected MetadataStoreExtended createConfigurationMetadataStore() throws MetadataStoreException { - return MetadataStoreExtended.create(testZKServer.getConnectionString(), MetadataStoreConfig.builder().build()); - } - @Test public void shouldElectOneLeader() { int leaders = 0; diff --git a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/api/MetadataCacheConfig.java b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/api/MetadataCacheConfig.java new file mode 100644 index 0000000000000..1671fd39956d8 --- /dev/null +++ b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/api/MetadataCacheConfig.java @@ -0,0 +1,50 @@ +/** + * 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.pulsar.metadata.api; + +import java.util.concurrent.TimeUnit; +import lombok.Builder; +import lombok.Getter; +import lombok.ToString; + +/** + * The configuration builder for a {@link MetadataCache} config. + */ +@Builder +@Getter +@ToString +public class MetadataCacheConfig { + private static final long DEFAULT_CACHE_REFRESH_TIME_MILLIS = TimeUnit.MINUTES.toMillis(5); + + /** + * Specifies that active entries are eligible for automatic refresh once a fixed duration has + * elapsed after the entry's creation, or the most recent replacement of its value. + * A negative or zero value disables automatic refresh. + */ + @Builder.Default + private final long refreshAfterWriteMillis = DEFAULT_CACHE_REFRESH_TIME_MILLIS; + + /** + * Specifies that each entry should be automatically removed from the cache once a fixed duration + * has elapsed after the entry's creation, or the most recent replacement of its value. + * A negative or zero value disables automatic expiration. + */ + @Builder.Default + private final long expireAfterWriteMillis = 2 * DEFAULT_CACHE_REFRESH_TIME_MILLIS; +} diff --git a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/api/MetadataStore.java b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/api/MetadataStore.java index fbebc7b1741ae..8c0f24a571c95 100644 --- a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/api/MetadataStore.java +++ b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/api/MetadataStore.java @@ -138,9 +138,23 @@ public interface MetadataStore extends AutoCloseable { * @param * @param clazz * the class type to be used for serialization/deserialization + * @param cacheConfig + * the cache configuration to be used * @return the metadata cache object */ - MetadataCache getMetadataCache(Class clazz); + MetadataCache getMetadataCache(Class clazz, MetadataCacheConfig cacheConfig); + + /** + * Create a metadata cache specialized for a specific class. + * + * @param + * @param clazz + * the class type to be used for serialization/deserialization + * @return the metadata cache object + */ + default MetadataCache getMetadataCache(Class clazz) { + return getMetadataCache(clazz, getDefaultMetadataCacheConfig()); + } /** * Create a metadata cache specialized for a specific class. @@ -148,9 +162,23 @@ public interface MetadataStore extends AutoCloseable { * @param * @param typeRef * the type ref description to be used for serialization/deserialization + * @param cacheConfig + * the cache configuration to be used * @return the metadata cache object */ - MetadataCache getMetadataCache(TypeReference typeRef); + MetadataCache getMetadataCache(TypeReference typeRef, MetadataCacheConfig cacheConfig); + + /** + * Create a metadata cache specialized for a specific class. + * + * @param + * @param typeRef + * the type ref description to be used for serialization/deserialization + * @return the metadata cache object + */ + default MetadataCache getMetadataCache(TypeReference typeRef) { + return getMetadataCache(typeRef, getDefaultMetadataCacheConfig()); + } /** * Create a metadata cache that uses a particular serde object. @@ -158,7 +186,30 @@ public interface MetadataStore extends AutoCloseable { * @param * @param serde * the custom serialization/deserialization object + * @param cacheConfig + * the cache configuration to be used * @return the metadata cache object */ - MetadataCache getMetadataCache(MetadataSerde serde); + MetadataCache getMetadataCache(MetadataSerde serde, MetadataCacheConfig cacheConfig); + + /** + * Create a metadata cache that uses a particular serde object. + * + * @param + * @param serde + * the custom serialization/deserialization object + * @return the metadata cache object + */ + default MetadataCache getMetadataCache(MetadataSerde serde) { + return getMetadataCache(serde, getDefaultMetadataCacheConfig()); + } + + /** + * Returns the default metadata cache config. + * + * @return default metadata cache config + */ + default MetadataCacheConfig getDefaultMetadataCacheConfig() { + return MetadataCacheConfig.builder().build(); + } } diff --git a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/cache/impl/MetadataCacheImpl.java b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/cache/impl/MetadataCacheImpl.java index 8bf5a729b7694..0ad556e852421 100644 --- a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/cache/impl/MetadataCacheImpl.java +++ b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/cache/impl/MetadataCacheImpl.java @@ -39,6 +39,7 @@ import org.apache.pulsar.metadata.api.CacheGetResult; import org.apache.pulsar.metadata.api.GetResult; import org.apache.pulsar.metadata.api.MetadataCache; +import org.apache.pulsar.metadata.api.MetadataCacheConfig; import org.apache.pulsar.metadata.api.MetadataSerde; import org.apache.pulsar.metadata.api.MetadataStore; import org.apache.pulsar.metadata.api.MetadataStoreException.AlreadyExistsException; @@ -48,31 +49,36 @@ import org.apache.pulsar.metadata.api.Notification; import org.apache.pulsar.metadata.impl.AbstractMetadataStore; +import javax.ws.rs.HEAD; + @Slf4j public class MetadataCacheImpl implements MetadataCache, Consumer { - - private static final long CACHE_REFRESH_TIME_MILLIS = TimeUnit.MINUTES.toMillis(5); - @Getter private final MetadataStore store; private final MetadataSerde serde; private final AsyncLoadingCache>> objCache; - public MetadataCacheImpl(MetadataStore store, TypeReference typeRef) { - this(store, new JSONMetadataSerdeTypeRef<>(typeRef)); + public MetadataCacheImpl(MetadataStore store, TypeReference typeRef, MetadataCacheConfig cacheConfig) { + this(store, new JSONMetadataSerdeTypeRef<>(typeRef), cacheConfig); } - public MetadataCacheImpl(MetadataStore store, JavaType type) { - this(store, new JSONMetadataSerdeSimpleType<>(type)); + public MetadataCacheImpl(MetadataStore store, JavaType type, MetadataCacheConfig cacheConfig) { + this(store, new JSONMetadataSerdeSimpleType<>(type), cacheConfig); } - public MetadataCacheImpl(MetadataStore store, MetadataSerde serde) { + public MetadataCacheImpl(MetadataStore store, MetadataSerde serde, MetadataCacheConfig cacheConfig) { this.store = store; this.serde = serde; - this.objCache = Caffeine.newBuilder() - .refreshAfterWrite(CACHE_REFRESH_TIME_MILLIS, TimeUnit.MILLISECONDS) + Caffeine cacheBuilder = Caffeine.newBuilder(); + if (cacheConfig.getRefreshAfterWriteMillis() > 0) { + cacheBuilder.refreshAfterWrite(cacheConfig.getRefreshAfterWriteMillis(), TimeUnit.MILLISECONDS); + } + if (cacheConfig.getExpireAfterWriteMillis() > 0) { + cacheBuilder.expireAfterWrite(cacheConfig.getExpireAfterWriteMillis(), TimeUnit.MILLISECONDS); + } + this.objCache = cacheBuilder .buildAsync(new AsyncCacheLoader>>() { @Override public CompletableFuture>> asyncLoad(String key, Executor executor) { diff --git a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/coordination/impl/LeaderElectionImpl.java b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/coordination/impl/LeaderElectionImpl.java index 6599c625f788b..86fac33f7c2c9 100644 --- a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/coordination/impl/LeaderElectionImpl.java +++ b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/coordination/impl/LeaderElectionImpl.java @@ -19,28 +19,22 @@ package org.apache.pulsar.metadata.coordination.impl; import com.fasterxml.jackson.databind.type.TypeFactory; - -import io.netty.util.concurrent.DefaultThreadFactory; - -import java.util.ArrayList; import java.util.EnumSet; -import java.util.List; import java.util.Optional; import java.util.concurrent.CompletableFuture; import java.util.concurrent.CompletionException; import java.util.concurrent.ExecutionException; -import java.util.concurrent.Executors; import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.ScheduledFuture; import java.util.concurrent.TimeUnit; import java.util.function.Consumer; - import lombok.extern.slf4j.Slf4j; - import org.apache.bookkeeper.common.concurrent.FutureUtils; import org.apache.bookkeeper.common.util.SafeRunnable; -import org.apache.pulsar.common.util.FutureUtil; import org.apache.pulsar.metadata.api.GetResult; import org.apache.pulsar.metadata.api.MetadataCache; +import org.apache.pulsar.metadata.api.MetadataCacheConfig; +import org.apache.pulsar.metadata.api.MetadataSerde; import org.apache.pulsar.metadata.api.MetadataStoreException; import org.apache.pulsar.metadata.api.MetadataStoreException.AlreadyClosedException; import org.apache.pulsar.metadata.api.MetadataStoreException.BadVersionException; @@ -52,7 +46,6 @@ import org.apache.pulsar.metadata.api.extended.MetadataStoreExtended; import org.apache.pulsar.metadata.api.extended.SessionEvent; import org.apache.pulsar.metadata.cache.impl.JSONMetadataSerdeSimpleType; -import org.apache.pulsar.metadata.api.MetadataSerde; @Slf4j class LeaderElectionImpl implements LeaderElection { @@ -61,6 +54,7 @@ class LeaderElectionImpl implements LeaderElection { private final MetadataStoreExtended store; private final MetadataCache cache; private final Consumer stateChangesListener; + private final ScheduledFuture updateCachedValueFuture; private LeaderElectionState leaderElectionState; private Optional version = Optional.empty(); @@ -82,7 +76,10 @@ private enum InternalState { this.path = path; this.serde = new JSONMetadataSerdeSimpleType<>(TypeFactory.defaultInstance().constructSimpleType(clazz, null)); this.store = store; - this.cache = store.getMetadataCache(clazz); + MetadataCacheConfig metadataCacheConfig = MetadataCacheConfig.builder() + .expireAfterWriteMillis(-1L) + .build(); + this.cache = store.getMetadataCache(clazz, metadataCacheConfig); this.leaderElectionState = LeaderElectionState.NoLeader; this.internalState = InternalState.Init; this.stateChangesListener = stateChangesListener; @@ -90,6 +87,9 @@ private enum InternalState { store.registerListener(this::handlePathNotification); store.registerSessionListener(this::handleSessionNotification); + updateCachedValueFuture = executor.scheduleWithFixedDelay(SafeRunnable.safeRun(this::getLeaderValue), + metadataCacheConfig.getRefreshAfterWriteMillis() / 2, + metadataCacheConfig.getRefreshAfterWriteMillis(), TimeUnit.MILLISECONDS); } @Override @@ -111,10 +111,13 @@ private synchronized CompletableFuture elect() { } else { return tryToBecomeLeader(); } - }).thenCompose(leaderElectionState -> - // make sure that the cache contains the current leader - // so that getLeaderValueIfPresent works on all brokers - cache.get(path).thenApply(__ -> leaderElectionState)); + }).thenComposeAsync(leaderElectionState -> { + // make sure that the cache contains the current leader + // so that getLeaderValueIfPresent works on all brokers + cache.refresh(path); + return cache.get(path) + .thenApply(__ -> leaderElectionState); + }, executor); } private synchronized CompletableFuture handleExistingLeaderValue(GetResult res) { @@ -216,11 +219,6 @@ private synchronized CompletableFuture tryToBecomeLeader() // There was a conflict between 2 participants trying to become leaders at same time. Retry // to fetch info on new leader. - // We force the invalidation of the cache entry. Since we received a BadVersion error, we - // already know that the entry is out of date. If we don't invalidate, we'd be retrying the - // leader election many times until we finally receive the notification that invalidates the - // cache. - cache.invalidate(path); elect() .thenAccept(lse -> result.complete(lse)) .exceptionally(ex2 -> { @@ -238,6 +236,7 @@ private synchronized CompletableFuture tryToBecomeLeader() @Override public void close() throws Exception { + updateCachedValueFuture.cancel(true); try { asyncClose().join(); } catch (CompletionException e) { diff --git a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/AbstractMetadataStore.java b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/AbstractMetadataStore.java index ac7feb4747f67..08b259c5b87d1 100644 --- a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/AbstractMetadataStore.java +++ b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/AbstractMetadataStore.java @@ -40,12 +40,11 @@ import java.util.function.Consumer; import java.util.stream.Collectors; -import lombok.AccessLevel; import lombok.Getter; import lombok.extern.slf4j.Slf4j; - import org.apache.pulsar.common.util.FutureUtil; import org.apache.pulsar.metadata.api.MetadataCache; +import org.apache.pulsar.metadata.api.MetadataCacheConfig; import org.apache.pulsar.metadata.api.MetadataSerde; import org.apache.pulsar.metadata.api.Notification; import org.apache.pulsar.metadata.api.NotificationType; @@ -123,23 +122,23 @@ public CompletableFuture asyncReload(String key, Boolean oldValue, } @Override - public MetadataCache getMetadataCache(Class clazz) { + public MetadataCache getMetadataCache(Class clazz, MetadataCacheConfig cacheConfig) { MetadataCacheImpl metadataCache = new MetadataCacheImpl(this, - TypeFactory.defaultInstance().constructSimpleType(clazz, null)); + TypeFactory.defaultInstance().constructSimpleType(clazz, null), cacheConfig); metadataCaches.add(metadataCache); return metadataCache; } @Override - public MetadataCache getMetadataCache(TypeReference typeRef) { - MetadataCacheImpl metadataCache = new MetadataCacheImpl(this, typeRef); + public MetadataCache getMetadataCache(TypeReference typeRef, MetadataCacheConfig cacheConfig) { + MetadataCacheImpl metadataCache = new MetadataCacheImpl(this, typeRef, cacheConfig); metadataCaches.add(metadataCache); return metadataCache; } @Override - public MetadataCache getMetadataCache(MetadataSerde serde) { - MetadataCacheImpl metadataCache = new MetadataCacheImpl<>(this, serde); + public MetadataCache getMetadataCache(MetadataSerde serde, MetadataCacheConfig cacheConfig) { + MetadataCacheImpl metadataCache = new MetadataCacheImpl<>(this, serde, cacheConfig); metadataCaches.add(metadataCache); return metadataCache; } diff --git a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/FaultInjectionMetadataStore.java b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/FaultInjectionMetadataStore.java index 65731c0fed8fa..e351433feb923 100644 --- a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/FaultInjectionMetadataStore.java +++ b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/FaultInjectionMetadataStore.java @@ -31,6 +31,7 @@ import org.apache.pulsar.common.util.FutureUtil; import org.apache.pulsar.metadata.api.GetResult; import org.apache.pulsar.metadata.api.MetadataCache; +import org.apache.pulsar.metadata.api.MetadataCacheConfig; import org.apache.pulsar.metadata.api.MetadataSerde; import org.apache.pulsar.metadata.api.MetadataStore; import org.apache.pulsar.metadata.api.MetadataStoreException; @@ -147,18 +148,18 @@ public void registerListener(Consumer listener) { } @Override - public MetadataCache getMetadataCache(Class clazz) { - return store.getMetadataCache(clazz); + public MetadataCache getMetadataCache(Class clazz, MetadataCacheConfig cacheConfig) { + return store.getMetadataCache(clazz, cacheConfig); } @Override - public MetadataCache getMetadataCache(TypeReference typeRef) { - return store.getMetadataCache(typeRef); + public MetadataCache getMetadataCache(TypeReference typeRef, MetadataCacheConfig cacheConfig) { + return store.getMetadataCache(typeRef, cacheConfig); } @Override - public MetadataCache getMetadataCache(MetadataSerde serde) { - return store.getMetadataCache(serde); + public MetadataCache getMetadataCache(MetadataSerde serde, MetadataCacheConfig cacheConfig) { + return store.getMetadataCache(serde, cacheConfig); } @Override