diff --git a/CHANGES.txt b/CHANGES.txt index 0f9c4f156a..b7b0e42949 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -1,4 +1,5 @@ 2.0 + * Improve LeveledScanner work estimation (CASSANDRA-5250) * Replace compaction lock with runWithCompactionsDisabled (CASSANDRA-3430) * Change Message IDs to ints (CASSANDRA-5307) * Move sstable level information into the Stats component, removing the diff --git a/src/java/org/apache/cassandra/db/compaction/LeveledCompactionStrategy.java b/src/java/org/apache/cassandra/db/compaction/LeveledCompactionStrategy.java index ffe45ad475..d08154285a 100644 --- a/src/java/org/apache/cassandra/db/compaction/LeveledCompactionStrategy.java +++ b/src/java/org/apache/cassandra/db/compaction/LeveledCompactionStrategy.java @@ -34,7 +34,6 @@ import org.apache.cassandra.dht.Token; import org.apache.cassandra.exceptions.ConfigurationException; import org.apache.cassandra.io.sstable.SSTable; import org.apache.cassandra.io.sstable.SSTableReader; -import org.apache.cassandra.io.sstable.SSTableScanner; import org.apache.cassandra.notifications.INotification; import org.apache.cassandra.notifications.INotificationConsumer; import org.apache.cassandra.notifications.SSTableAddedNotification; @@ -172,7 +171,9 @@ public class LeveledCompactionStrategy extends AbstractCompactionStrategy implem else { // Create a LeveledScanner that only opens one sstable at a time, in sorted order - scanners.add(new LeveledScanner(byLevel.get(level), range)); + List intersecting = LeveledScanner.intersecting(byLevel.get(level), range); + if (!intersecting.isEmpty()) + scanners.add(new LeveledScanner(intersecting, range)); } } @@ -194,19 +195,46 @@ public class LeveledCompactionStrategy extends AbstractCompactionStrategy implem public LeveledScanner(Collection sstables, Range range) { this.range = range; - this.sstables = new ArrayList(sstables); - Collections.sort(this.sstables, SSTable.sstableComparator); - sstableIterator = this.sstables.iterator(); - currentScanner = sstableIterator.next().getDirectScanner(range); + // add only sstables that intersect our range, and estimate how much data that involves + this.sstables = new ArrayList(sstables.size()); long length = 0; for (SSTableReader sstable : sstables) - length += sstable.uncompressedLength(); + { + this.sstables.add(sstable); + long estimatedKeys = sstable.estimatedKeys(); + double estKeysInRangeRatio = 1.0; + + if (estimatedKeys > 0 && range != null) + estKeysInRangeRatio = ((double) sstable.estimatedKeysForRanges(Collections.singleton(range))) / estimatedKeys; + + length += sstable.uncompressedLength() * estKeysInRangeRatio; + } + totalLength = length; + Collections.sort(this.sstables, SSTable.sstableComparator); + sstableIterator = this.sstables.iterator(); + assert sstableIterator.hasNext(); // caller should check intersecting first + currentScanner = sstableIterator.next().getDirectScanner(range); + } + + public static List intersecting(Collection sstables, Range range) + { + ArrayList filtered = new ArrayList(); + for (SSTableReader sstable : sstables) + { + Range sstableRange = new Range(sstable.first.getToken(), sstable.last.getToken(), sstable.partitioner); + if (range == null || sstableRange.intersects(range)) + filtered.add(sstable); + } + return filtered; } protected OnDiskAtomIterator computeNext() { + if (currentScanner == null) + return endOfData(); + try { while (true) diff --git a/src/java/org/apache/cassandra/io/sstable/SSTableBoundedScanner.java b/src/java/org/apache/cassandra/io/sstable/SSTableBoundedScanner.java index a5719011a2..0e31896b8a 100644 --- a/src/java/org/apache/cassandra/io/sstable/SSTableBoundedScanner.java +++ b/src/java/org/apache/cassandra/io/sstable/SSTableBoundedScanner.java @@ -35,11 +35,10 @@ public class SSTableBoundedScanner extends SSTableScanner private final Iterator> rangeIterator; private Pair currentRange; - SSTableBoundedScanner(SSTableReader sstable, boolean skipCache, Iterator> rangeIterator) + SSTableBoundedScanner(SSTableReader sstable, boolean skipCache, Range range) { super(sstable, skipCache); - this.rangeIterator = rangeIterator; - assert rangeIterator.hasNext(); // use EmptyCompactionScanner otherwise + this.rangeIterator = sstable.getPositionsForRanges(Collections.singletonList(range)).iterator(); currentRange = rangeIterator.next(); dfile.seek(currentRange.left); } diff --git a/src/java/org/apache/cassandra/io/sstable/SSTableReader.java b/src/java/org/apache/cassandra/io/sstable/SSTableReader.java index e2ef70c882..bbef4ec34f 100644 --- a/src/java/org/apache/cassandra/io/sstable/SSTableReader.java +++ b/src/java/org/apache/cassandra/io/sstable/SSTableReader.java @@ -973,10 +973,7 @@ public class SSTableReader extends SSTable if (range == null) return getDirectScanner(); - Iterator> rangeIterator = getPositionsForRanges(Collections.singletonList(range)).iterator(); - return rangeIterator.hasNext() - ? new SSTableBoundedScanner(this, true, rangeIterator) - : new EmptyCompactionScanner(getFilename()); + return new SSTableBoundedScanner(this, true, range); } public FileDataInput getFileDataInput(long position) @@ -1210,48 +1207,4 @@ public class SSTableReader extends SSTable sstable.releaseReference(); } } - - private static class EmptyCompactionScanner implements ICompactionScanner - { - private final String filename; - - private EmptyCompactionScanner(String filename) - { - this.filename = filename; - } - - public long getLengthInBytes() - { - return 0; - } - - public long getCurrentPosition() - { - return 0; - } - - public String getBackingFiles() - { - return filename; - } - - public void close() - { - } - - public boolean hasNext() - { - return false; - } - - public OnDiskAtomIterator next() - { - throw new IndexOutOfBoundsException(); - } - - public void remove() - { - throw new UnsupportedOperationException(); - } - } }