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..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; @@ -1137,7 +1141,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) @@ -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/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..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 @@ -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,173 @@ 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); + // 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) { + 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)) { + 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; + 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..fad54b51a2c --- /dev/null +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/tools/cli/commands/bookies/CorrectEnsemblePlacementCommand.java @@ -0,0 +1,228 @@ +/* + * 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.util.concurrent.UncheckedExecutionException; +import java.io.IOException; +import java.net.URI; +import java.util.List; +import java.util.NavigableSet; +import java.util.TreeSet; +import java.util.concurrent.ConcurrentSkipListSet; +import java.util.concurrent.CountDownLatch; +import java.util.stream.Collectors; +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.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.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(bookKeeper); + 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 { + @Cleanup + final MetadataBookieDriver metadataDriver = instantiateMetadataDriver(conf); + @Cleanup + final LedgerManagerFactory lmf = metadataDriver.getLedgerManagerFactory(); + @Cleanup + final LedgerUnderreplicationManager lum = lmf.newLedgerUnderreplicationManager(); + + final List targetLedgers = + flags.ledgerIds.stream().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 (LOG.isDebugEnabled()) { + LOG.debug("Start relocation of the ledger {}.", ledgerId); + } + 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; + } + 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.error("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); + } + } +} 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 new file mode 100644 index 00000000000..e5c472751ae --- /dev/null +++ b/bookkeeper-server/src/test/java/org/apache/bookkeeper/client/CorrectEnsemblePlacementCmdTest.java @@ -0,0 +1,91 @@ +/* + * 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 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.test.BookKeeperClusterTestCase; +import org.apache.bookkeeper.util.EntryFormatter; +import org.apache.bookkeeper.util.LedgerIdFormatter; +import org.junit.Test; +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(1); + 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 { + @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)); + } + + 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)))); + } + } }