From 114e9bfe2b2a46fc8b36b2afd55b3b539d383d96 Mon Sep 17 00:00:00 2001 From: horizonzy Date: Thu, 31 Aug 2023 01:21:26 +0800 Subject: [PATCH 01/12] Fix isolated group not work problem. --- ...IsolatedBookieEnsemblePlacementPolicy.java | 21 ++++++++++--------- 1 file changed, 11 insertions(+), 10 deletions(-) diff --git a/pulsar-broker-common/src/main/java/org/apache/pulsar/bookie/rackawareness/IsolatedBookieEnsemblePlacementPolicy.java b/pulsar-broker-common/src/main/java/org/apache/pulsar/bookie/rackawareness/IsolatedBookieEnsemblePlacementPolicy.java index ae288fcf7f238..97a537199dc0b 100644 --- a/pulsar-broker-common/src/main/java/org/apache/pulsar/bookie/rackawareness/IsolatedBookieEnsemblePlacementPolicy.java +++ b/pulsar-broker-common/src/main/java/org/apache/pulsar/bookie/rackawareness/IsolatedBookieEnsemblePlacementPolicy.java @@ -28,7 +28,8 @@ import java.util.Map; import java.util.Optional; import java.util.Set; -import java.util.concurrent.CompletableFuture; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.TimeoutException; import lombok.extern.slf4j.Slf4j; import org.apache.bookkeeper.client.BKException.BKNotEnoughBookiesException; import org.apache.bookkeeper.client.RackawareEnsemblePlacementPolicy; @@ -48,6 +49,8 @@ import org.apache.pulsar.common.policies.data.EnsemblePlacementPolicyConfig; import org.apache.pulsar.metadata.api.MetadataCache; import org.apache.pulsar.metadata.api.MetadataStore; +import org.apache.pulsar.zookeeper.ZkBookieRackAffinityMapping; +import org.apache.zookeeper.KeeperException; @Slf4j public class IsolatedBookieEnsemblePlacementPolicy extends RackawareEnsemblePlacementPolicy { @@ -187,17 +190,13 @@ private Set getBlacklistedBookiesWithIsolationGroups(int ensembleSize, Set blacklistedBookies = new HashSet<>(); try { if (bookieMappingCache != null) { - CompletableFuture> future = - bookieMappingCache.get(BookieRackAffinityMapping.BOOKIE_INFO_ROOT_PATH); - - Optional optRes = (future.isDone() && !future.isCompletedExceptionally()) - ? future.join() : Optional.empty(); - - if (!optRes.isPresent()) { - return blacklistedBookies; + Optional optional = + bookieMappingCache.get(BookieRackAffinityMapping.BOOKIE_INFO_ROOT_PATH).get(30, TimeUnit.SECONDS); + if (!optional.isPresent()) { + throw new KeeperException.NoNodeException(ZkBookieRackAffinityMapping.BOOKIE_INFO_ROOT_PATH); } - BookiesRackConfiguration allGroupsBookieMapping = optRes.get(); + BookiesRackConfiguration allGroupsBookieMapping = optional.get(); Set allBookies = allGroupsBookieMapping.keySet(); int totalAvailableBookiesInPrimaryGroup = 0; Set primaryIsolationGroup = Collections.emptySet(); @@ -246,6 +245,8 @@ private Set getBlacklistedBookiesWithIsolationGroups(int ensembleSize, } } } + } catch (TimeoutException e) { + log.warn("Getting bookie isolation info from metadata store timeout."); } catch (Exception e) { log.warn("Error getting bookie isolation info from metadata store: {}", e.getMessage()); } From f750c30421c14bf1f9603494ad420b1dbf5e2fbb Mon Sep 17 00:00:00 2001 From: horizonzy Date: Thu, 31 Aug 2023 01:34:23 +0800 Subject: [PATCH 02/12] Fix checkstyle. --- .../rackawareness/IsolatedBookieEnsemblePlacementPolicy.java | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/pulsar-broker-common/src/main/java/org/apache/pulsar/bookie/rackawareness/IsolatedBookieEnsemblePlacementPolicy.java b/pulsar-broker-common/src/main/java/org/apache/pulsar/bookie/rackawareness/IsolatedBookieEnsemblePlacementPolicy.java index 97a537199dc0b..9bcd467811978 100644 --- a/pulsar-broker-common/src/main/java/org/apache/pulsar/bookie/rackawareness/IsolatedBookieEnsemblePlacementPolicy.java +++ b/pulsar-broker-common/src/main/java/org/apache/pulsar/bookie/rackawareness/IsolatedBookieEnsemblePlacementPolicy.java @@ -191,7 +191,8 @@ private Set getBlacklistedBookiesWithIsolationGroups(int ensembleSize, try { if (bookieMappingCache != null) { Optional optional = - bookieMappingCache.get(BookieRackAffinityMapping.BOOKIE_INFO_ROOT_PATH).get(30, TimeUnit.SECONDS); + bookieMappingCache.get(BookieRackAffinityMapping.BOOKIE_INFO_ROOT_PATH) + .get(30, TimeUnit.SECONDS); if (!optional.isPresent()) { throw new KeeperException.NoNodeException(ZkBookieRackAffinityMapping.BOOKIE_INFO_ROOT_PATH); } From ee4d1acadc4c494a8d3cd27a2dfcc257548133c9 Mon Sep 17 00:00:00 2001 From: horizonzy Date: Thu, 31 Aug 2023 11:40:07 +0800 Subject: [PATCH 03/12] Add test case. --- ...IsolatedBookieEnsemblePlacementPolicy.java | 9 ++- ...atedBookieEnsemblePlacementPolicyTest.java | 75 +++++++++++++++++++ 2 files changed, 81 insertions(+), 3 deletions(-) diff --git a/pulsar-broker-common/src/main/java/org/apache/pulsar/bookie/rackawareness/IsolatedBookieEnsemblePlacementPolicy.java b/pulsar-broker-common/src/main/java/org/apache/pulsar/bookie/rackawareness/IsolatedBookieEnsemblePlacementPolicy.java index 9bcd467811978..67fc9d24164a1 100644 --- a/pulsar-broker-common/src/main/java/org/apache/pulsar/bookie/rackawareness/IsolatedBookieEnsemblePlacementPolicy.java +++ b/pulsar-broker-common/src/main/java/org/apache/pulsar/bookie/rackawareness/IsolatedBookieEnsemblePlacementPolicy.java @@ -19,6 +19,7 @@ package org.apache.pulsar.bookie.rackawareness; import static org.apache.pulsar.bookie.rackawareness.BookieRackAffinityMapping.METADATA_STORE_INSTANCE; +import com.google.common.annotations.VisibleForTesting; import io.netty.util.HashedWheelTimer; import java.util.Arrays; import java.util.Collections; @@ -62,7 +63,8 @@ public class IsolatedBookieEnsemblePlacementPolicy extends RackawareEnsemblePlac private ImmutablePair, Set> defaultIsolationGroups; private MetadataCache bookieMappingCache; - + //For test. + long metaOpTimeout = TimeUnit.SECONDS.toMillis(30); public IsolatedBookieEnsemblePlacementPolicy() { super(); @@ -185,14 +187,15 @@ private static Pair, Set> getIsolationGroup( return pair; } - private Set getBlacklistedBookiesWithIsolationGroups(int ensembleSize, + @VisibleForTesting + Set getBlacklistedBookiesWithIsolationGroups(int ensembleSize, Pair, Set> isolationGroups) { Set blacklistedBookies = new HashSet<>(); try { if (bookieMappingCache != null) { Optional optional = bookieMappingCache.get(BookieRackAffinityMapping.BOOKIE_INFO_ROOT_PATH) - .get(30, TimeUnit.SECONDS); + .get(metaOpTimeout, TimeUnit.MILLISECONDS); if (!optional.isPresent()) { throw new KeeperException.NoNodeException(ZkBookieRackAffinityMapping.BOOKIE_INFO_ROOT_PATH); } diff --git a/pulsar-broker-common/src/test/java/org/apache/pulsar/bookie/rackawareness/IsolatedBookieEnsemblePlacementPolicyTest.java b/pulsar-broker-common/src/test/java/org/apache/pulsar/bookie/rackawareness/IsolatedBookieEnsemblePlacementPolicyTest.java index 85feeaecfdd78..8e1cf4b64076c 100644 --- a/pulsar-broker-common/src/test/java/org/apache/pulsar/bookie/rackawareness/IsolatedBookieEnsemblePlacementPolicyTest.java +++ b/pulsar-broker-common/src/test/java/org/apache/pulsar/bookie/rackawareness/IsolatedBookieEnsemblePlacementPolicyTest.java @@ -18,11 +18,14 @@ */ package org.apache.pulsar.bookie.rackawareness; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.when; import static org.testng.Assert.assertEquals; import static org.testng.Assert.assertFalse; import static org.testng.Assert.assertTrue; import static org.testng.Assert.fail; import com.fasterxml.jackson.databind.ObjectMapper; +import com.google.common.collect.Sets; import io.netty.util.HashedWheelTimer; import java.nio.charset.StandardCharsets; import java.util.ArrayList; @@ -34,18 +37,24 @@ import java.util.Map; import java.util.Optional; import java.util.Set; +import java.util.concurrent.CompletableFuture; import org.apache.bookkeeper.client.BKException.BKNotEnoughBookiesException; import org.apache.bookkeeper.conf.ClientConfiguration; import org.apache.bookkeeper.feature.SettableFeatureProvider; import org.apache.bookkeeper.net.BookieId; import org.apache.bookkeeper.net.BookieSocketAddress; import org.apache.bookkeeper.stats.NullStatsLogger; +import org.apache.commons.lang3.tuple.MutablePair; +import org.apache.commons.lang3.tuple.Pair; import org.apache.pulsar.common.policies.data.BookieInfo; +import org.apache.pulsar.common.policies.data.BookiesRackConfiguration; import org.apache.pulsar.common.policies.data.EnsemblePlacementPolicyConfig; import org.apache.pulsar.common.util.ObjectMapperFactory; import org.apache.pulsar.metadata.api.MetadataStore; import org.apache.pulsar.metadata.api.MetadataStoreConfig; import org.apache.pulsar.metadata.api.MetadataStoreFactory; +import org.apache.pulsar.metadata.api.extended.MetadataStoreExtended; +import org.apache.pulsar.metadata.cache.impl.MetadataCacheImpl; import org.awaitility.Awaitility; import org.testng.annotations.AfterMethod; import org.testng.annotations.BeforeMethod; @@ -87,6 +96,72 @@ void teardown() throws Exception { timer.stop(); } + @Test + public void testGetDataFromMetadataFailed() throws Exception { + Map mainBookieGroup = new HashMap<>(); + mainBookieGroup.put(BOOKIE1, BookieInfo.builder().rack("rack0").build()); + mainBookieGroup.put(BOOKIE2, BookieInfo.builder().rack("rack1").build()); + + Map secondaryBookieGroup = new HashMap<>(); + secondaryBookieGroup.put(BOOKIE3, BookieInfo.builder().rack("rack0").build()); + + store = mock(MetadataStoreExtended.class); + MetadataCacheImpl cache = mock(MetadataCacheImpl.class); + when(store.getMetadataCache(BookiesRackConfiguration.class)).thenReturn(cache); + CompletableFuture completableFuture = CompletableFuture.completedFuture(null); + long metaOpTimeout = 3000; + CompletableFuture> waitingCompleteFuture = new CompletableFuture<>(); + new Thread(() -> { + try { + Thread.sleep(metaOpTimeout - 1000); + } catch (InterruptedException e) { + throw new RuntimeException(e); + } + BookiesRackConfiguration rackConfiguration = new BookiesRackConfiguration(); + rackConfiguration.put("group1", mainBookieGroup); + rackConfiguration.put("group2", secondaryBookieGroup); + waitingCompleteFuture.complete(Optional.of(rackConfiguration)); + }).start(); + + CompletableFuture> timeoutFuture = new CompletableFuture<>(); + new Thread(() -> { + try { + Thread.sleep(metaOpTimeout + 1000); + } catch (InterruptedException e) { + throw new RuntimeException(e); + } + BookiesRackConfiguration rackConfiguration = new BookiesRackConfiguration(); + rackConfiguration.put("group1", mainBookieGroup); + rackConfiguration.put("group2", secondaryBookieGroup); + waitingCompleteFuture.complete(Optional.of(rackConfiguration)); + }).start(); + + + when(cache.get(BookieRackAffinityMapping.BOOKIE_INFO_ROOT_PATH)).thenReturn(completableFuture) + .thenReturn(waitingCompleteFuture).thenReturn(timeoutFuture); + + IsolatedBookieEnsemblePlacementPolicy isolationPolicy = new IsolatedBookieEnsemblePlacementPolicy(); + isolationPolicy.metaOpTimeout = metaOpTimeout; + ClientConfiguration bkClientConf = new ClientConfiguration(); + bkClientConf.setProperty(BookieRackAffinityMapping.METADATA_STORE_INSTANCE, store); + bkClientConf.setProperty(IsolatedBookieEnsemblePlacementPolicy.ISOLATION_BOOKIE_GROUPS, isolationGroups); + isolationPolicy.initialize(bkClientConf, Optional.empty(), timer, SettableFeatureProvider.DISABLE_ALL, NullStatsLogger.INSTANCE, BookieSocketAddress.LEGACY_BOOKIEID_RESOLVER); + isolationPolicy.onClusterChanged(writableBookies, readOnlyBookies); + + MutablePair, Set> groups = new MutablePair<>(); + groups.setLeft(Sets.newHashSet("group1")); + groups.setRight(Sets.newHashSet()); + + Set blacklist = + isolationPolicy.getBlacklistedBookiesWithIsolationGroups(2, groups); + assertFalse(blacklist.isEmpty()); + + blacklist = + isolationPolicy.getBlacklistedBookiesWithIsolationGroups(2, groups); + assertTrue(blacklist.isEmpty()); + + } + @Test public void testBasic() throws Exception { Map> bookieMapping = new HashMap<>(); From 0129a59cb3012ecc662c34d788eb744f3ac04e84 Mon Sep 17 00:00:00 2001 From: horizonzy Date: Thu, 31 Aug 2023 11:40:58 +0800 Subject: [PATCH 04/12] Code clean. --- .../rackawareness/IsolatedBookieEnsemblePlacementPolicy.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pulsar-broker-common/src/main/java/org/apache/pulsar/bookie/rackawareness/IsolatedBookieEnsemblePlacementPolicy.java b/pulsar-broker-common/src/main/java/org/apache/pulsar/bookie/rackawareness/IsolatedBookieEnsemblePlacementPolicy.java index 67fc9d24164a1..02b39a0950bed 100644 --- a/pulsar-broker-common/src/main/java/org/apache/pulsar/bookie/rackawareness/IsolatedBookieEnsemblePlacementPolicy.java +++ b/pulsar-broker-common/src/main/java/org/apache/pulsar/bookie/rackawareness/IsolatedBookieEnsemblePlacementPolicy.java @@ -63,7 +63,7 @@ public class IsolatedBookieEnsemblePlacementPolicy extends RackawareEnsemblePlac private ImmutablePair, Set> defaultIsolationGroups; private MetadataCache bookieMappingCache; - //For test. + @VisibleForTesting long metaOpTimeout = TimeUnit.SECONDS.toMillis(30); public IsolatedBookieEnsemblePlacementPolicy() { From 69ac89f28e42e1d20f9e47b2137660bf989684b7 Mon Sep 17 00:00:00 2001 From: horizonzy Date: Thu, 31 Aug 2023 11:51:29 +0800 Subject: [PATCH 05/12] fix ci. --- .../IsolatedBookieEnsemblePlacementPolicyTest.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/pulsar-broker-common/src/test/java/org/apache/pulsar/bookie/rackawareness/IsolatedBookieEnsemblePlacementPolicyTest.java b/pulsar-broker-common/src/test/java/org/apache/pulsar/bookie/rackawareness/IsolatedBookieEnsemblePlacementPolicyTest.java index 8e1cf4b64076c..a92c9954d28d4 100644 --- a/pulsar-broker-common/src/test/java/org/apache/pulsar/bookie/rackawareness/IsolatedBookieEnsemblePlacementPolicyTest.java +++ b/pulsar-broker-common/src/test/java/org/apache/pulsar/bookie/rackawareness/IsolatedBookieEnsemblePlacementPolicyTest.java @@ -97,7 +97,7 @@ void teardown() throws Exception { } @Test - public void testGetDataFromMetadataFailed() throws Exception { + public void testMetadataStoreCases() throws Exception { Map mainBookieGroup = new HashMap<>(); mainBookieGroup.put(BOOKIE1, BookieInfo.builder().rack("rack0").build()); mainBookieGroup.put(BOOKIE2, BookieInfo.builder().rack("rack1").build()); @@ -150,7 +150,7 @@ public void testGetDataFromMetadataFailed() throws Exception { MutablePair, Set> groups = new MutablePair<>(); groups.setLeft(Sets.newHashSet("group1")); - groups.setRight(Sets.newHashSet()); + groups.setRight(new HashSet<>()); Set blacklist = isolationPolicy.getBlacklistedBookiesWithIsolationGroups(2, groups); From 7775182d0b2e1f18c86b71760c63242ca2cb0e8b Mon Sep 17 00:00:00 2001 From: horizonzy Date: Thu, 31 Aug 2023 12:34:06 +0800 Subject: [PATCH 06/12] fix ci. --- .../rackawareness/IsolatedBookieEnsemblePlacementPolicyTest.java | 1 - 1 file changed, 1 deletion(-) diff --git a/pulsar-broker-common/src/test/java/org/apache/pulsar/bookie/rackawareness/IsolatedBookieEnsemblePlacementPolicyTest.java b/pulsar-broker-common/src/test/java/org/apache/pulsar/bookie/rackawareness/IsolatedBookieEnsemblePlacementPolicyTest.java index a92c9954d28d4..0f86ace56fd51 100644 --- a/pulsar-broker-common/src/test/java/org/apache/pulsar/bookie/rackawareness/IsolatedBookieEnsemblePlacementPolicyTest.java +++ b/pulsar-broker-common/src/test/java/org/apache/pulsar/bookie/rackawareness/IsolatedBookieEnsemblePlacementPolicyTest.java @@ -45,7 +45,6 @@ import org.apache.bookkeeper.net.BookieSocketAddress; import org.apache.bookkeeper.stats.NullStatsLogger; import org.apache.commons.lang3.tuple.MutablePair; -import org.apache.commons.lang3.tuple.Pair; import org.apache.pulsar.common.policies.data.BookieInfo; import org.apache.pulsar.common.policies.data.BookiesRackConfiguration; import org.apache.pulsar.common.policies.data.EnsemblePlacementPolicyConfig; From 9bb1024b59fda8a411bff40094b107abef29ef22 Mon Sep 17 00:00:00 2001 From: horizonzy Date: Fri, 1 Sep 2023 00:41:51 +0800 Subject: [PATCH 07/12] tune code. --- .../IsolatedBookieEnsemblePlacementPolicyTest.java | 5 ++--- 1 file changed, 2 insertions(+), 3 deletions(-) diff --git a/pulsar-broker-common/src/test/java/org/apache/pulsar/bookie/rackawareness/IsolatedBookieEnsemblePlacementPolicyTest.java b/pulsar-broker-common/src/test/java/org/apache/pulsar/bookie/rackawareness/IsolatedBookieEnsemblePlacementPolicyTest.java index 0f86ace56fd51..e4e6f93806978 100644 --- a/pulsar-broker-common/src/test/java/org/apache/pulsar/bookie/rackawareness/IsolatedBookieEnsemblePlacementPolicyTest.java +++ b/pulsar-broker-common/src/test/java/org/apache/pulsar/bookie/rackawareness/IsolatedBookieEnsemblePlacementPolicyTest.java @@ -125,14 +125,14 @@ public void testMetadataStoreCases() throws Exception { CompletableFuture> timeoutFuture = new CompletableFuture<>(); new Thread(() -> { try { - Thread.sleep(metaOpTimeout + 1000); + Thread.sleep(metaOpTimeout + 5000); } catch (InterruptedException e) { throw new RuntimeException(e); } BookiesRackConfiguration rackConfiguration = new BookiesRackConfiguration(); rackConfiguration.put("group1", mainBookieGroup); rackConfiguration.put("group2", secondaryBookieGroup); - waitingCompleteFuture.complete(Optional.of(rackConfiguration)); + timeoutFuture.complete(Optional.of(rackConfiguration)); }).start(); @@ -158,7 +158,6 @@ public void testMetadataStoreCases() throws Exception { blacklist = isolationPolicy.getBlacklistedBookiesWithIsolationGroups(2, groups); assertTrue(blacklist.isEmpty()); - } @Test From 51e8589de2e65fc1d28e2b50f56bf3d771f4b656 Mon Sep 17 00:00:00 2001 From: horizonzy Date: Fri, 1 Sep 2023 23:37:34 +0800 Subject: [PATCH 08/12] cache the RackConfiguration to avoid sync operation. --- ...IsolatedBookieEnsemblePlacementPolicy.java | 39 +++++++++------ ...atedBookieEnsemblePlacementPolicyTest.java | 50 +++++++++---------- 2 files changed, 48 insertions(+), 41 deletions(-) diff --git a/pulsar-broker-common/src/main/java/org/apache/pulsar/bookie/rackawareness/IsolatedBookieEnsemblePlacementPolicy.java b/pulsar-broker-common/src/main/java/org/apache/pulsar/bookie/rackawareness/IsolatedBookieEnsemblePlacementPolicy.java index 02b39a0950bed..ffe66d5d5651a 100644 --- a/pulsar-broker-common/src/main/java/org/apache/pulsar/bookie/rackawareness/IsolatedBookieEnsemblePlacementPolicy.java +++ b/pulsar-broker-common/src/main/java/org/apache/pulsar/bookie/rackawareness/IsolatedBookieEnsemblePlacementPolicy.java @@ -23,14 +23,12 @@ import io.netty.util.HashedWheelTimer; import java.util.Arrays; import java.util.Collections; -import java.util.HashMap; import java.util.HashSet; import java.util.List; import java.util.Map; import java.util.Optional; import java.util.Set; -import java.util.concurrent.TimeUnit; -import java.util.concurrent.TimeoutException; +import java.util.concurrent.CompletableFuture; import lombok.extern.slf4j.Slf4j; import org.apache.bookkeeper.client.BKException.BKNotEnoughBookiesException; import org.apache.bookkeeper.client.RackawareEnsemblePlacementPolicy; @@ -63,8 +61,8 @@ public class IsolatedBookieEnsemblePlacementPolicy extends RackawareEnsemblePlac private ImmutablePair, Set> defaultIsolationGroups; private MetadataCache bookieMappingCache; - @VisibleForTesting - long metaOpTimeout = TimeUnit.SECONDS.toMillis(30); + + private BookiesRackConfiguration cachedRackConfiguration = null; public IsolatedBookieEnsemblePlacementPolicy() { super(); @@ -92,7 +90,9 @@ public RackawareEnsemblePlacementPolicyImpl initialize(ClientConfiguration conf, // Only add the bookieMappingCache if we have defined an isolation group bookieMappingCache = store.getMetadataCache(BookiesRackConfiguration.class); - bookieMappingCache.get(BookieRackAffinityMapping.BOOKIE_INFO_ROOT_PATH).join(); + Optional optional = + bookieMappingCache.get(BookieRackAffinityMapping.BOOKIE_INFO_ROOT_PATH).join(); + optional.ifPresent(bookiesRackConfiguration -> cachedRackConfiguration = bookiesRackConfiguration); } } if (conf.getProperty(SECONDARY_ISOLATION_BOOKIE_GROUPS) != null) { @@ -112,7 +112,6 @@ public RackawareEnsemblePlacementPolicyImpl initialize(ClientConfiguration conf, public PlacementResult> newEnsemble(int ensembleSize, int writeQuorumSize, int ackQuorumSize, Map customMetadata, Set excludeBookies) throws BKNotEnoughBookiesException { - Map> isolationGroup = new HashMap<>(); Set blacklistedBookies = getBlacklistedBookiesWithIsolationGroups( ensembleSize, defaultIsolationGroups); if (excludeBookies == null) { @@ -193,14 +192,26 @@ Set getBlacklistedBookiesWithIsolationGroups(int ensembleSize, Set blacklistedBookies = new HashSet<>(); try { if (bookieMappingCache != null) { - Optional optional = - bookieMappingCache.get(BookieRackAffinityMapping.BOOKIE_INFO_ROOT_PATH) - .get(metaOpTimeout, TimeUnit.MILLISECONDS); - if (!optional.isPresent()) { - throw new KeeperException.NoNodeException(ZkBookieRackAffinityMapping.BOOKIE_INFO_ROOT_PATH); + CompletableFuture> future = + bookieMappingCache.get(BookieRackAffinityMapping.BOOKIE_INFO_ROOT_PATH); + + BookiesRackConfiguration allGroupsBookieMapping; + Optional optRes = (future.isDone() && !future.isCompletedExceptionally()) + ? future.join() : Optional.empty(); + + if (!optRes.isPresent()) { + if (cachedRackConfiguration != null) { + log.debug("The newest rack config is not available now, use the cached rack config : {}", + cachedRackConfiguration); + allGroupsBookieMapping = cachedRackConfiguration; + } else { + throw new KeeperException.NoNodeException(ZkBookieRackAffinityMapping.BOOKIE_INFO_ROOT_PATH); + } + } else { + cachedRackConfiguration = optRes.get(); + allGroupsBookieMapping = optRes.get(); } - BookiesRackConfiguration allGroupsBookieMapping = optional.get(); Set allBookies = allGroupsBookieMapping.keySet(); int totalAvailableBookiesInPrimaryGroup = 0; Set primaryIsolationGroup = Collections.emptySet(); @@ -249,8 +260,6 @@ Set getBlacklistedBookiesWithIsolationGroups(int ensembleSize, } } } - } catch (TimeoutException e) { - log.warn("Getting bookie isolation info from metadata store timeout."); } catch (Exception e) { log.warn("Error getting bookie isolation info from metadata store: {}", e.getMessage()); } diff --git a/pulsar-broker-common/src/test/java/org/apache/pulsar/bookie/rackawareness/IsolatedBookieEnsemblePlacementPolicyTest.java b/pulsar-broker-common/src/test/java/org/apache/pulsar/bookie/rackawareness/IsolatedBookieEnsemblePlacementPolicyTest.java index e4e6f93806978..73ddd91c231b0 100644 --- a/pulsar-broker-common/src/test/java/org/apache/pulsar/bookie/rackawareness/IsolatedBookieEnsemblePlacementPolicyTest.java +++ b/pulsar-broker-common/src/test/java/org/apache/pulsar/bookie/rackawareness/IsolatedBookieEnsemblePlacementPolicyTest.java @@ -107,40 +107,31 @@ public void testMetadataStoreCases() throws Exception { store = mock(MetadataStoreExtended.class); MetadataCacheImpl cache = mock(MetadataCacheImpl.class); when(store.getMetadataCache(BookiesRackConfiguration.class)).thenReturn(cache); - CompletableFuture completableFuture = CompletableFuture.completedFuture(null); - long metaOpTimeout = 3000; - CompletableFuture> waitingCompleteFuture = new CompletableFuture<>(); - new Thread(() -> { - try { - Thread.sleep(metaOpTimeout - 1000); - } catch (InterruptedException e) { - throw new RuntimeException(e); - } - BookiesRackConfiguration rackConfiguration = new BookiesRackConfiguration(); - rackConfiguration.put("group1", mainBookieGroup); - rackConfiguration.put("group2", secondaryBookieGroup); - waitingCompleteFuture.complete(Optional.of(rackConfiguration)); - }).start(); + CompletableFuture> initialFuture = new CompletableFuture<>(); + //The initialFuture only has group1. + BookiesRackConfiguration rackConfiguration1 = new BookiesRackConfiguration(); + rackConfiguration1.put("group1", mainBookieGroup); + initialFuture.complete(Optional.of(rackConfiguration1)); - CompletableFuture> timeoutFuture = new CompletableFuture<>(); + long waitTime = 2000; + CompletableFuture> waitingCompleteFuture = new CompletableFuture<>(); new Thread(() -> { try { - Thread.sleep(metaOpTimeout + 5000); + Thread.sleep(waitTime); } catch (InterruptedException e) { throw new RuntimeException(e); } - BookiesRackConfiguration rackConfiguration = new BookiesRackConfiguration(); - rackConfiguration.put("group1", mainBookieGroup); - rackConfiguration.put("group2", secondaryBookieGroup); - timeoutFuture.complete(Optional.of(rackConfiguration)); + //The waitingCompleteFuture has group1 and group2. + BookiesRackConfiguration rackConfiguration2 = new BookiesRackConfiguration(); + rackConfiguration2.put("group1", mainBookieGroup); + rackConfiguration2.put("group2", secondaryBookieGroup); + waitingCompleteFuture.complete(Optional.of(rackConfiguration2)); }).start(); - - when(cache.get(BookieRackAffinityMapping.BOOKIE_INFO_ROOT_PATH)).thenReturn(completableFuture) - .thenReturn(waitingCompleteFuture).thenReturn(timeoutFuture); + when(cache.get(BookieRackAffinityMapping.BOOKIE_INFO_ROOT_PATH)).thenReturn(initialFuture) + .thenReturn(waitingCompleteFuture); IsolatedBookieEnsemblePlacementPolicy isolationPolicy = new IsolatedBookieEnsemblePlacementPolicy(); - isolationPolicy.metaOpTimeout = metaOpTimeout; ClientConfiguration bkClientConf = new ClientConfiguration(); bkClientConf.setProperty(BookieRackAffinityMapping.METADATA_STORE_INSTANCE, store); bkClientConf.setProperty(IsolatedBookieEnsemblePlacementPolicy.ISOLATION_BOOKIE_GROUPS, isolationGroups); @@ -151,13 +142,20 @@ public void testMetadataStoreCases() throws Exception { groups.setLeft(Sets.newHashSet("group1")); groups.setRight(new HashSet<>()); + //The future is waiting done, so use the cached rack config. Set blacklist = isolationPolicy.getBlacklistedBookiesWithIsolationGroups(2, groups); - assertFalse(blacklist.isEmpty()); + assertTrue(blacklist.isEmpty()); + Thread.sleep(waitTime); + + //The future is already done, use the newest blacklist = isolationPolicy.getBlacklistedBookiesWithIsolationGroups(2, groups); - assertTrue(blacklist.isEmpty()); + assertFalse(blacklist.isEmpty()); + assertEquals(blacklist.size(), 1); + BookieId excludeBookie = blacklist.iterator().next(); + assertEquals(excludeBookie.toString(), BOOKIE3); } @Test From d3c118eaead9bd0eed66fe5eb352746af4133cc5 Mon Sep 17 00:00:00 2001 From: horizonzy Date: Fri, 1 Sep 2023 23:44:41 +0800 Subject: [PATCH 09/12] tune comment. --- .../IsolatedBookieEnsemblePlacementPolicyTest.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pulsar-broker-common/src/test/java/org/apache/pulsar/bookie/rackawareness/IsolatedBookieEnsemblePlacementPolicyTest.java b/pulsar-broker-common/src/test/java/org/apache/pulsar/bookie/rackawareness/IsolatedBookieEnsemblePlacementPolicyTest.java index 73ddd91c231b0..eda28878c167f 100644 --- a/pulsar-broker-common/src/test/java/org/apache/pulsar/bookie/rackawareness/IsolatedBookieEnsemblePlacementPolicyTest.java +++ b/pulsar-broker-common/src/test/java/org/apache/pulsar/bookie/rackawareness/IsolatedBookieEnsemblePlacementPolicyTest.java @@ -149,7 +149,7 @@ public void testMetadataStoreCases() throws Exception { Thread.sleep(waitTime); - //The future is already done, use the newest + //The future is already done, use the newest rack config. blacklist = isolationPolicy.getBlacklistedBookiesWithIsolationGroups(2, groups); assertFalse(blacklist.isEmpty()); From 57debf9afb527c9ecb66afa26fc443b9906bc679 Mon Sep 17 00:00:00 2001 From: horizonzy Date: Mon, 4 Sep 2023 18:49:09 +0800 Subject: [PATCH 10/12] Address the comments. --- ...IsolatedBookieEnsemblePlacementPolicy.java | 39 +++++++---------- ...atedBookieEnsemblePlacementPolicyTest.java | 43 +++++++++++++++++-- 2 files changed, 55 insertions(+), 27 deletions(-) diff --git a/pulsar-broker-common/src/main/java/org/apache/pulsar/bookie/rackawareness/IsolatedBookieEnsemblePlacementPolicy.java b/pulsar-broker-common/src/main/java/org/apache/pulsar/bookie/rackawareness/IsolatedBookieEnsemblePlacementPolicy.java index ffe66d5d5651a..0b62b6a55a460 100644 --- a/pulsar-broker-common/src/main/java/org/apache/pulsar/bookie/rackawareness/IsolatedBookieEnsemblePlacementPolicy.java +++ b/pulsar-broker-common/src/main/java/org/apache/pulsar/bookie/rackawareness/IsolatedBookieEnsemblePlacementPolicy.java @@ -28,7 +28,6 @@ import java.util.Map; import java.util.Optional; import java.util.Set; -import java.util.concurrent.CompletableFuture; import lombok.extern.slf4j.Slf4j; import org.apache.bookkeeper.client.BKException.BKNotEnoughBookiesException; import org.apache.bookkeeper.client.RackawareEnsemblePlacementPolicy; @@ -62,7 +61,7 @@ public class IsolatedBookieEnsemblePlacementPolicy extends RackawareEnsemblePlac private MetadataCache bookieMappingCache; - private BookiesRackConfiguration cachedRackConfiguration = null; + private volatile BookiesRackConfiguration cachedRackConfiguration = null; public IsolatedBookieEnsemblePlacementPolicy() { super(); @@ -90,9 +89,12 @@ public RackawareEnsemblePlacementPolicyImpl initialize(ClientConfiguration conf, // Only add the bookieMappingCache if we have defined an isolation group bookieMappingCache = store.getMetadataCache(BookiesRackConfiguration.class); - Optional optional = - bookieMappingCache.get(BookieRackAffinityMapping.BOOKIE_INFO_ROOT_PATH).join(); - optional.ifPresent(bookiesRackConfiguration -> cachedRackConfiguration = bookiesRackConfiguration); + bookieMappingCache.get(BookieRackAffinityMapping.BOOKIE_INFO_ROOT_PATH).thenAccept(opt -> opt.ifPresent( + bookiesRackConfiguration -> cachedRackConfiguration = bookiesRackConfiguration)) + .exceptionally(e -> { + log.warn("Failed to load bookies rack configuration while initialize the PlacementPolicy."); + return null; + }); } } if (conf.getProperty(SECONDARY_ISOLATION_BOOKIE_GROUPS) != null) { @@ -192,26 +194,17 @@ Set getBlacklistedBookiesWithIsolationGroups(int ensembleSize, Set blacklistedBookies = new HashSet<>(); try { if (bookieMappingCache != null) { - CompletableFuture> future = - bookieMappingCache.get(BookieRackAffinityMapping.BOOKIE_INFO_ROOT_PATH); + bookieMappingCache.get(BookieRackAffinityMapping.BOOKIE_INFO_ROOT_PATH) + .thenAccept(opt -> cachedRackConfiguration = opt.orElse(null)).exceptionally(e -> { + log.warn("Failed to update the newest bookies rack config."); + return null; + }); - BookiesRackConfiguration allGroupsBookieMapping; - Optional optRes = (future.isDone() && !future.isCompletedExceptionally()) - ? future.join() : Optional.empty(); - - if (!optRes.isPresent()) { - if (cachedRackConfiguration != null) { - log.debug("The newest rack config is not available now, use the cached rack config : {}", - cachedRackConfiguration); - allGroupsBookieMapping = cachedRackConfiguration; - } else { - throw new KeeperException.NoNodeException(ZkBookieRackAffinityMapping.BOOKIE_INFO_ROOT_PATH); - } - } else { - cachedRackConfiguration = optRes.get(); - allGroupsBookieMapping = optRes.get(); + BookiesRackConfiguration allGroupsBookieMapping = cachedRackConfiguration; + if (allGroupsBookieMapping == null) { + log.debug("The bookies rack config is not available at now."); + return blacklistedBookies; } - Set allBookies = allGroupsBookieMapping.keySet(); int totalAvailableBookiesInPrimaryGroup = 0; Set primaryIsolationGroup = Collections.emptySet(); diff --git a/pulsar-broker-common/src/test/java/org/apache/pulsar/bookie/rackawareness/IsolatedBookieEnsemblePlacementPolicyTest.java b/pulsar-broker-common/src/test/java/org/apache/pulsar/bookie/rackawareness/IsolatedBookieEnsemblePlacementPolicyTest.java index eda28878c167f..e8a1568354d26 100644 --- a/pulsar-broker-common/src/test/java/org/apache/pulsar/bookie/rackawareness/IsolatedBookieEnsemblePlacementPolicyTest.java +++ b/pulsar-broker-common/src/test/java/org/apache/pulsar/bookie/rackawareness/IsolatedBookieEnsemblePlacementPolicyTest.java @@ -128,8 +128,23 @@ public void testMetadataStoreCases() throws Exception { waitingCompleteFuture.complete(Optional.of(rackConfiguration2)); }).start(); - when(cache.get(BookieRackAffinityMapping.BOOKIE_INFO_ROOT_PATH)).thenReturn(initialFuture) - .thenReturn(waitingCompleteFuture); + long longWaitTime = 4000; + CompletableFuture> emptyFuture = new CompletableFuture<>(); + new Thread(() -> { + try { + Thread.sleep(longWaitTime); + } catch (InterruptedException e) { + throw new RuntimeException(e); + } + //The emptyFuture means that the zk node /bookies already be removed. + emptyFuture.complete(Optional.empty()); + }).start(); + + //Return different future means that cache expire. + when(cache.get(BookieRackAffinityMapping.BOOKIE_INFO_ROOT_PATH)) + .thenReturn(initialFuture).thenReturn(initialFuture) + .thenReturn(waitingCompleteFuture).thenReturn(waitingCompleteFuture) + .thenReturn(emptyFuture).thenReturn(emptyFuture); IsolatedBookieEnsemblePlacementPolicy isolationPolicy = new IsolatedBookieEnsemblePlacementPolicy(); ClientConfiguration bkClientConf = new ClientConfiguration(); @@ -142,20 +157,40 @@ public void testMetadataStoreCases() throws Exception { groups.setLeft(Sets.newHashSet("group1")); groups.setRight(new HashSet<>()); - //The future is waiting done, so use the cached rack config. + //initialFuture, the future is waiting done. Set blacklist = isolationPolicy.getBlacklistedBookiesWithIsolationGroups(2, groups); assertTrue(blacklist.isEmpty()); + //waitingCompleteFuture, the future is waiting done. + blacklist = + isolationPolicy.getBlacklistedBookiesWithIsolationGroups(2, groups); + assertTrue(blacklist.isEmpty()); + Thread.sleep(waitTime); - //The future is already done, use the newest rack config. + //waitingCompleteFuture, the future is already done. blacklist = isolationPolicy.getBlacklistedBookiesWithIsolationGroups(2, groups); assertFalse(blacklist.isEmpty()); assertEquals(blacklist.size(), 1); BookieId excludeBookie = blacklist.iterator().next(); assertEquals(excludeBookie.toString(), BOOKIE3); + + //emptyFuture, the future is waiting done. + blacklist = + isolationPolicy.getBlacklistedBookiesWithIsolationGroups(2, groups); + assertFalse(blacklist.isEmpty()); + assertEquals(blacklist.size(), 1); + excludeBookie = blacklist.iterator().next(); + assertEquals(excludeBookie.toString(), BOOKIE3); + + Thread.sleep(longWaitTime - waitTime); + + //emptyFuture, the future is already done. + blacklist = + isolationPolicy.getBlacklistedBookiesWithIsolationGroups(2, groups); + assertTrue(blacklist.isEmpty()); } @Test From d91cf3eb7cb3d074b5a2f1d34cd64887b2823e67 Mon Sep 17 00:00:00 2001 From: horizonzy Date: Mon, 4 Sep 2023 19:13:33 +0800 Subject: [PATCH 11/12] Fix checkstyle. --- .../rackawareness/IsolatedBookieEnsemblePlacementPolicy.java | 4 +--- 1 file changed, 1 insertion(+), 3 deletions(-) diff --git a/pulsar-broker-common/src/main/java/org/apache/pulsar/bookie/rackawareness/IsolatedBookieEnsemblePlacementPolicy.java b/pulsar-broker-common/src/main/java/org/apache/pulsar/bookie/rackawareness/IsolatedBookieEnsemblePlacementPolicy.java index 0b62b6a55a460..158e219eca38c 100644 --- a/pulsar-broker-common/src/main/java/org/apache/pulsar/bookie/rackawareness/IsolatedBookieEnsemblePlacementPolicy.java +++ b/pulsar-broker-common/src/main/java/org/apache/pulsar/bookie/rackawareness/IsolatedBookieEnsemblePlacementPolicy.java @@ -47,8 +47,6 @@ import org.apache.pulsar.common.policies.data.EnsemblePlacementPolicyConfig; import org.apache.pulsar.metadata.api.MetadataCache; import org.apache.pulsar.metadata.api.MetadataStore; -import org.apache.pulsar.zookeeper.ZkBookieRackAffinityMapping; -import org.apache.zookeeper.KeeperException; @Slf4j public class IsolatedBookieEnsemblePlacementPolicy extends RackawareEnsemblePlacementPolicy { @@ -76,7 +74,7 @@ public RackawareEnsemblePlacementPolicyImpl initialize(ClientConfiguration conf, store = BookieRackAffinityMapping.createMetadataStore(conf); } catch (MetadataException e) { throw new RuntimeException(METADATA_STORE_INSTANCE + " failed initialized"); - } + Set primaryIsolationGroups = new HashSet<>(); Set secondaryIsolationGroups = new HashSet<>(); if (conf.getProperty(ISOLATION_BOOKIE_GROUPS) != null) { From cc8f08af88940b95736aff2cc6f60acba8ea76c5 Mon Sep 17 00:00:00 2001 From: horizonzy Date: Mon, 4 Sep 2023 22:49:22 +0800 Subject: [PATCH 12/12] fix compile problem. --- .../rackawareness/IsolatedBookieEnsemblePlacementPolicy.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pulsar-broker-common/src/main/java/org/apache/pulsar/bookie/rackawareness/IsolatedBookieEnsemblePlacementPolicy.java b/pulsar-broker-common/src/main/java/org/apache/pulsar/bookie/rackawareness/IsolatedBookieEnsemblePlacementPolicy.java index 158e219eca38c..64bc59057d0f8 100644 --- a/pulsar-broker-common/src/main/java/org/apache/pulsar/bookie/rackawareness/IsolatedBookieEnsemblePlacementPolicy.java +++ b/pulsar-broker-common/src/main/java/org/apache/pulsar/bookie/rackawareness/IsolatedBookieEnsemblePlacementPolicy.java @@ -74,7 +74,7 @@ public RackawareEnsemblePlacementPolicyImpl initialize(ClientConfiguration conf, store = BookieRackAffinityMapping.createMetadataStore(conf); } catch (MetadataException e) { throw new RuntimeException(METADATA_STORE_INSTANCE + " failed initialized"); - + } Set primaryIsolationGroups = new HashSet<>(); Set secondaryIsolationGroups = new HashSet<>(); if (conf.getProperty(ISOLATION_BOOKIE_GROUPS) != null) {