From 3aaefa390bf81c8d5ce60e7e4743a9adaf02715e Mon Sep 17 00:00:00 2001 From: Dmitry Konstantinov Date: Sun, 30 Aug 2026 21:20:15 +0100 Subject: [PATCH] Avoid double chunk read for BTI small partitions during a select query execution patch by Dmitry Konstantinov; reviewed by Branimir Lambov for CASSANDRA-21637 --- .../io/sstable/AbstractSSTableIterator.java | 17 +- .../io/sstable/format/bti/BtiTableReader.java | 123 +++++-- .../sstable/format/bti/SSTableIterator.java | 3 +- .../format/bti/SSTableReversedIterator.java | 3 +- .../format/bti/BtiTableReaderTest.java | 302 ++++++++++++++++++ 5 files changed, 416 insertions(+), 32 deletions(-) create mode 100644 test/unit/org/apache/cassandra/io/sstable/format/bti/BtiTableReaderTest.java diff --git a/src/java/org/apache/cassandra/io/sstable/AbstractSSTableIterator.java b/src/java/org/apache/cassandra/io/sstable/AbstractSSTableIterator.java index b156806e6284..d78f8df54e48 100644 --- a/src/java/org/apache/cassandra/io/sstable/AbstractSSTableIterator.java +++ b/src/java/org/apache/cassandra/io/sstable/AbstractSSTableIterator.java @@ -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; @@ -72,7 +73,6 @@ public abstract class AbstractSSTableIterator protected final Slices slices; - // file on every path where we created it. protected AbstractSSTableIterator(SSTableReader sstable, FileDataInput file, DecoratedKey key, @@ -80,6 +80,18 @@ protected AbstractSSTableIterator(SSTableReader sstable, 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(); @@ -91,6 +103,8 @@ 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; @@ -98,7 +112,6 @@ protected AbstractSSTableIterator(SSTableReader sstable, else { Reader reader = null; - boolean shouldCloseFile = file == null; try { // We seek to the beginning to the partition if either: diff --git a/src/java/org/apache/cassandra/io/sstable/format/bti/BtiTableReader.java b/src/java/org/apache/cassandra/io/sstable/format/bti/BtiTableReader.java index 1eee4515e342..9f04caec8da3 100644 --- a/src/java/org/apache/cassandra/io/sstable/format/bti/BtiTableReader.java +++ b/src/java/org/apache/cassandra/io/sstable/format/bti/BtiTableReader.java @@ -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; @@ -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. + * + * 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; + + 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()) @@ -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) @@ -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, @@ -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 diff --git a/src/java/org/apache/cassandra/io/sstable/format/bti/SSTableIterator.java b/src/java/org/apache/cassandra/io/sstable/format/bti/SSTableIterator.java index be62f70945f0..4e38346f7aac 100644 --- a/src/java/org/apache/cassandra/io/sstable/format/bti/SSTableIterator.java +++ b/src/java/org/apache/cassandra/io/sstable/format/bti/SSTableIterator.java @@ -42,13 +42,14 @@ class SSTableIterator extends AbstractSSTableIterator 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) diff --git a/src/java/org/apache/cassandra/io/sstable/format/bti/SSTableReversedIterator.java b/src/java/org/apache/cassandra/io/sstable/format/bti/SSTableReversedIterator.java index 8e69054d5263..36f4cc9ffdd3 100644 --- a/src/java/org/apache/cassandra/io/sstable/format/bti/SSTableReversedIterator.java +++ b/src/java/org/apache/cassandra/io/sstable/format/bti/SSTableReversedIterator.java @@ -51,13 +51,14 @@ class SSTableReversedIterator extends AbstractSSTableIterator 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) diff --git a/test/unit/org/apache/cassandra/io/sstable/format/bti/BtiTableReaderTest.java b/test/unit/org/apache/cassandra/io/sstable/format/bti/BtiTableReaderTest.java new file mode 100644 index 000000000000..843f07e14bb0 --- /dev/null +++ b/test/unit/org/apache/cassandra/io/sstable/format/bti/BtiTableReaderTest.java @@ -0,0 +1,302 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.cassandra.io.sstable.format.bti; + +import java.util.Collections; +import java.util.Set; + +import org.junit.Before; +import org.junit.Test; + +import org.apache.cassandra.config.TestDatabaseDescriptor; +import org.apache.cassandra.cql3.CQLTester; +import org.apache.cassandra.db.ColumnFamilyStore; +import org.apache.cassandra.db.DecoratedKey; +import org.apache.cassandra.db.Slices; +import org.apache.cassandra.db.filter.ColumnFilter; +import org.apache.cassandra.db.rows.Unfiltered; +import org.apache.cassandra.db.rows.UnfilteredRowIterator; +import org.apache.cassandra.io.sstable.SSTableReadsListener; +import org.apache.cassandra.io.sstable.format.SSTableReader; +import org.apache.cassandra.io.util.FileDataInput; +import org.apache.cassandra.io.util.RandomAccessReader; +import org.apache.cassandra.utils.ByteBufferUtil; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertFalse; +import static org.junit.Assert.assertNotNull; +import static org.junit.Assert.assertNull; +import static org.junit.Assert.assertSame; +import static org.junit.Assert.assertTrue; + +/** + * Covers the reuse of the data file reader between partition key verification and the row iterator. + *

+ * A partition small enough to have no row index entry has its key verified by reading the data file at the position + * the trie points to, which is where the row iterator then has to start. The reader is handed on rather than closed, + * so the same position - and on a compressed table the same decompression and checksum - is not paid for twice. + */ +public class BtiTableReaderTest extends CQLTester +{ + /** Over the 4KiB default column_index_size, so this partition gets an entry in the row index. */ + private static final int ROWS_IN_WIDE_PARTITION = 400; + private static final int NARROW_PARTITIONS = 100; + private static final int WIDE_PARTITION_KEY = -1; + + @Before + public void selectBtiFormatAndPopulate() + { + TestDatabaseDescriptor.setUnsafeSelectedSSTableFormat(new BtiFormat.BtiFormatFactory().getInstance(Collections.emptyMap())); + + createTable("CREATE TABLE %s (pk int, ck int, v text, PRIMARY KEY (pk, ck))"); + + String value = new String(new char[64]).replace('\0', 'x'); + for (int pk = 0; pk < NARROW_PARTITIONS; pk++) + execute("INSERT INTO %s (pk, ck, v) VALUES (?, ?, ?)", pk, 0, value); + + for (int ck = 0; ck < ROWS_IN_WIDE_PARTITION; ck++) + execute("INSERT INTO %s (pk, ck, v) VALUES (?, ?, ?)", WIDE_PARTITION_KEY, ck, value); + + flush(); + } + + private BtiTableReader sstable() + { + ColumnFamilyStore cfs = getCurrentColumnFamilyStore(); + Set live = cfs.getLiveSSTables(); + assertEquals("expected the test data to be in a single sstable", 1, live.size()); + SSTableReader reader = live.iterator().next(); + assertTrue("expected a BTI sstable, got " + reader.getClass().getSimpleName(), reader instanceof BtiTableReader); + return (BtiTableReader) reader; + } + + private DecoratedKey key(int pk) + { + return sstable().decorateKey(ByteBufferUtil.bytes(pk)); + } + + private static void assertClosed(FileDataInput in) + { + assertFalse("reader should have been closed", isOpen(in)); + } + + private static void assertOpen(FileDataInput in) + { + assertTrue("reader should still be open", isOpen(in)); + } + + /** + * A closed {@link RandomAccessReader} has dropped its buffer and refuses to seek; that is the externally visible + * sign of closing it. + */ + private static boolean isOpen(FileDataInput in) + { + RandomAccessReader reader = (RandomAccessReader) in; + try + { + long position = reader.getFilePointer(); + reader.seek(0); // accepted at any position while open, throws once closed + reader.seek(position); // leave the reader where we found it + return true; + } + catch (IllegalStateException e) + { + return false; + } + } + + @Test + public void testRetainsDataFileReaderForPartitionWithoutRowIndex() throws Exception + { + BtiTableReader sstable = sstable(); + DecoratedKey key = key(0); + + BtiTableReader.ExactPosition pos = sstable.getExactPosition(key, SSTableReadsListener.NOOP_LISTENER, false, true); + + assertNotNull(pos.entry); + assertFalse("a narrow partition should have no row index entry", pos.entry.isIndexed()); + assertNotNull("the data file reader used for key verification should have been handed over", pos.dataInput); + assertOpen(pos.dataInput); + + // it is left where the verification put it, just past the partition key, which is where the iterator seeks from + assertTrue(pos.dataInput.getFilePointer() > pos.entry.position); + + pos.dataInput.close(); + } + + @Test + public void testDoesNotRetainReaderForRowIndexedPartition() + { + BtiTableReader sstable = sstable(); + + BtiTableReader.ExactPosition pos = sstable.getExactPosition(key(WIDE_PARTITION_KEY), + SSTableReadsListener.NOOP_LISTENER, + false, + true); + + assertNotNull(pos.entry); + assertTrue("a wide partition should have a row index entry", pos.entry.isIndexed()); + // that key was verified against the row index file, which is of no use to a data file iterator + assertNull(pos.dataInput); + } + + @Test + public void testDoesNotRetainWhenNotRequested() + { + BtiTableReader sstable = sstable(); + + BtiTableReader.ExactPosition pos = sstable.getExactPosition(key(0), SSTableReadsListener.NOOP_LISTENER, false, false); + + assertNotNull(pos.entry); + assertNull(pos.dataInput); + } + + @Test + public void testMissingKeyRetainsNothing() + { + BtiTableReader sstable = sstable(); + + BtiTableReader.ExactPosition pos = sstable.getExactPosition(key(NARROW_PARTITIONS + 1000), + SSTableReadsListener.NOOP_LISTENER, + false, + true); + + assertNull(pos.entry); + assertNull(pos.dataInput); + assertSame(BtiTableReader.ExactPosition.NOT_FOUND, pos); + } + + /** + * The ownership half of the change: the iterator must close a reader handed to it, or every point read leaks a + * reader and its buffer. + */ + @Test + public void testIteratorClosesRetainedReader() + { + for (boolean reversed : new boolean[]{ false, true }) + { + BtiTableReader sstable = sstable(); + DecoratedKey key = key(1); + + BtiTableReader.ExactPosition pos = sstable.getExactPosition(key, SSTableReadsListener.NOOP_LISTENER, false, true); + assertNotNull(pos.dataInput); + + int rows = 0; + try (UnfilteredRowIterator iter = sstable.rowIterator(pos.dataInput, true, key, pos.entry, + Slices.ALL, ColumnFilter.all(sstable.metadata()), reversed)) + { + while (iter.hasNext()) + { + Unfiltered u = iter.next(); + assertNotNull(u); + rows++; + } + } + + assertEquals("reversed=" + reversed, 1, rows); + assertClosed(pos.dataInput); + } + } + + /** + * A reader the caller keeps must survive the iterator (a scanner scenario). + */ + @Test + public void testIteratorLeavesCallerOwnedReaderOpen() throws Exception + { + BtiTableReader sstable = sstable(); + DecoratedKey key = key(2); + + BtiTableReader.ExactPosition pos = sstable.getExactPosition(key, SSTableReadsListener.NOOP_LISTENER, false, true); + assertNotNull(pos.dataInput); + + try (UnfilteredRowIterator iter = sstable.rowIterator(pos.dataInput, false, key, pos.entry, + Slices.ALL, ColumnFilter.all(sstable.metadata()), false)) + { + while (iter.hasNext()) + iter.next(); + } + + assertOpen(pos.dataInput); + pos.dataInput.close(); + } + + @Test + public void testNarrowPartitionsReadCorrectly() + { + String value = new String(new char[64]).replace('\0', 'x'); + for (int pk = 0; pk < NARROW_PARTITIONS; pk++) + assertRows(execute("SELECT pk, ck, v FROM %s WHERE pk = ?", pk), row(pk, 0, value)); + } + + /** + * A row-indexed partition is the case where nothing is handed over, so the iterator has to open its own reader + * and close it again. Reversed as well as forward, since they build different readers. + */ + @Test + public void testWidePartitionReadsCorrectly() + { + assertEquals(ROWS_IN_WIDE_PARTITION, + execute("SELECT ck FROM %s WHERE pk = ?", WIDE_PARTITION_KEY).size()); + assertRows(execute("SELECT ck FROM %s WHERE pk = ? ORDER BY ck DESC LIMIT 1", WIDE_PARTITION_KEY), + row(ROWS_IN_WIDE_PARTITION - 1)); + assertRows(execute("SELECT ck FROM %s WHERE pk = ? ORDER BY ck ASC LIMIT 1", WIDE_PARTITION_KEY), + row(0)); + } + + @Test + public void testMissingPartitionReturnsNothing() + { + assertEmpty(execute("SELECT * FROM %s WHERE pk = ?", NARROW_PARTITIONS + 1000)); + } + + @Test + public void testFullScanStillWorks() + { + assertEquals(NARROW_PARTITIONS + ROWS_IN_WIDE_PARTITION, execute("SELECT pk, ck FROM %s").size()); + } + + /** + * Repeated point reads through the normal entry point - the path that now hands the reader over every time. + */ + @Test + public void testRepeatedPointReadsThroughRowIterator() + { + BtiTableReader sstable = sstable(); + + for (int round = 0; round < 5; round++) + { + for (int pk = 0; pk < NARROW_PARTITIONS; pk++) + { + DecoratedKey key = key(pk); + int rows = 0; + try (UnfilteredRowIterator iter = sstable.rowIterator(key, Slices.ALL, ColumnFilter.all(sstable.metadata()), + false, SSTableReadsListener.NOOP_LISTENER)) + { + while (iter.hasNext()) + { + iter.next(); + rows++; + } + } + assertEquals("pk=" + pk, 1, rows); + } + } + } +}