diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/loadbalance/BundleSplitStrategy.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/loadbalance/BundleSplitStrategy.java index d3d76b8e92db1..484eee77f898f 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/loadbalance/BundleSplitStrategy.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/loadbalance/BundleSplitStrategy.java @@ -18,7 +18,7 @@ */ package org.apache.pulsar.broker.loadbalance; -import java.util.Set; +import java.util.Map; import org.apache.pulsar.broker.PulsarService; /** @@ -33,7 +33,7 @@ public interface BundleSplitStrategy { * leader broker). * @param pulsar * Service to use. - * @return A set of the bundles that should be split. + * @return A map of the bundles that should be split and the brokers on which they reside. */ - Set findBundlesToSplit(LoadData loadData, PulsarService pulsar); + Map findBundlesToSplit(LoadData loadData, PulsarService pulsar); } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/loadbalance/impl/BundleSplitterTask.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/loadbalance/impl/BundleSplitterTask.java index 310f80d6f74fc..7f3e43d352346 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/loadbalance/impl/BundleSplitterTask.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/loadbalance/impl/BundleSplitterTask.java @@ -19,9 +19,7 @@ package org.apache.pulsar.broker.loadbalance.impl; import java.util.HashMap; -import java.util.HashSet; import java.util.Map; -import java.util.Set; import org.apache.pulsar.broker.PulsarService; import org.apache.pulsar.broker.ServiceConfiguration; import org.apache.pulsar.broker.loadbalance.BundleSplitStrategy; @@ -38,7 +36,7 @@ */ public class BundleSplitterTask implements BundleSplitStrategy { private static final Logger log = LoggerFactory.getLogger(BundleSplitStrategy.class); - private final Set bundleCache; + private final Map bundleCache; private final Map namespaceBundleCount; @@ -48,7 +46,7 @@ public class BundleSplitterTask implements BundleSplitStrategy { * */ public BundleSplitterTask() { - bundleCache = new HashSet<>(); + bundleCache = new HashMap<>(); namespaceBundleCount = new HashMap<>(); } @@ -61,10 +59,10 @@ public BundleSplitterTask() { * @param pulsar * Service to use. * @return All bundles who have exceeded configured thresholds in number of topics, number of sessions, total - * message rates, or total throughput. + * message rates, or total throughput and the brokers on which they reside. */ @Override - public Set findBundlesToSplit(final LoadData loadData, final PulsarService pulsar) { + public Map findBundlesToSplit(final LoadData loadData, final PulsarService pulsar) { bundleCache.clear(); namespaceBundleCount.clear(); final ServiceConfiguration conf = pulsar.getConfiguration(); @@ -108,7 +106,7 @@ public Set findBundlesToSplit(final LoadData loadData, final PulsarServi maxBundleSessions, totalMessageRate, maxBundleMsgRate, totalMessageThroughput / LoadManagerShared.MIBI, maxBundleBandwidth / LoadManagerShared.MIBI); - bundleCache.add(bundle); + bundleCache.put(bundle, broker); int bundleNum = namespaceBundleCount.getOrDefault(namespace, 0); namespaceBundleCount.put(namespace, bundleNum + 1); } else { diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/loadbalance/impl/ModularLoadManagerImpl.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/loadbalance/impl/ModularLoadManagerImpl.java index 6e63643a859d0..e25ec981e7705 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/loadbalance/impl/ModularLoadManagerImpl.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/loadbalance/impl/ModularLoadManagerImpl.java @@ -748,10 +748,10 @@ public void checkNamespaceBundleSplit() { } final boolean unloadSplitBundles = pulsar.getConfiguration().isLoadBalancerAutoUnloadSplitBundlesEnabled(); synchronized (bundleSplitStrategy) { - final Set bundlesToBeSplit = bundleSplitStrategy.findBundlesToSplit(loadData, pulsar); + final Map bundlesToBeSplit = bundleSplitStrategy.findBundlesToSplit(loadData, pulsar); NamespaceBundleFactory namespaceBundleFactory = pulsar.getNamespaceService().getNamespaceBundleFactory(); int splitCount = 0; - for (String bundleName : bundlesToBeSplit) { + for (String bundleName : bundlesToBeSplit.keySet()) { try { final String namespaceName = LoadManagerShared.getNamespaceNameFromBundleName(bundleName); final String bundleRange = LoadManagerShared.getBundleRangeFromBundleName(bundleName); @@ -768,9 +768,17 @@ public void checkNamespaceBundleSplit() { .invalidateBundleCache(NamespaceName.get(namespaceName)); deleteBundleDataFromMetadataStore(bundleName); - log.info("Load-manager splitting bundle {} and unloading {}", bundleName, unloadSplitBundles); + // Check NamespacePolicies and AntiAffinityNamespace support unload bundle. + boolean isUnload = false; + String broker = bundlesToBeSplit.get(bundleName); + if (unloadSplitBundles + && shouldNamespacePoliciesUnload(namespaceName, bundleRange, broker) + && shouldAntiAffinityNamespaceUnload(namespaceName, bundleRange, broker)) { + isUnload = true; + } + log.info("Load-manager splitting bundle {} and unloading {}", bundleName, isUnload); pulsar.getAdminClient().namespaces().splitNamespaceBundle(namespaceName, bundleRange, - unloadSplitBundles, null); + isUnload, null); splitCount++; log.info("Successfully split namespace bundle {}", bundleName); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/loadbalance/impl/BundleSplitterTaskTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/loadbalance/impl/BundleSplitterTaskTest.java index fbefdb74d41fc..3173987a3c8a8 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/loadbalance/impl/BundleSplitterTaskTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/loadbalance/impl/BundleSplitterTaskTest.java @@ -37,7 +37,6 @@ import java.util.HashMap; import java.util.Map; import java.util.Optional; -import java.util.Set; /** * @author hezhangjian @@ -92,7 +91,7 @@ public void testSplitTaskWhenTopicJustOne() { bundleData.setLongTermData(averageMessageData); loadData.getBundleData().put("ten/ns/0x00000000_0x80000000", bundleData); - final Set bundlesToSplit = bundleSplitterTask.findBundlesToSplit(loadData, pulsar); + final Map bundlesToSplit = bundleSplitterTask.findBundlesToSplit(loadData, pulsar); Assert.assertEquals(bundlesToSplit.size(), 0); } @@ -142,7 +141,7 @@ public void testLoadBalancerNamespaceMaximumBundles() throws Exception { loadData.getBundleData().put("ten/ns/0x40000000_0x60000000", bundleData3); int currentBundleCount = pulsar.getNamespaceService().getBundleCount(NamespaceName.get("ten/ns")); - final Set bundlesToSplit = bundleSplitterTask.findBundlesToSplit(loadData, pulsar); + final Map bundlesToSplit = bundleSplitterTask.findBundlesToSplit(loadData, pulsar); Assert.assertEquals(bundlesToSplit.size() + currentBundleCount, pulsar.getConfiguration().getLoadBalancerNamespaceMaximumBundles()); } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/BrokerServiceLookupTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/BrokerServiceLookupTest.java index 53da440736208..33f0d6ee05a41 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/BrokerServiceLookupTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/BrokerServiceLookupTest.java @@ -752,6 +752,11 @@ public void testModularLoadManagerSplitBundle() throws Exception { assertNotEquals(pulsar2.getNamespaceService().getBundle(topicName), bundleInBroker2); }); + // Unload the NamespacePolicies and AntiAffinity check. + String currentBroker = String.format("%s:%d", "localhost", pulsar.getListenPortHTTP().get()); + assertTrue(loadManager.shouldNamespacePoliciesUnload(namespace,"0x00000000_0xffffffff", currentBroker)); + assertTrue(loadManager.shouldAntiAffinityNamespaceUnload(namespace,"0x00000000_0xffffffff", currentBroker)); + // (7) Make lookup request again to Broker-2 which should succeed. final String topic3 = "persistent://" + namespace + "/topic3"; @Cleanup