Refactor CompactionStrategyManager Patch by Blake Eggleston; reviewed by Marcus Eriksson for CASSANDRA-14621
Project: http://git-wip-us.apache.org/repos/asf/cassandra/repo Commit: http://git-wip-us.apache.org/repos/asf/cassandra/commit/cba5d513 Tree: http://git-wip-us.apache.org/repos/asf/cassandra/tree/cba5d513 Diff: http://git-wip-us.apache.org/repos/asf/cassandra/diff/cba5d513 Branch: refs/heads/trunk Commit: cba5d5131e75c0631b2f8845ad63b413ab4f420f Parents: 186a860 Author: Blake Eggleston <[email protected]> Authored: Fri Aug 3 08:07:50 2018 -0700 Committer: Blake Eggleston <[email protected]> Committed: Fri Aug 24 08:21:12 2018 -0700 ---------------------------------------------------------------------- CHANGES.txt | 1 + .../apache/cassandra/db/ColumnFamilyStore.java | 6 - .../compaction/AbstractCompactionStrategy.java | 6 + .../db/compaction/AbstractStrategyHolder.java | 206 +++++++ .../db/compaction/CompactionStrategyHolder.java | 240 ++++++++ .../compaction/CompactionStrategyManager.java | 547 ++++++++----------- .../db/compaction/PendingRepairHolder.java | 252 +++++++++ .../db/compaction/PendingRepairManager.java | 30 +- .../CompactionStrategyManagerTest.java | 153 +++++- .../cassandra/io/sstable/LegacySSTableTest.java | 16 +- 10 files changed, 1100 insertions(+), 357 deletions(-) ---------------------------------------------------------------------- http://git-wip-us.apache.org/repos/asf/cassandra/blob/cba5d513/CHANGES.txt ---------------------------------------------------------------------- diff --git a/CHANGES.txt b/CHANGES.txt index cd96555..0e97f9a 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -1,4 +1,5 @@ 4.0 + * Refactor CompactionStrategyManager (CASSANDRA-14621) * Flush netty client messages immediately by default (CASSANDRA-13651) * Improve read repair blocking behavior (CASSANDRA-10726) * Add a virtual table to expose settings (CASSANDRA-14573) http://git-wip-us.apache.org/repos/asf/cassandra/blob/cba5d513/src/java/org/apache/cassandra/db/ColumnFamilyStore.java ---------------------------------------------------------------------- diff --git a/src/java/org/apache/cassandra/db/ColumnFamilyStore.java b/src/java/org/apache/cassandra/db/ColumnFamilyStore.java index f03ffe6..496fc10 100644 --- a/src/java/org/apache/cassandra/db/ColumnFamilyStore.java +++ b/src/java/org/apache/cassandra/db/ColumnFamilyStore.java @@ -2553,12 +2553,6 @@ public class ColumnFamilyStore implements ColumnFamilyStoreMBean return getDirectories().trueSnapshotsSize(); } - @VisibleForTesting - void resetFileIndexGenerator() - { - fileIndexGenerator.set(0); - } - /** * Returns a ColumnFamilyStore by id if it exists, null otherwise * Differently from others, this method does not throw exception if the table does not exist. http://git-wip-us.apache.org/repos/asf/cassandra/blob/cba5d513/src/java/org/apache/cassandra/db/compaction/AbstractCompactionStrategy.java ---------------------------------------------------------------------- diff --git a/src/java/org/apache/cassandra/db/compaction/AbstractCompactionStrategy.java b/src/java/org/apache/cassandra/db/compaction/AbstractCompactionStrategy.java index 3410f13..59bdce6 100644 --- a/src/java/org/apache/cassandra/db/compaction/AbstractCompactionStrategy.java +++ b/src/java/org/apache/cassandra/db/compaction/AbstractCompactionStrategy.java @@ -281,6 +281,12 @@ public abstract class AbstractCompactionStrategy public abstract void removeSSTable(SSTableReader sstable); + public void removeSSTables(Iterable<SSTableReader> removed) + { + for (SSTableReader sstable : removed) + removeSSTable(sstable); + } + /** * Returns the sstables managed by this strategy instance */ http://git-wip-us.apache.org/repos/asf/cassandra/blob/cba5d513/src/java/org/apache/cassandra/db/compaction/AbstractStrategyHolder.java ---------------------------------------------------------------------- diff --git a/src/java/org/apache/cassandra/db/compaction/AbstractStrategyHolder.java b/src/java/org/apache/cassandra/db/compaction/AbstractStrategyHolder.java new file mode 100644 index 0000000..dc16261 --- /dev/null +++ b/src/java/org/apache/cassandra/db/compaction/AbstractStrategyHolder.java @@ -0,0 +1,206 @@ +/* + * 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.cassandra.db.compaction; + +import java.util.Collection; +import java.util.Collections; +import java.util.HashSet; +import java.util.List; +import java.util.Set; +import java.util.UUID; +import java.util.function.Supplier; + +import com.google.common.base.Preconditions; + +import org.apache.cassandra.db.ColumnFamilyStore; +import org.apache.cassandra.db.SerializationHeader; +import org.apache.cassandra.db.lifecycle.LifecycleTransaction; +import org.apache.cassandra.dht.Range; +import org.apache.cassandra.dht.Token; +import org.apache.cassandra.index.Index; +import org.apache.cassandra.io.sstable.Descriptor; +import org.apache.cassandra.io.sstable.ISSTableScanner; +import org.apache.cassandra.io.sstable.SSTableMultiWriter; +import org.apache.cassandra.io.sstable.format.SSTableReader; +import org.apache.cassandra.io.sstable.metadata.MetadataCollector; +import org.apache.cassandra.schema.CompactionParams; + +/** + * Wrapper that's aware of how sstables are divided between separate strategies, + * and provides a standard interface to them + * + * not threadsafe, calls must be synchronized by caller + */ +public abstract class AbstractStrategyHolder +{ + public static class TaskSupplier implements Comparable<TaskSupplier> + { + private final int numRemaining; + private final Supplier<AbstractCompactionTask> supplier; + + TaskSupplier(int numRemaining, Supplier<AbstractCompactionTask> supplier) + { + this.numRemaining = numRemaining; + this.supplier = supplier; + } + + public AbstractCompactionTask getTask() + { + return supplier.get(); + } + + public int compareTo(TaskSupplier o) + { + return o.numRemaining - numRemaining; + } + } + + public static interface DestinationRouter + { + int getIndexForSSTable(SSTableReader sstable); + int getIndexForSSTableDirectory(Descriptor descriptor); + } + + /** + * Maps sstables to their token partition bucket + */ + static class GroupedSSTableContainer + { + private final AbstractStrategyHolder holder; + private final Set<SSTableReader>[] groups; + + private GroupedSSTableContainer(AbstractStrategyHolder holder) + { + this.holder = holder; + Preconditions.checkArgument(holder.numTokenPartitions > 0, "numTokenPartitions not set"); + groups = new Set[holder.numTokenPartitions]; + } + + void add(SSTableReader sstable) + { + Preconditions.checkArgument(holder.managesSSTable(sstable), "this strategy holder doesn't manage %s", sstable); + int idx = holder.router.getIndexForSSTable(sstable); + Preconditions.checkState(idx >= 0 && idx < holder.numTokenPartitions, "Invalid sstable index (%s) for %s", idx, sstable); + if (groups[idx] == null) + groups[idx] = new HashSet<>(); + groups[idx].add(sstable); + } + + int numGroups() + { + return groups.length; + } + + Set<SSTableReader> getGroup(int i) + { + Preconditions.checkArgument(i >= 0 && i < groups.length); + Set<SSTableReader> group = groups[i]; + return group != null ? group : Collections.emptySet(); + } + + boolean isGroupEmpty(int i) + { + return getGroup(i).isEmpty(); + } + + boolean isEmpty() + { + for (int i = 0; i < groups.length; i++) + if (!isGroupEmpty(i)) + return false; + return true; + } + } + + protected final ColumnFamilyStore cfs; + final DestinationRouter router; + private int numTokenPartitions = -1; + + AbstractStrategyHolder(ColumnFamilyStore cfs, DestinationRouter router) + { + this.cfs = cfs; + this.router = router; + } + + public abstract void startup(); + + public abstract void shutdown(); + + final void setStrategy(CompactionParams params, int numTokenPartitions) + { + Preconditions.checkArgument(numTokenPartitions > 0, "at least one token partition required"); + shutdown(); + this.numTokenPartitions = numTokenPartitions; + setStrategyInternal(params, numTokenPartitions); + } + + protected abstract void setStrategyInternal(CompactionParams params, int numTokenPartitions); + + /** + * SSTables are grouped by their repaired and pending repair status. This method determines if this holder + * holds the sstable for the given repaired/grouped statuses. Holders should be mutually exclusive in the + * groups they deal with. IOW, if one holder returns true for a given isRepaired/isPendingRepair combo, + * none of the others should. + */ + public abstract boolean managesRepairedGroup(boolean isRepaired, boolean isPendingRepair); + + public boolean managesSSTable(SSTableReader sstable) + { + return managesRepairedGroup(sstable.isRepaired(), sstable.isPendingRepair()); + } + + public abstract AbstractCompactionStrategy getStrategyFor(SSTableReader sstable); + + public abstract Iterable<AbstractCompactionStrategy> allStrategies(); + + public abstract Collection<TaskSupplier> getBackgroundTaskSuppliers(int gcBefore); + + public abstract Collection<AbstractCompactionTask> getMaximalTasks(int gcBefore, boolean splitOutput); + + public abstract Collection<AbstractCompactionTask> getUserDefinedTasks(GroupedSSTableContainer sstables, int gcBefore); + + public GroupedSSTableContainer createGroupedSSTableContainer() + { + return new GroupedSSTableContainer(this); + } + + public abstract void addSSTables(GroupedSSTableContainer sstables); + + public abstract void removeSSTables(GroupedSSTableContainer sstables); + + public abstract void replaceSSTables(GroupedSSTableContainer removed, GroupedSSTableContainer added); + + public abstract List<ISSTableScanner> getScanners(GroupedSSTableContainer sstables, Collection<Range<Token>> ranges); + + + public abstract SSTableMultiWriter createSSTableMultiWriter(Descriptor descriptor, + long keyCount, + long repairedAt, + UUID pendingRepair, + MetadataCollector collector, + SerializationHeader header, + Collection<Index> indexes, + LifecycleTransaction txn); + + /** + * Return the directory index the given compaction strategy belongs to, or -1 + * if it's not held by this holder + */ + public abstract int getStrategyIndex(AbstractCompactionStrategy strategy); +} http://git-wip-us.apache.org/repos/asf/cassandra/blob/cba5d513/src/java/org/apache/cassandra/db/compaction/CompactionStrategyHolder.java ---------------------------------------------------------------------- diff --git a/src/java/org/apache/cassandra/db/compaction/CompactionStrategyHolder.java b/src/java/org/apache/cassandra/db/compaction/CompactionStrategyHolder.java new file mode 100644 index 0000000..8fba121 --- /dev/null +++ b/src/java/org/apache/cassandra/db/compaction/CompactionStrategyHolder.java @@ -0,0 +1,240 @@ +/* + * 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.cassandra.db.compaction; + +import java.util.ArrayList; +import java.util.Collection; +import java.util.List; +import java.util.UUID; + +import com.google.common.base.Preconditions; + +import org.apache.cassandra.db.ColumnFamilyStore; +import org.apache.cassandra.db.SerializationHeader; +import org.apache.cassandra.db.lifecycle.LifecycleTransaction; +import org.apache.cassandra.dht.Range; +import org.apache.cassandra.dht.Token; +import org.apache.cassandra.index.Index; +import org.apache.cassandra.io.sstable.Descriptor; +import org.apache.cassandra.io.sstable.ISSTableScanner; +import org.apache.cassandra.io.sstable.SSTableMultiWriter; +import org.apache.cassandra.io.sstable.format.SSTableReader; +import org.apache.cassandra.io.sstable.metadata.MetadataCollector; +import org.apache.cassandra.schema.CompactionParams; +import org.apache.cassandra.service.ActiveRepairService; + +public class CompactionStrategyHolder extends AbstractStrategyHolder +{ + private final List<AbstractCompactionStrategy> strategies = new ArrayList<>(); + private final boolean isRepaired; + + public CompactionStrategyHolder(ColumnFamilyStore cfs, DestinationRouter router, boolean isRepaired) + { + super(cfs, router); + this.isRepaired = isRepaired; + } + + @Override + public void startup() + { + strategies.forEach(AbstractCompactionStrategy::startup); + } + + @Override + public void shutdown() + { + strategies.forEach(AbstractCompactionStrategy::shutdown); + } + + @Override + public void setStrategyInternal(CompactionParams params, int numTokenPartitions) + { + strategies.clear(); + for (int i = 0; i < numTokenPartitions; i++) + strategies.add(cfs.createCompactionStrategyInstance(params)); + } + + @Override + public boolean managesRepairedGroup(boolean isRepaired, boolean isPendingRepair) + { + Preconditions.checkArgument(!isPendingRepair || !isRepaired, + "SSTables cannot be both repaired and pending repair"); + return !isPendingRepair && (isRepaired == this.isRepaired); + } + + @Override + public AbstractCompactionStrategy getStrategyFor(SSTableReader sstable) + { + Preconditions.checkArgument(managesSSTable(sstable), "Attempting to get compaction strategy from wrong holder"); + return strategies.get(router.getIndexForSSTable(sstable)); + } + + @Override + public Iterable<AbstractCompactionStrategy> allStrategies() + { + return strategies; + } + + @Override + public Collection<TaskSupplier> getBackgroundTaskSuppliers(int gcBefore) + { + List<TaskSupplier> suppliers = new ArrayList<>(strategies.size()); + for (AbstractCompactionStrategy strategy : strategies) + suppliers.add(new TaskSupplier(strategy.getEstimatedRemainingTasks(), () -> strategy.getNextBackgroundTask(gcBefore))); + + return suppliers; + } + + @Override + public Collection<AbstractCompactionTask> getMaximalTasks(int gcBefore, boolean splitOutput) + { + List<AbstractCompactionTask> tasks = new ArrayList<>(strategies.size()); + for (AbstractCompactionStrategy strategy : strategies) + { + Collection<AbstractCompactionTask> task = strategy.getMaximalTask(gcBefore, splitOutput); + if (task != null) + tasks.addAll(task); + } + return tasks; + } + + @Override + public Collection<AbstractCompactionTask> getUserDefinedTasks(GroupedSSTableContainer sstables, int gcBefore) + { + List<AbstractCompactionTask> tasks = new ArrayList<>(strategies.size()); + for (int i = 0; i < strategies.size(); i++) + { + if (sstables.isGroupEmpty(i)) + continue; + + tasks.add(strategies.get(i).getUserDefinedTask(sstables.getGroup(i), gcBefore)); + } + return tasks; + } + + @Override + public void addSSTables(GroupedSSTableContainer sstables) + { + Preconditions.checkArgument(sstables.numGroups() == strategies.size()); + for (int i = 0; i < strategies.size(); i++) + { + if (!sstables.isGroupEmpty(i)) + strategies.get(i).addSSTables(sstables.getGroup(i)); + } + } + + @Override + public void removeSSTables(GroupedSSTableContainer sstables) + { + Preconditions.checkArgument(sstables.numGroups() == strategies.size()); + for (int i = 0; i < strategies.size(); i++) + { + if (!sstables.isGroupEmpty(i)) + strategies.get(i).removeSSTables(sstables.getGroup(i)); + } + } + + @Override + public void replaceSSTables(GroupedSSTableContainer removed, GroupedSSTableContainer added) + { + Preconditions.checkArgument(removed.numGroups() == strategies.size()); + Preconditions.checkArgument(added.numGroups() == strategies.size()); + for (int i = 0; i < strategies.size(); i++) + { + if (removed.isGroupEmpty(i) && added.isGroupEmpty(i)) + continue; + + if (removed.isGroupEmpty(i)) + strategies.get(i).addSSTables(added.getGroup(i)); + else + strategies.get(i).replaceSSTables(removed.getGroup(i), added.getGroup(i)); + } + } + + public AbstractCompactionStrategy first() + { + return strategies.get(0); + } + + @Override + @SuppressWarnings("resource") + public List<ISSTableScanner> getScanners(GroupedSSTableContainer sstables, Collection<Range<Token>> ranges) + { + List<ISSTableScanner> scanners = new ArrayList<>(strategies.size()); + for (int i = 0; i < strategies.size(); i++) + { + if (sstables.isGroupEmpty(i)) + continue; + + scanners.addAll(strategies.get(i).getScanners(sstables.getGroup(i), ranges).scanners); + } + return scanners; + } + + Collection<Collection<SSTableReader>> groupForAnticompaction(Iterable<SSTableReader> sstables) + { + Preconditions.checkState(!isRepaired); + GroupedSSTableContainer group = createGroupedSSTableContainer(); + sstables.forEach(group::add); + + Collection<Collection<SSTableReader>> anticompactionGroups = new ArrayList<>(); + for (int i = 0; i < strategies.size(); i++) + { + if (group.isGroupEmpty(i)) + continue; + + anticompactionGroups.addAll(strategies.get(i).groupSSTablesForAntiCompaction(group.getGroup(i))); + } + + return anticompactionGroups; + } + + @Override + public SSTableMultiWriter createSSTableMultiWriter(Descriptor descriptor, long keyCount, long repairedAt, UUID pendingRepair, MetadataCollector collector, SerializationHeader header, Collection<Index> indexes, LifecycleTransaction txn) + { + if (isRepaired) + { + Preconditions.checkArgument(repairedAt != ActiveRepairService.UNREPAIRED_SSTABLE, + "Repaired CompactionStrategyHolder can't create unrepaired sstable writers"); + } + else + { + Preconditions.checkArgument(repairedAt == ActiveRepairService.UNREPAIRED_SSTABLE, + "Unrepaired CompactionStrategyHolder can't create repaired sstable writers"); + } + Preconditions.checkArgument(pendingRepair == null, + "CompactionStrategyHolder can't create sstable writer with pendingRepair id"); + // to avoid creating a compaction strategy for the wrong pending repair manager, we get the index based on where the sstable is to be written + AbstractCompactionStrategy strategy = strategies.get(router.getIndexForSSTableDirectory(descriptor)); + return strategy.createSSTableMultiWriter(descriptor, + keyCount, + repairedAt, + pendingRepair, + collector, + header, + indexes, + txn); + } + + @Override + public int getStrategyIndex(AbstractCompactionStrategy strategy) + { + return strategies.indexOf(strategy); + } +} http://git-wip-us.apache.org/repos/asf/cassandra/blob/cba5d513/src/java/org/apache/cassandra/db/compaction/CompactionStrategyManager.java ---------------------------------------------------------------------- diff --git a/src/java/org/apache/cassandra/db/compaction/CompactionStrategyManager.java b/src/java/org/apache/cassandra/db/compaction/CompactionStrategyManager.java index 81b7c7e..9766454 100644 --- a/src/java/org/apache/cassandra/db/compaction/CompactionStrategyManager.java +++ b/src/java/org/apache/cassandra/db/compaction/CompactionStrategyManager.java @@ -20,50 +20,67 @@ package org.apache.cassandra.db.compaction; import java.io.File; import java.io.IOException; -import java.util.*; +import java.util.ArrayList; +import java.util.Arrays; +import java.util.Collection; +import java.util.Collections; +import java.util.ConcurrentModificationException; +import java.util.HashSet; +import java.util.List; +import java.util.Set; +import java.util.UUID; import java.util.concurrent.Callable; import java.util.concurrent.locks.ReentrantReadWriteLock; import java.util.function.Supplier; import java.util.stream.Collectors; -import java.util.stream.Stream; import com.google.common.annotations.VisibleForTesting; +import com.google.common.collect.ImmutableList; import com.google.common.collect.Iterables; import com.google.common.collect.Lists; import com.google.common.primitives.Longs; - -import org.apache.cassandra.config.DatabaseDescriptor; -import org.apache.cassandra.db.DiskBoundaries; -import org.apache.cassandra.index.Index; - import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import org.apache.cassandra.config.DatabaseDescriptor; import org.apache.cassandra.db.ColumnFamilyStore; import org.apache.cassandra.db.Directories; +import org.apache.cassandra.db.DiskBoundaries; import org.apache.cassandra.db.SerializationHeader; +import org.apache.cassandra.db.compaction.AbstractStrategyHolder.TaskSupplier; import org.apache.cassandra.db.lifecycle.LifecycleTransaction; import org.apache.cassandra.db.lifecycle.SSTableSet; import org.apache.cassandra.dht.Range; import org.apache.cassandra.dht.Token; +import org.apache.cassandra.index.Index; import org.apache.cassandra.io.sstable.Component; import org.apache.cassandra.io.sstable.Descriptor; +import org.apache.cassandra.io.sstable.ISSTableScanner; import org.apache.cassandra.io.sstable.SSTableMultiWriter; import org.apache.cassandra.io.sstable.format.SSTableReader; -import org.apache.cassandra.io.sstable.ISSTableScanner; import org.apache.cassandra.io.sstable.metadata.MetadataCollector; import org.apache.cassandra.io.sstable.metadata.StatsMetadata; -import org.apache.cassandra.notifications.*; +import org.apache.cassandra.notifications.INotification; +import org.apache.cassandra.notifications.INotificationConsumer; +import org.apache.cassandra.notifications.SSTableAddedNotification; +import org.apache.cassandra.notifications.SSTableDeletingNotification; +import org.apache.cassandra.notifications.SSTableListChangedNotification; +import org.apache.cassandra.notifications.SSTableMetadataChanged; +import org.apache.cassandra.notifications.SSTableRepairStatusChanged; import org.apache.cassandra.schema.CompactionParams; import org.apache.cassandra.schema.TableMetadata; import org.apache.cassandra.service.ActiveRepairService; -import org.apache.cassandra.utils.Pair; + +import static org.apache.cassandra.db.compaction.AbstractStrategyHolder.GroupedSSTableContainer; /** * Manages the compaction strategies. * - * For each directory, a separate compaction strategy instance for both repaired and unrepaired data, and also one instance - * for each pending repair. This is done to keep the different sets of sstables completely separate. + * SSTables are isolated from each other based on their incremental repair status (repaired, unrepaired, or pending repair) + * and directory (determined by their starting token). This class handles the routing between {@link AbstractStrategyHolder} + * instances based on repair status, and the {@link AbstractStrategyHolder} instances have separate compaction strategies + * for each directory, which it routes sstables to. Note that {@link PendingRepairHolder} also divides sstables on their + * pending repair id. * * Operations on this class are guarded by a {@link ReentrantReadWriteLock}. This lock performs mutual exclusion on * reads and writes to the following variables: {@link this#repaired}, {@link this#unrepaired}, {@link this#isActive}, @@ -95,9 +112,12 @@ public class CompactionStrategyManager implements INotificationConsumer /** * Variables guarded by read and write lock above */ - private final List<AbstractCompactionStrategy> repaired = new ArrayList<>(); - private final List<AbstractCompactionStrategy> unrepaired = new ArrayList<>(); - private final List<PendingRepairManager> pendingRepairs = new ArrayList<>(); + private final PendingRepairHolder pendingRepairs; + private final CompactionStrategyHolder repaired; + private final CompactionStrategyHolder unrepaired; + + private final ImmutableList<AbstractStrategyHolder> holders; + private volatile CompactionParams params; private DiskBoundaries currentBoundaries; private volatile boolean enabled = true; @@ -124,6 +144,23 @@ public class CompactionStrategyManager implements INotificationConsumer public CompactionStrategyManager(ColumnFamilyStore cfs, Supplier<DiskBoundaries> boundariesSupplier, boolean partitionSSTablesByTokenRange) { + AbstractStrategyHolder.DestinationRouter router = new AbstractStrategyHolder.DestinationRouter() + { + public int getIndexForSSTable(SSTableReader sstable) + { + return compactionStrategyIndexFor(sstable); + } + + public int getIndexForSSTableDirectory(Descriptor descriptor) + { + return compactionStrategyIndexForDirectory(descriptor); + } + }; + pendingRepairs = new PendingRepairHolder(cfs, router); + repaired = new CompactionStrategyHolder(cfs, router, true); + unrepaired = new CompactionStrategyHolder(cfs, router, false); + holders = ImmutableList.of(pendingRepairs, repaired, unrepaired); + cfs.getTracker().subscribe(this); logger.trace("{} subscribed to the data tracker.", this); this.cfs = cfs; @@ -150,55 +187,36 @@ public class CompactionStrategyManager implements INotificationConsumer if (!isEnabled()) return null; + int numPartitions = getNumTokenPartitions(); // first try to promote/demote sstables from completed repairs - ArrayList<Pair<Integer, PendingRepairManager>> pendingRepairManagers = new ArrayList<>(pendingRepairs.size()); - for (PendingRepairManager pendingRepair : pendingRepairs) + List<TaskSupplier> repairFinishedSuppliers = pendingRepairs.getRepairFinishedTaskSuppliers(); + if (!repairFinishedSuppliers.isEmpty()) { - int numPending = pendingRepair.getNumPendingRepairFinishedTasks(); - if (numPending > 0) + Collections.sort(repairFinishedSuppliers); + for (TaskSupplier supplier : repairFinishedSuppliers) { - pendingRepairManagers.add(Pair.create(numPending, pendingRepair)); - } - } - if (!pendingRepairManagers.isEmpty()) - { - pendingRepairManagers.sort((x, y) -> y.left - x.left); - for (Pair<Integer, PendingRepairManager> pair : pendingRepairManagers) - { - AbstractCompactionTask task = pair.right.getNextRepairFinishedTask(); + AbstractCompactionTask task = supplier.getTask(); if (task != null) - { return task; - } } } // sort compaction task suppliers by remaining tasks descending - ArrayList<Pair<Integer, Supplier<AbstractCompactionTask>>> sortedSuppliers = new ArrayList<>(repaired.size() + unrepaired.size() + 1); + List<TaskSupplier> suppliers = new ArrayList<>(numPartitions * holders.size()); + for (AbstractStrategyHolder holder : holders) + suppliers.addAll(holder.getBackgroundTaskSuppliers(gcBefore)); - for (AbstractCompactionStrategy strategy : repaired) - sortedSuppliers.add(Pair.create(strategy.getEstimatedRemainingTasks(), () -> strategy.getNextBackgroundTask(gcBefore))); - - for (AbstractCompactionStrategy strategy : unrepaired) - sortedSuppliers.add(Pair.create(strategy.getEstimatedRemainingTasks(), () -> strategy.getNextBackgroundTask(gcBefore))); - - for (PendingRepairManager pending : pendingRepairs) - sortedSuppliers.add(Pair.create(pending.getMaxEstimatedRemainingTasks(), () -> pending.getNextBackgroundTask(gcBefore))); - - sortedSuppliers.sort((x, y) -> y.left - x.left); + Collections.sort(suppliers); // return the first non-null task - AbstractCompactionTask task; - Iterator<Supplier<AbstractCompactionTask>> suppliers = Iterables.transform(sortedSuppliers, p -> p.right).iterator(); - assert suppliers.hasNext(); - - do + for (TaskSupplier supplier : suppliers) { - task = suppliers.next().get(); + AbstractCompactionTask task = supplier.getTask(); + if (task != null) + return task; } - while (suppliers.hasNext() && task == null); - return task; + return null; } finally { @@ -289,21 +307,17 @@ public class CompactionStrategyManager implements INotificationConsumer if (sstable.openReason != SSTableReader.OpenReason.EARLY) compactionStrategyFor(sstable).addSSTable(sstable); } - repaired.forEach(AbstractCompactionStrategy::startup); - unrepaired.forEach(AbstractCompactionStrategy::startup); - pendingRepairs.forEach(PendingRepairManager::startup); - shouldDefragment = repaired.get(0).shouldDefragment(); - supportsEarlyOpen = repaired.get(0).supportsEarlyOpen(); - fanout = (repaired.get(0) instanceof LeveledCompactionStrategy) ? ((LeveledCompactionStrategy) repaired.get(0)).getLevelFanoutSize() : LeveledCompactionStrategy.DEFAULT_LEVEL_FANOUT_SIZE; + holders.forEach(AbstractStrategyHolder::startup); + shouldDefragment = repaired.first().shouldDefragment(); + supportsEarlyOpen = repaired.first().supportsEarlyOpen(); + fanout = (repaired.first() instanceof LeveledCompactionStrategy) ? ((LeveledCompactionStrategy) repaired.first()).getLevelFanoutSize() : LeveledCompactionStrategy.DEFAULT_LEVEL_FANOUT_SIZE; } finally { writeLock.unlock(); } - repaired.forEach(AbstractCompactionStrategy::startup); - unrepaired.forEach(AbstractCompactionStrategy::startup); - pendingRepairs.forEach(PendingRepairManager::startup); - if (Stream.concat(repaired.stream(), unrepaired.stream()).anyMatch(cs -> cs.logAll)) + + if (repaired.first().logAll) compactionLogger.enable(); } @@ -321,19 +335,13 @@ public class CompactionStrategyManager implements INotificationConsumer } @VisibleForTesting - protected AbstractCompactionStrategy compactionStrategyFor(SSTableReader sstable) + AbstractCompactionStrategy compactionStrategyFor(SSTableReader sstable) { // should not call maybeReloadDiskBoundaries because it may be called from within lock readLock.lock(); try { - int index = compactionStrategyIndexFor(sstable); - if (sstable.isPendingRepair()) - return pendingRepairs.get(index).getOrCreate(sstable); - else if (sstable.isRepaired()) - return repaired.get(index); - else - return unrepaired.get(index); + return getHolder(sstable).getStrategyFor(sstable); } finally { @@ -353,7 +361,7 @@ public class CompactionStrategyManager implements INotificationConsumer * @return */ @VisibleForTesting - protected int compactionStrategyIndexFor(SSTableReader sstable) + int compactionStrategyIndexFor(SSTableReader sstable) { // should not call maybeReloadDiskBoundaries because it may be called from within lock readLock.lock(); @@ -372,13 +380,28 @@ public class CompactionStrategyManager implements INotificationConsumer } } + private int compactionStrategyIndexForDirectory(Descriptor descriptor) + { + readLock.lock(); + try + { + return partitionSSTablesByTokenRange ? currentBoundaries.getBoundariesFromSSTableDirectory(descriptor) : 0; + } + finally + { + readLock.unlock(); + } + } + + + @VisibleForTesting List<AbstractCompactionStrategy> getRepaired() { readLock.lock(); try { - return Lists.newArrayList(repaired); + return Lists.newArrayList(repaired.allStrategies()); } finally { @@ -392,7 +415,7 @@ public class CompactionStrategyManager implements INotificationConsumer readLock.lock(); try { - return Lists.newArrayList(unrepaired); + return Lists.newArrayList(unrepaired.allStrategies()); } finally { @@ -401,14 +424,12 @@ public class CompactionStrategyManager implements INotificationConsumer } @VisibleForTesting - List<AbstractCompactionStrategy> getForPendingRepair(UUID sessionID) + Iterable<AbstractCompactionStrategy> getForPendingRepair(UUID sessionID) { readLock.lock(); try { - List<AbstractCompactionStrategy> strategies = new ArrayList<>(pendingRepairs.size()); - pendingRepairs.forEach(p -> strategies.add(p.get(sessionID))); - return strategies; + return pendingRepairs.getStrategiesFor(sessionID); } finally { @@ -423,7 +444,7 @@ public class CompactionStrategyManager implements INotificationConsumer try { Set<UUID> ids = new HashSet<>(); - pendingRepairs.forEach(p -> ids.addAll(p.getSessions())); + pendingRepairs.getManagers().forEach(p -> ids.addAll(p.getSessions())); return ids; } finally @@ -437,7 +458,8 @@ public class CompactionStrategyManager implements INotificationConsumer readLock.lock(); try { - return Iterables.any(pendingRepairs, prm -> prm.hasDataForSession(sessionID)); + return Iterables.any(pendingRepairs.getManagers(), + prm -> prm.hasDataForSession(sessionID)); } finally { @@ -451,9 +473,7 @@ public class CompactionStrategyManager implements INotificationConsumer try { isActive = false; - repaired.forEach(AbstractCompactionStrategy::shutdown); - unrepaired.forEach(AbstractCompactionStrategy::shutdown); - pendingRepairs.forEach(PendingRepairManager::shutdown); + holders.forEach(AbstractStrategyHolder::shutdown); compactionLogger.disable(); } finally @@ -546,22 +566,22 @@ public class CompactionStrategyManager implements INotificationConsumer startup(); } + private Iterable<AbstractCompactionStrategy> getAllStrategies() + { + return Iterables.concat(Iterables.transform(holders, AbstractStrategyHolder::allStrategies)); + } + public int getUnleveledSSTables() { maybeReloadDiskBoundaries(); readLock.lock(); try { - if (repaired.get(0) instanceof LeveledCompactionStrategy && unrepaired.get(0) instanceof LeveledCompactionStrategy) + if (repaired.first() instanceof LeveledCompactionStrategy) { int count = 0; - for (AbstractCompactionStrategy strategy : repaired) + for (AbstractCompactionStrategy strategy : getAllStrategies()) count += ((LeveledCompactionStrategy) strategy).getLevelSize(0); - for (AbstractCompactionStrategy strategy : unrepaired) - count += ((LeveledCompactionStrategy) strategy).getLevelSize(0); - for (PendingRepairManager pendingManager : pendingRepairs) - for (AbstractCompactionStrategy strategy : pendingManager.getStrategies()) - count += ((LeveledCompactionStrategy) strategy).getLevelSize(0); return count; } } @@ -583,24 +603,14 @@ public class CompactionStrategyManager implements INotificationConsumer readLock.lock(); try { - if (repaired.get(0) instanceof LeveledCompactionStrategy && unrepaired.get(0) instanceof LeveledCompactionStrategy) + if (repaired.first() instanceof LeveledCompactionStrategy) { int[] res = new int[LeveledManifest.MAX_LEVEL_COUNT]; - for (AbstractCompactionStrategy strategy : repaired) + for (AbstractCompactionStrategy strategy : getAllStrategies()) { int[] repairedCountPerLevel = ((LeveledCompactionStrategy) strategy).getAllLevelSize(); res = sumArrays(res, repairedCountPerLevel); } - for (AbstractCompactionStrategy strategy : unrepaired) - { - int[] unrepairedCountPerLevel = ((LeveledCompactionStrategy) strategy).getAllLevelSize(); - res = sumArrays(res, unrepairedCountPerLevel); - } - for (PendingRepairManager pending : pendingRepairs) - { - int[] pendingRepairCountPerLevel = pending.getSSTableCountPerLevel(); - res = sumArrays(res, pendingRepairCountPerLevel); - } return res; } } @@ -650,6 +660,75 @@ public class CompactionStrategyManager implements INotificationConsumer } } + private int getHolderIndex(SSTableReader sstable) + { + for (int i = 0; i < holders.size(); i++) + { + if (holders.get(i).managesSSTable(sstable)) + return i; + } + + throw new IllegalStateException("No holder claimed " + sstable); + } + + private AbstractStrategyHolder getHolder(SSTableReader sstable) + { + for (AbstractStrategyHolder holder : holders) + { + if (holder.managesSSTable(sstable)) + return holder; + } + + throw new IllegalStateException("No holder claimed " + sstable); + } + + private AbstractStrategyHolder getHolder(long repairedAt, UUID pendingRepair) + { + return getHolder(repairedAt != ActiveRepairService.UNREPAIRED_SSTABLE, + pendingRepair != ActiveRepairService.NO_PENDING_REPAIR); + } + + @VisibleForTesting + AbstractStrategyHolder getHolder(boolean isRepaired, boolean isPendingRepair) + { + for (AbstractStrategyHolder holder : holders) + { + if (holder.managesRepairedGroup(isRepaired, isPendingRepair)) + return holder; + } + + throw new IllegalStateException(String.format("No holder claimed isPendingRepair: %s, isPendingRepair %s", + isRepaired, isPendingRepair)); + } + + @VisibleForTesting + ImmutableList<AbstractStrategyHolder> getHolders() + { + return holders; + } + + /** + * Split sstables into a list of grouped sstable containers, the list index an sstable + * + * lives in matches the list index of the holder that's responsible for it + */ + @VisibleForTesting + List<GroupedSSTableContainer> groupSSTables(Iterable<SSTableReader> sstables) + { + List<GroupedSSTableContainer> classified = new ArrayList<>(holders.size()); + for (AbstractStrategyHolder holder : holders) + { + classified.add(holder.createGroupedSSTableContainer()); + } + + for (SSTableReader sstable : sstables) + { + classified.get(getHolderIndex(sstable)).add(sstable); + } + + return classified; + } + private void handleListChangedNotification(Iterable<SSTableReader> added, Iterable<SSTableReader> removed) { // If reloaded, SSTables will be placed in their correct locations @@ -660,70 +739,11 @@ public class CompactionStrategyManager implements INotificationConsumer readLock.lock(); try { - // a bit of gymnastics to be able to replace sstables in compaction strategies - // we use this to know that a compaction finished and where to start the next compaction in LCS - Directories.DataDirectory [] locations = cfs.getDirectories().getWriteableLocations(); - int locationSize = cfs.getPartitioner().splitter().isPresent() ? locations.length : 1; - - List<Set<SSTableReader>> pendingRemoved = new ArrayList<>(locationSize); - List<Set<SSTableReader>> pendingAdded = new ArrayList<>(locationSize); - List<Set<SSTableReader>> repairedRemoved = new ArrayList<>(locationSize); - List<Set<SSTableReader>> repairedAdded = new ArrayList<>(locationSize); - List<Set<SSTableReader>> unrepairedRemoved = new ArrayList<>(locationSize); - List<Set<SSTableReader>> unrepairedAdded = new ArrayList<>(locationSize); - - for (int i = 0; i < locationSize; i++) - { - pendingRemoved.add(new HashSet<>()); - pendingAdded.add(new HashSet<>()); - repairedRemoved.add(new HashSet<>()); - repairedAdded.add(new HashSet<>()); - unrepairedRemoved.add(new HashSet<>()); - unrepairedAdded.add(new HashSet<>()); - } - - for (SSTableReader sstable : removed) - { - int i = compactionStrategyIndexFor(sstable); - if (sstable.isPendingRepair()) - pendingRemoved.get(i).add(sstable); - else if (sstable.isRepaired()) - repairedRemoved.get(i).add(sstable); - else - unrepairedRemoved.get(i).add(sstable); - } - for (SSTableReader sstable : added) + List<GroupedSSTableContainer> addedGroups = groupSSTables(added); + List<GroupedSSTableContainer> removedGroups = groupSSTables(removed); + for (int i=0; i<holders.size(); i++) { - int i = compactionStrategyIndexFor(sstable); - if (sstable.isPendingRepair()) - pendingAdded.get(i).add(sstable); - else if (sstable.isRepaired()) - repairedAdded.get(i).add(sstable); - else - unrepairedAdded.get(i).add(sstable); - } - for (int i = 0; i < locationSize; i++) - { - - if (!pendingRemoved.get(i).isEmpty()) - { - pendingRepairs.get(i).replaceSSTables(pendingRemoved.get(i), pendingAdded.get(i)); - } - else - { - PendingRepairManager pendingManager = pendingRepairs.get(i); - pendingAdded.get(i).forEach(s -> pendingManager.addSSTable(s)); - } - - if (!repairedRemoved.get(i).isEmpty()) - repaired.get(i).replaceSSTables(repairedRemoved.get(i), repairedAdded.get(i)); - else - repaired.get(i).addSSTables(repairedAdded.get(i)); - - if (!unrepairedRemoved.get(i).isEmpty()) - unrepaired.get(i).replaceSSTables(unrepairedRemoved.get(i), unrepairedAdded.get(i)); - else - unrepaired.get(i).addSSTables(unrepairedAdded.get(i)); + holders.get(i).replaceSSTables(removedGroups.get(i), addedGroups.get(i)); } } finally @@ -742,26 +762,21 @@ public class CompactionStrategyManager implements INotificationConsumer readLock.lock(); try { - for (SSTableReader sstable : sstables) + List<GroupedSSTableContainer> groups = groupSSTables(sstables); + for (int i = 0; i < holders.size(); i++) { - int index = compactionStrategyIndexFor(sstable); - if (sstable.isPendingRepair()) - { - pendingRepairs.get(index).addSSTable(sstable); - unrepaired.get(index).removeSSTable(sstable); - repaired.get(index).removeSSTable(sstable); - } - else if (sstable.isRepaired()) - { - pendingRepairs.get(index).removeSSTable(sstable); - unrepaired.get(index).removeSSTable(sstable); - repaired.get(index).addSSTable(sstable); - } - else + GroupedSSTableContainer group = groups.get(i); + + if (group.isEmpty()) + continue; + + AbstractStrategyHolder dstHolder = holders.get(i); + dstHolder.addSSTables(group); + + for (AbstractStrategyHolder holder : holders) { - pendingRepairs.get(index).removeSSTable(sstable); - repaired.get(index).removeSSTable(sstable); - unrepaired.get(index).addSSTable(sstable); + if (holder != dstHolder) + holder.removeSSTables(group); } } } @@ -863,46 +878,13 @@ public class CompactionStrategyManager implements INotificationConsumer List<ISSTableScanner> scanners = new ArrayList<>(sstables.size()); try { - assert repaired.size() == unrepaired.size(); - assert repaired.size() == pendingRepairs.size(); - - int numRepaired = repaired.size(); - List<Set<SSTableReader>> pendingSSTables = new ArrayList<>(numRepaired); - List<Set<SSTableReader>> repairedSSTables = new ArrayList<>(numRepaired); - List<Set<SSTableReader>> unrepairedSSTables = new ArrayList<>(numRepaired); - - for (int i = 0; i < numRepaired; i++) - { - pendingSSTables.add(new HashSet<>()); - repairedSSTables.add(new HashSet<>()); - unrepairedSSTables.add(new HashSet<>()); - } + List<GroupedSSTableContainer> sstableGroups = groupSSTables(sstables); - for (SSTableReader sstable : sstables) + for (int i = 0; i < holders.size(); i++) { - int idx = compactionStrategyIndexFor(sstable); - if (sstable.isPendingRepair()) - pendingSSTables.get(idx).add(sstable); - else if (sstable.isRepaired()) - repairedSSTables.get(idx).add(sstable); - else - unrepairedSSTables.get(idx).add(sstable); - } - - for (int i = 0; i < pendingSSTables.size(); i++) - { - if (!pendingSSTables.get(i).isEmpty()) - scanners.addAll(pendingRepairs.get(i).getScanners(pendingSSTables.get(i), ranges)); - } - for (int i = 0; i < repairedSSTables.size(); i++) - { - if (!repairedSSTables.get(i).isEmpty()) - scanners.addAll(repaired.get(i).getScanners(repairedSSTables.get(i), ranges).scanners); - } - for (int i = 0; i < unrepairedSSTables.size(); i++) - { - if (!unrepairedSSTables.get(i).isEmpty()) - scanners.addAll(unrepaired.get(i).getScanners(unrepairedSSTables.get(i), ranges).scanners); + AbstractStrategyHolder holder = holders.get(i); + GroupedSSTableContainer group = sstableGroups.get(i); + scanners.addAll(holder.getScanners(group, ranges)); } } catch (PendingRepairManager.IllegalSSTableArgumentException e) @@ -942,12 +924,7 @@ public class CompactionStrategyManager implements INotificationConsumer readLock.lock(); try { - Map<Integer, List<SSTableReader>> groups = sstablesToGroup.stream().collect(Collectors.groupingBy((s) -> compactionStrategyIndexFor(s))); - Collection<Collection<SSTableReader>> anticompactionGroups = new ArrayList<>(); - - for (Map.Entry<Integer, List<SSTableReader>> group : groups.entrySet()) - anticompactionGroups.addAll(unrepaired.get(group.getKey()).groupSSTablesForAntiCompaction(group.getValue())); - return anticompactionGroups; + return unrepaired.groupForAnticompaction(sstablesToGroup); } finally { @@ -960,7 +937,7 @@ public class CompactionStrategyManager implements INotificationConsumer readLock.lock(); try { - return unrepaired.get(0).getMaxSSTableBytes(); + return unrepaired.first().getMaxSSTableBytes(); } finally { @@ -1026,24 +1003,9 @@ public class CompactionStrategyManager implements INotificationConsumer readLock.lock(); try { - for (AbstractCompactionStrategy strategy : repaired) - { - Collection<AbstractCompactionTask> task = strategy.getMaximalTask(gcBefore, splitOutput); - if (task != null) - tasks.addAll(task); - } - for (AbstractCompactionStrategy strategy : unrepaired) + for (AbstractStrategyHolder holder : holders) { - Collection<AbstractCompactionTask> task = strategy.getMaximalTask(gcBefore, splitOutput); - if (task != null) - tasks.addAll(task); - } - - for (PendingRepairManager pending : pendingRepairs) - { - Collection<AbstractCompactionTask> pendingRepairTasks = pending.getMaximalTasks(gcBefore, splitOutput); - if (pendingRepairTasks != null) - tasks.addAll(pendingRepairTasks); + tasks.addAll(holder.getMaximalTasks(gcBefore, splitOutput)); } } finally @@ -1073,27 +1035,11 @@ public class CompactionStrategyManager implements INotificationConsumer readLock.lock(); try { - Map<Integer, List<SSTableReader>> repairedSSTables = sstables.stream() - .filter(s -> !s.isMarkedSuspect() && s.isRepaired() && !s.isPendingRepair()) - .collect(Collectors.groupingBy((s) -> compactionStrategyIndexFor(s))); - - Map<Integer, List<SSTableReader>> unrepairedSSTables = sstables.stream() - .filter(s -> !s.isMarkedSuspect() && !s.isRepaired() && !s.isPendingRepair()) - .collect(Collectors.groupingBy((s) -> compactionStrategyIndexFor(s))); - - Map<Integer, List<SSTableReader>> pendingSSTables = sstables.stream() - .filter(s -> !s.isMarkedSuspect() && s.isPendingRepair()) - .collect(Collectors.groupingBy((s) -> compactionStrategyIndexFor(s))); - - for (Map.Entry<Integer, List<SSTableReader>> group : repairedSSTables.entrySet()) - ret.add(repaired.get(group.getKey()).getUserDefinedTask(group.getValue(), gcBefore)); - - for (Map.Entry<Integer, List<SSTableReader>> group : unrepairedSSTables.entrySet()) - ret.add(unrepaired.get(group.getKey()).getUserDefinedTask(group.getValue(), gcBefore)); - - for (Map.Entry<Integer, List<SSTableReader>> group : pendingSSTables.entrySet()) - ret.addAll(pendingRepairs.get(group.getKey()).createUserDefinedTasks(group.getValue(), gcBefore)); - + List<GroupedSSTableContainer> groupedSSTables = groupSSTables(sstables); + for (int i = 0; i < holders.size(); i++) + { + ret.addAll(holders.get(i).getUserDefinedTasks(groupedSSTables.get(i), gcBefore)); + } return ret; } finally @@ -1109,12 +1055,8 @@ public class CompactionStrategyManager implements INotificationConsumer readLock.lock(); try { - for (AbstractCompactionStrategy strategy : repaired) - tasks += strategy.getEstimatedRemainingTasks(); - for (AbstractCompactionStrategy strategy : unrepaired) + for (AbstractCompactionStrategy strategy : getAllStrategies()) tasks += strategy.getEstimatedRemainingTasks(); - for (PendingRepairManager pending : pendingRepairs) - tasks += pending.getEstimatedRemainingTasks(); } finally { @@ -1134,7 +1076,7 @@ public class CompactionStrategyManager implements INotificationConsumer readLock.lock(); try { - return unrepaired.get(0).getName(); + return unrepaired.first().getName(); } finally { @@ -1148,9 +1090,9 @@ public class CompactionStrategyManager implements INotificationConsumer readLock.lock(); try { - List<AbstractCompactionStrategy> pending = new ArrayList<>(); - pendingRepairs.forEach(p -> pending.addAll(p.getStrategies())); - return Arrays.asList(repaired, unrepaired, pending); + return Arrays.asList(Lists.newArrayList(repaired.allStrategies()), + Lists.newArrayList(unrepaired.allStrategies()), + Lists.newArrayList(pendingRepairs.allStrategies())); } finally { @@ -1177,30 +1119,16 @@ public class CompactionStrategyManager implements INotificationConsumer } } - private void setStrategy(CompactionParams params) + private int getNumTokenPartitions() { - repaired.forEach(AbstractCompactionStrategy::shutdown); - unrepaired.forEach(AbstractCompactionStrategy::shutdown); - pendingRepairs.forEach(PendingRepairManager::shutdown); - repaired.clear(); - unrepaired.clear(); - pendingRepairs.clear(); + return partitionSSTablesByTokenRange ? currentBoundaries.directories.size() : 1; + } - if (partitionSSTablesByTokenRange) - { - for (int i = 0; i < currentBoundaries.directories.size(); i++) - { - repaired.add(cfs.createCompactionStrategyInstance(params)); - unrepaired.add(cfs.createCompactionStrategyInstance(params)); - pendingRepairs.add(new PendingRepairManager(cfs, params)); - } - } - else - { - repaired.add(cfs.createCompactionStrategyInstance(params)); - unrepaired.add(cfs.createCompactionStrategyInstance(params)); - pendingRepairs.add(new PendingRepairManager(cfs, params)); - } + private void setStrategy(CompactionParams params) + { + int numPartitions = getNumTokenPartitions(); + for (AbstractStrategyHolder holder : holders) + holder.setStrategy(params, numPartitions); this.params = params; } @@ -1227,14 +1155,7 @@ public class CompactionStrategyManager implements INotificationConsumer readLock.lock(); try { - // to avoid creating a compaction strategy for the wrong pending repair manager, we get the index based on where the sstable is to be written - int index = partitionSSTablesByTokenRange? currentBoundaries.getBoundariesFromSSTableDirectory(descriptor) : 0; - if (pendingRepair != ActiveRepairService.NO_PENDING_REPAIR) - return pendingRepairs.get(index).getOrCreate(pendingRepair).createSSTableMultiWriter(descriptor, keyCount, ActiveRepairService.UNREPAIRED_SSTABLE, pendingRepair, collector, header, indexes, txn); - else if (repairedAt == ActiveRepairService.UNREPAIRED_SSTABLE) - return unrepaired.get(index).createSSTableMultiWriter(descriptor, keyCount, repairedAt, ActiveRepairService.NO_PENDING_REPAIR, collector, header, indexes, txn); - else - return repaired.get(index).createSSTableMultiWriter(descriptor, keyCount, repairedAt, ActiveRepairService.NO_PENDING_REPAIR, collector, header, indexes, txn); + return getHolder(repairedAt, pendingRepair).createSSTableMultiWriter(descriptor, keyCount, repairedAt, pendingRepair, collector, header, indexes, txn); } finally { @@ -1244,7 +1165,7 @@ public class CompactionStrategyManager implements INotificationConsumer public boolean isRepaired(AbstractCompactionStrategy strategy) { - return repaired.contains(strategy); + return repaired.getStrategyIndex(strategy) >= 0; } public List<String> getStrategyFolders(AbstractCompactionStrategy strategy) @@ -1255,23 +1176,11 @@ public class CompactionStrategyManager implements INotificationConsumer Directories.DataDirectory[] locations = cfs.getDirectories().getWriteableLocations(); if (partitionSSTablesByTokenRange) { - int unrepairedIndex = unrepaired.indexOf(strategy); - if (unrepairedIndex > 0) - { - return Collections.singletonList(locations[unrepairedIndex].location.getAbsolutePath()); - } - int repairedIndex = repaired.indexOf(strategy); - if (repairedIndex > 0) - { - return Collections.singletonList(locations[repairedIndex].location.getAbsolutePath()); - } - for (int i = 0; i < pendingRepairs.size(); i++) + for (AbstractStrategyHolder holder : holders) { - PendingRepairManager pending = pendingRepairs.get(i); - if (pending.hasStrategy(strategy)) - { - return Collections.singletonList(locations[i].location.getAbsolutePath()); - } + int idx = holder.getStrategyIndex(strategy); + if (idx >= 0) + return Collections.singletonList(locations[idx].location.getAbsolutePath()); } } List<String> folders = new ArrayList<>(locations.length); @@ -1299,7 +1208,7 @@ public class CompactionStrategyManager implements INotificationConsumer readLock.lock(); try { - return pendingRepairs; + return Lists.newArrayList(pendingRepairs.getManagers()); } finally { http://git-wip-us.apache.org/repos/asf/cassandra/blob/cba5d513/src/java/org/apache/cassandra/db/compaction/PendingRepairHolder.java ---------------------------------------------------------------------- diff --git a/src/java/org/apache/cassandra/db/compaction/PendingRepairHolder.java b/src/java/org/apache/cassandra/db/compaction/PendingRepairHolder.java new file mode 100644 index 0000000..7b9123f --- /dev/null +++ b/src/java/org/apache/cassandra/db/compaction/PendingRepairHolder.java @@ -0,0 +1,252 @@ +/* + * 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.cassandra.db.compaction; + +import java.util.ArrayList; +import java.util.Collection; +import java.util.List; +import java.util.UUID; + +import com.google.common.base.Preconditions; +import com.google.common.collect.Iterables; + +import org.apache.cassandra.db.ColumnFamilyStore; +import org.apache.cassandra.db.SerializationHeader; +import org.apache.cassandra.db.lifecycle.LifecycleTransaction; +import org.apache.cassandra.dht.Range; +import org.apache.cassandra.dht.Token; +import org.apache.cassandra.index.Index; +import org.apache.cassandra.io.sstable.Descriptor; +import org.apache.cassandra.io.sstable.ISSTableScanner; +import org.apache.cassandra.io.sstable.SSTableMultiWriter; +import org.apache.cassandra.io.sstable.format.SSTableReader; +import org.apache.cassandra.io.sstable.metadata.MetadataCollector; +import org.apache.cassandra.schema.CompactionParams; +import org.apache.cassandra.service.ActiveRepairService; + +public class PendingRepairHolder extends AbstractStrategyHolder +{ + private final List<PendingRepairManager> managers = new ArrayList<>(); + + public PendingRepairHolder(ColumnFamilyStore cfs, DestinationRouter router) + { + super(cfs, router); + } + + @Override + public void startup() + { + managers.forEach(PendingRepairManager::startup); + } + + @Override + public void shutdown() + { + managers.forEach(PendingRepairManager::shutdown); + } + + @Override + public void setStrategyInternal(CompactionParams params, int numTokenPartitions) + { + managers.clear(); + for (int i = 0; i < numTokenPartitions; i++) + managers.add(new PendingRepairManager(cfs, params)); + } + + @Override + public boolean managesRepairedGroup(boolean isRepaired, boolean isPendingRepair) + { + Preconditions.checkArgument(!isPendingRepair || !isRepaired, + "SSTables cannot be both repaired and pending repair"); + return isPendingRepair; + } + + @Override + public AbstractCompactionStrategy getStrategyFor(SSTableReader sstable) + { + Preconditions.checkArgument(managesSSTable(sstable), "Attempting to get compaction strategy from wrong holder"); + return managers.get(router.getIndexForSSTable(sstable)).getOrCreate(sstable); + } + + @Override + public Iterable<AbstractCompactionStrategy> allStrategies() + { + return Iterables.concat(Iterables.transform(managers, PendingRepairManager::getStrategies)); + } + + Iterable<AbstractCompactionStrategy> getStrategiesFor(UUID session) + { + List<AbstractCompactionStrategy> strategies = new ArrayList<>(managers.size()); + for (PendingRepairManager manager : managers) + { + AbstractCompactionStrategy strategy = manager.get(session); + if (strategy != null) + strategies.add(strategy); + } + return strategies; + } + + public Iterable<PendingRepairManager> getManagers() + { + return managers; + } + + @Override + public Collection<TaskSupplier> getBackgroundTaskSuppliers(int gcBefore) + { + List<TaskSupplier> suppliers = new ArrayList<>(managers.size()); + for (PendingRepairManager manager : managers) + suppliers.add(new TaskSupplier(manager.getMaxEstimatedRemainingTasks(), () -> manager.getNextBackgroundTask(gcBefore))); + + return suppliers; + } + + @Override + public Collection<AbstractCompactionTask> getMaximalTasks(int gcBefore, boolean splitOutput) + { + List<AbstractCompactionTask> tasks = new ArrayList<>(managers.size()); + for (PendingRepairManager manager : managers) + { + Collection<AbstractCompactionTask> task = manager.getMaximalTasks(gcBefore, splitOutput); + if (task != null) + tasks.addAll(task); + } + return tasks; + } + + @Override + public Collection<AbstractCompactionTask> getUserDefinedTasks(GroupedSSTableContainer sstables, int gcBefore) + { + List<AbstractCompactionTask> tasks = new ArrayList<>(managers.size()); + + for (int i = 0; i < managers.size(); i++) + { + if (sstables.isGroupEmpty(i)) + continue; + + tasks.addAll(managers.get(i).createUserDefinedTasks(sstables.getGroup(i), gcBefore)); + } + return tasks; + } + + public ArrayList<TaskSupplier> getRepairFinishedTaskSuppliers() + { + ArrayList<TaskSupplier> suppliers = new ArrayList<>(managers.size()); + for (PendingRepairManager manager : managers) + { + int numPending = manager.getNumPendingRepairFinishedTasks(); + if (numPending > 0) + { + suppliers.add(new TaskSupplier(numPending, manager::getNextRepairFinishedTask)); + } + } + + return suppliers; + } + + @Override + public void addSSTables(GroupedSSTableContainer sstables) + { + Preconditions.checkArgument(sstables.numGroups() == managers.size()); + for (int i = 0; i < managers.size(); i++) + { + if (!sstables.isGroupEmpty(i)) + managers.get(i).addSSTables(sstables.getGroup(i)); + } + } + + @Override + public void removeSSTables(GroupedSSTableContainer sstables) + { + Preconditions.checkArgument(sstables.numGroups() == managers.size()); + for (int i = 0; i < managers.size(); i++) + { + if (!sstables.isGroupEmpty(i)) + managers.get(i).removeSSTables(sstables.getGroup(i)); + } + } + + @Override + public void replaceSSTables(GroupedSSTableContainer removed, GroupedSSTableContainer added) + { + Preconditions.checkArgument(removed.numGroups() == managers.size()); + Preconditions.checkArgument(added.numGroups() == managers.size()); + for (int i = 0; i < managers.size(); i++) + { + if (removed.isGroupEmpty(i) && added.isGroupEmpty(i)) + continue; + + if (removed.isGroupEmpty(i)) + managers.get(i).addSSTables(added.getGroup(i)); + else + managers.get(i).replaceSSTables(removed.getGroup(i), added.getGroup(i)); + } + } + + @Override + public List<ISSTableScanner> getScanners(GroupedSSTableContainer sstables, Collection<Range<Token>> ranges) + { + List<ISSTableScanner> scanners = new ArrayList<>(managers.size()); + for (int i = 0; i < managers.size(); i++) + { + if (sstables.isGroupEmpty(i)) + continue; + + scanners.addAll(managers.get(i).getScanners(sstables.getGroup(i), ranges)); + } + return scanners; + } + + @Override + public SSTableMultiWriter createSSTableMultiWriter(Descriptor descriptor, + long keyCount, + long repairedAt, + UUID pendingRepair, + MetadataCollector collector, + SerializationHeader header, + Collection<Index> indexes, + LifecycleTransaction txn) + { + Preconditions.checkArgument(repairedAt == ActiveRepairService.UNREPAIRED_SSTABLE, + "PendingRepairHolder can't create sstablewriter with repaired at set"); + Preconditions.checkArgument(pendingRepair != null, + "PendingRepairHolder can't create sstable writer without pendingRepair id"); + // to avoid creating a compaction strategy for the wrong pending repair manager, we get the index based on where the sstable is to be written + AbstractCompactionStrategy strategy = managers.get(router.getIndexForSSTableDirectory(descriptor)).getOrCreate(pendingRepair); + return strategy.createSSTableMultiWriter(descriptor, + keyCount, + repairedAt, + pendingRepair, + collector, + header, + indexes, + txn); + } + + @Override + public int getStrategyIndex(AbstractCompactionStrategy strategy) + { + for (int i = 0; i < managers.size(); i++) + { + if (managers.get(i).hasStrategy(strategy)) + return i; + } + return -1; + } +} http://git-wip-us.apache.org/repos/asf/cassandra/blob/cba5d513/src/java/org/apache/cassandra/db/compaction/PendingRepairManager.java ---------------------------------------------------------------------- diff --git a/src/java/org/apache/cassandra/db/compaction/PendingRepairManager.java b/src/java/org/apache/cassandra/db/compaction/PendingRepairManager.java index 64aa3f2..edc9a2f 100644 --- a/src/java/org/apache/cassandra/db/compaction/PendingRepairManager.java +++ b/src/java/org/apache/cassandra/db/compaction/PendingRepairManager.java @@ -147,11 +147,24 @@ class PendingRepairManager strategy.removeSSTable(sstable); } + + void removeSSTables(Iterable<SSTableReader> removed) + { + for (SSTableReader sstable : removed) + removeSSTable(sstable); + } + synchronized void addSSTable(SSTableReader sstable) { getOrCreate(sstable).addSSTable(sstable); } + void addSSTables(Iterable<SSTableReader> added) + { + for (SSTableReader sstable : added) + addSSTable(sstable); + } + synchronized void replaceSSTables(Set<SSTableReader> removed, Set<SSTableReader> added) { if (removed.isEmpty() && added.isEmpty()) @@ -335,21 +348,6 @@ class PendingRepairManager return !ActiveRepairService.instance.consistent.local.isSessionInProgress(sessionID); } - /** - * calling this when underlying strategy is not LeveledCompactionStrategy is an error - */ - synchronized int[] getSSTableCountPerLevel() - { - int [] res = new int[LeveledManifest.MAX_LEVEL_COUNT]; - for (AbstractCompactionStrategy strategy : strategies.values()) - { - assert strategy instanceof LeveledCompactionStrategy; - int[] counts = ((LeveledCompactionStrategy) strategy).getAllLevelSize(); - res = CompactionStrategyManager.sumArrays(res, counts); - } - return res; - } - @SuppressWarnings("resource") synchronized Set<ISSTableScanner> getScanners(Collection<SSTableReader> sstables, Collection<Range<Token>> ranges) { @@ -391,7 +389,7 @@ class PendingRepairManager return strategies.keySet().contains(sessionID); } - public Collection<AbstractCompactionTask> createUserDefinedTasks(List<SSTableReader> sstables, int gcBefore) + public Collection<AbstractCompactionTask> createUserDefinedTasks(Collection<SSTableReader> sstables, int gcBefore) { Map<UUID, List<SSTableReader>> group = sstables.stream().collect(Collectors.groupingBy(s -> s.getSSTableMetadata().pendingRepair)); return group.entrySet().stream().map(g -> strategies.get(g.getKey()).getUserDefinedTask(g.getValue(), gcBefore)).collect(Collectors.toList()); http://git-wip-us.apache.org/repos/asf/cassandra/blob/cba5d513/test/unit/org/apache/cassandra/db/compaction/CompactionStrategyManagerTest.java ---------------------------------------------------------------------- diff --git a/test/unit/org/apache/cassandra/db/compaction/CompactionStrategyManagerTest.java b/test/unit/org/apache/cassandra/db/compaction/CompactionStrategyManagerTest.java index 549ba3f..eeaaf5b 100644 --- a/test/unit/org/apache/cassandra/db/compaction/CompactionStrategyManagerTest.java +++ b/test/unit/org/apache/cassandra/db/compaction/CompactionStrategyManagerTest.java @@ -20,20 +20,28 @@ package org.apache.cassandra.db.compaction; import java.io.File; import java.io.IOException; +import java.util.ArrayList; import java.util.Arrays; import java.util.Collections; import java.util.List; import java.util.Set; +import java.util.UUID; import java.util.concurrent.CountDownLatch; import java.util.concurrent.atomic.AtomicInteger; import java.util.stream.Collectors; +import com.google.common.collect.Iterables; +import com.google.common.collect.Lists; import com.google.common.collect.Sets; import com.google.common.io.Files; import org.junit.AfterClass; +import org.junit.Before; import org.junit.BeforeClass; import org.junit.Test; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + import org.apache.cassandra.SchemaLoader; import org.apache.cassandra.Util; import org.apache.cassandra.config.DatabaseDescriptor; @@ -41,10 +49,10 @@ import org.apache.cassandra.db.ColumnFamilyStore; import org.apache.cassandra.db.DecoratedKey; import org.apache.cassandra.db.Directories; import org.apache.cassandra.db.DiskBoundaries; -import org.apache.cassandra.db.DiskBoundaryManager; import org.apache.cassandra.db.Keyspace; import org.apache.cassandra.db.PartitionPosition; import org.apache.cassandra.db.RowUpdateBuilder; +import org.apache.cassandra.db.compaction.AbstractStrategyHolder.GroupedSSTableContainer; import org.apache.cassandra.dht.ByteOrderedPartitioner; import org.apache.cassandra.dht.IPartitioner; import org.apache.cassandra.io.sstable.format.SSTableReader; @@ -53,14 +61,20 @@ import org.apache.cassandra.notifications.SSTableDeletingNotification; import org.apache.cassandra.schema.CompactionParams; import org.apache.cassandra.schema.KeyspaceParams; import org.apache.cassandra.service.StorageService; +import org.apache.cassandra.utils.ByteBufferUtil; import org.apache.cassandra.utils.UUIDGen; import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertFalse; +import static org.junit.Assert.assertNotNull; +import static org.junit.Assert.assertSame; import static org.junit.Assert.assertTrue; +import static org.junit.Assert.fail; public class CompactionStrategyManagerTest { + private static final Logger logger = LoggerFactory.getLogger(CompactionStrategyManagerTest.class); + private static final String KS_PREFIX = "Keyspace1"; private static final String TABLE_PREFIX = "CF_STANDARD"; @@ -84,6 +98,13 @@ public class CompactionStrategyManagerTest .compaction(CompactionParams.stcs(Collections.emptyMap()))); } + @Before + public void setUp() throws Exception + { + ColumnFamilyStore cfs = Keyspace.open(KS_PREFIX).getColumnFamilyStore(TABLE_PREFIX); + cfs.truncateBlocking(); + } + @AfterClass public static void afterClass() { @@ -137,7 +158,7 @@ public class CompactionStrategyManagerTest final Integer[] boundaries = computeBoundaries(numSSTables, numDisks); MockBoundaryManager mockBoundaryManager = new MockBoundaryManager(cfs, boundaries); - System.out.println("Boundaries for " + numDisks + " disks is " + Arrays.toString(boundaries)); + logger.debug("Boundaries for {} disks is {}", numDisks, Arrays.toString(boundaries)); CompactionStrategyManager csm = new CompactionStrategyManager(cfs, mockBoundaryManager::getBoundaries, true); @@ -154,7 +175,7 @@ public class CompactionStrategyManagerTest updateBoundaries(mockBoundaryManager, boundaries, delta); // Check that SSTables are still assigned to the previous boundary layout - System.out.println("Old boundaries: " + Arrays.toString(previousBoundaries) + " New boundaries: " + Arrays.toString(boundaries)); + logger.debug("Old boundaries: {} New boundaries: {}", Arrays.toString(previousBoundaries), Arrays.toString(boundaries)); for (SSTableReader reader : cfs.getLiveSSTables()) { verifySSTableIsAssignedToCorrectStrategy(previousBoundaries, csm, reader); @@ -179,8 +200,6 @@ public class CompactionStrategyManagerTest } } - - @Test public void testAutomaticUpgradeConcurrency() throws Exception { @@ -253,6 +272,114 @@ public class CompactionStrategyManagerTest DatabaseDescriptor.setAutomaticSSTableUpgradeEnabled(false); } + private static void assertHolderExclusivity(boolean isRepaired, boolean isPendingRepair, Class<? extends AbstractStrategyHolder> expectedType) + { + ColumnFamilyStore cfs = Keyspace.open(KS_PREFIX).getColumnFamilyStore(TABLE_PREFIX); + CompactionStrategyManager csm = cfs.getCompactionStrategyManager(); + + AbstractStrategyHolder holder = csm.getHolder(isRepaired, isPendingRepair); + assertNotNull(holder); + assertSame(expectedType, holder.getClass()); + + int matches = 0; + for (AbstractStrategyHolder other : csm.getHolders()) + { + if (other.managesRepairedGroup(isRepaired, isPendingRepair)) + { + assertSame("holder assignment should be mutually exclusive", holder, other); + matches++; + } + } + assertEquals(1, matches); + } + + private static void assertInvalieHolderConfig(boolean isRepaired, boolean isPendingRepair) + { + ColumnFamilyStore cfs = Keyspace.open(KS_PREFIX).getColumnFamilyStore(TABLE_PREFIX); + CompactionStrategyManager csm = cfs.getCompactionStrategyManager(); + try + { + csm.getHolder(isRepaired, isPendingRepair); + fail("Expected IllegalArgumentException"); + } + catch (IllegalArgumentException e) + { + // expected + } + } + + /** + * If an sstable can be be assigned to a strategy holder, it shouldn't be possibly to + * assign it to any of the other holders. + */ + @Test + public void testMutualExclusiveHolderClassification() throws Exception + { + assertHolderExclusivity(false, false, CompactionStrategyHolder.class); + assertHolderExclusivity(true, false, CompactionStrategyHolder.class); + assertHolderExclusivity(false, true, PendingRepairHolder.class); + assertInvalieHolderConfig(true, true); + } + + PartitionPosition forKey(int key) + { + DecoratedKey dk = Util.dk(String.format("%04d", key)); + return dk.getToken().minKeyBound(); + } + + /** + * Test that csm.groupSSTables correctly groups sstables by repaired status and directory + */ + @Test + public void groupSSTables() throws Exception + { + final int numDir = 4; + ColumnFamilyStore cfs = createJBODMockCFS(numDir); + Keyspace.open(cfs.keyspace.getName()).getColumnFamilyStore(cfs.name).disableAutoCompaction(); + assertTrue(cfs.getLiveSSTables().isEmpty()); + List<SSTableReader> unrepaired = new ArrayList<>(); + List<SSTableReader> pendingRepair = new ArrayList<>(); + List<SSTableReader> repaired = new ArrayList<>(); + + for (int i = 0; i < numDir; i++) + { + int key = 100 * i; + unrepaired.add(createSSTableWithKey(cfs.keyspace.getName(), cfs.name, key++)); + pendingRepair.add(createSSTableWithKey(cfs.keyspace.getName(), cfs.name, key++)); + repaired.add(createSSTableWithKey(cfs.keyspace.getName(), cfs.name, key++)); + } + + cfs.getCompactionStrategyManager().mutateRepaired(pendingRepair, 0, UUID.randomUUID()); + cfs.getCompactionStrategyManager().mutateRepaired(repaired, 1000, null); + + DiskBoundaries boundaries = new DiskBoundaries(cfs.getDirectories().getWriteableLocations(), + Lists.newArrayList(forKey(100), forKey(200), forKey(300)), + 10, 10); + + CompactionStrategyManager csm = new CompactionStrategyManager(cfs, () -> boundaries, true); + + List<GroupedSSTableContainer> grouped = csm.groupSSTables(Iterables.concat(repaired, pendingRepair, unrepaired)); + + for (int x=0; x<grouped.size(); x++) + { + GroupedSSTableContainer group = grouped.get(x); + AbstractStrategyHolder holder = csm.getHolders().get(x); + for (int y=0; y<numDir; y++) + { + SSTableReader sstable = Iterables.getOnlyElement(group.getGroup(y)); + assertTrue(holder.managesSSTable(sstable)); + SSTableReader expected; + if (sstable.isRepaired()) + expected = repaired.get(y); + else if (sstable.isPendingRepair()) + expected = pendingRepair.get(y); + else + expected = unrepaired.get(y); + + assertSame(expected, sstable); + } + } + } private MockCFS createJBODMockCFS(int disks) { @@ -313,18 +440,13 @@ public class CompactionStrategyManagerTest return result; } - /** - * Since each SSTable contains keys from 0-99, and each sstable - * generation is numbered from 1-100, since we are using ByteOrderedPartitioner - * we can compute the sstable position in the disk boundaries by finding - * the generation position relative to the boundaries - */ private int getSSTableIndex(Integer[] boundaries, SSTableReader reader) { int index = 0; - while (boundaries[index] < reader.descriptor.generation) + int firstKey = Integer.parseInt(new String(ByteBufferUtil.getArray(reader.first.getKey()))); + while (boundaries[index] <= firstKey) index++; - System.out.println("Index for SSTable " + reader.descriptor.generation + " on boundary " + Arrays.toString(boundaries) + " is " + index); + logger.debug("Index for SSTable {} on boundary {} is {}", reader.descriptor.generation, Arrays.toString(boundaries), index); return index; } @@ -362,7 +484,7 @@ public class CompactionStrategyManagerTest } } - private static void createSSTableWithKey(String keyspace, String table, int key) + private static SSTableReader createSSTableWithKey(String keyspace, String table, int key) { long timestamp = System.currentTimeMillis(); DecoratedKey dk = Util.dk(String.format("%04d", key)); @@ -372,7 +494,10 @@ public class CompactionStrategyManagerTest .add("val", "val") .build() .applyUnsafe(); + Set<SSTableReader> before = cfs.getLiveSSTables(); cfs.forceBlockingFlush(); + Set<SSTableReader> after = cfs.getLiveSSTables(); + return Iterables.getOnlyElement(Sets.difference(after, before)); } // just to be able to override the data directories --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
