Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -223,6 +223,8 @@ Set<BookieId> 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)
Expand Down Expand Up @@ -315,6 +317,8 @@ Set<BookieId> 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;
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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;
Expand All @@ -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 {
Expand Down Expand Up @@ -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<String, Map<String, BookieInfo>> 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<Thread> 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<Set<BookieId>> clusterChange = new FutureTask<>(() -> {
clusterChangeThread.set(Thread.currentThread());
return isolationPolicy.onClusterChanged(Set.of(primaryBookie, joiningPrimaryBookie, ungroupedBookie),
Set.of());
});
Set<BookieId> excludedBookies = new HashSet<>();
FutureTask<List<BookieId>> 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<BookieId> 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 {

Expand Down