From 151171b653ed34a6cf925584f77028db2438cf40 Mon Sep 17 00:00:00 2001 From: Yuri Mizushima Date: Mon, 15 Nov 2021 16:57:41 +0900 Subject: [PATCH 1/4] feat: add relocate placement command --- .../apache/bookkeeper/bookie/BookieShell.java | 55 +++ .../apache/bookkeeper/client/BookKeeper.java | 3 +- .../bookkeeper/client/BookKeeperAdmin.java | 2 +- .../client/EnsemblePlacementPolicy.java | 25 ++ .../bookkeeper/client/LedgerFragment.java | 4 +- .../RackawareEnsemblePlacementPolicy.java | 22 ++ .../RackawareEnsemblePlacementPolicyImpl.java | 174 ++++++++++ .../TopologyAwareEnsemblePlacementPolicy.java | 14 +- .../CorrectEnsemblePlacementCommand.java | 321 ++++++++++++++++++ 9 files changed, 610 insertions(+), 10 deletions(-) create mode 100644 bookkeeper-server/src/main/java/org/apache/bookkeeper/tools/cli/commands/bookies/CorrectEnsemblePlacementCommand.java diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/BookieShell.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/BookieShell.java index ad7e8dbc15e..21c24b209aa 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/BookieShell.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/BookieShell.java @@ -71,6 +71,7 @@ import org.apache.bookkeeper.tools.cli.commands.bookie.RegenerateInterleavedStorageIndexFileCommand; import org.apache.bookkeeper.tools.cli.commands.bookie.SanityTestCommand; import org.apache.bookkeeper.tools.cli.commands.bookie.UpdateBookieInLedgerCommand; +import org.apache.bookkeeper.tools.cli.commands.bookies.CorrectEnsemblePlacementCommand; import org.apache.bookkeeper.tools.cli.commands.bookies.DecommissionCommand; import org.apache.bookkeeper.tools.cli.commands.bookies.EndpointInfoCommand; import org.apache.bookkeeper.tools.cli.commands.bookies.InfoCommand; @@ -160,6 +161,7 @@ public class BookieShell implements Tool { static final String CMD_CHECK_DB_LEDGERS_INDEX = "check-db-ledgers-index"; static final String CMD_REGENERATE_INTERLEAVED_STORAGE_INDEX_FILE = "regenerate-interleaved-storage-index-file"; static final String CMD_QUERY_AUTORECOVERY_STATUS = "queryrecoverystatus"; + static final String CMD_CORRECT_ENSEMBLE_PLACEMENT = "correct-ensemble-placement"; // cookie commands static final String CMD_CREATE_COOKIE = "cookie_create"; @@ -2257,6 +2259,58 @@ int runCmd(CommandLine cmdLine) throws Exception { } } + class CorrectEnsemblePlacementCmd extends MyCommand { + final Options opts = new Options(); + + public CorrectEnsemblePlacementCmd() { + super(CMD_CORRECT_ENSEMBLE_PLACEMENT); + Option ledgerOption = new Option("l", "ledgerids", true, + "Target ledger IDs to relocate." + + " Multiple can be specified, comma separated."); + ledgerOption.setRequired(true); + ledgerOption.setValueSeparator(','); + ledgerOption.setArgs(Option.UNLIMITED_VALUES); + + opts.addOption(ledgerOption); + opts.addOption("dr", "dryrun", false, + "Printing the relocation plan w/o doing actual relocation"); + opts.addOption("f", "force", false, + "Force relocation without confirmation"); + opts.addOption("sk", "skipOpenLedgers", false, + "Skip relocating open ledgers"); + } + + @Override + Options getOptions() { + return opts; + } + + @Override + String getDescription() { + return "Relocate ledgers to adhere ensemble placement policy."; + } + + @Override + String getUsage() { + return CMD_CORRECT_ENSEMBLE_PLACEMENT + + " --ledgerids [--dryrun] [--force] [--skipOpenLedgers]"; + } + + @Override + int runCmd(CommandLine cmdLine) throws Exception { + final CorrectEnsemblePlacementCommand cmd = new CorrectEnsemblePlacementCommand(); + final CorrectEnsemblePlacementCommand.CorrectEnsemblePlacementFlags + flags = new CorrectEnsemblePlacementCommand.CorrectEnsemblePlacementFlags(); + final List ledgerIds = Arrays.stream(cmdLine.getOptionValues("ledgerids")).map(Long::parseLong) + .collect(Collectors.toList()); + flags.ledgerIds(ledgerIds); + flags.dryRun(cmdLine.hasOption("dryrun")); + flags.force(cmdLine.hasOption("force")); + flags.skipOpenLedgers(cmdLine.hasOption("skipOpenLedgers")); + return cmd.apply(bkConf, flags) ? 0 : 1; + } + } + final Map commands = new HashMap<>(); { @@ -2302,6 +2356,7 @@ int runCmd(CommandLine cmdLine) throws Exception { commands.put(CMD_LOSTBOOKIERECOVERYDELAY, new LostBookieRecoveryDelayCmd()); commands.put(CMD_TRIGGERAUDIT, new TriggerAuditCmd()); commands.put(CMD_FORCEAUDITCHECKS, new ForceAuditorChecksCmd()); + commands.put(CMD_CORRECT_ENSEMBLE_PLACEMENT, new CorrectEnsemblePlacementCmd()); // cookie related commands commands.put(CMD_CREATE_COOKIE, new CreateCookieCommand().asShellCommand(CMD_CREATE_COOKIE, bkConf)); diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/client/BookKeeper.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/client/BookKeeper.java index f08f47b98c6..96ec06e3a58 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/client/BookKeeper.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/client/BookKeeper.java @@ -696,8 +696,7 @@ OrderedScheduler getScheduler() { return scheduler; } - @VisibleForTesting - EnsemblePlacementPolicy getPlacementPolicy() { + public EnsemblePlacementPolicy getPlacementPolicy() { return placementPolicy; } diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/client/BookKeeperAdmin.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/client/BookKeeperAdmin.java index 1dde401853e..55c1ccc3ff7 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/client/BookKeeperAdmin.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/client/BookKeeperAdmin.java @@ -1137,7 +1137,7 @@ public void replicateLedgerFragment(LedgerHandle lh, replicateLedgerFragment(lh, ledgerFragment, targetBookieAddresses, onReadEntryFailureCallback); } - private void replicateLedgerFragment(LedgerHandle lh, + public void replicateLedgerFragment(LedgerHandle lh, final LedgerFragment ledgerFragment, final Map targetBookieAddresses, final BiConsumer onReadEntryFailureCallback) diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/client/EnsemblePlacementPolicy.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/client/EnsemblePlacementPolicy.java index fcd38f2a92f..1922885d88c 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/client/EnsemblePlacementPolicy.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/client/EnsemblePlacementPolicy.java @@ -441,6 +441,31 @@ default boolean areAckedBookiesAdheringToPlacementPolicy(Set ackedBook return true; } + /** + * Returns placement result. If the currentEnsemble is not adhering placement policy, returns new ensemble that + * adheres placement policy. It should be implemented so as to minify the number of bookies replaced. + * + * @param ensembleSize + * ensemble size + * @param writeQuorumSize + * writeQuorumSize of the ensemble + * @param ackQuorumSize + * ackQuorumSize of the ensemble + * @param excludeBookies + * bookies that should not be considered as targets + * @param currentEnsemble + * current ensemble + * @return a placement result + */ + default PlacementResult> replaceToAdherePlacementPolicy( + int ensembleSize, + int writeQuorumSize, + int ackQuorumSize, + Set excludeBookies, + List currentEnsemble) { + throw new UnsupportedOperationException(); + } + /** * enum for PlacementPolicyAdherence. Currently we are supporting tri-value * enum for PlacementPolicyAdherence. If placement policy is met strictly diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/client/LedgerFragment.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/client/LedgerFragment.java index 1fb1e50cb02..a18e944e342 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/client/LedgerFragment.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/client/LedgerFragment.java @@ -40,7 +40,7 @@ public class LedgerFragment { private final DistributionSchedule schedule; private final boolean isLedgerClosed; - LedgerFragment(LedgerHandle lh, + public LedgerFragment(LedgerHandle lh, long firstEntryId, long lastKnownEntryId, Set bookieIndexes) { @@ -56,7 +56,7 @@ public class LedgerFragment { || !ensemble.equals(ensembles.get(ensembles.lastKey())); } - LedgerFragment(LedgerFragment lf, Set subset) { + public LedgerFragment(LedgerFragment lf, Set subset) { this.ledgerId = lf.ledgerId; this.firstEntryId = lf.firstEntryId; this.lastKnownEntryId = lf.lastKnownEntryId; diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/client/RackawareEnsemblePlacementPolicy.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/client/RackawareEnsemblePlacementPolicy.java index 626b7cc2e60..dad2f514d1c 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/client/RackawareEnsemblePlacementPolicy.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/client/RackawareEnsemblePlacementPolicy.java @@ -236,6 +236,28 @@ public BookieNode selectFromNetworkLocation( } } + @Override + public PlacementResult> replaceToAdherePlacementPolicy( + int ensembleSize, + int writeQuorumSize, + int ackQuorumSize, + Set excludeBookies, + List currentEnsemble) { + final PlacementResult> placementResult = + super.replaceToAdherePlacementPolicy(ensembleSize, writeQuorumSize, ackQuorumSize, + excludeBookies, currentEnsemble); + if (placementResult.isAdheringToPolicy() != PlacementPolicyAdherence.FAIL) { + return placementResult; + } else { + if (slave == null) { + return placementResult; + } else { + return slave.replaceToAdherePlacementPolicy(ensembleSize, writeQuorumSize, ackQuorumSize, + excludeBookies, currentEnsemble); + } + } + } + @Override public void handleBookiesThatLeft(Set leftBookies) { super.handleBookiesThatLeft(leftBookies); diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/client/RackawareEnsemblePlacementPolicyImpl.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/client/RackawareEnsemblePlacementPolicyImpl.java index 46c5a10786b..d9167c321d4 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/client/RackawareEnsemblePlacementPolicyImpl.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/client/RackawareEnsemblePlacementPolicyImpl.java @@ -46,6 +46,8 @@ import java.util.concurrent.TimeUnit; import java.util.concurrent.locks.ReentrantReadWriteLock; +import java.util.stream.Collectors; +import java.util.stream.Stream; import org.apache.bookkeeper.client.BKException.BKNotEnoughBookiesException; import org.apache.bookkeeper.client.BookieInfoReader.BookieInfo; import org.apache.bookkeeper.client.WeightedRandomSelection.WeightedObject; @@ -69,6 +71,7 @@ import org.apache.bookkeeper.stats.StatsLogger; import org.apache.bookkeeper.stats.annotations.StatsDoc; import org.apache.commons.collections4.CollectionUtils; +import org.apache.commons.lang3.tuple.Pair; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -1071,4 +1074,175 @@ public boolean areAckedBookiesAdheringToPlacementPolicy(Set ackedBooki } return rackCounter.size() >= minWriteQuorumNumRacksPerWriteQuorum; } + + @Override + public PlacementResult> replaceToAdherePlacementPolicy( + int ensembleSize, + int writeQuorumSize, + int ackQuorumSize, + Set excludeBookies, + List currentEnsemble) { + rwLock.readLock().lock(); + try { + final List provisionalEnsembleNodes = currentEnsemble.stream() + .map(this::convertBookieToNode).collect(Collectors.toList()); + final Set excludeNodes = convertBookiesToNodes( + addDefaultRackBookiesIfMinNumRacksIsEnforced(excludeBookies)); + int minNumRacksPerWriteQuorumForThisEnsemble = Math.min(writeQuorumSize, minNumRacksPerWriteQuorum); + final RRTopologyAwareCoverageEnsemble ensemble = + new RRTopologyAwareCoverageEnsemble( + ensembleSize, + writeQuorumSize, + ackQuorumSize, + RACKNAME_DISTANCE_FROM_LEAVES, + null, + null, + minNumRacksPerWriteQuorumForThisEnsemble); + + int numRacks = topology.getNumOfRacks(); + // only one rack or less than minNumRacksPerWriteQuorumForThisEnsemble, stop calculation to skip relocation + if (numRacks < 2 || numRacks < minNumRacksPerWriteQuorumForThisEnsemble) { + LOG.warn("Skip ensemble relocation because the cluster has only {} rack.", numRacks); + return PlacementResult.of(Collections.emptyList(), PlacementPolicyAdherence.FAIL); + } + + BookieNode prevNode = null; + final BookieNode firstNode = provisionalEnsembleNodes.get(0); + // use same bookie at first to reduce ledger replication + if (!excludeNodes.contains(firstNode) && ensemble.apply(firstNode, ensemble) + && ensemble.addNode(firstNode)) { + excludeNodes.add(firstNode); + prevNode = firstNode; + } + + for (int i = prevNode == null ? 0 : 1; i < ensembleSize; i++) { + final String curRack; + if (null == prevNode) { + if ((null == localNode) || defaultRack.equals(localNode.getNetworkLocation())) { + curRack = NodeBase.ROOT; + } else { + curRack = localNode.getNetworkLocation(); + } + } else { + curRack = "~" + prevNode.getNetworkLocation(); + } + + try { + prevNode = replaceToAdherePlacementPolicyInternal( + curRack, excludeNodes, ensemble, ensemble, + provisionalEnsembleNodes, i, ensembleSize, minNumRacksPerWriteQuorumForThisEnsemble); + // replace to newer node + provisionalEnsembleNodes.set(i, prevNode); + } catch (BKNotEnoughBookiesException e) { + LOG.warn("Skip ensemble relocation because the cluster has not enough bookies."); + return PlacementResult.of(Collections.emptyList(), PlacementPolicyAdherence.FAIL); + } + } + List bookieList = ensemble.toList(); + if (ensembleSize != bookieList.size()) { + LOG.warn("Not enough {} bookies are available to form an ensemble : {}.", + ensembleSize, bookieList); + return PlacementResult.of(Collections.emptyList(), PlacementPolicyAdherence.FAIL); + } + return PlacementResult.of(bookieList, + isEnsembleAdheringToPlacementPolicy( + bookieList, writeQuorumSize, ackQuorumSize)); + } finally { + rwLock.readLock().unlock(); + } + } + + private BookieNode replaceToAdherePlacementPolicyInternal( + String netPath, Set excludeBookies, Predicate predicate, + Ensemble ensemble, List provisionalEnsembleNodes, int ensembleIndex, + int ensembleSize, int minNumRacksPerWriteQuorumForThisEnsemble) throws BKNotEnoughBookiesException { + final BookieNode currentNode = provisionalEnsembleNodes.get(ensembleIndex); + // if the current bookie could be applied to the ensemble, apply it to minify the number of bookies replaced + if (!excludeBookies.contains(currentNode) && predicate.apply(currentNode, ensemble)) { + if (ensemble.addNode(currentNode)) { + // add the candidate to exclude set + excludeBookies.add(currentNode); + } + return currentNode; + } + + final List>> conditionList = new ArrayList<>(); + final Set preExcludeRacks = new HashSet<>(); + final Set postExcludeRacks = new HashSet<>(); + for (int i = 0; i < minNumRacksPerWriteQuorumForThisEnsemble - 1; i++) { + preExcludeRacks.add(provisionalEnsembleNodes.get(Math.floorMod((ensembleIndex - i - 1), ensembleSize)) + .getNetworkLocation()); + postExcludeRacks.add(provisionalEnsembleNodes.get(Math.floorMod((ensembleIndex + i + 1), ensembleSize)) + .getNetworkLocation()); + } + // adhere minNumRacksPerWriteQuorum by preExcludeRacks + // avoid additional replace from write quorum candidates by preExcludeRacks and postExcludeRacks + // avoid to use first candidate bookies for election by provisionalEnsembleNodes + conditionList.add(Pair.of( + "~" + String.join(",", + Stream.concat(preExcludeRacks.stream(), postExcludeRacks.stream()).collect(Collectors.toSet())), + provisionalEnsembleNodes + )); + // avoid to use same rack between previous index by netPath + // avoid to use first candidate bookies for election by provisionalEnsembleNodes + conditionList.add(Pair.of(netPath, provisionalEnsembleNodes)); + // avoid to use same rack between previous index by netPath + conditionList.add(Pair.of(netPath, Collections.emptyList())); + + for (Pair> condition : conditionList) { + WeightedRandomSelection wRSelection = null; + + final List leaves = new ArrayList<>(topology.getLeaves(condition.getLeft())); + if (!isWeighted) { + Collections.shuffle(leaves); + } else { + if (CollectionUtils.subtract(leaves, excludeBookies).size() < 1) { + throw new BKNotEnoughBookiesException(); + } + wRSelection = prepareForWeightedSelection(leaves); + if (wRSelection == null) { + throw new BKNotEnoughBookiesException(); + } + } + + final Iterator it = leaves.iterator(); + final Set bookiesSeenSoFar = new HashSet<>(); + while (true) { + Node n; + if (isWeighted) { + if (bookiesSeenSoFar.size() == leaves.size()) { + // Don't loop infinitely. + break; + } + n = wRSelection.getNextRandom(); + bookiesSeenSoFar.add(n); + } else { + if (it.hasNext()) { + n = it.next(); + } else { + break; + } + } + if (excludeBookies.contains(n)) { + continue; + } + if (!(n instanceof BookieNode) || !predicate.apply((BookieNode) n, ensemble)) { + continue; + } + // additional excludeBookies + if (condition.getRight().contains(n)) { + continue; + } + BookieNode bn = (BookieNode) n; + // got a good candidate + if (ensemble.addNode(bn)) { + // add the candidate to exclude set + excludeBookies.add(bn); + } + return bn; + } + } + + throw new BKNotEnoughBookiesException(); + } } diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/client/TopologyAwareEnsemblePlacementPolicy.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/client/TopologyAwareEnsemblePlacementPolicy.java index 438053f5449..dc136fd83b5 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/client/TopologyAwareEnsemblePlacementPolicy.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/client/TopologyAwareEnsemblePlacementPolicy.java @@ -802,12 +802,16 @@ protected String resolveNetworkLocation(BookieId addr) { protected Set convertBookiesToNodes(Collection excludeBookies) { Set nodes = new HashSet(); for (BookieId addr : excludeBookies) { - BookieNode bn = knownBookies.get(addr); - if (null == bn) { - bn = createBookieNode(addr); - } - nodes.add(bn); + nodes.add(convertBookieToNode(addr)); } return nodes; } + + protected BookieNode convertBookieToNode(BookieId addr) { + BookieNode bn = knownBookies.get(addr); + if (null == bn) { + bn = createBookieNode(addr); + } + return bn; + } } diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/tools/cli/commands/bookies/CorrectEnsemblePlacementCommand.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/tools/cli/commands/bookies/CorrectEnsemblePlacementCommand.java new file mode 100644 index 00000000000..36cddb7f11f --- /dev/null +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/tools/cli/commands/bookies/CorrectEnsemblePlacementCommand.java @@ -0,0 +1,321 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.bookkeeper.tools.cli.commands.bookies; + +import com.beust.jcommander.Parameter; +import com.beust.jcommander.converters.CommaParameterSplitter; +import com.google.common.annotations.VisibleForTesting; +import com.google.common.base.Functions; +import com.google.common.util.concurrent.UncheckedExecutionException; +import java.io.IOException; +import java.net.URI; +import java.util.Collections; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.NavigableSet; +import java.util.TreeSet; +import java.util.concurrent.ConcurrentSkipListSet; +import java.util.concurrent.CountDownLatch; +import java.util.stream.Collectors; +import java.util.stream.IntStream; +import lombok.Cleanup; +import lombok.Setter; +import lombok.experimental.Accessors; +import org.apache.bookkeeper.bookie.BookieException; +import org.apache.bookkeeper.client.BKException; +import org.apache.bookkeeper.client.BookKeeper; +import org.apache.bookkeeper.client.BookKeeperAdmin; +import org.apache.bookkeeper.client.EnsemblePlacementPolicy; +import org.apache.bookkeeper.client.LedgerFragment; +import org.apache.bookkeeper.client.api.LedgerMetadata; +import org.apache.bookkeeper.conf.ClientConfiguration; +import org.apache.bookkeeper.conf.ServerConfiguration; +import org.apache.bookkeeper.meta.LedgerManagerFactory; +import org.apache.bookkeeper.meta.LedgerUnderreplicationManager; +import org.apache.bookkeeper.meta.MetadataBookieDriver; +import org.apache.bookkeeper.meta.MetadataDrivers; +import org.apache.bookkeeper.meta.exceptions.MetadataException; +import org.apache.bookkeeper.net.BookieId; +import org.apache.bookkeeper.replication.ReplicationException; +import org.apache.bookkeeper.stats.NullStatsLogger; +import org.apache.bookkeeper.tools.cli.helpers.BookieCommand; +import org.apache.bookkeeper.tools.framework.CliFlags; +import org.apache.bookkeeper.tools.framework.CliSpec; +import org.apache.bookkeeper.util.IOUtils; +import org.apache.commons.configuration.ConfigurationException; +import org.apache.commons.lang3.tuple.Pair; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * Command to relocate ledgers to adhere ensemble placement policy. + */ +public class CorrectEnsemblePlacementCommand extends + BookieCommand { + private static final Logger LOG = LoggerFactory.getLogger(CorrectEnsemblePlacementCommand.class); + + private static final String NAME = "correct-ensemble-placement"; + private static final String DESC = "Relocate ledgers to adhere ensemble placement policy."; + + public CorrectEnsemblePlacementCommand() { + this(new CorrectEnsemblePlacementFlags()); + } + + private CorrectEnsemblePlacementCommand(CorrectEnsemblePlacementFlags flags) { + super(CliSpec.newBuilder() + .withName(NAME) + .withDescription(DESC) + .withFlags(flags) + .build()); + } + + /** + * Flags for correct-ensemble-placement command. + */ + @Accessors(fluent = true) + @Setter + public static class CorrectEnsemblePlacementFlags extends CliFlags { + @Parameter(names = { "-l", "--ledgerids" }, + description = "Target ledger IDs to relocate. Multiple can be specified, comma separated.", + splitter = CommaParameterSplitter.class, required = true) + private List ledgerIds; + + @Parameter(names = { "-dr", "--dryrun" }, + description = "Printing the relocation plan w/o doing actual relocation") + private boolean dryRun; + + @Parameter(names = { "-f", "--force" }, description = "Force relocation without confirmation") + private boolean force; + + @Parameter(names = {"-sk", "--skipOpenLedgers"}, description = "Skip relocating open ledgers") + private boolean skipOpenLedgers; + } + + @Override + public boolean apply(ServerConfiguration conf, CorrectEnsemblePlacementFlags flags) { + try { + if (flags.dryRun) { + LOG.info("The dry-run output could change every time you run" + + " since the selection of bookies replaced includes some randomness."); + } + if (!flags.skipOpenLedgers) { + LOG.warn("Try to relocate also open ledgers. It is not recommended."); + } + try { + if (!flags.force) { + final boolean confirm = + IOUtils.confirmPrompt("Are you sure to relocate target ledgers?"); + if (!confirm) { + LOG.error("Relocation is aborted."); + return false; + } + } + } catch (IOException e) { + LOG.error("Error during relocation", e); + return false; + } + + final ClientConfiguration clientConf = new ClientConfiguration(conf); + final BookKeeper bookKeeper = new BookKeeper(clientConf); + final BookKeeperAdmin admin = new BookKeeperAdmin(clientConf); + return relocate(conf, flags, bookKeeper, admin); + } catch (Exception e) { + throw new UncheckedExecutionException(e.getMessage(), e); + } + } + + @VisibleForTesting + public boolean relocate(ServerConfiguration conf, CorrectEnsemblePlacementFlags flags, + BookKeeper bookKeeper, BookKeeperAdmin admin) throws Exception { + final EnsemblePlacementPolicy placementPolicy = bookKeeper.getPlacementPolicy(); + + @Cleanup + final MetadataBookieDriver metadataDriver = instantiateMetadataDriver(conf); + @Cleanup + final LedgerManagerFactory lmf = metadataDriver.getLedgerManagerFactory(); + @Cleanup + final LedgerUnderreplicationManager lum = lmf.newLedgerUnderreplicationManager(); + + final List targetLedgers = + flags.ledgerIds.stream().parallel().distinct().filter(ledgerId -> { + try { + return (!flags.skipOpenLedgers || bookKeeper.isClosed(ledgerId)) + && !lum.isLedgerBeingReplicated(ledgerId); + } catch (BKException | InterruptedException | ReplicationException e) { + LOG.warn("Failed to add the ledger {} to target.", ledgerId, e); + return false; + } + }).collect(Collectors.toList()); + + if (targetLedgers.isEmpty()) { + LOG.info("None of ledgers are relocated."); + return true; + } + + final NavigableSet unAcquirableLedgers = new TreeSet<>(); + final NavigableSet> failedTargets = new ConcurrentSkipListSet<>(); + final CountDownLatch latch = new CountDownLatch(targetLedgers.size()); + for (long ledgerId : targetLedgers) { + if (!flags.dryRun) { + try { + lum.acquireUnderreplicatedLedger(ledgerId); + } catch (ReplicationException e) { + LOG.warn("Failed to acquire ledger to under replicated {}.", ledgerId); + unAcquirableLedgers.add(ledgerId); + latch.countDown(); + continue; + } + } + admin.asyncOpenLedger(ledgerId, (rc, lh, ctx) -> { + try { + if (rc != BKException.Code.OK) { + LOG.warn("Failed to open ledger {}", ledgerId); + return; + } + + final LedgerMetadata ledgerMeta = lh.getLedgerMetadata(); + final Map ledgerFragmentsRange = new HashMap<>(); + Long curEntryId = null; + for (Map.Entry> entry : + ledgerMeta.getAllEnsembles().entrySet()) { + if (curEntryId != null) { + ledgerFragmentsRange.put(curEntryId, entry.getKey() - 1); + } + curEntryId = entry.getKey(); + } + if (curEntryId != null) { + ledgerFragmentsRange.put(curEntryId, lh.getLastAddConfirmed()); + } + + for (Map.Entry> entry : ledgerMeta.getAllEnsembles().entrySet()) { + if (placementPolicy.isEnsembleAdheringToPlacementPolicy(entry.getValue(), + ledgerMeta.getWriteQuorumSize(), ledgerMeta.getAckQuorumSize()) + == EnsemblePlacementPolicy.PlacementPolicyAdherence.FAIL) { + try { + final List currentEnsemble = entry.getValue(); + // Currently, don't consider quarantinedBookies + final EnsemblePlacementPolicy.PlacementResult> placementResult = + placementPolicy.replaceToAdherePlacementPolicy( + ledgerMeta.getEnsembleSize(), + ledgerMeta.getWriteQuorumSize(), + ledgerMeta.getAckQuorumSize(), + Collections.emptySet(), + currentEnsemble); + + if (placementResult.isAdheringToPolicy() + == EnsemblePlacementPolicy.PlacementPolicyAdherence.FAIL) { + LOG.warn("Failed to relocate the ensemble. So, skip the operation." + + " ledgerId: {}, fragmentIndex: {}", + ledgerId, entry.getKey()); + failedTargets.add(Pair.of(ledgerId, entry.getKey())); + } else { + final List newEnsemble = placementResult.getResult(); + final Map replaceBookiesMap = IntStream + .range(0, ledgerMeta.getEnsembleSize()).boxed() + .filter(i -> !newEnsemble.get(i).equals(currentEnsemble.get(i))) + .collect(Collectors.toMap(Functions.identity(), newEnsemble::get)); + if (replaceBookiesMap.isEmpty()) { + LOG.warn("Failed to get bookies to replace. So, skip the operation." + + " ledgerId: {}, fragmentIndex: {}", + ledgerId, entry.getKey()); + failedTargets.add(Pair.of(ledgerId, entry.getKey())); + } else if (flags.dryRun) { + LOG.info("Would replace the ensemble. ledgerId: {}, fragmentIndex: {}," + + " currentEnsemble: {} replaceBookiesMap {}", + ledgerId, entry.getKey(), + currentEnsemble, replaceBookiesMap); + } else { + if (LOG.isDebugEnabled()) { + LOG.debug("Try to replace the ensemble. ledgerId: {}, fragmentIndex: {}," + + " replaceBookiesMap {}", + ledgerId, entry.getKey(), replaceBookiesMap); + } + final LedgerFragment fragment = new LedgerFragment(lh, entry.getKey(), + ledgerFragmentsRange.get(entry.getKey()), replaceBookiesMap.keySet()); + + try { + admin.replicateLedgerFragment(lh, fragment, replaceBookiesMap, + (lId, eId) -> { + // This consumer is already accepted before the method returns + // void. Therefore, use failedTargets in this consumer. + LOG.warn("Failed to read entry {}:{}", lId, eId); + failedTargets.add(Pair.of(ledgerId, entry.getKey())); + }); + LOG.info("Operation finished in the ensemble. ledgerId: {}," + + " fragmentIndex: {}, replaceBookiesMap {}", + ledgerId, entry.getKey(), replaceBookiesMap); + } catch (BKException | InterruptedException e) { + LOG.warn("Failed to replicate ledger fragment.", e); + failedTargets.add(Pair.of(ledgerId, entry.getKey())); + } + } + } + } catch (UnsupportedOperationException e) { + LOG.warn("UnsupportedOperationException caught. The placement policy might not support" + + " replaceToAdherePlacementPolicy method.", e); + failedTargets.add(Pair.of(ledgerId, entry.getKey())); + } + } else { + if (LOG.isDebugEnabled()) { + LOG.debug("The fragment is adhering to placement policy. So, skip the operation." + + " ledgerId: {}, fragmentIndex: {}", ledgerId, entry.getKey()); + } + } + } + } finally { + try { + if (!flags.dryRun) { + lum.releaseUnderreplicatedLedger(ledgerId); + } + } catch (ReplicationException e) { + LOG.warn("Failed to release under replicated ledger {}.", ledgerId, e); + } finally { + ((CountDownLatch) ctx).countDown(); + } + } + }, latch); + } + + // Currently, don't add timeout + latch.await(); + + if (unAcquirableLedgers.isEmpty() && failedTargets.isEmpty()) { + return true; + } else { + LOG.warn("Some ensembles couldn't be relocated to adhere placement policy." + + " Un-acquirable ledgers: {}" + + ", Failed targets [(ledgerId, fragmentIndex), ...]: {}", unAcquirableLedgers, failedTargets); + return false; + } + } + + private static MetadataBookieDriver instantiateMetadataDriver(ServerConfiguration conf) + throws BookieException { + try { + final String metadataServiceUriStr = conf.getMetadataServiceUri(); + final MetadataBookieDriver driver = MetadataDrivers.getBookieDriver(URI.create(metadataServiceUriStr)); + driver.initialize(conf, NullStatsLogger.INSTANCE); + return driver; + } catch (MetadataException me) { + throw new BookieException.MetadataStoreException("Failed to initialize metadata bookie driver", me); + } catch (ConfigurationException e) { + throw new BookieException.BookieIllegalOpException(e); + } + } +} From 4ab4f85cd2d8a12e0c81b1d7d8dd5ec63d54a000 Mon Sep 17 00:00:00 2001 From: Yuri Mizushima Date: Tue, 7 Dec 2021 14:07:02 +0900 Subject: [PATCH 2/4] test: add test for relocate placement feature --- .../CorrectEnsemblePlacementCmdTest.java | 235 ++++++++++++++++++ .../TestRackawareEnsemblePlacementPolicy.java | 144 +++++++++++ 2 files changed, 379 insertions(+) create mode 100644 bookkeeper-server/src/test/java/org/apache/bookkeeper/client/CorrectEnsemblePlacementCmdTest.java diff --git a/bookkeeper-server/src/test/java/org/apache/bookkeeper/client/CorrectEnsemblePlacementCmdTest.java b/bookkeeper-server/src/test/java/org/apache/bookkeeper/client/CorrectEnsemblePlacementCmdTest.java new file mode 100644 index 00000000000..75ea9613604 --- /dev/null +++ b/bookkeeper-server/src/test/java/org/apache/bookkeeper/client/CorrectEnsemblePlacementCmdTest.java @@ -0,0 +1,235 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.bookkeeper.client; + +import static junit.framework.TestCase.assertEquals; +import static org.apache.bookkeeper.client.RackawareEnsemblePlacementPolicyImpl.REPP_DNS_RESOLVER_CLASS; +import static org.mockito.ArgumentMatchers.eq; +import java.lang.reflect.Field; +import java.util.Arrays; +import java.util.Collections; +import java.util.HashMap; +import java.util.HashSet; +import java.util.List; +import java.util.Map; +import java.util.NavigableMap; +import java.util.Set; +import java.util.TreeMap; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicInteger; +import lombok.Cleanup; +import org.apache.bookkeeper.bookie.BookieShell; +import org.apache.bookkeeper.bookie.storage.ldb.DbLedgerStorage; +import org.apache.bookkeeper.client.api.LedgerMetadata; +import org.apache.bookkeeper.net.BookieId; +import org.apache.bookkeeper.net.BookieSocketAddress; +import org.apache.bookkeeper.net.NetworkTopology; +import org.apache.bookkeeper.proto.BookieAddressResolver; +import org.apache.bookkeeper.test.BookKeeperClusterTestCase; +import org.apache.bookkeeper.tools.cli.commands.bookies.CorrectEnsemblePlacementCommand; +import org.apache.bookkeeper.util.EntryFormatter; +import org.apache.bookkeeper.util.LedgerIdFormatter; +import org.apache.bookkeeper.util.StaticDNSResolver; +import org.junit.Test; +import org.mockito.Mockito; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * Tests of correct-ensemble-placement command. + */ +public class CorrectEnsemblePlacementCmdTest extends BookKeeperClusterTestCase { + + private static final Logger LOG = LoggerFactory.getLogger(CorrectEnsemblePlacementCmdTest.class); + private BookKeeper.DigestType digestType = BookKeeper.DigestType.CRC32; + private static final String PASSWORD = "testPasswd"; + + public CorrectEnsemblePlacementCmdTest() throws Exception { + super(0); + baseConf.setLedgerStorageClass(DbLedgerStorage.class.getName()); + baseConf.setGcWaitTime(60000); + baseConf.setFlushInterval(1); + } + + /** + * list of entry logger files that contains given ledgerId. + */ + @Test + public void testArgument() throws Exception { + startNewBookie(); + + @Cleanup + final BookKeeper bk = new BookKeeper(baseClientConf, zkc); + @Cleanup + final LedgerHandle lh = createLedgerWithEntries(bk, 10, 1, 1); + + final String[] argv1 = { "correct-ensemble-placement", "--ledgerids", String.valueOf(lh.getId()), + "--skipOpenLedgers", "--force"}; + final String[] argv2 = { "correct-ensemble-placement", "--ledgerids", String.valueOf(lh.getId()), + "--skipOpenLedgers", "--force", "--dryrun"}; + final BookieShell bkShell = + new BookieShell(LedgerIdFormatter.LONG_LEDGERID_FORMATTER, EntryFormatter.STRING_FORMATTER); + bkShell.setConf(baseClientConf); + + assertEquals("Failed to return exit code!", 0, bkShell.run(argv1)); + assertEquals("Failed to return exit code!", 0, bkShell.run(argv2)); + } + + @Test + public void testCorrectEnsemblePlacementByRackaware() throws Exception { + final int ensembleSize = 7; + final int quorumSize = 2; + + final BookieSocketAddress addr1 = new BookieSocketAddress("127.0.0.1", 3181); + final BookieSocketAddress addr2 = new BookieSocketAddress("127.0.0.2", 3181); + final BookieSocketAddress addr3 = new BookieSocketAddress("127.0.0.3", 3181); + final BookieSocketAddress addr4 = new BookieSocketAddress("127.0.0.4", 3181); + final BookieSocketAddress addr5 = new BookieSocketAddress("127.0.0.5", 3181); + final BookieSocketAddress addr6 = new BookieSocketAddress("127.0.0.6", 3181); + final BookieSocketAddress addr7 = new BookieSocketAddress("127.0.0.7", 3181); + final BookieSocketAddress addr8 = new BookieSocketAddress("127.0.0.8", 3181); + final BookieSocketAddress addr9 = new BookieSocketAddress("127.0.0.9", 3181); + + final Set writableBookies = new HashSet<>(); + writableBookies.add(addr1.toBookieId()); + writableBookies.add(addr2.toBookieId()); + writableBookies.add(addr3.toBookieId()); + writableBookies.add(addr4.toBookieId()); + writableBookies.add(addr5.toBookieId()); + writableBookies.add(addr6.toBookieId()); + writableBookies.add(addr7.toBookieId()); + writableBookies.add(addr8.toBookieId()); + writableBookies.add(addr9.toBookieId()); + + // add bookie node to resolver + StaticDNSResolver.reset(); + + final String rackName1 = NetworkTopology.DEFAULT_REGION + "/r1"; + final String rackName2 = NetworkTopology.DEFAULT_REGION + "/r2"; + final String rackName3 = NetworkTopology.DEFAULT_REGION + "/r3"; + + // update dns mapping + // add port for testing + StaticDNSResolver.addNodeToRack(addr1.getSocketAddress().getAddress().getHostAddress(), rackName1); + StaticDNSResolver.addNodeToRack(addr2.getSocketAddress().getAddress().getHostAddress(), rackName1); + StaticDNSResolver.addNodeToRack(addr3.getSocketAddress().getAddress().getHostAddress(), rackName1); + StaticDNSResolver.addNodeToRack(addr4.getSocketAddress().getAddress().getHostAddress(), rackName2); + StaticDNSResolver.addNodeToRack(addr5.getSocketAddress().getAddress().getHostAddress(), rackName2); + StaticDNSResolver.addNodeToRack(addr6.getSocketAddress().getAddress().getHostAddress(), rackName2); + StaticDNSResolver.addNodeToRack(addr7.getSocketAddress().getAddress().getHostAddress(), rackName3); + StaticDNSResolver.addNodeToRack(addr8.getSocketAddress().getAddress().getHostAddress(), rackName3); + StaticDNSResolver.addNodeToRack(addr9.getSocketAddress().getAddress().getHostAddress(), rackName3); + LOG.info("Set up static DNS Resolver."); + baseClientConf.setEnsemblePlacementPolicy(RackawareEnsemblePlacementPolicy.class); + baseClientConf.setProperty(REPP_DNS_RESOLVER_CLASS, StaticDNSResolver.class.getName()); + + startNewBookie(); + + final NavigableMap> ensemble = new TreeMap<>(); + // create failed ensemble + // expect that the ensemble will be replaced to + // [addr1, addr4, addr7, addr2, addr5, addr8, *addr6 (in /default-region/r2 bookies)*] + ensemble.put(0L, Arrays.asList( + addr1.toBookieId(), addr4.toBookieId(), + addr7.toBookieId(), addr2.toBookieId(), + addr5.toBookieId(), addr8.toBookieId(), + addr3.toBookieId())); + final LedgerMetadata lm1 = Mockito.mock(LedgerMetadata.class); + Mockito.doReturn(ensemble).when(lm1).getAllEnsembles(); + Mockito.doReturn(ensembleSize).when(lm1).getEnsembleSize(); + Mockito.doReturn(quorumSize).when(lm1).getWriteQuorumSize(); + Mockito.doReturn(quorumSize).when(lm1).getAckQuorumSize(); + final LedgerHandle lh1 = Mockito.mock(LedgerHandle.class); + Mockito.when(lh1.getLedgerMetadata()).thenReturn(lm1); + + @Cleanup + final BookKeeper bookKeeper = Mockito.spy(new BookKeeper(baseClientConf)); + Mockito.doReturn(true).when(bookKeeper).isClosed(Mockito.anyLong()); + + @Cleanup + final BookKeeperAdmin admin = Mockito.spy(new BookKeeperAdmin(baseClientConf)); + Mockito.doAnswer(invocationOnMock -> { + final AsyncCallback.OpenCallback op = invocationOnMock.getArgument(1); + final CountDownLatch ctx = invocationOnMock.getArgument(2); + op.openComplete(BKException.Code.OK, lh1, ctx); + return null; + }).when(admin).asyncOpenLedger(Mockito.anyLong(), Mockito.any(), Mockito.any()); + // expected return + final Map expectedMap = new HashMap<>(); + expectedMap.put(6, addr6.toBookieId()); + Mockito.doNothing().when(admin) + .replicateLedgerFragment(Mockito.any(), Mockito.any(), eq(expectedMap), Mockito.any()); + + final EnsemblePlacementPolicy policy = bookKeeper.getPlacementPolicy(); + final Field topologyField = TopologyAwareEnsemblePlacementPolicy.class + .getDeclaredField("topology"); + topologyField.setAccessible(true); + final NetworkTopology topology = Mockito.spy((NetworkTopology) topologyField.get(policy)); + Mockito.doReturn(2).when(topology).getNumOfRacks(); + topologyField.set(policy, topology); + final Field bookieAddressResolverField = TopologyAwareEnsemblePlacementPolicy.class + .getDeclaredField("bookieAddressResolver"); + bookieAddressResolverField.setAccessible(true); + final BookieAddressResolver bookieAddressResolver = + Mockito.spy((BookieAddressResolver) bookieAddressResolverField.get(policy)); + Mockito.doAnswer(invocationOnMock -> { + final BookieId bookieId = invocationOnMock.getArgument(0); + return new BookieSocketAddress(bookieId.getId()); + }).when(bookieAddressResolver).resolve(Mockito.any()); + bookieAddressResolverField.set(policy, bookieAddressResolver); + + // add mock bookies to knownBookies + policy.onClusterChanged(writableBookies, Collections.emptySet()); + + // make sure that the ensemble is FAIL state + assertEquals(EnsemblePlacementPolicy.PlacementPolicyAdherence.FAIL, + policy.isEnsembleAdheringToPlacementPolicy(ensemble.get(0L), quorumSize, quorumSize)); + + final CorrectEnsemblePlacementCommand.CorrectEnsemblePlacementFlags flag = + new CorrectEnsemblePlacementCommand.CorrectEnsemblePlacementFlags(); + flag.ledgerIds(Collections.singletonList(1L)); + final CorrectEnsemblePlacementCommand cmd = new CorrectEnsemblePlacementCommand(); + cmd.relocate(baseConf, flag, bookKeeper, admin); + + Mockito.verify(admin, Mockito.times(1)) + .replicateLedgerFragment(Mockito.any(), Mockito.any(), eq(expectedMap), Mockito.any()); + } + + private LedgerHandle createLedgerWithEntries(BookKeeper bk, int numOfEntries, + int ensembleSize, int quorumSize) throws Exception { + LedgerHandle lh = bk.createLedger(ensembleSize, quorumSize, digestType, PASSWORD.getBytes()); + final AtomicInteger rc = new AtomicInteger(BKException.Code.OK); + final CountDownLatch latch = new CountDownLatch(numOfEntries); + + final AsyncCallback.AddCallback cb = (rccb, lh1, entryId, ctx) -> { + rc.compareAndSet(BKException.Code.OK, rccb); + latch.countDown(); + }; + for (int i = 0; i < numOfEntries; i++) { + lh.asyncAddEntry(("foobar" + i).getBytes(), cb, null); + } + if (!latch.await(30, TimeUnit.SECONDS)) { + throw new Exception("Entries took too long to add"); + } + if (rc.get() != BKException.Code.OK) { + throw BKException.create(rc.get()); + } + return lh; + } +} diff --git a/bookkeeper-server/src/test/java/org/apache/bookkeeper/client/TestRackawareEnsemblePlacementPolicy.java b/bookkeeper-server/src/test/java/org/apache/bookkeeper/client/TestRackawareEnsemblePlacementPolicy.java index 28a27ee611b..37f7b292060 100644 --- a/bookkeeper-server/src/test/java/org/apache/bookkeeper/client/TestRackawareEnsemblePlacementPolicy.java +++ b/bookkeeper-server/src/test/java/org/apache/bookkeeper/client/TestRackawareEnsemblePlacementPolicy.java @@ -21,6 +21,9 @@ import static org.apache.bookkeeper.client.RackawareEnsemblePlacementPolicyImpl.shuffleWithMask; import static org.apache.bookkeeper.client.RoundRobinDistributionSchedule.writeSetFromValues; import static org.apache.bookkeeper.feature.SettableFeatureProvider.DISABLE_ALL; +import static org.hamcrest.Matchers.contains; +import static org.hamcrest.Matchers.is; +import static org.junit.Assert.assertThat; import com.google.common.util.concurrent.ThreadFactoryBuilder; import io.netty.util.HashedWheelTimer; @@ -36,6 +39,7 @@ import java.util.Optional; import java.util.Set; import java.util.concurrent.TimeUnit; +import java.util.function.Consumer; import junit.framework.TestCase; import org.apache.bookkeeper.client.BKException.BKNotEnoughBookiesException; import org.apache.bookkeeper.client.BookieInfoReader.BookieInfo; @@ -60,6 +64,10 @@ import org.apache.bookkeeper.test.TestStatsProvider.TestStatsLogger; import org.apache.bookkeeper.util.StaticDNSResolver; import org.apache.commons.collections4.CollectionUtils; +import org.apache.commons.lang3.tuple.Pair; +import org.hamcrest.Description; +import org.hamcrest.Matcher; +import org.hamcrest.TypeSafeMatcher; import org.junit.Test; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -2387,4 +2395,140 @@ public void testAreAckedBookiesAdheringToPlacementPolicy() throws Exception { testAreAckedBookiesAdheringToPlacementPolicyHelper(4, 6, 3, 2, 6, 3, 3); testAreAckedBookiesAdheringToPlacementPolicyHelper(5, 7, 5, 3, 7, 5, 2); } + + @SuppressWarnings("unchecked") + @Test + public void testReplaceToAdherePlacementPolicy() throws Exception { + final BookieSocketAddress addr1 = new BookieSocketAddress("127.0.0.1", 3181); + final BookieSocketAddress addr2 = new BookieSocketAddress("127.0.0.2", 3181); + final BookieSocketAddress addr3 = new BookieSocketAddress("127.0.0.3", 3181); + final BookieSocketAddress addr4 = new BookieSocketAddress("127.0.0.4", 3181); + final BookieSocketAddress addr5 = new BookieSocketAddress("127.0.0.5", 3181); + final BookieSocketAddress addr6 = new BookieSocketAddress("127.0.0.6", 3181); + final BookieSocketAddress addr7 = new BookieSocketAddress("127.0.0.7", 3181); + final BookieSocketAddress addr8 = new BookieSocketAddress("127.0.0.8", 3181); + final BookieSocketAddress addr9 = new BookieSocketAddress("127.0.0.9", 3181); + + final String rackName1 = NetworkTopology.DEFAULT_REGION + "/r1"; + final String rackName2 = NetworkTopology.DEFAULT_REGION + "/r2"; + final String rackName3 = NetworkTopology.DEFAULT_REGION + "/r3"; + + // update dns mapping + StaticDNSResolver.addNodeToRack(addr1.getSocketAddress().getAddress().getHostAddress(), rackName1); + StaticDNSResolver.addNodeToRack(addr2.getSocketAddress().getAddress().getHostAddress(), rackName1); + StaticDNSResolver.addNodeToRack(addr3.getSocketAddress().getAddress().getHostAddress(), rackName1); + StaticDNSResolver.addNodeToRack(addr4.getSocketAddress().getAddress().getHostAddress(), rackName2); + StaticDNSResolver.addNodeToRack(addr5.getSocketAddress().getAddress().getHostAddress(), rackName2); + StaticDNSResolver.addNodeToRack(addr6.getSocketAddress().getAddress().getHostAddress(), rackName2); + StaticDNSResolver.addNodeToRack(addr7.getSocketAddress().getAddress().getHostAddress(), rackName3); + StaticDNSResolver.addNodeToRack(addr8.getSocketAddress().getAddress().getHostAddress(), rackName3); + StaticDNSResolver.addNodeToRack(addr9.getSocketAddress().getAddress().getHostAddress(), rackName3); + + // Update cluster + final Set addrs = new HashSet<>(); + addrs.add(addr1.toBookieId()); + addrs.add(addr2.toBookieId()); + addrs.add(addr3.toBookieId()); + addrs.add(addr4.toBookieId()); + addrs.add(addr5.toBookieId()); + addrs.add(addr6.toBookieId()); + addrs.add(addr7.toBookieId()); + addrs.add(addr8.toBookieId()); + addrs.add(addr9.toBookieId()); + + final ClientConfiguration newConf = new ClientConfiguration(conf); + newConf.setDiskWeightBasedPlacementEnabled(false); + newConf.setMinNumRacksPerWriteQuorum(2); + newConf.setEnforceMinNumRacksPerWriteQuorum(true); + + repp.initialize(newConf, Optional.empty(), timer, + DISABLE_ALL, NullStatsLogger.INSTANCE, BookieSocketAddress.LEGACY_BOOKIEID_RESOLVER); + repp.withDefaultRack(NetworkTopology.DEFAULT_REGION_AND_RACK); + + repp.onClusterChanged(addrs, new HashSet<>()); + final Map bookieInfoMap = new HashMap<>(); + bookieInfoMap.put(addr1.toBookieId(), new BookieInfo(100L, 100L)); + bookieInfoMap.put(addr2.toBookieId(), new BookieInfo(100L, 100L)); + bookieInfoMap.put(addr3.toBookieId(), new BookieInfo(100L, 100L)); + bookieInfoMap.put(addr4.toBookieId(), new BookieInfo(100L, 100L)); + bookieInfoMap.put(addr5.toBookieId(), new BookieInfo(100L, 100L)); + bookieInfoMap.put(addr6.toBookieId(), new BookieInfo(100L, 100L)); + bookieInfoMap.put(addr7.toBookieId(), new BookieInfo(100L, 100L)); + bookieInfoMap.put(addr8.toBookieId(), new BookieInfo(100L, 100L)); + bookieInfoMap.put(addr9.toBookieId(), new BookieInfo(100L, 100L)); + + repp.updateBookieInfo(bookieInfoMap); + + final Set excludeList = new HashSet<>(); + final int ensembleSize = 7; + final int writeQuorumSize = 2; + final int ackQuorumSize = 2; + + class BookieRackMatcher extends TypeSafeMatcher { + final List expectedRacks; + + public BookieRackMatcher(String... expectedRacks) { + this.expectedRacks = Arrays.asList(expectedRacks); + } + + @Override + protected boolean matchesSafely(BookieId bookieId) { + return expectedRacks.contains(StaticDNSResolver.getRack(bookieId.toString().split(":")[0])); + } + + @Override + public void describeTo(Description description) { + description.appendText("expected racks " + expectedRacks); + } + } + + final BookieRackMatcher rack1 = new BookieRackMatcher(rackName1); + final BookieRackMatcher rack2 = new BookieRackMatcher(rackName2); + final BookieRackMatcher rack3 = new BookieRackMatcher(rackName3); + final BookieRackMatcher rack12 = new BookieRackMatcher(rackName1, rackName2); + final BookieRackMatcher rack13 = new BookieRackMatcher(rackName1, rackName3); + final BookieRackMatcher rack23 = new BookieRackMatcher(rackName2, rackName3); + final BookieRackMatcher rack123 = new BookieRackMatcher(rackName1, rackName2, rackName3); + final Consumer, Matcher>>> test = (pair) -> { + // RackawareEnsemblePlacementPolicyImpl#isEnsembleAdheringToPlacementPolicy + // is not scope of this test case. So, use the method in assertion for convenience. + assertEquals(PlacementPolicyAdherence.FAIL, + repp.isEnsembleAdheringToPlacementPolicy(pair.getLeft(), writeQuorumSize, ackQuorumSize)); + final EnsemblePlacementPolicy.PlacementResult> result = + repp.replaceToAdherePlacementPolicy(ensembleSize, writeQuorumSize, ackQuorumSize, + excludeList, pair.getLeft()); + if (LOG.isDebugEnabled()) { + LOG.debug("input: {}, result: {}", pair.getLeft(), result.getResult()); + } + assertEquals(PlacementPolicyAdherence.MEETS_STRICT, result.isAdheringToPolicy()); + assertThat(result.getResult(), pair.getRight()); + }; + + for (int i = 0; i < 1000; i++) { + test.accept(Pair.of(Arrays.asList(addr1.toBookieId(), addr4.toBookieId(), addr7.toBookieId(), + addr2.toBookieId(), addr5.toBookieId(), addr8.toBookieId(), addr9.toBookieId()), + // first, same, same, same, same, same, condition[0] + contains(is(addr1.toBookieId()), is(addr4.toBookieId()), is(addr7.toBookieId()), + is(addr2.toBookieId()), is(addr5.toBookieId()), is(addr8.toBookieId()), + is(addr6.toBookieId())))); + + test.accept(Pair.of(Arrays.asList(addr6.toBookieId(), addr4.toBookieId(), addr7.toBookieId(), + addr2.toBookieId(), addr5.toBookieId(), addr8.toBookieId(), addr3.toBookieId()), + // first, condition[0], same, same, same, same, same + contains(is(addr6.toBookieId()), is(addr1.toBookieId()), is(addr7.toBookieId()), + is(addr2.toBookieId()), is(addr5.toBookieId()), is(addr8.toBookieId()), + is(addr3.toBookieId())))); + + test.accept(Pair.of(Arrays.asList(addr1.toBookieId(), addr2.toBookieId(), addr3.toBookieId(), + addr4.toBookieId(), addr5.toBookieId(), addr6.toBookieId(), addr7.toBookieId()), + // first, candidate[0], same, same, candidate[0], same, same + contains(is(addr1.toBookieId()), is(rack3), is(addr3.toBookieId()), + is(addr4.toBookieId()), is(rack13), is(addr6.toBookieId()), is(addr7.toBookieId())))); + + test.accept(Pair.of(Arrays.asList(addr1.toBookieId(), addr2.toBookieId(), addr4.toBookieId(), + addr5.toBookieId(), addr7.toBookieId(), addr8.toBookieId(), addr9.toBookieId()), + contains(is(addr1.toBookieId()), is(rack23), is(rack123), is(rack123), + is(rack123), is(rack123), is(rack23)))); + } + } } From e6e4deb9ebdc3f1f5ad082965cf0e65cb7537c8f Mon Sep 17 00:00:00 2001 From: Yuri Mizushima Date: Thu, 17 Mar 2022 16:55:13 +0900 Subject: [PATCH 3/4] refactor: modify addNode call to outside of replaceToAdherePlacementPolicyInternal --- .../RackawareEnsemblePlacementPolicyImpl.java | 16 +++++++--------- .../client/CorrectEnsemblePlacementCmdTest.java | 7 +++++++ 2 files changed, 14 insertions(+), 9 deletions(-) diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/client/RackawareEnsemblePlacementPolicyImpl.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/client/RackawareEnsemblePlacementPolicyImpl.java index d9167c321d4..fa3e59cd65f 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/client/RackawareEnsemblePlacementPolicyImpl.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/client/RackawareEnsemblePlacementPolicyImpl.java @@ -1131,6 +1131,13 @@ public PlacementResult> replaceToAdherePlacementPolicy( prevNode = replaceToAdherePlacementPolicyInternal( curRack, excludeNodes, ensemble, ensemble, provisionalEnsembleNodes, i, ensembleSize, minNumRacksPerWriteQuorumForThisEnsemble); + // got a good candidate + if (ensemble.addNode(prevNode)) { + // add the candidate to exclude set + excludeNodes.add(prevNode); + } else { + throw new BKNotEnoughBookiesException(); + } // replace to newer node provisionalEnsembleNodes.set(i, prevNode); } catch (BKNotEnoughBookiesException e) { @@ -1159,10 +1166,6 @@ private BookieNode replaceToAdherePlacementPolicyInternal( final BookieNode currentNode = provisionalEnsembleNodes.get(ensembleIndex); // if the current bookie could be applied to the ensemble, apply it to minify the number of bookies replaced if (!excludeBookies.contains(currentNode) && predicate.apply(currentNode, ensemble)) { - if (ensemble.addNode(currentNode)) { - // add the candidate to exclude set - excludeBookies.add(currentNode); - } return currentNode; } @@ -1234,11 +1237,6 @@ private BookieNode replaceToAdherePlacementPolicyInternal( continue; } BookieNode bn = (BookieNode) n; - // got a good candidate - if (ensemble.addNode(bn)) { - // add the candidate to exclude set - excludeBookies.add(bn); - } return bn; } } diff --git a/bookkeeper-server/src/test/java/org/apache/bookkeeper/client/CorrectEnsemblePlacementCmdTest.java b/bookkeeper-server/src/test/java/org/apache/bookkeeper/client/CorrectEnsemblePlacementCmdTest.java index 75ea9613604..aad2b2d56ab 100644 --- a/bookkeeper-server/src/test/java/org/apache/bookkeeper/client/CorrectEnsemblePlacementCmdTest.java +++ b/bookkeeper-server/src/test/java/org/apache/bookkeeper/client/CorrectEnsemblePlacementCmdTest.java @@ -46,6 +46,7 @@ import org.apache.bookkeeper.util.EntryFormatter; import org.apache.bookkeeper.util.LedgerIdFormatter; import org.apache.bookkeeper.util.StaticDNSResolver; +import org.junit.After; import org.junit.Test; import org.mockito.Mockito; import org.slf4j.Logger; @@ -67,6 +68,12 @@ public CorrectEnsemblePlacementCmdTest() throws Exception { baseConf.setFlushInterval(1); } + @After + public void tearDown() throws Exception { + super.tearDown(); + StaticDNSResolver.reset(); + } + /** * list of entry logger files that contains given ledgerId. */ From 62970b3122fa1f8266763f0d3f6115a066fcabe9 Mon Sep 17 00:00:00 2001 From: Yuri Mizushima Date: Mon, 23 May 2022 16:13:16 +0900 Subject: [PATCH 4/4] fix: move relocation procecure to BookKeeperAdmin --- .../bookkeeper/client/BookKeeperAdmin.java | 106 ++++++++++++ .../CorrectEnsemblePlacementCommand.java | 115 ++----------- .../client/BookKeeperAdminTest.java | 139 +++++++++++++++ .../CorrectEnsemblePlacementCmdTest.java | 161 +----------------- 4 files changed, 261 insertions(+), 260 deletions(-) diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/client/BookKeeperAdmin.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/client/BookKeeperAdmin.java index 55c1ccc3ff7..1742f7ccdb9 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/client/BookKeeperAdmin.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/client/BookKeeperAdmin.java @@ -23,6 +23,7 @@ import static com.google.common.base.Preconditions.checkArgument; import static org.apache.bookkeeper.meta.MetadataDrivers.runFunctionWithMetadataBookieDriver; import static org.apache.bookkeeper.meta.MetadataDrivers.runFunctionWithRegistrationManager; +import com.google.common.base.Functions; import com.google.common.collect.Lists; import com.google.common.collect.Maps; import com.google.common.collect.Sets; @@ -31,6 +32,7 @@ import java.io.IOException; import java.util.ArrayList; import java.util.Collection; +import java.util.Collections; import java.util.Enumeration; import java.util.HashMap; import java.util.Iterator; @@ -50,6 +52,8 @@ import java.util.concurrent.atomic.AtomicReference; import java.util.function.BiConsumer; import java.util.function.Predicate; +import java.util.stream.Collectors; +import java.util.stream.IntStream; import lombok.SneakyThrows; import org.apache.bookkeeper.bookie.BookieException; import org.apache.bookkeeper.bookie.BookieImpl; @@ -1225,6 +1229,108 @@ public void processResult(int rc, String s, Object ctx) { } } + /** + * + * @param lh Ledger Handle + * @param dryRun if true, run it without any modification. + * @return failed ledger fragment indices + * @throws UnsupportedOperationException Default behavior of + * {@link EnsemblePlacementPolicy#replaceToAdherePlacementPolicy(int, int, int, java.util.Set, java.util.List)}. + */ + public List relocateLedgerToAdherePlacementPolicy(LedgerHandle lh, boolean dryRun) + throws UnsupportedOperationException { + final EnsemblePlacementPolicy placementPolicy = bkc.getPlacementPolicy(); + + final long ledgerId = lh.getId(); + final LedgerMetadata ledgerMeta = lh.getLedgerMetadata(); + final List failedFragmentIndexList = new ArrayList<>(); + final Map ledgerFragmentsRange = new HashMap<>(); + Long curEntryId = null; + for (Map.Entry> entry : + ledgerMeta.getAllEnsembles().entrySet()) { + if (curEntryId != null) { + ledgerFragmentsRange.put(curEntryId, entry.getKey() - 1); + } + curEntryId = entry.getKey(); + } + if (curEntryId != null) { + ledgerFragmentsRange.put(curEntryId, lh.getLastAddConfirmed()); + } + + for (Map.Entry> entry : ledgerMeta.getAllEnsembles().entrySet()) { + if (placementPolicy.isEnsembleAdheringToPlacementPolicy(entry.getValue(), + ledgerMeta.getWriteQuorumSize(), ledgerMeta.getAckQuorumSize()) + == EnsemblePlacementPolicy.PlacementPolicyAdherence.FAIL) { + final List currentEnsemble = entry.getValue(); + // Currently, don't consider quarantinedBookies + final EnsemblePlacementPolicy.PlacementResult> placementResult = + placementPolicy.replaceToAdherePlacementPolicy( + ledgerMeta.getEnsembleSize(), + ledgerMeta.getWriteQuorumSize(), + ledgerMeta.getAckQuorumSize(), + Collections.emptySet(), + currentEnsemble); + + if (placementResult.isAdheringToPolicy() + == EnsemblePlacementPolicy.PlacementPolicyAdherence.FAIL) { + LOG.warn("Failed to relocate the ensemble. So, skip the operation." + + " ledgerId: {}, fragmentIndex: {}", + ledgerId, entry.getKey()); + failedFragmentIndexList.add(entry.getKey()); + } else { + final List newEnsemble = placementResult.getResult(); + final Map replaceBookiesMap = IntStream + .range(0, ledgerMeta.getEnsembleSize()).boxed() + .filter(i -> !newEnsemble.get(i).equals(currentEnsemble.get(i))) + .collect(Collectors.toMap(Functions.identity(), newEnsemble::get)); + if (replaceBookiesMap.isEmpty()) { + LOG.warn("Failed to get bookies to replace. So, skip the operation." + + " ledgerId: {}, fragmentIndex: {}", + ledgerId, entry.getKey()); + failedFragmentIndexList.add(entry.getKey()); + } else if (dryRun) { + LOG.info("Would replace the ensemble. ledgerId: {}, fragmentIndex: {}," + + " currentEnsemble: {} replaceBookiesMap {}", + ledgerId, entry.getKey(), + currentEnsemble, replaceBookiesMap); + } else { + if (LOG.isDebugEnabled()) { + LOG.debug("Try to replace the ensemble. ledgerId: {}, fragmentIndex: {}," + + " replaceBookiesMap {}", + ledgerId, entry.getKey(), replaceBookiesMap); + } + final LedgerFragment fragment = new LedgerFragment(lh, entry.getKey(), + ledgerFragmentsRange.get(entry.getKey()), replaceBookiesMap.keySet()); + + try { + replicateLedgerFragment(lh, fragment, replaceBookiesMap, + (lId, eId) -> { + // This consumer is already accepted before the method returns + // void. Therefore, use failedFragmentIndexList in this consumer. + LOG.warn("Failed to read entry {}:{}", lId, eId); + failedFragmentIndexList.add(entry.getKey()); + }); + if (LOG.isDebugEnabled()) { + LOG.debug("Operation finished in the ensemble. ledgerId: {}," + + " fragmentIndex: {}, replaceBookiesMap {}", + ledgerId, entry.getKey(), replaceBookiesMap); + } + } catch (BKException | InterruptedException e) { + LOG.warn("Failed to replicate ledger fragment.", e); + failedFragmentIndexList.add(entry.getKey()); + } + } + } + } else { + if (LOG.isDebugEnabled()) { + LOG.debug("The fragment is adhering to placement policy. So, skip the operation." + + " ledgerId: {}, fragmentIndex: {}", ledgerId, entry.getKey()); + } + } + } + return failedFragmentIndexList; + } + /** * Format the BookKeeper metadata in zookeeper. * diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/tools/cli/commands/bookies/CorrectEnsemblePlacementCommand.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/tools/cli/commands/bookies/CorrectEnsemblePlacementCommand.java index 36cddb7f11f..fad54b51a2c 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/tools/cli/commands/bookies/CorrectEnsemblePlacementCommand.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/tools/cli/commands/bookies/CorrectEnsemblePlacementCommand.java @@ -20,20 +20,15 @@ import com.beust.jcommander.Parameter; import com.beust.jcommander.converters.CommaParameterSplitter; import com.google.common.annotations.VisibleForTesting; -import com.google.common.base.Functions; import com.google.common.util.concurrent.UncheckedExecutionException; import java.io.IOException; import java.net.URI; -import java.util.Collections; -import java.util.HashMap; import java.util.List; -import java.util.Map; import java.util.NavigableSet; import java.util.TreeSet; import java.util.concurrent.ConcurrentSkipListSet; import java.util.concurrent.CountDownLatch; import java.util.stream.Collectors; -import java.util.stream.IntStream; import lombok.Cleanup; import lombok.Setter; import lombok.experimental.Accessors; @@ -41,9 +36,6 @@ import org.apache.bookkeeper.client.BKException; import org.apache.bookkeeper.client.BookKeeper; import org.apache.bookkeeper.client.BookKeeperAdmin; -import org.apache.bookkeeper.client.EnsemblePlacementPolicy; -import org.apache.bookkeeper.client.LedgerFragment; -import org.apache.bookkeeper.client.api.LedgerMetadata; import org.apache.bookkeeper.conf.ClientConfiguration; import org.apache.bookkeeper.conf.ServerConfiguration; import org.apache.bookkeeper.meta.LedgerManagerFactory; @@ -51,7 +43,6 @@ import org.apache.bookkeeper.meta.MetadataBookieDriver; import org.apache.bookkeeper.meta.MetadataDrivers; import org.apache.bookkeeper.meta.exceptions.MetadataException; -import org.apache.bookkeeper.net.BookieId; import org.apache.bookkeeper.replication.ReplicationException; import org.apache.bookkeeper.stats.NullStatsLogger; import org.apache.bookkeeper.tools.cli.helpers.BookieCommand; @@ -133,7 +124,7 @@ public boolean apply(ServerConfiguration conf, CorrectEnsemblePlacementFlags fla final ClientConfiguration clientConf = new ClientConfiguration(conf); final BookKeeper bookKeeper = new BookKeeper(clientConf); - final BookKeeperAdmin admin = new BookKeeperAdmin(clientConf); + final BookKeeperAdmin admin = new BookKeeperAdmin(bookKeeper); return relocate(conf, flags, bookKeeper, admin); } catch (Exception e) { throw new UncheckedExecutionException(e.getMessage(), e); @@ -143,8 +134,6 @@ public boolean apply(ServerConfiguration conf, CorrectEnsemblePlacementFlags fla @VisibleForTesting public boolean relocate(ServerConfiguration conf, CorrectEnsemblePlacementFlags flags, BookKeeper bookKeeper, BookKeeperAdmin admin) throws Exception { - final EnsemblePlacementPolicy placementPolicy = bookKeeper.getPlacementPolicy(); - @Cleanup final MetadataBookieDriver metadataDriver = instantiateMetadataDriver(conf); @Cleanup @@ -153,7 +142,7 @@ public boolean relocate(ServerConfiguration conf, CorrectEnsemblePlacementFlags final LedgerUnderreplicationManager lum = lmf.newLedgerUnderreplicationManager(); final List targetLedgers = - flags.ledgerIds.stream().parallel().distinct().filter(ledgerId -> { + flags.ledgerIds.stream().distinct().filter(ledgerId -> { try { return (!flags.skipOpenLedgers || bookKeeper.isClosed(ledgerId)) && !lum.isLedgerBeingReplicated(ledgerId); @@ -172,6 +161,9 @@ public boolean relocate(ServerConfiguration conf, CorrectEnsemblePlacementFlags final NavigableSet> failedTargets = new ConcurrentSkipListSet<>(); final CountDownLatch latch = new CountDownLatch(targetLedgers.size()); for (long ledgerId : targetLedgers) { + if (LOG.isDebugEnabled()) { + LOG.debug("Start relocation of the ledger {}.", ledgerId); + } if (!flags.dryRun) { try { lum.acquireUnderreplicatedLedger(ledgerId); @@ -188,103 +180,18 @@ public boolean relocate(ServerConfiguration conf, CorrectEnsemblePlacementFlags LOG.warn("Failed to open ledger {}", ledgerId); return; } - - final LedgerMetadata ledgerMeta = lh.getLedgerMetadata(); - final Map ledgerFragmentsRange = new HashMap<>(); - Long curEntryId = null; - for (Map.Entry> entry : - ledgerMeta.getAllEnsembles().entrySet()) { - if (curEntryId != null) { - ledgerFragmentsRange.put(curEntryId, entry.getKey() - 1); - } - curEntryId = entry.getKey(); - } - if (curEntryId != null) { - ledgerFragmentsRange.put(curEntryId, lh.getLastAddConfirmed()); - } - - for (Map.Entry> entry : ledgerMeta.getAllEnsembles().entrySet()) { - if (placementPolicy.isEnsembleAdheringToPlacementPolicy(entry.getValue(), - ledgerMeta.getWriteQuorumSize(), ledgerMeta.getAckQuorumSize()) - == EnsemblePlacementPolicy.PlacementPolicyAdherence.FAIL) { - try { - final List currentEnsemble = entry.getValue(); - // Currently, don't consider quarantinedBookies - final EnsemblePlacementPolicy.PlacementResult> placementResult = - placementPolicy.replaceToAdherePlacementPolicy( - ledgerMeta.getEnsembleSize(), - ledgerMeta.getWriteQuorumSize(), - ledgerMeta.getAckQuorumSize(), - Collections.emptySet(), - currentEnsemble); - - if (placementResult.isAdheringToPolicy() - == EnsemblePlacementPolicy.PlacementPolicyAdherence.FAIL) { - LOG.warn("Failed to relocate the ensemble. So, skip the operation." - + " ledgerId: {}, fragmentIndex: {}", - ledgerId, entry.getKey()); - failedTargets.add(Pair.of(ledgerId, entry.getKey())); - } else { - final List newEnsemble = placementResult.getResult(); - final Map replaceBookiesMap = IntStream - .range(0, ledgerMeta.getEnsembleSize()).boxed() - .filter(i -> !newEnsemble.get(i).equals(currentEnsemble.get(i))) - .collect(Collectors.toMap(Functions.identity(), newEnsemble::get)); - if (replaceBookiesMap.isEmpty()) { - LOG.warn("Failed to get bookies to replace. So, skip the operation." - + " ledgerId: {}, fragmentIndex: {}", - ledgerId, entry.getKey()); - failedTargets.add(Pair.of(ledgerId, entry.getKey())); - } else if (flags.dryRun) { - LOG.info("Would replace the ensemble. ledgerId: {}, fragmentIndex: {}," - + " currentEnsemble: {} replaceBookiesMap {}", - ledgerId, entry.getKey(), - currentEnsemble, replaceBookiesMap); - } else { - if (LOG.isDebugEnabled()) { - LOG.debug("Try to replace the ensemble. ledgerId: {}, fragmentIndex: {}," - + " replaceBookiesMap {}", - ledgerId, entry.getKey(), replaceBookiesMap); - } - final LedgerFragment fragment = new LedgerFragment(lh, entry.getKey(), - ledgerFragmentsRange.get(entry.getKey()), replaceBookiesMap.keySet()); - - try { - admin.replicateLedgerFragment(lh, fragment, replaceBookiesMap, - (lId, eId) -> { - // This consumer is already accepted before the method returns - // void. Therefore, use failedTargets in this consumer. - LOG.warn("Failed to read entry {}:{}", lId, eId); - failedTargets.add(Pair.of(ledgerId, entry.getKey())); - }); - LOG.info("Operation finished in the ensemble. ledgerId: {}," - + " fragmentIndex: {}, replaceBookiesMap {}", - ledgerId, entry.getKey(), replaceBookiesMap); - } catch (BKException | InterruptedException e) { - LOG.warn("Failed to replicate ledger fragment.", e); - failedTargets.add(Pair.of(ledgerId, entry.getKey())); - } - } - } - } catch (UnsupportedOperationException e) { - LOG.warn("UnsupportedOperationException caught. The placement policy might not support" - + " replaceToAdherePlacementPolicy method.", e); - failedTargets.add(Pair.of(ledgerId, entry.getKey())); - } - } else { - if (LOG.isDebugEnabled()) { - LOG.debug("The fragment is adhering to placement policy. So, skip the operation." - + " ledgerId: {}, fragmentIndex: {}", ledgerId, entry.getKey()); - } - } - } + admin.relocateLedgerToAdherePlacementPolicy(lh, flags.dryRun) + .forEach(e -> failedTargets.add(Pair.of(ledgerId, e))); + } catch (UnsupportedOperationException e) { + LOG.warn("UnsupportedOperationException caught. The placement policy might not support" + + " replaceToAdherePlacementPolicy method.", e); } finally { try { if (!flags.dryRun) { lum.releaseUnderreplicatedLedger(ledgerId); } } catch (ReplicationException e) { - LOG.warn("Failed to release under replicated ledger {}.", ledgerId, e); + LOG.error("Failed to release under replicated ledger {}.", ledgerId, e); } finally { ((CountDownLatch) ctx).countDown(); } diff --git a/bookkeeper-server/src/test/java/org/apache/bookkeeper/client/BookKeeperAdminTest.java b/bookkeeper-server/src/test/java/org/apache/bookkeeper/client/BookKeeperAdminTest.java index 591761b1140..facf7037581 100644 --- a/bookkeeper-server/src/test/java/org/apache/bookkeeper/client/BookKeeperAdminTest.java +++ b/bookkeeper-server/src/test/java/org/apache/bookkeeper/client/BookKeeperAdminTest.java @@ -21,6 +21,7 @@ package org.apache.bookkeeper.client; import static java.nio.charset.StandardCharsets.UTF_8; +import static org.apache.bookkeeper.client.TopologyAwareEnsemblePlacementPolicy.REPP_DNS_RESOLVER_CLASS; import static org.apache.bookkeeper.util.BookKeeperConstants.AVAILABLE_NODE; import static org.apache.bookkeeper.util.BookKeeperConstants.READONLY; import static org.hamcrest.Matchers.is; @@ -30,21 +31,30 @@ import static org.junit.Assert.assertThat; import static org.junit.Assert.assertTrue; import static org.junit.Assert.fail; +import static org.mockito.ArgumentMatchers.eq; import com.google.common.net.InetAddresses; import java.io.File; +import java.lang.reflect.Field; import java.util.ArrayList; +import java.util.Arrays; import java.util.Collection; +import java.util.Collections; +import java.util.HashMap; import java.util.HashSet; import java.util.Iterator; import java.util.List; +import java.util.Map; +import java.util.NavigableMap; import java.util.Objects; import java.util.Random; import java.util.Set; +import java.util.TreeMap; import java.util.concurrent.CompletableFuture; import java.util.concurrent.CountDownLatch; import java.util.concurrent.ExecutionException; import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicInteger; +import lombok.Cleanup; import org.apache.bookkeeper.bookie.BookieImpl; import org.apache.bookkeeper.bookie.BookieResources; import org.apache.bookkeeper.bookie.CookieValidation; @@ -65,6 +75,9 @@ import org.apache.bookkeeper.meta.zk.ZKMetadataDriverBase; import org.apache.bookkeeper.net.BookieId; import org.apache.bookkeeper.net.BookieSocketAddress; +import org.apache.bookkeeper.net.NetworkTopology; +import org.apache.bookkeeper.net.NetworkTopologyImpl; +import org.apache.bookkeeper.proto.BookieAddressResolver; import org.apache.bookkeeper.proto.BookieServer; import org.apache.bookkeeper.replication.ReplicationException.UnavailableException; import org.apache.bookkeeper.server.Main; @@ -74,11 +87,14 @@ import org.apache.bookkeeper.util.AvailabilityOfEntriesOfLedger; import org.apache.bookkeeper.util.BookKeeperConstants; import org.apache.bookkeeper.util.PortManager; +import org.apache.bookkeeper.util.StaticDNSResolver; import org.apache.commons.io.FileUtils; import org.apache.zookeeper.CreateMode; import org.apache.zookeeper.ZooDefs.Ids; +import org.junit.After; import org.junit.Assert; import org.junit.Test; +import org.mockito.Mockito; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -100,6 +116,12 @@ public BookKeeperAdminTest() { setAutoRecoveryEnabled(true); } + @After + public void tearDown() throws Exception { + super.tearDown(); + StaticDNSResolver.reset(); + } + @Test public void testLostBookieRecoveryDelayValue() throws Exception { try (BookKeeperAdmin bkAdmin = new BookKeeperAdmin(zkUtil.getZooKeeperConnectString())) { @@ -735,4 +757,121 @@ public void testBookieServiceInfoReadonly() throws Exception { public void testLegacyBookieServiceInfo() throws Exception { testBookieServiceInfo(false, true); } + + @Test + public void testRelocateLedgerToAdherePlacementPolicyByRackaware() throws Exception { + final int ensembleSize = 7; + final int quorumSize = 2; + + final BookieSocketAddress addr1 = new BookieSocketAddress("127.0.0.1", 3181); + final BookieSocketAddress addr2 = new BookieSocketAddress("127.0.0.2", 3181); + final BookieSocketAddress addr3 = new BookieSocketAddress("127.0.0.3", 3181); + final BookieSocketAddress addr4 = new BookieSocketAddress("127.0.0.4", 3181); + final BookieSocketAddress addr5 = new BookieSocketAddress("127.0.0.5", 3181); + final BookieSocketAddress addr6 = new BookieSocketAddress("127.0.0.6", 3181); + final BookieSocketAddress addr7 = new BookieSocketAddress("127.0.0.7", 3181); + final BookieSocketAddress addr8 = new BookieSocketAddress("127.0.0.8", 3181); + final BookieSocketAddress addr9 = new BookieSocketAddress("127.0.0.9", 3181); + + final Set writableBookies = new HashSet<>(); + writableBookies.add(addr1.toBookieId()); + writableBookies.add(addr2.toBookieId()); + writableBookies.add(addr3.toBookieId()); + writableBookies.add(addr4.toBookieId()); + writableBookies.add(addr5.toBookieId()); + writableBookies.add(addr6.toBookieId()); + writableBookies.add(addr7.toBookieId()); + writableBookies.add(addr8.toBookieId()); + writableBookies.add(addr9.toBookieId()); + + // add bookie node to resolver + StaticDNSResolver.reset(); + + final String rackName1 = NetworkTopology.DEFAULT_REGION + "/r1"; + final String rackName2 = NetworkTopology.DEFAULT_REGION + "/r2"; + final String rackName3 = NetworkTopology.DEFAULT_REGION + "/r3"; + + // update dns mapping + // add port for testing + StaticDNSResolver.addNodeToRack(addr1.getSocketAddress().getAddress().getHostAddress(), rackName1); + StaticDNSResolver.addNodeToRack(addr2.getSocketAddress().getAddress().getHostAddress(), rackName1); + StaticDNSResolver.addNodeToRack(addr3.getSocketAddress().getAddress().getHostAddress(), rackName1); + StaticDNSResolver.addNodeToRack(addr4.getSocketAddress().getAddress().getHostAddress(), rackName2); + StaticDNSResolver.addNodeToRack(addr5.getSocketAddress().getAddress().getHostAddress(), rackName2); + StaticDNSResolver.addNodeToRack(addr6.getSocketAddress().getAddress().getHostAddress(), rackName2); + StaticDNSResolver.addNodeToRack(addr7.getSocketAddress().getAddress().getHostAddress(), rackName3); + StaticDNSResolver.addNodeToRack(addr8.getSocketAddress().getAddress().getHostAddress(), rackName3); + StaticDNSResolver.addNodeToRack(addr9.getSocketAddress().getAddress().getHostAddress(), rackName3); + LOG.info("Set up static DNS Resolver."); + baseClientConf.setEnsemblePlacementPolicy(RackawareEnsemblePlacementPolicy.class); + baseClientConf.setProperty(REPP_DNS_RESOLVER_CLASS, StaticDNSResolver.class.getName()); + + final NavigableMap> ensemble = new TreeMap<>(); + // create failed ensemble + // expect that the ensemble will be replaced to + // [addr1, addr4, addr7, addr2, addr5, addr8, *addr6 (in /default-region/r2 bookies)*] + ensemble.put(0L, Arrays.asList( + addr1.toBookieId(), addr4.toBookieId(), + addr7.toBookieId(), addr2.toBookieId(), + addr5.toBookieId(), addr8.toBookieId(), + addr3.toBookieId())); + final LedgerMetadata lm1 = Mockito.mock(LedgerMetadata.class); + Mockito.doReturn(ensemble).when(lm1).getAllEnsembles(); + Mockito.doReturn(ensembleSize).when(lm1).getEnsembleSize(); + Mockito.doReturn(quorumSize).when(lm1).getWriteQuorumSize(); + Mockito.doReturn(quorumSize).when(lm1).getAckQuorumSize(); + final LedgerHandle lh1 = Mockito.mock(LedgerHandle.class); + Mockito.when(lh1.getLedgerMetadata()).thenReturn(lm1); + + @Cleanup + final BookKeeper bookKeeper = Mockito.spy(new BookKeeper(baseClientConf)); + Mockito.doReturn(true).when(bookKeeper).isClosed(Mockito.anyLong()); + + @Cleanup + final BookKeeperAdmin admin = Mockito.spy(new BookKeeperAdmin(bookKeeper)); + Mockito.doAnswer(invocationOnMock -> { + final AsyncCallback.OpenCallback op = invocationOnMock.getArgument(1); + final CountDownLatch ctx = invocationOnMock.getArgument(2); + op.openComplete(BKException.Code.OK, lh1, ctx); + return null; + }).when(admin).asyncOpenLedger(Mockito.anyLong(), Mockito.any(), Mockito.any()); + // expected return + final Map expectedMap = new HashMap<>(); + expectedMap.put(6, addr6.toBookieId()); + Mockito.doNothing().when(admin) + .replicateLedgerFragment(Mockito.any(), Mockito.any(), eq(expectedMap), Mockito.any()); + + final EnsemblePlacementPolicy policy = bookKeeper.getPlacementPolicy(); + final Field depthOfAllLeavesField = NetworkTopologyImpl.class.getDeclaredField("depthOfAllLeaves"); + depthOfAllLeavesField.setAccessible(true); + final Field topologyField = TopologyAwareEnsemblePlacementPolicy.class + .getDeclaredField("topology"); + topologyField.setAccessible(true); + final NetworkTopology topology = Mockito.spy((NetworkTopology) topologyField.get(policy)); + Mockito.doReturn(2).when(topology).getNumOfRacks(); + depthOfAllLeavesField.set(topology, -1); + topologyField.set(policy, topology); + final Field bookieAddressResolverField = TopologyAwareEnsemblePlacementPolicy.class + .getDeclaredField("bookieAddressResolver"); + bookieAddressResolverField.setAccessible(true); + final BookieAddressResolver bookieAddressResolver = + Mockito.spy((BookieAddressResolver) bookieAddressResolverField.get(policy)); + Mockito.doAnswer(invocationOnMock -> { + final BookieId bookieId = invocationOnMock.getArgument(0); + return new BookieSocketAddress(bookieId.getId()); + }).when(bookieAddressResolver).resolve(Mockito.any()); + bookieAddressResolverField.set(policy, bookieAddressResolver); + + // add mock bookies to knownBookies + policy.onClusterChanged(writableBookies, Collections.emptySet()); + + // make sure that the ensemble is FAIL state + assertEquals(EnsemblePlacementPolicy.PlacementPolicyAdherence.FAIL, + policy.isEnsembleAdheringToPlacementPolicy(ensemble.get(0L), quorumSize, quorumSize)); + + assertEquals(Collections.emptyList(), admin.relocateLedgerToAdherePlacementPolicy(lh1, false)); + + Mockito.verify(admin, Mockito.times(1)) + .replicateLedgerFragment(Mockito.any(), Mockito.any(), eq(expectedMap), Mockito.any()); + } } diff --git a/bookkeeper-server/src/test/java/org/apache/bookkeeper/client/CorrectEnsemblePlacementCmdTest.java b/bookkeeper-server/src/test/java/org/apache/bookkeeper/client/CorrectEnsemblePlacementCmdTest.java index aad2b2d56ab..e5c472751ae 100644 --- a/bookkeeper-server/src/test/java/org/apache/bookkeeper/client/CorrectEnsemblePlacementCmdTest.java +++ b/bookkeeper-server/src/test/java/org/apache/bookkeeper/client/CorrectEnsemblePlacementCmdTest.java @@ -18,37 +18,16 @@ package org.apache.bookkeeper.client; import static junit.framework.TestCase.assertEquals; -import static org.apache.bookkeeper.client.RackawareEnsemblePlacementPolicyImpl.REPP_DNS_RESOLVER_CLASS; -import static org.mockito.ArgumentMatchers.eq; -import java.lang.reflect.Field; -import java.util.Arrays; -import java.util.Collections; -import java.util.HashMap; -import java.util.HashSet; -import java.util.List; -import java.util.Map; -import java.util.NavigableMap; -import java.util.Set; -import java.util.TreeMap; import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicInteger; import lombok.Cleanup; import org.apache.bookkeeper.bookie.BookieShell; import org.apache.bookkeeper.bookie.storage.ldb.DbLedgerStorage; -import org.apache.bookkeeper.client.api.LedgerMetadata; -import org.apache.bookkeeper.net.BookieId; -import org.apache.bookkeeper.net.BookieSocketAddress; -import org.apache.bookkeeper.net.NetworkTopology; -import org.apache.bookkeeper.proto.BookieAddressResolver; import org.apache.bookkeeper.test.BookKeeperClusterTestCase; -import org.apache.bookkeeper.tools.cli.commands.bookies.CorrectEnsemblePlacementCommand; import org.apache.bookkeeper.util.EntryFormatter; import org.apache.bookkeeper.util.LedgerIdFormatter; -import org.apache.bookkeeper.util.StaticDNSResolver; -import org.junit.After; import org.junit.Test; -import org.mockito.Mockito; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -62,33 +41,23 @@ public class CorrectEnsemblePlacementCmdTest extends BookKeeperClusterTestCase { private static final String PASSWORD = "testPasswd"; public CorrectEnsemblePlacementCmdTest() throws Exception { - super(0); + super(1); baseConf.setLedgerStorageClass(DbLedgerStorage.class.getName()); baseConf.setGcWaitTime(60000); baseConf.setFlushInterval(1); } - @After - public void tearDown() throws Exception { - super.tearDown(); - StaticDNSResolver.reset(); - } - /** * list of entry logger files that contains given ledgerId. */ @Test public void testArgument() throws Exception { - startNewBookie(); - - @Cleanup - final BookKeeper bk = new BookKeeper(baseClientConf, zkc); - @Cleanup - final LedgerHandle lh = createLedgerWithEntries(bk, 10, 1, 1); + @Cleanup final BookKeeper bk = new BookKeeper(baseClientConf, zkc); + @Cleanup final LedgerHandle lh = createLedgerWithEntries(bk, 10, 1, 1); - final String[] argv1 = { "correct-ensemble-placement", "--ledgerids", String.valueOf(lh.getId()), + final String[] argv1 = {"correct-ensemble-placement", "--ledgerids", String.valueOf(lh.getId()), "--skipOpenLedgers", "--force"}; - final String[] argv2 = { "correct-ensemble-placement", "--ledgerids", String.valueOf(lh.getId()), + final String[] argv2 = {"correct-ensemble-placement", "--ledgerids", String.valueOf(lh.getId()), "--skipOpenLedgers", "--force", "--dryrun"}; final BookieShell bkShell = new BookieShell(LedgerIdFormatter.LONG_LEDGERID_FORMATTER, EntryFormatter.STRING_FORMATTER); @@ -98,126 +67,6 @@ public void testArgument() throws Exception { assertEquals("Failed to return exit code!", 0, bkShell.run(argv2)); } - @Test - public void testCorrectEnsemblePlacementByRackaware() throws Exception { - final int ensembleSize = 7; - final int quorumSize = 2; - - final BookieSocketAddress addr1 = new BookieSocketAddress("127.0.0.1", 3181); - final BookieSocketAddress addr2 = new BookieSocketAddress("127.0.0.2", 3181); - final BookieSocketAddress addr3 = new BookieSocketAddress("127.0.0.3", 3181); - final BookieSocketAddress addr4 = new BookieSocketAddress("127.0.0.4", 3181); - final BookieSocketAddress addr5 = new BookieSocketAddress("127.0.0.5", 3181); - final BookieSocketAddress addr6 = new BookieSocketAddress("127.0.0.6", 3181); - final BookieSocketAddress addr7 = new BookieSocketAddress("127.0.0.7", 3181); - final BookieSocketAddress addr8 = new BookieSocketAddress("127.0.0.8", 3181); - final BookieSocketAddress addr9 = new BookieSocketAddress("127.0.0.9", 3181); - - final Set writableBookies = new HashSet<>(); - writableBookies.add(addr1.toBookieId()); - writableBookies.add(addr2.toBookieId()); - writableBookies.add(addr3.toBookieId()); - writableBookies.add(addr4.toBookieId()); - writableBookies.add(addr5.toBookieId()); - writableBookies.add(addr6.toBookieId()); - writableBookies.add(addr7.toBookieId()); - writableBookies.add(addr8.toBookieId()); - writableBookies.add(addr9.toBookieId()); - - // add bookie node to resolver - StaticDNSResolver.reset(); - - final String rackName1 = NetworkTopology.DEFAULT_REGION + "/r1"; - final String rackName2 = NetworkTopology.DEFAULT_REGION + "/r2"; - final String rackName3 = NetworkTopology.DEFAULT_REGION + "/r3"; - - // update dns mapping - // add port for testing - StaticDNSResolver.addNodeToRack(addr1.getSocketAddress().getAddress().getHostAddress(), rackName1); - StaticDNSResolver.addNodeToRack(addr2.getSocketAddress().getAddress().getHostAddress(), rackName1); - StaticDNSResolver.addNodeToRack(addr3.getSocketAddress().getAddress().getHostAddress(), rackName1); - StaticDNSResolver.addNodeToRack(addr4.getSocketAddress().getAddress().getHostAddress(), rackName2); - StaticDNSResolver.addNodeToRack(addr5.getSocketAddress().getAddress().getHostAddress(), rackName2); - StaticDNSResolver.addNodeToRack(addr6.getSocketAddress().getAddress().getHostAddress(), rackName2); - StaticDNSResolver.addNodeToRack(addr7.getSocketAddress().getAddress().getHostAddress(), rackName3); - StaticDNSResolver.addNodeToRack(addr8.getSocketAddress().getAddress().getHostAddress(), rackName3); - StaticDNSResolver.addNodeToRack(addr9.getSocketAddress().getAddress().getHostAddress(), rackName3); - LOG.info("Set up static DNS Resolver."); - baseClientConf.setEnsemblePlacementPolicy(RackawareEnsemblePlacementPolicy.class); - baseClientConf.setProperty(REPP_DNS_RESOLVER_CLASS, StaticDNSResolver.class.getName()); - - startNewBookie(); - - final NavigableMap> ensemble = new TreeMap<>(); - // create failed ensemble - // expect that the ensemble will be replaced to - // [addr1, addr4, addr7, addr2, addr5, addr8, *addr6 (in /default-region/r2 bookies)*] - ensemble.put(0L, Arrays.asList( - addr1.toBookieId(), addr4.toBookieId(), - addr7.toBookieId(), addr2.toBookieId(), - addr5.toBookieId(), addr8.toBookieId(), - addr3.toBookieId())); - final LedgerMetadata lm1 = Mockito.mock(LedgerMetadata.class); - Mockito.doReturn(ensemble).when(lm1).getAllEnsembles(); - Mockito.doReturn(ensembleSize).when(lm1).getEnsembleSize(); - Mockito.doReturn(quorumSize).when(lm1).getWriteQuorumSize(); - Mockito.doReturn(quorumSize).when(lm1).getAckQuorumSize(); - final LedgerHandle lh1 = Mockito.mock(LedgerHandle.class); - Mockito.when(lh1.getLedgerMetadata()).thenReturn(lm1); - - @Cleanup - final BookKeeper bookKeeper = Mockito.spy(new BookKeeper(baseClientConf)); - Mockito.doReturn(true).when(bookKeeper).isClosed(Mockito.anyLong()); - - @Cleanup - final BookKeeperAdmin admin = Mockito.spy(new BookKeeperAdmin(baseClientConf)); - Mockito.doAnswer(invocationOnMock -> { - final AsyncCallback.OpenCallback op = invocationOnMock.getArgument(1); - final CountDownLatch ctx = invocationOnMock.getArgument(2); - op.openComplete(BKException.Code.OK, lh1, ctx); - return null; - }).when(admin).asyncOpenLedger(Mockito.anyLong(), Mockito.any(), Mockito.any()); - // expected return - final Map expectedMap = new HashMap<>(); - expectedMap.put(6, addr6.toBookieId()); - Mockito.doNothing().when(admin) - .replicateLedgerFragment(Mockito.any(), Mockito.any(), eq(expectedMap), Mockito.any()); - - final EnsemblePlacementPolicy policy = bookKeeper.getPlacementPolicy(); - final Field topologyField = TopologyAwareEnsemblePlacementPolicy.class - .getDeclaredField("topology"); - topologyField.setAccessible(true); - final NetworkTopology topology = Mockito.spy((NetworkTopology) topologyField.get(policy)); - Mockito.doReturn(2).when(topology).getNumOfRacks(); - topologyField.set(policy, topology); - final Field bookieAddressResolverField = TopologyAwareEnsemblePlacementPolicy.class - .getDeclaredField("bookieAddressResolver"); - bookieAddressResolverField.setAccessible(true); - final BookieAddressResolver bookieAddressResolver = - Mockito.spy((BookieAddressResolver) bookieAddressResolverField.get(policy)); - Mockito.doAnswer(invocationOnMock -> { - final BookieId bookieId = invocationOnMock.getArgument(0); - return new BookieSocketAddress(bookieId.getId()); - }).when(bookieAddressResolver).resolve(Mockito.any()); - bookieAddressResolverField.set(policy, bookieAddressResolver); - - // add mock bookies to knownBookies - policy.onClusterChanged(writableBookies, Collections.emptySet()); - - // make sure that the ensemble is FAIL state - assertEquals(EnsemblePlacementPolicy.PlacementPolicyAdherence.FAIL, - policy.isEnsembleAdheringToPlacementPolicy(ensemble.get(0L), quorumSize, quorumSize)); - - final CorrectEnsemblePlacementCommand.CorrectEnsemblePlacementFlags flag = - new CorrectEnsemblePlacementCommand.CorrectEnsemblePlacementFlags(); - flag.ledgerIds(Collections.singletonList(1L)); - final CorrectEnsemblePlacementCommand cmd = new CorrectEnsemblePlacementCommand(); - cmd.relocate(baseConf, flag, bookKeeper, admin); - - Mockito.verify(admin, Mockito.times(1)) - .replicateLedgerFragment(Mockito.any(), Mockito.any(), eq(expectedMap), Mockito.any()); - } - private LedgerHandle createLedgerWithEntries(BookKeeper bk, int numOfEntries, int ensembleSize, int quorumSize) throws Exception { LedgerHandle lh = bk.createLedger(ensembleSize, quorumSize, digestType, PASSWORD.getBytes());