From 9c810580dd37dc3e724d80d0e40ca2172341512e Mon Sep 17 00:00:00 2001 From: Heesung Sohn Date: Tue, 22 Apr 2025 10:00:21 -0700 Subject: [PATCH] [fix][broker] fix NPE from the wrong iterator in the ownership cleanup job(ExtensibleLoadManagerImpl only) --- .../channel/ServiceUnitStateChannelImpl.java | 2 +- .../channel/ServiceUnitStateChannelTest.java | 20 +++++++++++++++++++ 2 files changed, 21 insertions(+), 1 deletion(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/loadbalance/extensions/channel/ServiceUnitStateChannelImpl.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/loadbalance/extensions/channel/ServiceUnitStateChannelImpl.java index ede15680512fb..f20fb70cffd21 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/loadbalance/extensions/channel/ServiceUnitStateChannelImpl.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/loadbalance/extensions/channel/ServiceUnitStateChannelImpl.java @@ -1600,7 +1600,7 @@ private void doCleanup(String broker, boolean gracefully) { // clean system bundles in the end var orphanSystemServiceUnitIter = orphanSystemServiceUnits.entrySet().iterator(); while (orphanSystemServiceUnitIter.hasNext()) { - var orphanSystemServiceUnit = iter.next(); + var orphanSystemServiceUnit = orphanSystemServiceUnitIter.next(); log.info("Overriding orphan system service unit:{}", orphanSystemServiceUnit.getKey()); overrideFutures.add( overrideOwnership(orphanSystemServiceUnit.getKey(), orphanSystemServiceUnit.getValue(), broker, diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/loadbalance/extensions/channel/ServiceUnitStateChannelTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/loadbalance/extensions/channel/ServiceUnitStateChannelTest.java index b6e38d4f6956c..7db0baabf6552 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/loadbalance/extensions/channel/ServiceUnitStateChannelTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/loadbalance/extensions/channel/ServiceUnitStateChannelTest.java @@ -1897,6 +1897,26 @@ public void unloadTimeoutCheckTest() TimeUnit.SECONDS).get(2, TimeUnit.SECONDS); } + @Test(priority = 23) + public void testCleanSystemTopicOwnership() + throws Exception { + String topic = "persistent://pulsar/system/test-system-topic"; + NamespaceBundle bundleName = pulsar.getNamespaceService().getBundle(TopicName.get(topic)); + var releasing = new ServiceUnitStateData(Releasing, pulsar2.getBrokerId(), pulsar1.getBrokerId(), 1); + doReturn(CompletableFuture.completedFuture(Optional.of(brokerId1))) + .when(loadManager).selectAsync(any(), any(), any()); + + try { + disableChannels(); + overrideTableView(channel1, bundleName.toString(), releasing); + } finally { + enableChannels(); + } + + channel1.cleanOwnerships(); + channel2.cleanOwnerships(); + } + private static ConcurrentHashMap>> getOwnerRequests( ServiceUnitStateChannel channel) throws IllegalAccessException {