Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -96,6 +96,13 @@ Set<SSTableReader> get(int level)
return levels[level - 1];
}

TreeSet<SSTableReader> 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;
Expand Down
32 changes: 30 additions & 2 deletions src/java/org/apache/cassandra/db/compaction/LeveledManifest.java
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -573,14 +574,17 @@ private Collection<SSTableReader> getCandidatesFor(int level)
return candidates;
}

// We know that because we are in level L0+, this is a tree set with disjoint SSTables
TreeSet<SSTableReader> 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<SSTableReader, Bounds<Token>> sstablesNextLevel = genBounds(generations.get(level + 1));
Iterator<SSTableReader> levelIterator = generations.wrappingIterator(level, lastCompactedSSTables[level]);
while (levelIterator.hasNext())
{
SSTableReader sstable = levelIterator.next();
Set<SSTableReader> candidates = Sets.union(Collections.singleton(sstable), overlappingWithBounds(sstable, sstablesNextLevel));
Set<SSTableReader> candidates = getIntersectingSSTablesFromTreeSet(sstable, sstablesNextLevel);
candidates.add(sstable);

if (Iterables.any(candidates, SSTableReader::isMarkedSuspect))
continue;
Expand All @@ -592,6 +596,30 @@ private Collection<SSTableReader> getCandidatesFor(int level)
return Collections.emptyList();
}

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

maybe add some javadoc to describe that sstable will be included in the result

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Refactored this so that we add the sstable outside of getIntersectingSSTablesFromTreeSet to make it less confusing.

/**
* Precondition: SSTables in sstablesNextLevel must be disjoint
*/
@VisibleForTesting
protected static Set<SSTableReader> getIntersectingSSTablesFromTreeSet(SSTableReader sstable, TreeSet<SSTableReader> sstablesNextLevel)
{
Set<SSTableReader> candidates = new HashSet<>();

if (sstablesNextLevel.isEmpty())
return candidates;

SSTableReader start = sstablesNextLevel.floor(sstable);
Iterator<SSTableReader> 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<SSTableReader> getCompactingL0()
{
Set<SSTableReader> sstables = new HashSet<>();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -595,6 +598,116 @@ public void testDisableSTCSInL0() throws IOException
}
}

@Test
public void testLinearScanTreeSetIntersectionEquivalence()
{
qt().withExamples(10).check(rs -> {
ColumnFamilyStore cfs = MockSchema.newCFS();
List<SSTableReader> 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<SSTableReader> treeSetIntersectingSSTables = getIntersectingSSTablesFromTreeSet(sstable, generations.getSortedLevel(1));
Collection<SSTableReader> 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<SSTableReader> 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;
Expand Down