From 31cb0051f4f9867c37d492ad6c5e6c9c8795e6e2 Mon Sep 17 00:00:00 2001 From: fengyubiao Date: Tue, 6 Sep 2022 17:13:07 +0800 Subject: [PATCH 1/4] [fix][flaky-test]NamespaceServiceTest.testModularLoadManagerRemoveBundleAndLoad --- .../namespace/NamespaceServiceTest.java | 56 +++++++++++++++---- 1 file changed, 46 insertions(+), 10 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/namespace/NamespaceServiceTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/namespace/NamespaceServiceTest.java index cd615f6a0ecf4..e4dfc55e2e3b6 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/namespace/NamespaceServiceTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/namespace/NamespaceServiceTest.java @@ -42,9 +42,11 @@ import java.util.Optional; import java.util.Set; import java.util.concurrent.CompletableFuture; +import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicReference; +import java.util.function.Predicate; import lombok.Cleanup; import org.apache.bookkeeper.mledger.ManagedLedger; import org.apache.commons.collections4.CollectionUtils; @@ -85,6 +87,7 @@ import org.apache.pulsar.policies.data.loadbalancer.BundleData; import org.awaitility.Awaitility; import org.mockito.stubbing.Answer; +import org.powermock.reflect.Whitebox; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.testng.Assert; @@ -746,19 +749,18 @@ public void testModularLoadManagerRemoveInactiveBundleFromLoadData() throws Exce public void testModularLoadManagerRemoveBundleAndLoad() throws Exception { final String BUNDLE_DATA_PATH = "/loadbalance/bundle-data"; final String namespace = "prop/ns-abc"; + final String bundleName = namespace + "/0x00000000_0xffffffff"; final String topic1 = "persistent://" + namespace + "/topic1"; final String topic2 = "persistent://" + namespace + "/topic2"; // configure broker with ModularLoadManager conf.setLoadManagerClassName(ModularLoadManagerImpl.class.getName()); conf.setForceDeleteNamespaceAllowed(true); + // Make sure LoadReportUpdaterTask has a 100% chance to write ZK. + conf.setLoadBalancerReportUpdateMaxIntervalMinutes(-1); restartBroker(); - LoadManager loadManager = spy(pulsar.getLoadManager().get()); - Field loadManagerField = NamespaceService.class.getDeclaredField("loadManager"); - loadManagerField.setAccessible(true); - doReturn(true).when(loadManager).isCentralized(); - loadManagerField.set(pulsar.getNamespaceService(), new AtomicReference<>(loadManager)); + LoadManager loadManager = pulsar.getLoadManager().get(); NamespaceName nsname = NamespaceName.get(namespace); @Cleanup @@ -778,10 +780,10 @@ public void testModularLoadManagerRemoveBundleAndLoad() throws Exception { //create znode for bundle-data pulsar.getBrokerService().updateRates(); - loadManager.writeLoadReportOnZookeeper(); - loadManager.writeResourceQuotasToZooKeeper(); - String path = BUNDLE_DATA_PATH + "/" + nsname.toString() + "/0x00000000_0xffffffff"; + waitResourceDataUpdateToZK(loadManager, + loadData -> loadData.getBundleData().containsKey(bundleName)); + String path = BUNDLE_DATA_PATH + "/" + bundleName; Optional getResult = pulsar.getLocalMetadataStore().get(path).get(); assertTrue(getResult.isPresent()); @@ -792,12 +794,46 @@ public void testModularLoadManagerRemoveBundleAndLoad() throws Exception { TimeUnit.SECONDS.sleep(5); // update broker bundle report to zk - loadManager.writeLoadReportOnZookeeper(); - loadManager.writeResourceQuotasToZooKeeper(); + waitResourceDataUpdateToZK(loadManager, + loadData -> !loadData.getBundleData().containsKey(bundleName)); getResult = pulsar.getLocalMetadataStore().get(path).get(); assertFalse(getResult.isPresent()); + } + private void waitResourceDataUpdateToZK( + LoadManager loadManager, Predicate utilChecker) throws Exception { + CompletableFuture waitForBrokerChangeNotice = registryBrokerDataChangeNotice(); + // Manually trigger "LoadReportUpdaterTask" + loadManager.writeLoadReportOnZookeeper(); + waitForBrokerChangeNotice.join(); + // Wait until "ModularLoadManager" completes processing the ZK notification. + ModularLoadManagerWrapper modularLoadManagerWrapper = (ModularLoadManagerWrapper) loadManager; + ModularLoadManagerImpl modularLoadManager = (ModularLoadManagerImpl) modularLoadManagerWrapper.getLoadManager(); + ScheduledExecutorService scheduler = Whitebox.getInternalState(modularLoadManager, "scheduler"); + CompletableFuture waitForNoticeHandleFinishByLoadManager = new CompletableFuture<>(); + scheduler.execute(() -> { + waitForNoticeHandleFinishByLoadManager.complete(null); + }); + waitForNoticeHandleFinishByLoadManager.join(); + // Manually trigger "LoadResourceQuotaUpdaterTask" + loadManager.writeResourceQuotasToZooKeeper(); + } + + public CompletableFuture registryBrokerDataChangeNotice() { + CompletableFuture completableFuture = new CompletableFuture<>(); + String lookupServiceAddress = pulsar.getAdvertisedAddress() + ":" + + (conf.getWebServicePort().isPresent() ? conf.getWebServicePort().get() + : conf.getWebServicePortTls().get()); + String brokerDataPath = LoadManager.LOADBALANCE_BROKERS_ROOT + "/" + lookupServiceAddress; + pulsar.getLocalMetadataStore().registerListener(notice -> { + if (brokerDataPath.equals(notice.getPath())){ + if (!completableFuture.isDone()) { + completableFuture.complete(null); + } + } + }); + return completableFuture; } @SuppressWarnings("unchecked") From e3f6862ee6417f64fd59680de5bff6b0eb476927 Mon Sep 17 00:00:00 2001 From: fengyubiao Date: Tue, 6 Sep 2022 20:16:22 +0800 Subject: [PATCH 2/4] remove unnecessary code --- .../pulsar/broker/namespace/NamespaceServiceTest.java | 9 +++------ 1 file changed, 3 insertions(+), 6 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/namespace/NamespaceServiceTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/namespace/NamespaceServiceTest.java index e4dfc55e2e3b6..bc0bf55b2a005 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/namespace/NamespaceServiceTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/namespace/NamespaceServiceTest.java @@ -781,8 +781,7 @@ public void testModularLoadManagerRemoveBundleAndLoad() throws Exception { //create znode for bundle-data pulsar.getBrokerService().updateRates(); - waitResourceDataUpdateToZK(loadManager, - loadData -> loadData.getBundleData().containsKey(bundleName)); + waitResourceDataUpdateToZK(loadManager); String path = BUNDLE_DATA_PATH + "/" + bundleName; Optional getResult = pulsar.getLocalMetadataStore().get(path).get(); @@ -794,15 +793,13 @@ public void testModularLoadManagerRemoveBundleAndLoad() throws Exception { TimeUnit.SECONDS.sleep(5); // update broker bundle report to zk - waitResourceDataUpdateToZK(loadManager, - loadData -> !loadData.getBundleData().containsKey(bundleName)); + waitResourceDataUpdateToZK(loadManager); getResult = pulsar.getLocalMetadataStore().get(path).get(); assertFalse(getResult.isPresent()); } - private void waitResourceDataUpdateToZK( - LoadManager loadManager, Predicate utilChecker) throws Exception { + private void waitResourceDataUpdateToZK(LoadManager loadManager) throws Exception { CompletableFuture waitForBrokerChangeNotice = registryBrokerDataChangeNotice(); // Manually trigger "LoadReportUpdaterTask" loadManager.writeLoadReportOnZookeeper(); From f89088cbfe60ac94f39e041aa10e1b185969f0f7 Mon Sep 17 00:00:00 2001 From: fengyubiao Date: Wed, 7 Sep 2022 09:34:10 +0800 Subject: [PATCH 3/4] remove unnecessary import --- .../org/apache/pulsar/broker/namespace/NamespaceServiceTest.java | 1 - 1 file changed, 1 deletion(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/namespace/NamespaceServiceTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/namespace/NamespaceServiceTest.java index bc0bf55b2a005..f3cfd33519c9b 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/namespace/NamespaceServiceTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/namespace/NamespaceServiceTest.java @@ -46,7 +46,6 @@ import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicReference; -import java.util.function.Predicate; import lombok.Cleanup; import org.apache.bookkeeper.mledger.ManagedLedger; import org.apache.commons.collections4.CollectionUtils; From abe8bff9ef2b8026e577adc72f42b28c9adb845f Mon Sep 17 00:00:00 2001 From: fengyubiao Date: Wed, 7 Sep 2022 10:21:21 +0800 Subject: [PATCH 4/4] add code comment --- .../broker/namespace/NamespaceServiceTest.java | 16 ++++++++++++++++ 1 file changed, 16 insertions(+) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/namespace/NamespaceServiceTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/namespace/NamespaceServiceTest.java index f3cfd33519c9b..03260aee46d7d 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/namespace/NamespaceServiceTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/namespace/NamespaceServiceTest.java @@ -78,6 +78,7 @@ import org.apache.pulsar.common.util.collections.ConcurrentOpenHashMap; import org.apache.pulsar.metadata.api.GetResult; import org.apache.pulsar.metadata.api.MetadataCache; +import org.apache.pulsar.metadata.api.Notification; import org.apache.pulsar.metadata.api.extended.CreateOption; import org.apache.pulsar.policies.data.loadbalancer.AdvertisedListener; import org.apache.pulsar.policies.data.loadbalancer.LoadReport; @@ -798,6 +799,21 @@ public void testModularLoadManagerRemoveBundleAndLoad() throws Exception { assertFalse(getResult.isPresent()); } + /** + * 1. Manually trigger "LoadReportUpdaterTask" + * 2. Registry another new zk-node-listener "waitForBrokerChangeNotice". + * 3. Wait "waitForBrokerChangeNotice" is done, this task will be executed after + * {@link ModularLoadManagerImpl#handleDataNotification(Notification)}, because it is registry later than + * {@link ModularLoadManagerImpl#handleDataNotification(Notification)}. So if "waitForBrokerChangeNotice" is done + * we can guarantee {@link ModularLoadManagerImpl#handleDataNotification(Notification)} is done. At this time + * we still could not guarantee {@link ModularLoadManagerImpl#handleDataNotification(Notification)} has + * finished all things, because there has a async task be submitted to "ModularLoadManagerImpl.scheduler" by + * {@link ModularLoadManagerImpl#handleDataNotification(Notification)}. + * 4. Submit a new task to "scheduler"(it is a singleton thread executor). + * 5. Wait the new task done, if the new task done, we can guarantee + * {@link ModularLoadManagerImpl#handleDataNotification(Notification)} has finished all things. + * 6. Manually trigger "LoadResourceQuotaUpdaterTask". + */ private void waitResourceDataUpdateToZK(LoadManager loadManager) throws Exception { CompletableFuture waitForBrokerChangeNotice = registryBrokerDataChangeNotice(); // Manually trigger "LoadReportUpdaterTask"