Skip to content
Closed
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 @@ -46,6 +46,7 @@
import org.apache.cassandra.io.sstable.format.Version;
import org.apache.cassandra.io.util.FileDataInput;
import org.apache.cassandra.io.util.FileHandle;
import org.apache.cassandra.io.util.FileUtils;
import org.apache.cassandra.schema.TableMetadata;
import org.apache.cassandra.utils.ByteBufferUtil;

Expand All @@ -72,14 +73,25 @@ public abstract class AbstractSSTableIterator<RIE extends AbstractRowIndexEntry>

protected final Slices slices;

// file on every path where we created it.
protected AbstractSSTableIterator(SSTableReader sstable,
FileDataInput file,
DecoratedKey key,
RIE indexEntry,
Slices slices,
ColumnFilter columnFilter,
FileHandle ifile)
{
this(sstable, file, file == null, key, indexEntry, slices, columnFilter, ifile);
}

protected AbstractSSTableIterator(SSTableReader sstable,
FileDataInput file,
boolean shouldCloseFile,
DecoratedKey key,
RIE indexEntry,
Slices slices,
ColumnFilter columnFilter,
FileHandle ifile)
{
this.sstable = sstable;
this.metadata = sstable.metadata();
Expand All @@ -91,14 +103,15 @@ protected AbstractSSTableIterator(SSTableReader sstable,

if (indexEntry == null)
{
if (shouldCloseFile)
FileUtils.closeQuietly(file);
this.partitionLevelDeletion = DeletionTime.LIVE;
this.reader = null;
this.staticRow = Rows.EMPTY_STATIC_ROW;
}
else
{
Reader reader = null;
boolean shouldCloseFile = file == null;
try
{
// We seek to the beginning to the partition if either:
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -55,6 +55,7 @@
import org.apache.cassandra.io.sstable.format.SSTableReaderWithFilter;
import org.apache.cassandra.io.util.FileDataInput;
import org.apache.cassandra.io.util.FileHandle;
import org.apache.cassandra.io.util.FileUtils;
import org.apache.cassandra.io.util.RandomAccessReader;
import org.apache.cassandra.utils.ByteBufferUtil;
import org.apache.cassandra.utils.IFilter;
Expand Down Expand Up @@ -232,20 +233,66 @@ public DecoratedKey keyAtPositionFromSecondaryIndex(long keyPositionFromSecondar
}
}

/**
* The result of an exact-position lookup: the index entry, and optionally the data file reader that was used to
* verify the partition key.
*
* When a partition is small enough to have no entry in the row index, the trie payload holds its data file
* position and the key is checked by reading the data file there - which is the same place where the row
* iterator has to start. {@link #dataInput} lets that reader be reused instead of closed, so the iterator does
* not have to open the same position again; on a compressed table that would decompress and checksum the same
* chunk twice. It is null whenever there is nothing to reuse, and when it is not null the caller owns it and
* must close it.
Comment thread
netudima marked this conversation as resolved.
*
* This only applies to the exact-match lookup used by {@link #rowIterator}. In the other cases there is either no
* reader worth keeping, or no one to give it to. A partition that has a row index entry has its key checked in the
* row index file, not in the data file.
* The generic {@link org.apache.cassandra.io.sstable.format.SSTableReader#getPosition} is shared with the BIG format,
* which never reads the data file during a lookup. Its callers either only use the returned position, or, like
* compaction's shadow-source iterators and the secondary index builders, already keep one long-lived data file
* reader that they seek themselves for many partitions.
*/
static final class ExactPosition
{
static final ExactPosition NOT_FOUND = new ExactPosition(null, null);

final TrieIndexEntry entry;
final FileDataInput dataInput;
Comment thread
netudima marked this conversation as resolved.

private ExactPosition(TrieIndexEntry entry, FileDataInput dataInput)
{
this.entry = entry;
this.dataInput = dataInput;
}
}

TrieIndexEntry getExactPosition(DecoratedKey dk,
SSTableReadsListener listener,
boolean updateStats)
{
return getExactPosition(dk, listener, updateStats, false).entry;
}

/**
* @param retainDataInput when true, and the key was verified by reading the data file, the reader used for that
* verification is returned still open in {@link ExactPosition#dataInput} and ownership
* passes to the caller. Otherwise every reader opened here is closed before returning.
*/
ExactPosition getExactPosition(DecoratedKey dk,
SSTableReadsListener listener,
boolean updateStats,
boolean retainDataInput)
{
if ((filterFirst() && getFirst().compareTo(dk) > 0) || (filterLast() && getLast().compareTo(dk) < 0))
{
notifySkipped(SkippingReason.MIN_MAX_KEYS, listener, EQ, updateStats);
return null;
return ExactPosition.NOT_FOUND;
}

if (!isPresentInFilter(dk))
{
notifySkipped(SkippingReason.BLOOM_FILTER, listener, EQ, updateStats);
return null;
return ExactPosition.NOT_FOUND;
}

try (PartitionIndex.Reader reader = partitionIndex.openReader())
Expand All @@ -254,36 +301,35 @@ TrieIndexEntry getExactPosition(DecoratedKey dk,
if (indexPos == PartitionIndex.NOT_FOUND)
{
notifySkipped(SkippingReason.PARTITION_INDEX_LOOKUP, listener, EQ, updateStats);
return null;
return ExactPosition.NOT_FOUND;
}

FileHandle fh;
long seekPosition;
if (indexPos >= 0)
{
fh = rowIndexFile;
seekPosition = indexPos;
}
else
{
fh = dfile;
seekPosition = ~indexPos;
}
boolean fromDataFile = indexPos < 0;
FileHandle fh = fromDataFile ? dfile : rowIndexFile;
long seekPosition = fromDataFile ? ~indexPos : indexPos;

try (FileDataInput in = fh.createReader(seekPosition))
boolean retain = retainDataInput && fromDataFile;

FileDataInput in = fh.createReader(seekPosition);
boolean handedOver = false;
try
{
if (ByteBufferUtil.equalsWithShortLength(in, dk.getKey()))
{
TrieIndexEntry rie = indexPos >= 0 ? TrieIndexEntry.deserialize(in, in.getFilePointer(), descriptor.version)
: new TrieIndexEntry(~indexPos);
notifySelected(SelectionReason.INDEX_ENTRY_FOUND, listener, EQ, updateStats, rie);
return rie;
}
else
if (!ByteBufferUtil.equalsWithShortLength(in, dk.getKey()))
{
notifySkipped(SkippingReason.INDEX_ENTRY_NOT_FOUND, listener, EQ, updateStats);
return null;
return ExactPosition.NOT_FOUND;
}

TrieIndexEntry rie = fromDataFile ? new TrieIndexEntry(~indexPos)
: TrieIndexEntry.deserialize(in, in.getFilePointer(), descriptor.version);
notifySelected(SelectionReason.INDEX_ENTRY_FOUND, listener, EQ, updateStats, rie);
handedOver = retain;
return new ExactPosition(rie, retain ? in : null);
}
finally
{
if (!handedOver)
in.close();
}
}
catch (IOException | IllegalArgumentException | ArrayIndexOutOfBoundsException | AssertionError e)
Expand Down Expand Up @@ -373,7 +419,10 @@ public UnfilteredRowIterator rowIterator(DecoratedKey key,
boolean reversed,
SSTableReadsListener listener)
{
return rowIterator(null, key, getExactPosition(key, listener, true), slices, selectedColumns, reversed);
// retain the reader that verified the key: for a partition with no row index entry it is already on the
// data file at the position the iterator has to start from
ExactPosition pos = getExactPosition(key, listener, true, true);
return rowIterator(pos.dataInput, pos.dataInput != null, key, pos.entry, slices, selectedColumns, reversed);
}

public UnfilteredRowIterator rowIterator(FileDataInput dataFileInput,
Expand All @@ -382,14 +431,32 @@ public UnfilteredRowIterator rowIterator(FileDataInput dataFileInput,
Slices slices,
ColumnFilter selectedColumns,
boolean reversed)
{
// callers of this overload - scanners walking many partitions with one reader - keep ownership of it
return rowIterator(dataFileInput, false, key, indexEntry, slices, selectedColumns, reversed);
}

UnfilteredRowIterator rowIterator(FileDataInput dataFileInput,
boolean ownsDataFileInput,
DecoratedKey key,
TrieIndexEntry indexEntry,
Slices slices,
ColumnFilter selectedColumns,
boolean reversed)
{
if (indexEntry == null)
{
if (ownsDataFileInput)
FileUtils.closeQuietly(dataFileInput);
return UnfilteredRowIterators.noRowsIterator(metadata(), key, Rows.EMPTY_STATIC_ROW, DeletionTime.LIVE, reversed);
}

// the iterator owns the reader either because we handed ours over, or because it has to open its own
boolean shouldCloseFile = ownsDataFileInput || dataFileInput == null;
if (reversed)
return new SSTableReversedIterator(this, dataFileInput, key, indexEntry, slices, selectedColumns, rowIndexFile);
return new SSTableReversedIterator(this, dataFileInput, shouldCloseFile, key, indexEntry, slices, selectedColumns, rowIndexFile);
else
return new SSTableIterator(this, dataFileInput, key, indexEntry, slices, selectedColumns, rowIndexFile);
return new SSTableIterator(this, dataFileInput, shouldCloseFile, key, indexEntry, slices, selectedColumns, rowIndexFile);
}

@VisibleForTesting
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -42,13 +42,14 @@ class SSTableIterator extends AbstractSSTableIterator<AbstractRowIndexEntry>

public SSTableIterator(BtiTableReader sstable,
FileDataInput file,
boolean shouldCloseFile,
DecoratedKey key,
AbstractRowIndexEntry indexEntry,
Slices slices,
ColumnFilter columns,
FileHandle ifile)
{
super(sstable, file, key, indexEntry, slices, columns, ifile);
super(sstable, file, shouldCloseFile, key, indexEntry, slices, columns, ifile);
}

protected Reader createReaderInternal(AbstractRowIndexEntry indexEntry, FileDataInput file, boolean shouldCloseFile, Version version)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -51,13 +51,14 @@ class SSTableReversedIterator extends AbstractSSTableIterator<TrieIndexEntry>

public SSTableReversedIterator(BtiTableReader sstable,
FileDataInput file,
boolean shouldCloseFile,
DecoratedKey key,
TrieIndexEntry indexEntry,
Slices slices,
ColumnFilter columns,
FileHandle ifile)
{
super(sstable, file, key, indexEntry, slices, columns, ifile);
super(sstable, file, shouldCloseFile, key, indexEntry, slices, columns, ifile);
}

protected Reader createReaderInternal(TrieIndexEntry indexEntry, FileDataInput file, boolean shouldCloseFile, Version version)
Expand Down
Loading