diff --git a/src/java/org/apache/cassandra/db/compaction/LeveledGenerations.java b/src/java/org/apache/cassandra/db/compaction/LeveledGenerations.java index 513e02aad99e..28ac9bbb9e85 100644 --- a/src/java/org/apache/cassandra/db/compaction/LeveledGenerations.java +++ b/src/java/org/apache/cassandra/db/compaction/LeveledGenerations.java @@ -96,6 +96,13 @@ Set get(int level) return levels[level - 1]; } + TreeSet getSortedLevel(int level) + { + if (level > levelCount() - 1 || level <= 0) + throw new ArrayIndexOutOfBoundsException("Invalid sorted generation " + level + " - maximum is " + (levelCount() - 1) + " and minimum is 1"); + return levels[level - 1]; + } + int levelCount() { return levels.length + 1; diff --git a/src/java/org/apache/cassandra/db/compaction/LeveledManifest.java b/src/java/org/apache/cassandra/db/compaction/LeveledManifest.java index 5e85511726a2..021604ea50a6 100644 --- a/src/java/org/apache/cassandra/db/compaction/LeveledManifest.java +++ b/src/java/org/apache/cassandra/db/compaction/LeveledManifest.java @@ -27,6 +27,7 @@ import java.util.List; import java.util.Map; import java.util.Set; +import java.util.TreeSet; import java.util.function.Function; import com.google.common.annotations.VisibleForTesting; @@ -573,14 +574,17 @@ private Collection getCandidatesFor(int level) return candidates; } + // We know that because we are in level L0+, this is a tree set with disjoint SSTables + TreeSet sstablesNextLevel = generations.getSortedLevel(level + 1); + // look for a non-suspect keyspace to compact with, starting with where we left off last time, // and wrapping back to the beginning of the generation if necessary - Map> sstablesNextLevel = genBounds(generations.get(level + 1)); Iterator levelIterator = generations.wrappingIterator(level, lastCompactedSSTables[level]); while (levelIterator.hasNext()) { SSTableReader sstable = levelIterator.next(); - Set candidates = Sets.union(Collections.singleton(sstable), overlappingWithBounds(sstable, sstablesNextLevel)); + Set candidates = getIntersectingSSTablesFromTreeSet(sstable, sstablesNextLevel); + candidates.add(sstable); if (Iterables.any(candidates, SSTableReader::isMarkedSuspect)) continue; @@ -592,6 +596,30 @@ private Collection getCandidatesFor(int level) return Collections.emptyList(); } + /** + * Precondition: SSTables in sstablesNextLevel must be disjoint + */ + @VisibleForTesting + protected static Set getIntersectingSSTablesFromTreeSet(SSTableReader sstable, TreeSet sstablesNextLevel) + { + Set candidates = new HashSet<>(); + + if (sstablesNextLevel.isEmpty()) + return candidates; + + SSTableReader start = sstablesNextLevel.floor(sstable); + Iterator it = sstablesNextLevel.tailSet(start != null ? start : sstablesNextLevel.first(), true).iterator(); + + while (it.hasNext()) + { + SSTableReader s = it.next(); + if (s.getFirst().compareTo(sstable.getLast()) > 0) break; + if (s.getLast().compareTo(sstable.getFirst()) >= 0) candidates.add(s); + } + + return candidates; + } + private Set getCompactingL0() { Set sstables = new HashSet<>(); diff --git a/test/unit/org/apache/cassandra/db/compaction/LeveledCompactionStrategyTest.java b/test/unit/org/apache/cassandra/db/compaction/LeveledCompactionStrategyTest.java index 0aeab98f7959..d352fd612560 100644 --- a/test/unit/org/apache/cassandra/db/compaction/LeveledCompactionStrategyTest.java +++ b/test/unit/org/apache/cassandra/db/compaction/LeveledCompactionStrategyTest.java @@ -75,7 +75,10 @@ import org.apache.cassandra.utils.Pair; import org.apache.cassandra.utils.TimeUUID; +import static accord.utils.Property.qt; import static java.util.Collections.singleton; +import static org.apache.cassandra.db.compaction.LeveledManifest.getIntersectingSSTablesFromTreeSet; +import static org.apache.cassandra.db.compaction.LeveledManifest.overlapping; import static org.apache.cassandra.schema.MockSchema.readerBounds; import static org.apache.cassandra.utils.TimeUUID.Generator.nextTimeUUID; import static org.assertj.core.api.Assertions.assertThat; @@ -595,6 +598,116 @@ public void testDisableSTCSInL0() throws IOException } } + @Test + public void testLinearScanTreeSetIntersectionEquivalence() + { + qt().withExamples(10).check(rs -> { + ColumnFamilyStore cfs = MockSchema.newCFS(); + List sstables = new ArrayList<>(); + + // Generates disjoint sorted SSTables that match what we see in L1 + int i = 0; + int start; + int end = 10; + while (i < 1000) + { + // Start with at least offset 1 to ensure that SSTables are disjoint + start = end + rs.nextInt(1, 15); + + // Include space in between SSTables, so we don't hit the case where when + // going through Overlap.SUBSET we generate a SSTable that has start > end + end = start + rs.nextInt(5, 20); + + sstables.add(MockSchema.sstableWithLevel(i++, start, end, 1, cfs)); + } + + LeveledGenerations generations = new LeveledGenerations(); + generations.addAll(sstables); + + // Every kind of overlap gets tested; randomness only selects the tokens + for (Overlap overlap : Overlap.values()) + { + SSTableReader sstableinL1 = rs.pick(sstables); + Token newStart; + Token newEnd; + + switch (overlap) + { + case NONE: + newStart = sstableinL1.getLast().getToken().increaseSlightly(); + newEnd = newStart.getToken().increaseSlightly(); + break; + case LEFT_END: + newStart = sstableinL1.getFirst().getToken().decreaseSlightly(); + newEnd = sstableinL1.getFirst().getToken().increaseSlightly(); + break; + case RIGHT_END: + newStart = sstableinL1.getLast().getToken().decreaseSlightly(); + newEnd = sstableinL1.getLast().getToken().increaseSlightly(); + break; + case SUBSET: + newStart = sstableinL1.getFirst().getToken().increaseSlightly(); + newEnd = sstableinL1.getLast().getToken().decreaseSlightly(); + break; + case SUPERSET: + newStart = sstableinL1.getFirst().getToken().decreaseSlightly(); + newEnd = sstableinL1.getLast().getToken().increaseSlightly(); + break; + case EXACT_MATCH: + newStart = sstableinL1.getFirst().getToken(); + newEnd = sstableinL1.getLast().getToken(); + break; + case WHOLE_LEVEL: + newStart = sstables.get(0).getFirst().getToken(); + newEnd = sstables.get(sstables.size() - 1).getLast().getToken(); + break; + case SINGLE_TOKEN: + newStart = sstableinL1.getFirst().getToken(); + newEnd = sstableinL1.getFirst().getToken(); + break; + default: + throw new IllegalStateException("Unhandled overlap " + overlap); + } + + SSTableReader sstable = MockSchema.sstableWithLevel(i++, newStart.getLongValue(), newEnd.getLongValue(), 0, cfs); + Collection treeSetIntersectingSSTables = getIntersectingSSTablesFromTreeSet(sstable, generations.getSortedLevel(1)); + Collection linearScanIntersectionSSTables = overlapping(sstable.getFirst().getToken(), sstable.getLast().getToken(), generations.getSortedLevel(1)); + + // The results should be equivalent to doing a linear scan + assertTrue("treeSet and linear scan produce different results for overlap " + overlap, + treeSetIntersectingSSTables.containsAll(linearScanIntersectionSSTables) && linearScanIntersectionSSTables.containsAll(treeSetIntersectingSSTables)); + } + }); + } + + @Test + public void testTreeSetIntersectionForEmptyNextLevelIsEmpty() + { + ColumnFamilyStore cfs = MockSchema.newCFS(); + LeveledGenerations generations = new LeveledGenerations(); + + SSTableReader sstable = MockSchema.sstableWithLevel(0, 1, 10, 0, cfs); + Collection treeSetIntersectingSSTables = getIntersectingSSTablesFromTreeSet(sstable, generations.getSortedLevel(1)); + + assertTrue("treeSet and linear scan produce different results", + treeSetIntersectingSSTables.isEmpty()); + } + + /** + * The ways a new L0 SSTable can overlap the SSTables in L1. + */ + private enum Overlap + { + NONE, + LEFT_END, + RIGHT_END, + SUBSET, + SUPERSET, + EXACT_MATCH, + WHOLE_LEVEL, + SINGLE_TOKEN + } + private int getTaskLevel(ColumnFamilyStore cfs) { int level = -1;