diff --git a/CHANGES.txt b/CHANGES.txt index c1406fee99..efac6cea44 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -1,4 +1,5 @@ 5.1 + * Add system_views.partition_key_statistics for querying SSTable metadata (CASSANDRA-20161) * CEP-42 - Add Constraints Framework (CASSANDRA-19947) * Add table metric PurgeableTombstoneScannedHistogram and a tracing event for scanned purgeable tombstones (CASSANDRA-20132) * Make sure we can parse the expanded CQL before writing it to the log or sending it to replicas (CASSANDRA-20218) diff --git a/doc/modules/cassandra/pages/managing/operating/virtualtables.adoc b/doc/modules/cassandra/pages/managing/operating/virtualtables.adoc index d3b948e3d1..362308372c 100644 --- a/doc/modules/cassandra/pages/managing/operating/virtualtables.adoc +++ b/doc/modules/cassandra/pages/managing/operating/virtualtables.adoc @@ -516,6 +516,53 @@ SELECT total - progress AS remaining FROM system_views.sstable_tasks; ---- +=== Virtual table for primary id's + +Since https://issues.apache.org/jira/browse/CASSANDRA-20161[CASSANDRA-20161], there is +a virtual table `system_views.partition_key_statistics` to allow users to query partition keys and related metadata for a specific table within a keyspace. This feature provides insights into SSTable-level details, such as token values, size estimates, and SSTable counts, without requiring expensive disk-based operations. + +[source,console] +---- +cassandra@cqlsh> select * from ks.tbl; + + id | cl1 | i +----+-----+----- + 1 | 202 | 999 + 1 | 200 | 999 + 1 | 101 | 600 + 2 | 101 | 600 +---- + +[source,console] +---- +cassandra@cqlsh> use system_views; +cassandra@cqlsh> select * from partition_key_statistics where keyspace_name = 'ks' and table_name = 'tbl' and key = '1'; + +@ Row 1 +---------------+---------------------- + keyspace_name | ks + table_name | tbl + token_value | -4069959284402364209 + key | 1 + size_estimate | 83 + sstables | 3 + +cassandra@cqlsh> select * from partition_key_statistics where keyspace_name = 'ks' and table_name = 'tbl' and key = '2'; + +@ Row 1 +---------------+---------------------- + keyspace_name | ks + table_name | tbl + token_value | -3248873570005575792 + key | 2 + size_estimate | 25 + sstables | 1 +---- + +The value in `sstables` column means how many SSTables a particular key is located in. + +If there is a composite partition key, you can separate the values by a colon for `key`. + === Other Virtual Tables Some examples of using other virtual tables are as follows. diff --git a/src/java/org/apache/cassandra/cql3/restrictions/StatementRestrictions.java b/src/java/org/apache/cassandra/cql3/restrictions/StatementRestrictions.java index 71958bc61b..e619993c6a 100644 --- a/src/java/org/apache/cassandra/cql3/restrictions/StatementRestrictions.java +++ b/src/java/org/apache/cassandra/cql3/restrictions/StatementRestrictions.java @@ -345,7 +345,7 @@ public final class StatementRestrictions } else { - if (!allowFiltering && requiresAllowFilteringIfNotSpecified()) + if (!allowFiltering && requiresAllowFilteringIfNotSpecified(table)) throw invalidRequest(allowFilteringMessage(state)); } @@ -356,12 +356,12 @@ public final class StatementRestrictions validateSecondaryIndexSelections(); } - public boolean requiresAllowFilteringIfNotSpecified() + public static boolean requiresAllowFilteringIfNotSpecified(TableMetadata metadata) { - if (!table.isVirtual()) + if (!metadata.isVirtual()) return true; - VirtualTable tableNullable = VirtualKeyspaceRegistry.instance.getTableNullable(table.id); + VirtualTable tableNullable = VirtualKeyspaceRegistry.instance.getTableNullable(metadata.id); assert tableNullable != null; return !tableNullable.allowFilteringImplicitly(); } @@ -568,7 +568,7 @@ public final class StatementRestrictions // components must have a EQ. Only the last partition key component can be in IN relation. if (partitionKeyRestrictions.needFiltering()) { - if (!allowFiltering && !forView && !hasQueriableIndex && requiresAllowFilteringIfNotSpecified()) + if (!allowFiltering && !forView && !hasQueriableIndex && requiresAllowFilteringIfNotSpecified(table)) throw new InvalidRequestException(allowFilteringMessage(state)); isKeyRange = true; diff --git a/src/java/org/apache/cassandra/cql3/statements/SelectStatement.java b/src/java/org/apache/cassandra/cql3/statements/SelectStatement.java index 03ff708ec5..1762b776a7 100644 --- a/src/java/org/apache/cassandra/cql3/statements/SelectStatement.java +++ b/src/java/org/apache/cassandra/cql3/statements/SelectStatement.java @@ -88,6 +88,7 @@ import org.apache.commons.lang3.builder.ToStringBuilder; import org.apache.commons.lang3.builder.ToStringStyle; import static java.lang.String.format; +import static org.apache.cassandra.cql3.restrictions.StatementRestrictions.requiresAllowFilteringIfNotSpecified; import static org.apache.cassandra.cql3.statements.RequestValidations.checkFalse; import static org.apache.cassandra.cql3.statements.RequestValidations.checkNotNull; import static org.apache.cassandra.cql3.statements.RequestValidations.checkNull; @@ -1357,7 +1358,7 @@ public class SelectStatement implements CQLStatement.SingleKeyspaceCqlStatement boundNames, orderings, selectsOnlyStaticColumns, - parameters.allowFiltering, + parameters.allowFiltering || !requiresAllowFilteringIfNotSpecified(metadata), forView); } @@ -1583,7 +1584,7 @@ public class SelectStatement implements CQLStatement.SingleKeyspaceCqlStatement { // We will potentially filter data if the row filter is not the identity and there isn't any index group // supporting all the expressions in the filter. - if (restrictions.requiresAllowFilteringIfNotSpecified()) + if (requiresAllowFilteringIfNotSpecified(table)) checkFalse(restrictions.needFiltering(table), StatementRestrictions.REQUIRES_ALLOW_FILTERING_MESSAGE); } } diff --git a/src/java/org/apache/cassandra/db/PartitionRangeReadCommand.java b/src/java/org/apache/cassandra/db/PartitionRangeReadCommand.java index c39f2cb2f3..4926061cb8 100644 --- a/src/java/org/apache/cassandra/db/PartitionRangeReadCommand.java +++ b/src/java/org/apache/cassandra/db/PartitionRangeReadCommand.java @@ -565,7 +565,7 @@ public class PartitionRangeReadCommand extends ReadCommand implements PartitionR public UnfilteredPartitionIterator executeLocally(ReadExecutionController executionController) { VirtualTable view = VirtualKeyspaceRegistry.instance.getTableNullable(metadata().id); - UnfilteredPartitionIterator resultIterator = view.select(dataRange, columnFilter()); + UnfilteredPartitionIterator resultIterator = view.select(dataRange, columnFilter(), rowFilter()); return limits().filter(rowFilter().filter(resultIterator, nowInSec()), nowInSec(), selectsFullPartition()); } diff --git a/src/java/org/apache/cassandra/db/SinglePartitionReadCommand.java b/src/java/org/apache/cassandra/db/SinglePartitionReadCommand.java index be1d5a2322..ad6da1e88a 100644 --- a/src/java/org/apache/cassandra/db/SinglePartitionReadCommand.java +++ b/src/java/org/apache/cassandra/db/SinglePartitionReadCommand.java @@ -1395,7 +1395,7 @@ public class SinglePartitionReadCommand extends ReadCommand implements SinglePar public UnfilteredPartitionIterator executeLocally(ReadExecutionController executionController) { VirtualTable view = VirtualKeyspaceRegistry.instance.getTableNullable(metadata().id); - UnfilteredPartitionIterator resultIterator = view.select(partitionKey, clusteringIndexFilter, columnFilter()); + UnfilteredPartitionIterator resultIterator = view.select(partitionKey, clusteringIndexFilter, columnFilter(), rowFilter()); return limits().filter(rowFilter().filter(resultIterator, nowInSec()), nowInSec(), selectsFullPartition()); } diff --git a/src/java/org/apache/cassandra/db/virtual/AbstractVirtualTable.java b/src/java/org/apache/cassandra/db/virtual/AbstractVirtualTable.java index 344369f7d7..df2e4bc7cc 100644 --- a/src/java/org/apache/cassandra/db/virtual/AbstractVirtualTable.java +++ b/src/java/org/apache/cassandra/db/virtual/AbstractVirtualTable.java @@ -29,6 +29,7 @@ import org.apache.cassandra.db.EmptyIterators; import org.apache.cassandra.db.PartitionPosition; import org.apache.cassandra.db.filter.ClusteringIndexFilter; import org.apache.cassandra.db.filter.ColumnFilter; +import org.apache.cassandra.db.filter.RowFilter; import org.apache.cassandra.db.partitions.AbstractUnfilteredPartitionIterator; import org.apache.cassandra.db.partitions.PartitionUpdate; import org.apache.cassandra.db.partitions.SingletonUnfilteredPartitionIterator; @@ -75,7 +76,7 @@ public abstract class AbstractVirtualTable implements VirtualTable } @Override - public final UnfilteredPartitionIterator select(DecoratedKey partitionKey, ClusteringIndexFilter clusteringIndexFilter, ColumnFilter columnFilter) + public final UnfilteredPartitionIterator select(DecoratedKey partitionKey, ClusteringIndexFilter clusteringIndexFilter, ColumnFilter columnFilter, RowFilter rowFilter) { Partition partition = data(partitionKey).getPartition(partitionKey); @@ -88,7 +89,7 @@ public abstract class AbstractVirtualTable implements VirtualTable } @Override - public final UnfilteredPartitionIterator select(DataRange dataRange, ColumnFilter columnFilter) + public final UnfilteredPartitionIterator select(DataRange dataRange, ColumnFilter columnFilter, RowFilter rowFilter) { DataSet data = data(); diff --git a/src/java/org/apache/cassandra/db/virtual/CollectionVirtualTableAdapter.java b/src/java/org/apache/cassandra/db/virtual/CollectionVirtualTableAdapter.java index c5079d73c0..47aa3bd5c4 100644 --- a/src/java/org/apache/cassandra/db/virtual/CollectionVirtualTableAdapter.java +++ b/src/java/org/apache/cassandra/db/virtual/CollectionVirtualTableAdapter.java @@ -51,6 +51,7 @@ import org.apache.cassandra.db.DeletionTime; import org.apache.cassandra.db.EmptyIterators; import org.apache.cassandra.db.filter.ClusteringIndexFilter; import org.apache.cassandra.db.filter.ColumnFilter; +import org.apache.cassandra.db.filter.RowFilter; import org.apache.cassandra.db.marshal.AbstractType; import org.apache.cassandra.db.marshal.BooleanType; import org.apache.cassandra.db.marshal.ByteType; @@ -307,7 +308,8 @@ public class CollectionVirtualTableAdapter implements VirtualTable @Override public UnfilteredPartitionIterator select(DecoratedKey partitionKey, ClusteringIndexFilter clusteringFilter, - ColumnFilter columnFilter) + ColumnFilter columnFilter, + RowFilter rowFilter) { if (!data.iterator().hasNext()) return EmptyIterators.unfilteredPartition(metadata); @@ -348,7 +350,7 @@ public class CollectionVirtualTableAdapter implements VirtualTable } @Override - public UnfilteredPartitionIterator select(DataRange dataRange, ColumnFilter columnFilter) + public UnfilteredPartitionIterator select(DataRange dataRange, ColumnFilter columnFilter, RowFilter rowFilter) { return createPartitionIterator(metadata, new AbstractIterator<>() { diff --git a/src/java/org/apache/cassandra/db/virtual/PartitionKeyStatsTable.java b/src/java/org/apache/cassandra/db/virtual/PartitionKeyStatsTable.java new file mode 100644 index 0000000000..d114e5faa7 --- /dev/null +++ b/src/java/org/apache/cassandra/db/virtual/PartitionKeyStatsTable.java @@ -0,0 +1,364 @@ +/* + * 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.db.virtual; + +import java.io.IOException; +import java.math.BigInteger; +import java.nio.ByteBuffer; +import java.util.List; +import java.util.concurrent.CopyOnWriteArrayList; +import java.util.function.Consumer; + +import com.google.common.annotations.VisibleForTesting; +import com.google.common.collect.Lists; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import org.apache.cassandra.cql3.Operator; +import org.apache.cassandra.db.Clustering; +import org.apache.cassandra.db.ColumnFamilyStore; +import org.apache.cassandra.db.DataRange; +import org.apache.cassandra.db.DecoratedKey; +import org.apache.cassandra.db.DeletionTime; +import org.apache.cassandra.db.PartitionPosition; +import org.apache.cassandra.db.Slice; +import org.apache.cassandra.db.Slices; +import org.apache.cassandra.db.context.CounterContext; +import org.apache.cassandra.db.filter.ClusteringIndexFilter; +import org.apache.cassandra.db.filter.ColumnFilter; +import org.apache.cassandra.db.filter.RowFilter; +import org.apache.cassandra.db.marshal.CompositeType; +import org.apache.cassandra.db.marshal.CounterColumnType; +import org.apache.cassandra.db.marshal.IntegerType; +import org.apache.cassandra.db.marshal.UTF8Type; +import org.apache.cassandra.db.partitions.PartitionUpdate; +import org.apache.cassandra.db.partitions.SingletonUnfilteredPartitionIterator; +import org.apache.cassandra.db.partitions.UnfilteredPartitionIterator; +import org.apache.cassandra.db.rows.AbstractUnfilteredRowIterator; +import org.apache.cassandra.db.rows.BTreeRow; +import org.apache.cassandra.db.rows.BufferCell; +import org.apache.cassandra.db.rows.Cell; +import org.apache.cassandra.db.rows.EncodingStats; +import org.apache.cassandra.db.rows.Row; +import org.apache.cassandra.db.rows.Rows; +import org.apache.cassandra.db.rows.Unfiltered; +import org.apache.cassandra.db.rows.UnfilteredRowIterator; +import org.apache.cassandra.db.rows.UnfilteredRowIterators; +import org.apache.cassandra.dht.AbstractBounds; +import org.apache.cassandra.dht.Bounds; +import org.apache.cassandra.dht.LocalPartitioner; +import org.apache.cassandra.dht.Range; +import org.apache.cassandra.dht.Token; +import org.apache.cassandra.exceptions.InvalidRequestException; +import org.apache.cassandra.io.sstable.KeyReader; +import org.apache.cassandra.io.sstable.format.SSTableReader; +import org.apache.cassandra.schema.ColumnMetadata; +import org.apache.cassandra.schema.KeyspaceMetadata; +import org.apache.cassandra.schema.Schema; +import org.apache.cassandra.schema.TableMetadata; +import org.apache.cassandra.serializers.MarshalException; + +import static org.apache.cassandra.cql3.statements.RequestValidations.invalidRequest; + +/** + * A virtual table for querying partition keys of SSTables in a specific keyspace. + * + *

This table is implemented as a virtual table in Cassandra, meaning it does not + * store data persistently on disk but instead derives its data from live metadata. + * + *

The CQL equivalent of this virtual table is: + *

+ * CREATE TABLE system_views.partition_key_statistics (
+ *     keyspace_name TEXT,
+ *     table_name TEXT,
+ *     token_value INT,
+ *     key TEXT,
+ *     size_estimate COUNTER,
+ *     sstables COUNTER,
+ *     PRIMARY KEY ((keyspace_name, table_name), token_value, key)
+ * );
+ * 
+ * + *

Note: + *

+ */ +public class PartitionKeyStatsTable implements VirtualTable +{ + private static final Logger logger = LoggerFactory.getLogger(PartitionKeyStatsTable.class); + public static final String NAME = "partition_key_statistics"; + + private static final String TABLE_READ_ONLY_ERROR = "The specified table is read-only."; + private static final String UNSUPPORTED_RANGE_QUERY_ERROR = "Range queries are not supported. Please provide both a keyspace and a table name."; + private static final String REVERSED_QUERY_ERROR = "Reversed queries are not supported."; + private static final String KEYSPACE_NOT_EXIST_ERROR = "The keyspace '%s' does not exist."; + private static final String TABLE_NOT_EXIST_ERROR = "The table '%s' does not exist in the keyspace '%s'."; + private static final String KEY_ONLY_EQUALS_ERROR = "The 'key' column can only be used in an equality query for this virtual table."; + private static final String KEY_NOT_WITHIN_BOUNDS_ERROR = "The specified 'key' is not within the provided token value bounds."; + private static final String PARTITIONER_NOT_SUPPORTED = "Partitioner '%s' for table '%s' in keyspace '%s' is not supported."; + + private static final String COLUMN_KEYSPACE_NAME = "keyspace_name"; + private static final String COLUMN_TABLE_NAME = "table_name"; + private static final String COLUMN_TOKEN_VALUE = "token_value"; + private static final String COLUMN_KEY = "key"; + private static final String COLUMN_SIZE_ESTIMATE = "size_estimate"; + private static final String COLUMN_SSTABLES = "sstables"; + + private final TableMetadata metadata; + private final ColumnMetadata sizeEstimateColumn; + private final ColumnMetadata sstablesColumn; + + @VisibleForTesting + final CopyOnWriteArrayList> readListener = new CopyOnWriteArrayList<>(); + + public PartitionKeyStatsTable(String keyspace) + { + this.metadata = TableMetadata.builder(keyspace, NAME) + .kind(TableMetadata.Kind.VIRTUAL) + .partitioner(new LocalPartitioner(CompositeType.getInstance(UTF8Type.instance, UTF8Type.instance))) + .addPartitionKeyColumn(COLUMN_KEYSPACE_NAME, UTF8Type.instance) + .addPartitionKeyColumn(COLUMN_TABLE_NAME, UTF8Type.instance) + .addClusteringColumn(COLUMN_TOKEN_VALUE, IntegerType.instance) + .addClusteringColumn(COLUMN_KEY, UTF8Type.instance) + .addRegularColumn(COLUMN_SIZE_ESTIMATE, CounterColumnType.instance) + .addRegularColumn(COLUMN_SSTABLES, CounterColumnType.instance) + .build(); + sizeEstimateColumn = metadata.regularColumns().getSimple(0); + sstablesColumn = metadata.regularColumns().getSimple(1); + } + + @Override + public UnfilteredPartitionIterator select(DecoratedKey partitionKey, ClusteringIndexFilter clusteringIndexFilter, ColumnFilter columnFilter, RowFilter rowFilter) + { + if (clusteringIndexFilter.isReversed()) + throw new InvalidRequestException(REVERSED_QUERY_ERROR); + + ByteBuffer[] key = ((CompositeType) this.metadata.partitionKeyType).split(partitionKey.getKey()); + String keyspace = UTF8Type.instance.getString(key[0]); + String table = UTF8Type.instance.getString(key[1]); + + KeyspaceMetadata ksm = Schema.instance.getKeyspaceMetadata(keyspace); + if (ksm == null) + throw invalidRequest(KEYSPACE_NOT_EXIST_ERROR, keyspace); + + TableMetadata metadata = ksm.getTableOrViewNullable(table); + if (metadata == null) + throw invalidRequest(TABLE_NOT_EXIST_ERROR, table, keyspace); + + if (!metadata.partitioner.supportsSplitting()) + throw invalidRequest(PARTITIONER_NOT_SUPPORTED, metadata.partitioner.getClass().getName(), table, keyspace); + + AbstractBounds range = getBounds(metadata, clusteringIndexFilter, rowFilter); + return new SingletonUnfilteredPartitionIterator(select(partitionKey, metadata, clusteringIndexFilter, range)); + } + + private List getSStables(TableMetadata metadata, AbstractBounds range) + { + return Lists.newArrayList(ColumnFamilyStore.getIfExists(metadata).getTracker().getView().liveSSTablesInBounds(range.left, range.right)); + } + + private UnfilteredRowIterator select(DecoratedKey partitionKey, TableMetadata metadata, ClusteringIndexFilter clusteringIndexFilter, AbstractBounds range) + { + List sstables = getSStables(metadata, range); + if (sstables.isEmpty()) + return UnfilteredRowIterators.noRowsIterator(metadata, partitionKey, Rows.EMPTY_STATIC_ROW, DeletionTime.LIVE, false); + + List sstableIterators = Lists.newArrayList(); + for (SSTableReader sstable : sstables) + sstableIterators.add(getSStableRowIterator(metadata, partitionKey, sstable, clusteringIndexFilter, range)); + + return UnfilteredRowIterators.merge(sstableIterators); + } + + private UnfilteredRowIterator getSStableRowIterator(TableMetadata target, DecoratedKey partitionKey, SSTableReader sstable, ClusteringIndexFilter filter, AbstractBounds range) + { + final KeyReader reader; + try + { + // ignore warning on try-with-resources, the reader will be closed on endOfData or close + reader = sstable.keyReader(range.left); + } + catch (IOException e) + { + logger.error("Error generating keyReader for SSTable: {}", sstable, e); + throw new RuntimeException(e); + } + + return new AbstractUnfilteredRowIterator(metadata, partitionKey, DeletionTime.LIVE, + metadata.regularAndStaticColumns(), Rows.EMPTY_STATIC_ROW, + false, EncodingStats.NO_STATS) + { + public Unfiltered endOfData() + { + reader.close(); + return super.endOfData(); + } + + public void close() + { + reader.close(); + } + + private Row buildRow(Clustering clustering, long size) + { + Row.Builder row = BTreeRow.sortedBuilder(); + row.newRow(clustering); + row.addCell(cell(sizeEstimateColumn, CounterContext.instance().createUpdate(size))); + row.addCell(cell(sstablesColumn, CounterContext.instance().createUpdate(1))); + return row.build(); + } + + @Override + protected Unfiltered computeNext() + { + while (!reader.isExhausted()) + { + DecoratedKey key = target.partitioner.decorateKey(reader.key()); + + for (Consumer listener : readListener) + listener.accept(key); + + // Store the reader's current data position to calculate size later + long lastPosition = reader.dataPosition(); + try + { + // Advance the reader to the next key for the next iteration. Also by moving to next key + // we move the dataPosition to the start of the next key for calculating size + reader.advance(); + } + catch (IOException e) + { + logger.error("Error advancing reader for SSTable: {}", sstable, e); + return endOfData(); + } + + // Calculate the size of the current key. If EOF use the length of the file + long current = reader.dataPosition() == -1 ? sstable.uncompressedLength() : reader.dataPosition(); + long size = current - lastPosition; + + String keyString = target.partitionKeyType.getString(key.getKey()); + + // Check if the current key is outside the queried range; if so, stop + if (range.right.compareTo(key) < 0) + return endOfData(); + + // Convert the token to a string and create a clustering object + String tokenString = key.getToken().toString(); + Clustering clustering = Clustering.make( + IntegerType.instance.decompose(new BigInteger(tokenString)), + UTF8Type.instance.decompose(keyString) + ); + + // Check if the current clustering matches the filter; if so, return the row + if (filter.selects(clustering)) + return buildRow(clustering, size); + } + return endOfData(); + } + }; + } + + /** + * This converts the clustering token/key into the partition level token/key for the target table. Also provides an + * optimization from RowFilter when a `key` is specified with or without the clustering `token` being set. + */ + private AbstractBounds getBounds(TableMetadata target, ClusteringIndexFilter clusteringIndexFilter, RowFilter rowFilter) + { + Slices s = clusteringIndexFilter.getSlices(target); + Token startToken = target.partitioner.getMinimumToken(); + Token endToken = target.partitioner.getMaximumToken(); + BigInteger startTokenValue = new BigInteger(endToken.getTokenValue().toString(), 10); + BigInteger endTokenValue = new BigInteger(startToken.getTokenValue().toString(), 10); + + // find min/max token values from the clustering key + for (int i = 0; i < s.size(); i++) + { + Slice slice = s.get(i); + if (!slice.start().isEmpty()) + { + startTokenValue = startTokenValue.min(IntegerType.instance.compose(slice.start().bufferAt(0))); + startToken = target.partitioner.getTokenFactory().fromString(startTokenValue.toString()); + } + if (!slice.end().isEmpty()) + { + endTokenValue = endTokenValue.max(IntegerType.instance.compose(slice.end().bufferAt(0))); + endToken = target.partitioner.getTokenFactory().fromString(endTokenValue.toString()); + } + } + + // override min/max of token if the `key` is specified + for (RowFilter.Expression expression : rowFilter.getExpressions()) + { + if (expression.column().name.toString().equals(COLUMN_KEY)) + { + if (expression.operator() != Operator.EQ) + throw new InvalidRequestException(KEY_ONLY_EQUALS_ERROR); + + String keyString = UTF8Type.instance.compose(expression.getIndexValue()); + ByteBuffer keyAsBB; + try + { + keyAsBB = target.partitionKeyType.fromString(keyString); + } + catch (MarshalException ex) + { + throw new InvalidRequestException(ex.getMessage()); + } + DecoratedKey decoratedKey = target.partitioner.decorateKey(keyAsBB); + + if (!DataRange.forKeyRange(new Range<>(startToken.minKeyBound(), endToken.maxKeyBound())).contains(decoratedKey.getToken().minKeyBound())) + throw new InvalidRequestException(KEY_NOT_WITHIN_BOUNDS_ERROR); + + return Bounds.bounds(decoratedKey, true, decoratedKey, true); + } + } + return Bounds.bounds(startToken.minKeyBound(), true, endToken.maxKeyBound(), true); + } + + private static Cell cell(ColumnMetadata column, ByteBuffer value) + { + return BufferCell.live(column, 1L, value); + } + + @Override + public TableMetadata metadata() + { + return this.metadata; + } + + @Override + public UnfilteredPartitionIterator select(DataRange dataRange, ColumnFilter columnFilter, RowFilter rowFilter) + { + throw new InvalidRequestException(UNSUPPORTED_RANGE_QUERY_ERROR); + } + + @Override + public void truncate() + { + throw new InvalidRequestException(TABLE_READ_ONLY_ERROR); + } + + @Override + public void apply(PartitionUpdate update) + { + throw new InvalidRequestException(TABLE_READ_ONLY_ERROR); + } +} diff --git a/src/java/org/apache/cassandra/db/virtual/SystemViewsKeyspace.java b/src/java/org/apache/cassandra/db/virtual/SystemViewsKeyspace.java index 8c1412e08c..dacf9f643a 100644 --- a/src/java/org/apache/cassandra/db/virtual/SystemViewsKeyspace.java +++ b/src/java/org/apache/cassandra/db/virtual/SystemViewsKeyspace.java @@ -56,6 +56,7 @@ public final class SystemViewsKeyspace extends VirtualKeyspace .add(new RolesCacheKeysTable(VIRTUAL_VIEWS)) .add(new CQLMetricsTable(VIRTUAL_VIEWS)) .add(new BatchMetricsTable(VIRTUAL_VIEWS)) + .add(new PartitionKeyStatsTable(VIRTUAL_VIEWS)) .add(new StreamingVirtualTable(VIRTUAL_VIEWS)) .add(new GossipInfoTable(VIRTUAL_VIEWS)) .add(new QueriesTable(VIRTUAL_VIEWS)) diff --git a/src/java/org/apache/cassandra/db/virtual/VirtualTable.java b/src/java/org/apache/cassandra/db/virtual/VirtualTable.java index 53a9f2ac7f..770cb13983 100644 --- a/src/java/org/apache/cassandra/db/virtual/VirtualTable.java +++ b/src/java/org/apache/cassandra/db/virtual/VirtualTable.java @@ -21,6 +21,7 @@ import org.apache.cassandra.db.DataRange; import org.apache.cassandra.db.DecoratedKey; import org.apache.cassandra.db.filter.ClusteringIndexFilter; import org.apache.cassandra.db.filter.ColumnFilter; +import org.apache.cassandra.db.filter.RowFilter; import org.apache.cassandra.db.partitions.PartitionUpdate; import org.apache.cassandra.db.partitions.UnfilteredPartitionIterator; import org.apache.cassandra.schema.TableMetadata; @@ -59,18 +60,20 @@ public interface VirtualTable * @param partitionKey the partition key * @param clusteringIndexFilter the clustering columns to selected * @param columnFilter the selected columns + * @param rowFilter filter on which rows a given query should include or exclude * @return the rows corresponding to the requested data. */ - UnfilteredPartitionIterator select(DecoratedKey partitionKey, ClusteringIndexFilter clusteringIndexFilter, ColumnFilter columnFilter); + UnfilteredPartitionIterator select(DecoratedKey partitionKey, ClusteringIndexFilter clusteringIndexFilter, ColumnFilter columnFilter, RowFilter rowFilter); /** * Selects the rows from a range of partitions. * * @param dataRange the range of data to retrieve * @param columnFilter the selected columns + * @param rowFilter filter on which rows a given query should include or exclude * @return the rows corresponding to the requested data. */ - UnfilteredPartitionIterator select(DataRange dataRange, ColumnFilter columnFilter); + UnfilteredPartitionIterator select(DataRange dataRange, ColumnFilter columnFilter, RowFilter rowFilter); /** * Truncates data from the underlying source, if supported. diff --git a/src/java/org/apache/cassandra/dht/IPartitioner.java b/src/java/org/apache/cassandra/dht/IPartitioner.java index 9dba3e0d9d..341ebc47f1 100644 --- a/src/java/org/apache/cassandra/dht/IPartitioner.java +++ b/src/java/org/apache/cassandra/dht/IPartitioner.java @@ -74,6 +74,15 @@ public interface IPartitioner throw new UnsupportedOperationException("If you are using a splitting partitioner, getMaximumToken has to be implemented"); } + /** + * + * @return true if supports splitting as per {@link IPartitioner#split(Token, Token, double)}, false otherwise. Defaults to false. + */ + default boolean supportsSplitting() + { + return false; + } + /** * @return a Token that can be used to route a given key * (This is NOT a method to create a Token from its string representation; diff --git a/src/java/org/apache/cassandra/dht/Murmur3Partitioner.java b/src/java/org/apache/cassandra/dht/Murmur3Partitioner.java index e2371c0937..dfe0971f7a 100644 --- a/src/java/org/apache/cassandra/dht/Murmur3Partitioner.java +++ b/src/java/org/apache/cassandra/dht/Murmur3Partitioner.java @@ -138,6 +138,11 @@ public class Murmur3Partitioner implements IPartitioner return new LongToken(newToken); } + public boolean supportsSplitting() + { + return true; + } + public LongToken getMinimumToken() { return MINIMUM; diff --git a/src/java/org/apache/cassandra/dht/RandomPartitioner.java b/src/java/org/apache/cassandra/dht/RandomPartitioner.java index a8fbe764d4..9b833e3868 100644 --- a/src/java/org/apache/cassandra/dht/RandomPartitioner.java +++ b/src/java/org/apache/cassandra/dht/RandomPartitioner.java @@ -135,6 +135,11 @@ public class RandomPartitioner implements IPartitioner return new BigIntegerToken(newToken); } + public boolean supportsSplitting() + { + return true; + } + public BigIntegerToken getMinimumToken() { return MINIMUM; diff --git a/src/java/org/apache/cassandra/io/sstable/format/SSTableReader.java b/src/java/org/apache/cassandra/io/sstable/format/SSTableReader.java index 8e2db41f6c..bf3b203e3e 100644 --- a/src/java/org/apache/cassandra/io/sstable/format/SSTableReader.java +++ b/src/java/org/apache/cassandra/io/sstable/format/SSTableReader.java @@ -903,6 +903,14 @@ public abstract class SSTableReader extends SSTable implements UnfilteredSource, */ public abstract KeyReader keyReader() throws IOException; + /** + * Returns a {@link KeyReader} over all keys in the sstable after a given key. + * @param key + * @return + * @throws IOException + */ + public abstract KeyReader keyReader(PartitionPosition key) throws IOException; + /** * Returns a {@link KeyIterator} over all keys in the sstable. */ diff --git a/src/java/org/apache/cassandra/io/sstable/format/big/BigTableKeyReader.java b/src/java/org/apache/cassandra/io/sstable/format/big/BigTableKeyReader.java index e0b965057f..04b07af2ce 100644 --- a/src/java/org/apache/cassandra/io/sstable/format/big/BigTableKeyReader.java +++ b/src/java/org/apache/cassandra/io/sstable/format/big/BigTableKeyReader.java @@ -53,7 +53,12 @@ public class BigTableKeyReader implements KeyReader public static BigTableKeyReader create(RandomAccessReader indexFileReader, IndexSerializer serializer) throws IOException { - BigTableKeyReader iterator = new BigTableKeyReader(null, indexFileReader, serializer); + return create(null, indexFileReader, serializer); + } + + public static BigTableKeyReader create(FileHandle indexFile, RandomAccessReader indexFileReader, IndexSerializer serializer) throws IOException + { + BigTableKeyReader iterator = new BigTableKeyReader(indexFile, indexFileReader, serializer); try { iterator.advance(); diff --git a/src/java/org/apache/cassandra/io/sstable/format/big/BigTableReader.java b/src/java/org/apache/cassandra/io/sstable/format/big/BigTableReader.java index 692cadf34d..0864a64cee 100644 --- a/src/java/org/apache/cassandra/io/sstable/format/big/BigTableReader.java +++ b/src/java/org/apache/cassandra/io/sstable/format/big/BigTableReader.java @@ -155,6 +155,28 @@ public class BigTableReader extends SSTableReaderWithFilter implements IndexSumm return BigTableKeyReader.create(ifile, rowIndexEntrySerializer); } + @Override + public KeyReader keyReader(PartitionPosition key) throws IOException + { + FileHandle iFile = ifile.sharedCopy(); + RandomAccessReader reader = iFile.createReader(); + reader.seek(getIndexScanPosition(key)); + KeyReader keys = BigTableKeyReader.create(iFile, reader, rowIndexEntrySerializer); + + boolean hasMoreKeys = true; + while (hasMoreKeys) + { + ByteBuffer indexKey = keys.key(); + DecoratedKey indexDecoratedKey = decorateKey(indexKey); + if (indexDecoratedKey.compareTo(key) >= 0) + break; + + // Advance the iterator and check if more keys are available + hasMoreKeys = keys.advance(); + } + return keys; + } + /** * Finds and returns the first key beyond a given token in this SSTable or null if no such key exists. */ 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 9a65be1137..3916063958 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 @@ -45,6 +45,7 @@ import org.apache.cassandra.dht.Token; import org.apache.cassandra.io.sstable.CorruptSSTableException; import org.apache.cassandra.io.sstable.Descriptor; import org.apache.cassandra.io.sstable.IVerifier; +import org.apache.cassandra.io.sstable.KeyReader; import org.apache.cassandra.io.sstable.SSTable; import org.apache.cassandra.io.sstable.SSTableReadsListener; import org.apache.cassandra.io.sstable.SSTableReadsListener.SelectionReason; @@ -124,6 +125,15 @@ public class BtiTableReader extends SSTableReaderWithFilter return partitionIndex == null ? 0 : partitionIndex.size(); } + @Override + public KeyReader keyReader(PartitionPosition key) throws IOException + { + return PartitionIterator.create(partitionIndex, metadata().partitioner, rowIndexFile, dfile, + key, -1, + metadata().partitioner.getMaximumToken().maxKeyBound(), 0, + descriptor.version); + } + @Override protected TrieIndexEntry getRowIndexEntry(PartitionPosition key, Operator operator, diff --git a/test/distributed/org/apache/cassandra/io/sstable/format/ForwardingSSTableReader.java b/test/distributed/org/apache/cassandra/io/sstable/format/ForwardingSSTableReader.java index 710a426647..60168f6519 100644 --- a/test/distributed/org/apache/cassandra/io/sstable/format/ForwardingSSTableReader.java +++ b/test/distributed/org/apache/cassandra/io/sstable/format/ForwardingSSTableReader.java @@ -243,6 +243,11 @@ public abstract class ForwardingSSTableReader extends SSTableReader return delegate.keyReader(); } + public KeyReader keyReader(PartitionPosition key) throws IOException + { + return delegate.keyReader(key); + } + @Override public KeyIterator keyIterator() throws IOException { diff --git a/test/unit/org/apache/cassandra/db/virtual/PartitionKeyStatsTableTest.java b/test/unit/org/apache/cassandra/db/virtual/PartitionKeyStatsTableTest.java new file mode 100644 index 0000000000..1d0a9e43e1 --- /dev/null +++ b/test/unit/org/apache/cassandra/db/virtual/PartitionKeyStatsTableTest.java @@ -0,0 +1,351 @@ +/* + * 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.db.virtual; + +import java.math.BigInteger; +import java.net.InetAddress; +import java.net.UnknownHostException; +import java.nio.ByteBuffer; +import java.util.Arrays; +import java.util.Collection; +import java.util.Collections; +import java.util.List; +import java.util.concurrent.atomic.AtomicInteger; + +import com.google.common.collect.ImmutableList; +import org.junit.Before; +import org.junit.Test; +import org.junit.runner.RunWith; +import org.junit.runners.Parameterized; +import org.junit.runners.Parameterized.Parameters; + +import com.datastax.driver.core.ResultSet; +import com.datastax.driver.core.Row; +import com.datastax.driver.core.exceptions.InvalidQueryException; +import org.apache.cassandra.Util; +import org.apache.cassandra.config.DatabaseDescriptor; +import org.apache.cassandra.cql3.CQLTester; +import org.apache.cassandra.dht.Murmur3Partitioner; +import org.apache.cassandra.io.sstable.format.bti.BtiFormat; +import org.bouncycastle.util.encoders.Hex; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertTrue; + +@RunWith(Parameterized.class) +public class PartitionKeyStatsTableTest extends CQLTester +{ + private static final String KS_NAME = "vts"; + private String table; + private AtomicInteger scanned; + + private final boolean useBtiFormat; + + @Parameters(name = "Use BtiFormat = {0}") + public static Collection parameters() + { + return Arrays.asList(new Object[][]{ { false }, { true } }); + } + + public PartitionKeyStatsTableTest(boolean useBtiFormat) + { + this.useBtiFormat = useBtiFormat; + } + + @Before + public void before() + { + if (useBtiFormat) + DatabaseDescriptor.setSelectedSSTableFormat(new BtiFormat.BtiFormatFactory().getInstance(Collections.emptyMap())); + + PartitionKeyStatsTable primaryIdTable = new PartitionKeyStatsTable(KS_NAME); + scanned = new AtomicInteger(); + VirtualKeyspaceRegistry.instance.register(new VirtualKeyspace(KS_NAME, ImmutableList.of(primaryIdTable))); + + table = createTable("CREATE TABLE %s (key blob PRIMARY KEY, value blob)"); + + ByteBuffer value = ByteBuffer.wrap(new byte[1]); + for (int i = -10; i < 1000; i++) + { + ByteBuffer key = Murmur3Partitioner.LongToken.keyForToken(i); + execute("INSERT INTO %s (key, value) VALUES (?, ?)", key, value); + } + Util.flushTable(KEYSPACE, table); + primaryIdTable.readListener.add(unused -> scanned.incrementAndGet()); + } + + @Test + public void testPrimaryIdTable() + { + ResultSet rs = executeNetWithPaging("SELECT * FROM vts.partition_key_statistics WHERE keyspace_name = ? AND table_name = ?", + 10, KEYSPACE, table); + List all = rs.all(); + assertEquals(1010, all.size()); + assertResults(all, -10, 1000); + // 1010 + 100 for the 1 per 10 page, +1 for the last + assertEquals(1111, scanned.get()); + } + + @Test + public void testTokenValueGreaterThanZero() + { + ResultSet rs = executeNetWithPaging("SELECT * FROM vts.partition_key_statistics WHERE keyspace_name = ? AND table_name = ? AND token_value > 0", + 10, KEYSPACE, table); + List all = rs.all(); + assertEquals(999, all.size()); + assertResults(all, 1, 1000); + assertEquals(1099, scanned.get()); + } + + @Test + public void testTokenValueGreaterThanNegativeFive() + { + ResultSet rs = executeNetWithPaging("SELECT * FROM vts.partition_key_statistics WHERE keyspace_name = ? AND table_name = ? AND token_value > -5", + 10, KEYSPACE, table); + List all = rs.all(); + assertEquals(1004, all.size()); + assertResults(all, -4, 1000); + // 1004 + 100 for the 1 per 10 page, +1 for the last + assertEquals(1105, scanned.get()); + } + + @Test + public void testTokenValueLessThanOrEqualToFive() + { + ResultSet rs = executeNetWithPaging("SELECT * FROM vts.partition_key_statistics WHERE keyspace_name = ? AND table_name = ? AND token_value <= 5", + 10, KEYSPACE, table); + List all = rs.all(); + assertEquals(16, all.size()); + assertResults(all, -10, 5); + assertEquals(18, scanned.get()); + } + + @Test + public void testTokenValueEqualToZero() + { + ResultSet rs = executeNetWithPaging("SELECT * FROM vts.partition_key_statistics WHERE keyspace_name = ? AND table_name = ? AND token_value = 0", + 10, KEYSPACE, table); + List all = rs.all(); + assertEquals(1, all.size()); + Row row = all.get(0); + assertEquals(BigInteger.valueOf(0), row.get("token_value", BigInteger.class)); + assertEquals(2, scanned.get()); + } + + @Test + public void testTokenValueBounds() + { + ResultSet rs = executeNetWithPaging("SELECT * FROM vts.partition_key_statistics WHERE keyspace_name = ? AND table_name = ? AND token_value > 0 AND token_value < 15", + 10, KEYSPACE, table); + List all = rs.all(); + assertEquals(14, all.size()); + assertResults(all, 1, 14); + // 0->10 = 11, 10->16 = 7 + assertEquals(18, scanned.get()); + } + + @Test + public void testTokenValueBoundsWithBetween() + { + ResultSet rs = executeNetWithPaging("SELECT * FROM vts.partition_key_statistics WHERE keyspace_name = ? AND table_name = ? AND token_value BETWEEN 0 AND 15", + 10, KEYSPACE, table); + List all = rs.all(); + assertEquals(16, all.size()); + assertResults(all, 0, 15); + assertEquals(18, scanned.get()); + } + + @Test + public void testTokenValueBoundsWithIn() + { + ResultSet rs = executeNetWithPaging("SELECT * FROM vts.partition_key_statistics WHERE keyspace_name = ? AND table_name = ? AND token_value IN (1,3,6)", + 10, KEYSPACE, table); + List all = rs.all(); + assertEquals(3, all.size()); + assertEquals(BigInteger.valueOf(1), all.get(0).get("token_value", BigInteger.class)); + assertEquals(BigInteger.valueOf(3), all.get(1).get("token_value", BigInteger.class)); + assertEquals(BigInteger.valueOf(6), all.get(2).get("token_value", BigInteger.class)); + assertEquals(7, scanned.get()); + } + + @Test + public void testTokenValueBoundsWithKey() + { + ByteBuffer ten = Murmur3Partitioner.LongToken.keyForToken(10); + ResultSet rs = executeNetWithPaging("SELECT * FROM vts.partition_key_statistics WHERE keyspace_name = ? AND table_name = ? AND token_value > 0 AND token_value < 15 AND key = ?", + 10, KEYSPACE, table, Hex.toHexString(ten.array())); + List all = rs.all(); + assertEquals(1, all.size()); + Row row = all.get(0); + assertEquals(BigInteger.valueOf(10), row.get("token_value", BigInteger.class)); + assertEquals(2, scanned.get()); + } + + @Test + public void testByKey() + { + ByteBuffer ten = Murmur3Partitioner.LongToken.keyForToken(10); + ResultSet rs = executeNetWithPaging("SELECT * FROM vts.partition_key_statistics WHERE keyspace_name = ? AND table_name = ? AND key = ?", + 10, KEYSPACE, table, Hex.toHexString(ten.array())); + List all = rs.all(); + assertEquals(1, all.size()); + Row row = all.get(0); + assertEquals(BigInteger.valueOf(10), row.get("token_value", BigInteger.class)); + assertEquals(2, scanned.get()); + } + + @Test + public void testIgnoreSStableOutOfRange() + { + ByteBuffer twok = Murmur3Partitioner.LongToken.keyForToken(2000); + execute("INSERT INTO %s (key, value) VALUES (?, ?)", twok, ByteBuffer.wrap(new byte[1])); + Util.flushTable(KEYSPACE, table); + ResultSet rs = executeNetWithPaging("SELECT * FROM vts.partition_key_statistics WHERE keyspace_name = ? AND table_name = ? AND token_value > 1500", + 10, KEYSPACE, table); + List all = rs.all(); + assertEquals(1, all.size()); + Row row = all.get(0); + assertEquals(BigInteger.valueOf(2000), row.get("token_value", BigInteger.class)); + assertEquals(1L, row.get("sstables", Long.class).longValue()); + assertEquals(1, scanned.get()); + } + + @Test + public void testNoResults() + { + ResultSet rs = executeNetWithPaging("SELECT * FROM vts.partition_key_statistics WHERE keyspace_name = ? AND table_name = ? AND token_value < -1000", + 10, KEYSPACE, table); + List all = rs.all(); + assertEquals(0, all.size()); + assertEquals(0, scanned.get()); // sstables shouldn't even of been touched + } + + @Test(expected = InvalidQueryException.class) + public void testNonExistantKeyspace() + { + ResultSet rs = executeNetWithPaging("SELECT * FROM vts.partition_key_statistics WHERE keyspace_name = 'non_existent' AND table_name = ?", + 10, table); + List all = rs.all(); + assertEquals(0, all.size()); + assertEquals(0, scanned.get()); + } + + @Test + public void testNoResultsWithSSTables() + { + ByteBuffer o1 = Murmur3Partitioner.LongToken.keyForToken(10000); + ByteBuffer o2 = Murmur3Partitioner.LongToken.keyForToken(10002); + ByteBuffer value = ByteBuffer.wrap(new byte[10]); + execute("INSERT INTO %s (key, value) VALUES (?, ?)", o1, value); + execute("INSERT INTO %s (key, value) VALUES (?, ?)", o2, value); + Util.flushTable(KEYSPACE, table); + + ResultSet rs = executeNetWithPaging("SELECT * FROM vts.partition_key_statistics WHERE keyspace_name = ? AND table_name = ? AND token_value = 10001", + 10, KEYSPACE, table); + List all = rs.all(); + assertEquals(0, all.size()); + assertEquals(1, scanned.get()); + } + + @Test + public void testPrimaryIdTableDuplicates() + { + // 0xc25f118f072d6ba5cab7fb1468ace617 hashes to 1563004846366 + ByteBuffer dup = Murmur3Partitioner.LongToken.keyForToken(1563004846366L); + // -19, 68, -61 (0xed44c3) hashes to 1563004846366 + ByteBuffer dup2 = ByteBuffer.wrap(new byte[]{ -19, 68, -61 }); + ByteBuffer value = ByteBuffer.wrap(new byte[10]); + execute("INSERT INTO %s (key, value) VALUES (?, ?)", dup, value); + execute("INSERT INTO %s (key, value) VALUES (?, ?)", dup2, value); + Util.flushTable(KEYSPACE, table); + + ResultSet rs = executeNetWithPaging("SELECT * FROM vts.partition_key_statistics WHERE keyspace_name = ? AND table_name = ? AND token_value = 1563004846366", + 10, KEYSPACE, table); + List all = rs.all(); + assertEquals(2, all.size()); + assertEquals(BigInteger.valueOf(1563004846366L), all.get(0).get("token_value", BigInteger.class)); + assertEquals(BigInteger.valueOf(1563004846366L), all.get(1).get("token_value", BigInteger.class)); + assertEquals("c25f118f072d6ba5cab7fb1468ace617", all.get(0).getString("key")); + assertEquals("ed44c3", all.get(1).getString("key")); + assertEquals(2, scanned.get()); + } + + @Test + public void testCompositeType() throws UnknownHostException + { + String table = createTable("CREATE TABLE %s (key text, keytwo inet, value text, primary key ((key, keytwo)))"); + + execute("INSERT INTO %s (key, keytwo, value) VALUES (?, ?, ?)", "testkey", InetAddress.getByName("127.0.0.1"), "value"); + Util.flushTable(KEYSPACE, table); + + ResultSet rs = executeNetWithPaging("SELECT * FROM vts.partition_key_statistics WHERE keyspace_name = ? AND table_name = ? AND key = 'testkey:127.0.0.1'", + 10, KEYSPACE, table); + List all = rs.all(); + assertEquals(1, all.size()); + } + + @Test + public void testTextType() + { + String table = createTable("CREATE TABLE %s (key text PRIMARY KEY, value text)"); + + execute("INSERT INTO %s (key, value) VALUES (?, ?)", "testkey", "value"); + Util.flushTable(KEYSPACE, table); + + ResultSet rs = executeNetWithPaging("SELECT * FROM vts.partition_key_statistics WHERE keyspace_name = ? AND table_name = ? AND key = 'testkey'", + 10, KEYSPACE, table); + List all = rs.all(); + assertEquals(1, all.size()); + } + + @Test + public void testSameKeyInMultipleSSTables() + { + String table = createTable("CREATE TABLE %s (key blob PRIMARY KEY, value blob)"); + + ByteBuffer key = Murmur3Partitioner.LongToken.keyForToken(1); + ByteBuffer value = ByteBuffer.wrap(new byte[10]); + execute("INSERT INTO %s (key, value) VALUES (?, ?)", key, value); + Util.flushTable(KEYSPACE, table); + value = ByteBuffer.wrap(new byte[100]); + execute("INSERT INTO %s (key, value) VALUES (?, ?)", key, value); + Util.flushTable(KEYSPACE, table); + + ResultSet rs = executeNetWithPaging("SELECT * FROM vts.partition_key_statistics WHERE keyspace_name = ? AND table_name = ?", + 10, KEYSPACE, table); + List all = rs.all(); + assertEquals(1, all.size()); + Row row = all.get(0); + assertEquals(BigInteger.valueOf(1), row.get("token_value", BigInteger.class)); + long size = row.get("size_estimate", Long.class); + // providing a range since with timestamp delta vint encoding worried this may drift with time or in wierd + // VMs so just want to make sure it's in the right ballpark + assertTrue(size >= 110 && size < 200); + assertEquals(2L, row.get("sstables", Long.class).longValue()); + assertEquals(2, scanned.get()); + } + + private static void assertResults(List all, int start, int end) + { + for (int i = start, offset = 0; i < end; i++, offset++) + { + Row row = all.get(offset); + } + } +}