mirror of https://github.com/apache/cassandra
Improve LeveledScanner work estimation
patch by Marcus Eriksson; reviewed by jbellis for CASSANDRA-5250
This commit is contained in:
parent
d72e9381fa
commit
effdb08b34
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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<SSTableReader> 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<SSTableReader> sstables, Range<Token> range)
|
||||
{
|
||||
this.range = range;
|
||||
this.sstables = new ArrayList<SSTableReader>(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<SSTableReader>(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<SSTableReader> intersecting(Collection<SSTableReader> sstables, Range<Token> range)
|
||||
{
|
||||
ArrayList<SSTableReader> filtered = new ArrayList<SSTableReader>();
|
||||
for (SSTableReader sstable : sstables)
|
||||
{
|
||||
Range<Token> sstableRange = new Range<Token>(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)
|
||||
|
|
|
|||
|
|
@ -35,11 +35,10 @@ public class SSTableBoundedScanner extends SSTableScanner
|
|||
private final Iterator<Pair<Long, Long>> rangeIterator;
|
||||
private Pair<Long, Long> currentRange;
|
||||
|
||||
SSTableBoundedScanner(SSTableReader sstable, boolean skipCache, Iterator<Pair<Long, Long>> rangeIterator)
|
||||
SSTableBoundedScanner(SSTableReader sstable, boolean skipCache, Range<Token> 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);
|
||||
}
|
||||
|
|
|
|||
|
|
@ -973,10 +973,7 @@ public class SSTableReader extends SSTable
|
|||
if (range == null)
|
||||
return getDirectScanner();
|
||||
|
||||
Iterator<Pair<Long, Long>> 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();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
Loading…
Reference in New Issue