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 6e49804905893..9a2b70d6996af 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 @@ -38,7 +38,6 @@ import java.util.concurrent.Executors; import java.util.concurrent.RejectedExecutionException; import java.util.concurrent.TimeUnit; -import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicReference; import java.util.concurrent.locks.Lock; import java.util.concurrent.locks.ReentrantLock; @@ -301,11 +300,26 @@ private void handleMetadataSessionEvent(SessionEvent e) { lastMetadataSessionEvent = e; } + @VisibleForTesting + boolean isLeader() { + return pulsar.getLeaderElectionService() != null && pulsar.getLeaderElectionService().isLeader(); + } + private LoadSheddingStrategy createLoadSheddingStrategy() { return Reflections.createInstance(conf.getLoadBalancerLoadSheddingStrategy(), LoadSheddingStrategy.class, Thread.currentThread().getContextClassLoader()); } + @VisibleForTesting + void setLoadSheddingStrategy(LoadSheddingStrategy loadSheddingStrategy) { + this.loadSheddingStrategy = loadSheddingStrategy; + } + + @VisibleForTesting + LoadData getLoadData() { + return loadData; + } + /** * Initialize this load manager. * @@ -635,6 +649,10 @@ public void disableBroker() throws PulsarServerException { */ @Override public synchronized void doLoadShedding() { + if (!isLeader()) { + log.debug().log("Skipping load shedding because this broker is not the leader"); + return; + } if (!LoadManagerShared.isLoadSheddingEnabled(pulsar)) { return; } @@ -651,23 +669,27 @@ public synchronized void doLoadShedding() { Set sheddingExcludedNamespaces = conf.getLoadBalancerSheddingExcludedNamespaces(); final Multimap bundlesToUnload = loadSheddingStrategy.findBundlesForUnloading(loadData, conf); - bundlesToUnload.asMap().forEach((broker, bundles) -> { - AtomicBoolean unloadBundleForBroker = new AtomicBoolean(false); - bundles.forEach(bundle -> { + for (Map.Entry> entry : bundlesToUnload.asMap().entrySet()) { + String broker = entry.getKey(); + boolean unloadBundleForBroker = false; + for (String bundle : entry.getValue()) { + if (!isLeader()) { + return; + } final String namespaceName = LoadManagerShared.getNamespaceNameFromBundleName(bundle); final String bundleRange = LoadManagerShared.getBundleRangeFromBundleName(bundle); if (sheddingExcludedNamespaces.contains(namespaceName)) { log.debug().attr("class", loadSheddingStrategy.getClass().getSimpleName()) .attr("namespace", namespaceName) .log("Skipping load shedding for namespace"); - return; + continue; } if (!shouldNamespacePoliciesUnload(namespaceName, bundleRange, broker)) { - return; + continue; } if (!shouldAntiAffinityNamespaceUnload(namespaceName, bundleRange, broker)) { - return; + continue; } NamespaceBundle bundleToUnload = LoadManagerShared.getNamespaceBundle(pulsar, bundle); Optional destBroker = this.selectBroker(bundleToUnload); @@ -675,37 +697,45 @@ public synchronized void doLoadShedding() { log.info().attr("class", loadSheddingStrategy.getClass().getSimpleName()) .attr("bundle", bundle).attr("broker", broker) .log("No broker available to unload bundle from broker"); - return; + continue; } if (destBroker.get().equals(broker)) { log.warn().attr("class", loadSheddingStrategy.getClass().getSimpleName()) .attr("broker", destBroker.get()).attr("bundle", bundle) .log("The destination broker is the same as the current owner broker for bundle"); - return; + continue; } - log.info().attr("class", loadSheddingStrategy.getClass().getSimpleName()) - .attr("bundle", bundle).attr("sourceBroker", broker).attr("destBroker", destBroker.get()) - .log("Unloading bundle from source broker to dest broker"); try { - pulsar.getAdminClient().namespaces() - .unloadNamespaceBundle(namespaceName, bundleRange, destBroker.get()); + if (!isLeader()) { + return; + } + log.info().attr("class", loadSheddingStrategy.getClass().getSimpleName()) + .attr("bundle", bundle).attr("sourceBroker", broker).attr("destBroker", destBroker.get()) + .log("Unloading bundle from source broker to dest broker"); + unloadNamespaceBundle(namespaceName, bundleRange, destBroker.get()); loadData.getRecentlyUnloadedBundles().put(bundle, System.currentTimeMillis()); unloadBundleCount++; - unloadBundleForBroker.set(true); + unloadBundleForBroker = true; } catch (PulsarServerException | PulsarAdminException e) { log.warn().attr("bundle", bundle).attr("broker", broker).exception(e) .log("Error when trying to perform load shedding on for broker"); } - }); - if (unloadBundleForBroker.get()) { + } + if (unloadBundleForBroker) { unloadBrokerCount++; } - }); + } updateBundleUnloadingMetrics(); } + @VisibleForTesting + void unloadNamespaceBundle(String namespaceName, String bundleRange, String destinationBroker) + throws PulsarServerException, PulsarAdminException { + pulsar.getAdminClient().namespaces().unloadNamespaceBundle(namespaceName, bundleRange, destinationBroker); + } + /** * As leader broker, update bundle unloading metrics. */ @@ -1204,20 +1234,34 @@ private int selectTopKBundle() { */ @Override public void writeBundleDataOnZooKeeper() { + if (!isLeader()) { + log.debug().log("Skipping bundle data write because this broker is not the leader"); + return; + } updateBundleData(); + if (!isLeader()) { + return; + } // Write the bundle data to metadata store. List> futures = new ArrayList<>(); // use synchronized to protect bundleArr. synchronized (bundleArr) { int updateBundleCount = selectTopKBundle(); - bundleArr.stream().limit(updateBundleCount).forEach(entry -> futures.add( - pulsarResources.getLoadBalanceResources().getBundleDataResources().updateBundleData( - entry.getKey(), (BundleData) entry.getValue()))); + for (Map.Entry entry : bundleArr.subList(0, updateBundleCount)) { + if (!isLeader()) { + break; + } + futures.add(pulsarResources.getLoadBalanceResources().getBundleDataResources().updateBundleData( + entry.getKey(), (BundleData) entry.getValue())); + } } // Write the time average broker data to metadata store. for (Map.Entry entry : loadData.getBrokerData().entrySet()) { + if (!isLeader()) { + break; + } final String broker = entry.getKey(); final TimeAverageBrokerData data = entry.getValue().getTimeAverageData(); futures.add(pulsarResources.getLoadBalanceResources() diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/loadbalance/impl/ModularLoadManagerImplTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/loadbalance/impl/ModularLoadManagerImplTest.java index c42ca4fefbb29..55a4b121e711d 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/loadbalance/impl/ModularLoadManagerImplTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/loadbalance/impl/ModularLoadManagerImplTest.java @@ -26,6 +26,7 @@ import static org.mockito.Mockito.doNothing; import static org.mockito.Mockito.spy; import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.verifyNoInteractions; import static org.mockito.Mockito.when; import static org.testng.Assert.assertEquals; import static org.testng.Assert.assertFalse; @@ -34,6 +35,7 @@ import static org.testng.Assert.fail; import com.fasterxml.jackson.databind.ObjectMapper; import com.google.common.collect.BoundType; +import com.google.common.collect.ImmutableMultimap; import com.google.common.collect.Range; import com.google.common.collect.Sets; import com.google.common.hash.Hashing; @@ -56,6 +58,8 @@ import java.util.concurrent.LinkedBlockingQueue; import java.util.concurrent.ThreadPoolExecutor; import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicInteger; import java.util.concurrent.atomic.AtomicReference; import java.util.function.Supplier; import lombok.Cleanup; @@ -66,6 +70,7 @@ import org.apache.pulsar.broker.loadbalance.LoadBalancerTestingUtils; import org.apache.pulsar.broker.loadbalance.LoadData; import org.apache.pulsar.broker.loadbalance.LoadManager; +import org.apache.pulsar.broker.loadbalance.LoadSheddingStrategy; import org.apache.pulsar.broker.loadbalance.ResourceUnit; import org.apache.pulsar.broker.loadbalance.impl.LoadManagerShared.BrokerTopicLoadingPredicate; import org.apache.pulsar.client.admin.Namespaces; @@ -295,6 +300,76 @@ private String mockBundleName(final int i) { return String.format("%d/%d/0x00000000_0xffffffff", i, i); } + @Test + public void testFollowerSkipsLoadShedding() { + Awaitility.await().until(() -> pulsar1.getLeaderElectionService().isLeader() + || pulsar2.getLeaderElectionService().isLeader()); + + ModularLoadManagerImpl followerLoadManager = pulsar1.getLeaderElectionService().isLeader() + ? secondaryLoadManager : primaryLoadManager; + + LoadSheddingStrategy loadSheddingStrategy = Mockito.mock(LoadSheddingStrategy.class); + followerLoadManager.setLoadSheddingStrategy(loadSheddingStrategy); + + followerLoadManager.doLoadShedding(); + + verifyNoInteractions(loadSheddingStrategy); + } + + @Test + public void testLoadSheddingStopsWhenLeadershipChangesBeforeUnload() throws Exception { + Awaitility.await().until(() -> primaryLoadManager.getAvailableBrokers().size() > 1); + + AtomicBoolean leader = new AtomicBoolean(true); + ModularLoadManagerImpl loadManagerSpy = spy(primaryLoadManager); + doAnswer(invocation -> leader.get()).when(loadManagerSpy).isLeader(); + + LoadSheddingStrategy loadSheddingStrategy = Mockito.mock(LoadSheddingStrategy.class); + loadManagerSpy.setLoadSheddingStrategy(loadSheddingStrategy); + when(loadSheddingStrategy.findBundlesForUnloading(any(), any())) + .thenReturn(ImmutableMultimap.of(primaryBrokerId, mockBundleName(1), + primaryBrokerId, mockBundleName(2))); + doAnswer(invocation -> true).when(loadManagerSpy).shouldNamespacePoliciesUnload( + Mockito.anyString(), Mockito.anyString(), Mockito.anyString()); + doAnswer(invocation -> true).when(loadManagerSpy).shouldAntiAffinityNamespaceUnload( + Mockito.anyString(), Mockito.anyString(), Mockito.anyString()); + doAnswer(invocation -> { + leader.set(false); + return Optional.of(secondaryBrokerId); + }).when(loadManagerSpy).selectBroker(any()); + doNothing().when(loadManagerSpy).unloadNamespaceBundle( + Mockito.anyString(), Mockito.anyString(), Mockito.anyString()); + + loadManagerSpy.doLoadShedding(); + + verify(loadManagerSpy, Mockito.times(1)).selectBroker(any()); + verify(loadManagerSpy, Mockito.never()).unloadNamespaceBundle( + Mockito.anyString(), Mockito.anyString(), Mockito.anyString()); + } + + @Test + public void testBundleDataWriteStopsWhenLeadershipChangesBeforeMetadataWrite() throws Exception { + String bundle = mockBundleName(99); + BundleData bundleData = new BundleData(10, 1000); + String bundleDataPath = String.format("%s/%s", BUNDLE_DATA_BASE_PATH, bundle); + MetadataCache metadataCache = pulsar1.getLocalMetadataStore().getMetadataCache(BundleData.class); + metadataCache.create(bundleDataPath, bundleData).join(); + + Awaitility.await().until(() -> primaryLoadManager.getLoadData().getBrokerData().containsKey(primaryBrokerId)); + AtomicInteger leaderChecks = new AtomicInteger(); + ModularLoadManagerImpl loadManagerSpy = spy(primaryLoadManager); + LoadData loadData = loadManagerSpy.getLoadData(); + loadData.getBundleData().clear(); + loadData.getBundleData().put(bundle, bundleData); + loadData.getBrokerData().get(primaryBrokerId).getLocalData().getLastStats() + .put(bundle, new NamespaceBundleStats()); + doAnswer(invocation -> leaderChecks.getAndIncrement() < 2).when(loadManagerSpy).isLeader(); + + loadManagerSpy.writeBundleDataOnZooKeeper(); + + assertEquals(metadataCache.getWithStats(bundleDataPath).get().get().getStat().getVersion(), 0); + } + // Test disabled since it's depending on CPU usage in the machine @Test(enabled = false) public void testCandidateConsistency() throws Exception {