diff --git a/CHANGES.txt b/CHANGES.txt index c9ec2e1ea0..970933fab5 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -1,10 +1,12 @@ 5.0-beta1 + * Fix incorrect seeking through the sstable iterator by IndexState (CASSANDRA-18932) * Upgrade Python driver to 3.28.0 (CASSANDRA-18960) * Add retries to IR messages (CASSANDRA-18962) * Add metrics and logging to repair retries (CASSANDRA-18952) * Remove deprecated code in Cassandra 1.x and 2.x (CASSANDRA-18959) * ClientRequestSize metrics should not treat CONTAINS restrictions as being equality-based (CASSANDRA-18896) + 5.0-alpha2 * Add support for vector search in SAI (CASSANDRA-18715) * Remove crc_check_chance from CompressionParams (CASSANDRA-18872) diff --git a/src/java/org/apache/cassandra/io/sstable/AbstractSSTableIterator.java b/src/java/org/apache/cassandra/io/sstable/AbstractSSTableIterator.java index 61dbda52fe..2690be69d7 100644 --- a/src/java/org/apache/cassandra/io/sstable/AbstractSSTableIterator.java +++ b/src/java/org/apache/cassandra/io/sstable/AbstractSSTableIterator.java @@ -340,6 +340,15 @@ public abstract class AbstractSSTableIterator deserializer = UnfilteredDeserializer.create(metadata, file, sstable.header, helper); } + /** + * Seek to the given position in the file. Initializes the file reader along with a deserializer if they are + * not initialized yet. Otherwise, the deserializer is reset by calling its {@link UnfilteredDeserializer#clearState()} + * method. + *

+ * Note that the only valid usage of this method is to seek to the beginning of the serialized record so that + * the deserializer can read the next unfiltered from that position. Setting an arbitrary position will lead + * to unexpected results and/or corrupted reads. + */ public void seekToPosition(long position) throws IOException { // This may be the first time we're actually looking into the file diff --git a/src/java/org/apache/cassandra/io/sstable/format/big/IndexState.java b/src/java/org/apache/cassandra/io/sstable/format/big/IndexState.java index c7697cfec0..f738f8ab7c 100644 --- a/src/java/org/apache/cassandra/io/sstable/format/big/IndexState.java +++ b/src/java/org/apache/cassandra/io/sstable/format/big/IndexState.java @@ -66,7 +66,6 @@ public class IndexState implements AutoCloseable { reader.seekToPosition(columnOffset(blockIdx)); mark = reader.file.mark(); - reader.deserializer.clearState(); } currentIndexIdx = blockIdx; @@ -114,9 +113,9 @@ public class IndexState implements AutoCloseable } else { - reader.seekToPosition(startOfBlock); + reader.file.seek(startOfBlock); mark = reader.file.mark(); - reader.seekToPosition(currentFilePointer); + reader.file.seek(currentFilePointer); } } } diff --git a/test/unit/org/apache/cassandra/cql3/CQLTester.java b/test/unit/org/apache/cassandra/cql3/CQLTester.java index 4fc35cede6..c9b43e4cd2 100644 --- a/test/unit/org/apache/cassandra/cql3/CQLTester.java +++ b/test/unit/org/apache/cassandra/cql3/CQLTester.java @@ -31,6 +31,7 @@ import java.util.ArrayList; import java.util.Arrays; import java.util.Collection; import java.util.Collections; +import java.util.Comparator; import java.util.Date; import java.util.HashMap; import java.util.HashSet; @@ -1494,19 +1495,19 @@ public abstract class CQLTester return sessionNet(protocolVersion).execute(statement); } - protected com.datastax.driver.core.ResultSet executeNetWithPaging(ProtocolVersion version, String query, int pageSize) + protected com.datastax.driver.core.ResultSet executeNetWithPaging(ProtocolVersion version, String query, int pageSize, Object... values) { - return sessionNet(version).execute(new SimpleStatement(formatQuery(query)).setFetchSize(pageSize)); + return sessionNet(version).execute(new SimpleStatement(formatQuery(query), values).setFetchSize(pageSize)); } - protected com.datastax.driver.core.ResultSet executeNetWithPaging(ProtocolVersion version, String query, String KS, int pageSize) + protected com.datastax.driver.core.ResultSet executeNetWithPaging(ProtocolVersion version, String query, String KS, int pageSize, Object... values) { - return sessionNet(version).execute(new SimpleStatement(formatQuery(KS, query)).setKeyspace(KS).setFetchSize(pageSize)); + return sessionNet(version).execute(new SimpleStatement(formatQuery(KS, query), values).setKeyspace(KS).setFetchSize(pageSize)); } - protected com.datastax.driver.core.ResultSet executeNetWithPaging(String query, int pageSize) + protected com.datastax.driver.core.ResultSet executeNetWithPaging(String query, int pageSize, Object... values) { - return sessionNet().execute(new SimpleStatement(formatQuery(query)).setFetchSize(pageSize)); + return sessionNet().execute(new SimpleStatement(formatQuery(query), values).setFetchSize(pageSize)); } protected com.datastax.driver.core.ResultSet executeNetWithoutPaging(String query) @@ -1622,6 +1623,23 @@ public abstract class CQLTester return rs; } + public static int compareNetRows(Row r1, Row r2) + { + Comparator bufComp = Comparator.nullsFirst(Comparator.naturalOrder()); + for (int c = 0; c < Math.min(r1.getColumnDefinitions().size(), r2.getColumnDefinitions().size()); c++) + { + DataType t1 = r1.getColumnDefinitions().getType(c); + DataType t2 = r2.getColumnDefinitions().getType(c); + if (!t1.equals(t2)) + return t1.getName().toString().compareTo(t2.getName().toString()); + + int cmp = bufComp.compare(r1.getBytesUnsafe(c), r2.getBytesUnsafe(c)); + if (cmp != 0) + return cmp; + } + return Integer.compare(r1.getColumnDefinitions().size(), r2.getColumnDefinitions().size()); + } + protected void assertRowsNet(ResultSet result, Object[]... rows) { assertRowsNet(getDefaultVersion(), result, rows); diff --git a/test/unit/org/apache/cassandra/cql3/ViewAbstractParameterizedTest.java b/test/unit/org/apache/cassandra/cql3/ViewAbstractParameterizedTest.java index a8c73ff84c..3455bb674d 100644 --- a/test/unit/org/apache/cassandra/cql3/ViewAbstractParameterizedTest.java +++ b/test/unit/org/apache/cassandra/cql3/ViewAbstractParameterizedTest.java @@ -60,9 +60,9 @@ public abstract class ViewAbstractParameterizedTest extends ViewAbstractTest } @Override - protected com.datastax.driver.core.ResultSet executeNetWithPaging(String query, int pageSize) + protected com.datastax.driver.core.ResultSet executeNetWithPaging(String query, int pageSize, Object... values) { - return executeNetWithPaging(version, query, pageSize); + return executeNetWithPaging(version, query, pageSize, values); } @Override diff --git a/test/unit/org/apache/cassandra/io/sstable/format/ColumnIndexTest.java b/test/unit/org/apache/cassandra/io/sstable/format/ColumnIndexTest.java new file mode 100644 index 0000000000..442d4d35bf --- /dev/null +++ b/test/unit/org/apache/cassandra/io/sstable/format/ColumnIndexTest.java @@ -0,0 +1,165 @@ +/* + * 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; + +import java.util.Iterator; +import java.util.List; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.function.IntFunction; + +import org.apache.commons.lang3.StringUtils; +import org.junit.Before; +import org.junit.BeforeClass; +import org.junit.Test; +import org.junit.runner.RunWith; + +import com.datastax.driver.core.ResultSet; +import com.datastax.driver.core.Row; +import org.apache.cassandra.config.DatabaseDescriptor; +import org.apache.cassandra.cql3.CQLTester; +import org.jboss.byteman.contrib.bmunit.BMRule; +import org.jboss.byteman.contrib.bmunit.BMRules; +import org.jboss.byteman.contrib.bmunit.BMUnitRunner; + +import static org.assertj.core.api.Assertions.assertThat; + +/** + * Tests initially implemented for CASSANDRA-18932 and CASSANDRA-18993. They verify the behaviour of querying data + * from a row with a clustering column condition, where the last item covered by the column index block is a range + * tombstone boundary marker. + *

+ * The column index is actually a part of primary index in BIG format and row index in BTI format. + */ +@RunWith(BMUnitRunner.class) +public class ColumnIndexTest extends CQLTester +{ + private static final AtomicBoolean rtbmLastInIndexBlock = new AtomicBoolean(false); + + public static void setRTBMLastInIndexBlock() + { + rtbmLastInIndexBlock.set(true); + } + + @BeforeClass + public static void setUpClass() + { + DatabaseDescriptor.daemonInitialization(); + DatabaseDescriptor.setColumnIndexSizeInKiB(1); + CQLTester.setUpClass(); + requireNetwork(); + } + + @Before + public void beforeTest() throws Throwable + { + super.beforeTest(); + rtbmLastInIndexBlock.set(false); + } + + @Test + @BMRules(rules = { @BMRule(name = "force_add_index_block_after_range_tombstone_boundary_marker_big", + targetClass = "org.apache.cassandra.io.sstable.format.big.BigFormatPartitionWriter", + targetMethod = "addUnfiltered", + targetLocation = "AT EXIT", + condition = "$1 instanceof org.apache.cassandra.db.rows.RangeTombstoneBoundaryMarker", + action = "$this.addIndexBlock()"), + @BMRule(name = "force_add_index_block_after_range_tombstone_boundary_marker_bti", + targetClass = "org.apache.cassandra.io.sstable.format.bti.BtiFormatPartitionWriter", + targetMethod = "addUnfiltered", + targetLocation = "AT EXIT", + condition = "$1 instanceof org.apache.cassandra.db.rows.RangeTombstoneBoundaryMarker", + action = "$this.addIndexBlock()"), + }) + public void testBoundaryMarkerAtTheEndOfIndexBlock() throws Throwable + { + int expectedRows = 10; + IntFunction fmt = (i) -> String.format("%05d", i); + + createTable("CREATE TABLE %s (k int, c int, v text, PRIMARY KEY (k, c)) WITH CLUSTERING ORDER BY (c ASC)"); + + for (int i = 0; i < expectedRows; i++) + { + // delete (0; 1] and then insert 1 will make a range tombstone boundary marker (0; 0] + execute("DELETE FROM %s WHERE k = ? AND c > ? AND c <= ?", 0, i - 1, i); + execute("INSERT INTO %s (k, c, v) VALUES (?, ?, ?)", 0, i, fmt.apply(i)); + } + + flush(); + + ResultSet r = executeNetWithPaging("SELECT * FROM %s", 1); + Iterator iter = r.iterator(); + int n = 0; + while (iter.hasNext()) + { + Row row = iter.next(); + int c = row.getInt("c"); + String v = row.getString("v"); + logger.info("Read row: " + row); + assertThat(v).isEqualTo(fmt.apply(c)); + n++; + } + assertThat(n).isEqualTo(expectedRows); + } + + @Test + @BMRules(rules = { @BMRule(name = "test_big", + targetClass = "org.apache.cassandra.io.sstable.format.big.BigFormatPartitionWriter", + targetMethod = "addUnfiltered", + targetLocation = "AT INVOKE addIndexBlock ALL", + condition = "$1 instanceof org.apache.cassandra.db.rows.RangeTombstoneBoundaryMarker", + action = "org.apache.cassandra.io.sstable.format.ColumnIndexTest.setRTBMLastInIndexBlock()"), + @BMRule(name = "test_bti", + targetClass = "org.apache.cassandra.io.sstable.format.bti.BtiFormatPartitionWriter", + targetMethod = "addUnfiltered", + targetLocation = "AT INVOKE addIndexBlock ALL", + condition = "$1 instanceof org.apache.cassandra.db.rows.RangeTombstoneBoundaryMarker", + action = "org.apache.cassandra.io.sstable.format.ColumnIndexTest.setRTBMLastInIndexBlock()") }) + public void silentDataLossTest() throws Throwable + { + rtbmLastInIndexBlock.set(false); + + String tab1 = createTable("CREATE TABLE %s (pk1 bigint, ck1 bigint, v1 ascii, PRIMARY KEY (pk1, ck1))"); + String tab2 = createTable("CREATE TABLE %s (pk1 bigint, ck1 bigint, v1 ascii, PRIMARY KEY (pk1, ck1))"); + + for (int size = 0; size < 1000; size += 10) + { + long pk = size; + String longString = StringUtils.repeat("a", size); + + for (String tbl : new String[]{ tab1, tab2 }) + { + execute(String.format("INSERT INTO %s.%s (pk1, ck1, v1) VALUES (?, ?, ?)", keyspace(), tbl), pk, 0L, longString); + execute(String.format("INSERT INTO %s.%s (pk1, ck1) VALUES (?, ?)", keyspace(), tbl), pk, 1L); + execute(String.format("DELETE FROM %s.%s WHERE pk1 = ? AND ck1 > ? AND ck1 < ?", keyspace(), tbl), pk, 1L, 5L); + execute(String.format("INSERT INTO %s.%s (pk1, ck1, v1) VALUES (?, ?, ?)", keyspace(), tbl), pk, 2L, longString); + execute(String.format("DELETE FROM %s.%s WHERE pk1 = ? AND ck1 > ? AND ck1 < ?", keyspace(), tbl), pk, 2L, 5L); + } + flush(keyspace(), tab1); + List r1 = executeNetWithPaging(String.format("SELECT * FROM %s.%s WHERE pk1 = ?", keyspace(), tab1), 1, pk).all(); + List r2 = executeNetWithPaging(String.format("SELECT * FROM %s.%s WHERE pk1 = ?", keyspace(), tab2), 1, pk).all(); + + assertThat(r1).usingElementComparator(CQLTester::compareNetRows) + .describedAs("Rows for pk=%s", pk) + .isEqualTo(r2); + assertThat(r1).hasSize(3); + } + + assertThat(rtbmLastInIndexBlock.get()).isTrue(); + } +}