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]

Reply via email to