diff --git a/CHANGES.txt b/CHANGES.txt index 48a33eb6f4..fac063b45b 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -194,6 +194,7 @@ * Add the ability to disable bulk loading of SSTables (CASSANDRA-18781) * Clean up obsolete functions and simplify cql_version handling in cqlsh (CASSANDRA-18787) Merged from 5.0: + * Avoid lambda usage in TrieMemoryIndex range queries and ensure queue size tracking is per column (CASSANDRA-20668) * Avoid CQLSH throwing an exception loading .cqlshrc on non-supported platforms (CASSANDRA-20478) * Relax validation of snapshot name as a part of SSTable files path validation (CASSANDRA-20649) * Optimize initial skipping logic for SAI queries on large partitions (CASSANDRA-20191) diff --git a/src/java/org/apache/cassandra/index/sai/memory/TrieMemoryIndex.java b/src/java/org/apache/cassandra/index/sai/memory/TrieMemoryIndex.java index 4798393fd7..c531ba6696 100644 --- a/src/java/org/apache/cassandra/index/sai/memory/TrieMemoryIndex.java +++ b/src/java/org/apache/cassandra/index/sai/memory/TrieMemoryIndex.java @@ -23,13 +23,13 @@ import java.util.Iterator; import java.util.Map; import java.util.PriorityQueue; import java.util.SortedSet; +import java.util.concurrent.atomic.AtomicInteger; import java.util.concurrent.atomic.LongAdder; import java.util.function.Function; import org.slf4j.Logger; import org.slf4j.LoggerFactory; -import io.netty.util.concurrent.FastThreadLocal; import org.apache.cassandra.db.Clustering; import org.apache.cassandra.db.DecoratedKey; import org.apache.cassandra.db.PartitionPosition; @@ -59,6 +59,7 @@ public class TrieMemoryIndex extends MemoryIndex { private static final Logger logger = LoggerFactory.getLogger(TrieMemoryIndex.class); private static final int MAX_RECURSIVE_KEY_LENGTH = 128; + private static final int MINIMUM_PRIORITY_QUEUE_SIZE = 128; private final InMemoryTrie data; private final PrimaryKeysReducer primaryKeysReducer; @@ -66,6 +67,11 @@ public class TrieMemoryIndex extends MemoryIndex private ByteBuffer minTerm; private ByteBuffer maxTerm; + // Maintain the last queue size used on this index to use for the next range match. + // This allows for receiving a stream of wide range queries where the queue size + // is larger than we would want to default the size to. + private final AtomicInteger lastPriorityQueueSize = new AtomicInteger(MINIMUM_PRIORITY_QUEUE_SIZE); + public TrieMemoryIndex(StorageAttachedIndex index) { super(index); @@ -142,7 +148,11 @@ public class TrieMemoryIndex extends MemoryIndex case CONTAINS_VALUE: return exactMatch(expression, keyRange); case RANGE: - return rangeMatch(expression, keyRange); + KeyRangeIterator keyIterator = rangeMatch(expression, keyRange); + int keyCount = (int) keyIterator.getMaxKeys(); + if (keyCount > MINIMUM_PRIORITY_QUEUE_SIZE) + lastPriorityQueueSize.set(keyCount); + return keyIterator; default: throw new IllegalArgumentException("Unsupported expression: " + expression); } @@ -251,29 +261,15 @@ public class TrieMemoryIndex extends MemoryIndex private static class Collector { - private static final int MINIMUM_QUEUE_SIZE = 128; - - // Maintain the last queue size used on this index to use for the next range match. - // This allows for receiving a stream of wide range queries where the queue size - // is larger than we would want to default the size to. - // TODO Investigate using a decaying histogram here to avoid the effect of outliers. - private static final FastThreadLocal lastQueueSize = new FastThreadLocal<>() - { - protected Integer initialValue() - { - return MINIMUM_QUEUE_SIZE; - } - }; - - PrimaryKey minimumKey = null; - PrimaryKey maximumKey = null; - final PriorityQueue mergedKeys = new PriorityQueue<>(lastQueueSize.get()); - + final PriorityQueue mergedKeys; final AbstractBounds keyRange; - public Collector(AbstractBounds keyRange) + PrimaryKey maximumKey = null; + + public Collector(AbstractBounds keyRange, int expectedKeys) { this.keyRange = keyRange; + this.mergedKeys = new PriorityQueue<>(expectedKeys); } public void processContent(PrimaryKeys keys) @@ -295,12 +291,8 @@ public class TrieMemoryIndex extends MemoryIndex || primaryKeys.last().partitionKey().compareTo(keyRange.left) < 0) return; - primaryKeys.forEach(this::processKey); - } - - public void updateLastQueueSize() - { - lastQueueSize.set(Math.max(MINIMUM_QUEUE_SIZE, mergedKeys.size())); + for (PrimaryKey primaryKey : primaryKeys) + processKey(primaryKey); } private void processKey(PrimaryKey key) @@ -309,7 +301,7 @@ public class TrieMemoryIndex extends MemoryIndex { mergedKeys.add(key); - minimumKey = minimumKey == null ? key : key.compareTo(minimumKey) < 0 ? key : minimumKey; + // We only track the maximum key, as the minimum can be peeked in constant time on the PQ itself. maximumKey = maximumKey == null ? key : key.compareTo(maximumKey) > 0 ? key : maximumKey; } } @@ -341,20 +333,16 @@ public class TrieMemoryIndex extends MemoryIndex upperInclusive = false; } - Collector cd = new Collector(keyRange); + Collector cd = new Collector(keyRange, lastPriorityQueueSize.get()); + Iterator values = data.subtrie(lowerBound, lowerInclusive, upperBound, upperInclusive).valueIterator(); - data.subtrie(lowerBound, lowerInclusive, upperBound, upperInclusive) - .values() - .forEach(cd::processContent); + while (values.hasNext()) + cd.processContent(values.next()); if (cd.mergedKeys.isEmpty()) - { return KeyRangeIterator.empty(); - } - cd.updateLastQueueSize(); - - return new InMemoryKeyRangeIterator(cd.minimumKey, cd.maximumKey, cd.mergedKeys); + return new InMemoryKeyRangeIterator(cd.mergedKeys.peek(), cd.maximumKey, cd.mergedKeys); } private static class PrimaryKeysReducer implements InMemoryTrie.UpsertTransformer diff --git a/test/unit/org/apache/cassandra/index/sai/memory/TrieMemoryIndexTest.java b/test/unit/org/apache/cassandra/index/sai/memory/TrieMemoryIndexTest.java index 0ab4c846f7..963742c02c 100644 --- a/test/unit/org/apache/cassandra/index/sai/memory/TrieMemoryIndexTest.java +++ b/test/unit/org/apache/cassandra/index/sai/memory/TrieMemoryIndexTest.java @@ -29,8 +29,10 @@ import java.util.TreeMap; import java.util.function.IntFunction; import java.util.stream.Collectors; +import org.junit.Ignore; import org.junit.Test; +import org.HdrHistogram.Histogram; import org.apache.cassandra.cql3.Operator; import org.apache.cassandra.cql3.statements.schema.IndexTarget; import org.apache.cassandra.db.Clustering; @@ -243,4 +245,28 @@ public class TrieMemoryIndexTest extends SAIRandomizedTester index = new StorageAttachedIndex(cfs, indexMetadata); return new TrieMemoryIndex(index); } + + @Ignore + @Test + public void testMemtableRangeQueryPerformance() + { + createTable("CREATE TABLE %S (pk int, ck int, val int, PRIMARY KEY (pk, ck))"); + createIndex("CREATE INDEX ON %s(val) USING 'sai'"); + + for (int pk = 0; pk < 20; pk++) + for (int ck = 0; ck < 10000; ck++) + execute("INSERT INTO %s (pk, ck, val) VALUES (?, ?, ?)", pk, ck, ck); + + Histogram histogram = new Histogram(4); + + for (int i = 0; i < 20000; i++) + { + long start = System.nanoTime(); + execute("SELECT * FROM %s WHERE pk = 5 AND val > ? LIMIT 10", 4000); + histogram.recordValue(System.nanoTime() - start); + } + + System.out.println("50th: " + histogram.getValueAtPercentile(0.5)); + System.out.println("99th: " + histogram.getValueAtPercentile(0.99)); + } }