From 26a5052de8e32164c8136f03bf7867fa1bc0729b Mon Sep 17 00:00:00 2001 From: Denovo1998 Date: Thu, 1 Oct 2026 14:25:12 +0800 Subject: [PATCH] [fix][broker] Prevent incorrect isolation fallback during cluster updates --- ...IsolatedBookieEnsemblePlacementPolicy.java | 4 + ...atedBookieEnsemblePlacementPolicyTest.java | 97 +++++++++++++++++++ 2 files changed, 101 insertions(+) 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 5dbefdc66629b..4f07f2c81792d 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 @@ -223,6 +223,8 @@ Set getExcludedBookiesWithIsolationGroups(int ensembleSize, if (isolationGroups != null && isolationGroups.getLeft().contains(PULSAR_SYSTEM_TOPIC_ISOLATION_GROUP)) { return excludedBookies; } + // Keep isolation counts and exclusions on one cluster view while onClusterChanged updates knownBookies. + rwLock.readLock().lock(); try { if (bookieMappingCache != null) { bookieMappingCache.get(BookieRackAffinityMapping.BOOKIE_INFO_ROOT_PATH) @@ -315,6 +317,8 @@ Set getExcludedBookiesWithIsolationGroups(int ensembleSize, } } catch (Exception e) { log.warn().attr("store", e.getMessage()).log("Error getting bookie isolation info from metadata store"); + } finally { + rwLock.readLock().unlock(); } return excludedBookies; } 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 4ce2923318b65..46fd1444d634c 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,6 +18,7 @@ */ package org.apache.pulsar.bookie.rackawareness; +import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.Assertions.assertThatThrownBy; import static org.mockito.Mockito.doReturn; import static org.mockito.Mockito.mock; @@ -43,14 +44,19 @@ import java.util.Optional; import java.util.Set; import java.util.concurrent.CompletableFuture; +import java.util.concurrent.CountDownLatch; import java.util.concurrent.ExecutionException; +import java.util.concurrent.FutureTask; import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicReference; import org.apache.bookkeeper.client.BKException.BKNotEnoughBookiesException; import org.apache.bookkeeper.client.RackawareEnsemblePlacementPolicy; 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.proto.BookieAddressResolver; +import org.apache.bookkeeper.proto.BookieAddressResolver.BookieIdNotResolvedException; import org.apache.bookkeeper.stats.NullStatsLogger; import org.apache.commons.lang3.tuple.MutablePair; import org.apache.commons.lang3.tuple.Pair; @@ -69,6 +75,7 @@ import org.testng.Assert; import org.testng.annotations.AfterMethod; import org.testng.annotations.BeforeMethod; +import org.testng.annotations.DataProvider; import org.testng.annotations.Test; public class IsolatedBookieEnsemblePlacementPolicyTest { @@ -541,6 +548,96 @@ public void testSecondaryIsolationGroupsBookies() throws Exception { assertTrue(ensemble.contains(new BookieSocketAddress(BOOKIE4).toBookieId())); } + @DataProvider + public Object[][] placementOperations() { + return new Object[][]{{false}, {true}}; + } + + @Test(dataProvider = "placementOperations") + public void testPlacementDuringClusterChange(boolean replaceBookie) throws Exception { + BookieId primaryBookie = BookieId.parse(BOOKIE1); + BookieId leavingPrimaryBookie = BookieId.parse(BOOKIE2); + BookieId ungroupedBookie = BookieId.parse(BOOKIE3); + BookieId joiningPrimaryBookie = BookieId.parse(BOOKIE4); + Map> bookieMapping = Map.of("group1", Map.of( + BOOKIE1, BookieInfo.builder().rack("rack0").build(), + BOOKIE2, BookieInfo.builder().rack("rack1").build(), + BOOKIE4, BookieInfo.builder().rack("rack1").build())); + store.put(BookieRackAffinityMapping.BOOKIE_INFO_ROOT_PATH, jsonMapper.writeValueAsBytes(bookieMapping), + Optional.empty()).join(); + + CountDownLatch clusterChangePaused = new CountDownLatch(1); + CountDownLatch finishClusterChange = new CountDownLatch(1); + AtomicReference clusterChangeThread = new AtomicReference<>(); + BookieAddressResolver resolver = bookieId -> { + if (bookieId.equals(joiningPrimaryBookie) && Thread.currentThread() == clusterChangeThread.get()) { + // onClusterChanged holds the write lock here and has already removed the old primary bookie. + clusterChangePaused.countDown(); + try { + if (!finishClusterChange.await(30, TimeUnit.SECONDS)) { + throw new IllegalStateException("Timed out waiting to finish the cluster change"); + } + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + throw new BookieIdNotResolvedException(bookieId, e); + } + } + return BookieSocketAddress.LEGACY_BOOKIEID_RESOLVER.resolve(bookieId); + }; + ClientConfiguration bkClientConf = new ClientConfiguration(); + bkClientConf.setProperty(BookieRackAffinityMapping.METADATA_STORE_INSTANCE, store); + bkClientConf.setProperty(IsolatedBookieEnsemblePlacementPolicy.ISOLATION_BOOKIE_GROUPS, "group1"); + IsolatedBookieEnsemblePlacementPolicy isolationPolicy = new IsolatedBookieEnsemblePlacementPolicy(); + isolationPolicy.initialize(bkClientConf, Optional.empty(), timer, SettableFeatureProvider.DISABLE_ALL, + NullStatsLogger.INSTANCE, resolver); + isolationPolicy.getInitialRackConfigurationLoadFuture().get(30, TimeUnit.SECONDS); + isolationPolicy.onClusterChanged(Set.of(primaryBookie, leavingPrimaryBookie, ungroupedBookie), Set.of()); + + FutureTask> clusterChange = new FutureTask<>(() -> { + clusterChangeThread.set(Thread.currentThread()); + return isolationPolicy.onClusterChanged(Set.of(primaryBookie, joiningPrimaryBookie, ungroupedBookie), + Set.of()); + }); + Set excludedBookies = new HashSet<>(); + FutureTask> placement = new FutureTask<>(() -> { + if (replaceBookie) { + return List.of(isolationPolicy.replaceBookie(2, 2, 2, Map.of(), + List.of(primaryBookie, leavingPrimaryBookie), leavingPrimaryBookie, excludedBookies) + .getResult()); + } + return isolationPolicy.newEnsemble(2, 2, 2, Map.of(), excludedBookies).getResult(); + }); + Thread updateThread = new Thread(clusterChange, "isolation-cluster-change"); + Thread placementThread = new Thread(placement, "isolation-placement"); + try { + updateThread.start(); + assertThat(clusterChangePaused.await(10, TimeUnit.SECONDS)).as("cluster update reached address resolution") + .isTrue(); + placementThread.start(); + // Both versions eventually wait for the write lock. Before the fix, isolation has already used + // the incomplete topology by then and incorrectly allowed the ungrouped fallback. + Awaitility.await().atMost(Duration.ofSeconds(10)) + .until(() -> placementThread.getState() == Thread.State.WAITING); + finishClusterChange.countDown(); + assertThat(clusterChange.get(10, TimeUnit.SECONDS)).containsExactly(leavingPrimaryBookie); + List result = placement.get(10, TimeUnit.SECONDS); + assertThat(excludedBookies).as("a complete primary group must keep ungrouped bookies excluded") + .contains(ungroupedBookie); + if (replaceBookie) { + assertThat(result).containsExactly(joiningPrimaryBookie); + } else { + assertThat(result).containsExactlyInAnyOrder(primaryBookie, joiningPrimaryBookie); + } + } finally { + finishClusterChange.countDown(); + updateThread.join(TimeUnit.SECONDS.toMillis(10)); + placementThread.join(TimeUnit.SECONDS.toMillis(10)); + isolationPolicy.uninitalize(); + assertThat(updateThread.isAlive()).as("cluster update thread stopped").isFalse(); + assertThat(placementThread.isAlive()).as("placement thread stopped").isFalse(); + } + } + @Test public void testSecondaryIsolationGroupsBookiesNegative() throws Exception {