From 5a4253b6a17de9810fbc4e1c3b8d4980e26adcca Mon Sep 17 00:00:00 2001 From: Tyler Hobbs Date: Sat, 19 Sep 2015 10:06:45 -0500 Subject: [PATCH] Allow MV's SELECT to restrict PK columns Patch by Tyler Hobbs; reviewed by Sylvain Lebresne for CASSANDRA-9664 --- CHANGES.txt | 2 + .../cassandra/config/ViewDefinition.java | 74 +- .../apache/cassandra/cql3/AbstractMarker.java | 31 +- .../cassandra/cql3/ColumnIdentifier.java | 35 + .../org/apache/cassandra/cql3/Constants.java | 26 +- src/java/org/apache/cassandra/cql3/Cql.g | 12 +- src/java/org/apache/cassandra/cql3/Json.java | 42 +- src/java/org/apache/cassandra/cql3/Lists.java | 8 +- src/java/org/apache/cassandra/cql3/Maps.java | 18 +- .../cassandra/cql3/MultiColumnRelation.java | 28 +- .../org/apache/cassandra/cql3/Operator.java | 58 +- .../org/apache/cassandra/cql3/Relation.java | 23 + src/java/org/apache/cassandra/cql3/Sets.java | 8 +- .../cassandra/cql3/SingleColumnRelation.java | 30 + src/java/org/apache/cassandra/cql3/Term.java | 19 +- .../apache/cassandra/cql3/TokenRelation.java | 26 + .../org/apache/cassandra/cql3/Tuples.java | 24 +- .../org/apache/cassandra/cql3/TypeCast.java | 5 +- .../org/apache/cassandra/cql3/UserTypes.java | 7 +- .../cql3/functions/FunctionCall.java | 16 +- .../restrictions/AbstractRestriction.java | 14 +- .../ForwardingPrimaryKeyRestrictions.java | 6 + .../restrictions/MultiColumnRestriction.java | 55 + .../cql3/restrictions/Restriction.java | 1 + .../restrictions/SingleColumnRestriction.java | 58 + .../restrictions/StatementRestrictions.java | 69 +- .../cql3/statements/AlterTableStatement.java | 2 +- .../cql3/statements/CreateViewStatement.java | 74 +- .../cql3/statements/IndexTarget.java | 14 +- .../statements/ModificationStatement.java | 2 +- .../cql3/statements/ParsedStatement.java | 5 + .../cql3/statements/SelectStatement.java | 27 +- .../cql3/statements/UpdateStatement.java | 2 + .../db/PartitionRangeReadCommand.java | 13 +- .../org/apache/cassandra/db/ReadCommand.java | 12 - .../org/apache/cassandra/db/ReadQuery.java | 24 +- .../db/SinglePartitionReadCommand.java | 25 +- .../apache/cassandra/db/filter/RowFilter.java | 39 + .../apache/cassandra/db/view/TemporalRow.java | 17 +- .../org/apache/cassandra/db/view/View.java | 156 +- .../apache/cassandra/db/view/ViewBuilder.java | 15 +- .../composites/CompositesSearcher.java | 2 +- .../cassandra/schema/SchemaKeyspace.java | 16 +- .../org/apache/cassandra/cql3/CQLTester.java | 78 + .../cassandra/cql3/ViewFilteringTest.java | 1292 +++++++++++++++++ .../org/apache/cassandra/cql3/ViewTest.java | 4 +- .../SelectSingleColumnRelationTest.java | 4 + 47 files changed, 2270 insertions(+), 248 deletions(-) create mode 100644 test/unit/org/apache/cassandra/cql3/ViewFilteringTest.java diff --git a/CHANGES.txt b/CHANGES.txt index e55fd0a613..e589626f73 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -1,4 +1,6 @@ 3.0.0-rc1 + * Allow MATERIALIZED VIEW's SELECT statement to restrict primary key + columns (CASSANDRA-9664) * Move crc_check_chance out of compression options (CASSANDRA-9839) * Fix descending iteration past end of BTreeSearchIterator (CASSANDRA-10301) * Transfer hints to a different node on decommission (CASSANDRA-10198) diff --git a/src/java/org/apache/cassandra/config/ViewDefinition.java b/src/java/org/apache/cassandra/config/ViewDefinition.java index 39695b98ce..02acc689c0 100644 --- a/src/java/org/apache/cassandra/config/ViewDefinition.java +++ b/src/java/org/apache/cassandra/config/ViewDefinition.java @@ -17,26 +17,35 @@ */ package org.apache.cassandra.config; +import java.util.List; import java.util.Objects; import java.util.UUID; +import java.util.stream.Collectors; +import org.antlr.runtime.*; +import org.apache.cassandra.cql3.*; +import org.apache.cassandra.cql3.statements.SelectStatement; +import org.apache.cassandra.db.view.View; +import org.apache.cassandra.exceptions.SyntaxException; import org.apache.commons.lang3.builder.HashCodeBuilder; import org.apache.commons.lang3.builder.ToStringBuilder; -import org.apache.cassandra.cql3.ColumnIdentifier; - public class ViewDefinition { public final String ksName; public final String viewName; public final UUID baseTableId; + public final String baseTableName; public final boolean includeAllColumns; // The order of partititon columns and clustering columns is important, so we cannot switch these two to sets public final CFMetaData metadata; + public SelectStatement.RawStatement select; + public String whereClause; + public ViewDefinition(ViewDefinition def) { - this(def.ksName, def.viewName, def.baseTableId, def.includeAllColumns, def.metadata); + this(def.ksName, def.viewName, def.baseTableId, def.baseTableName, def.includeAllColumns, def.select, def.whereClause, def.metadata); } /** @@ -44,12 +53,15 @@ public class ViewDefinition * @param baseTableId Internal ID of the table which this view is based off of * @param includeAllColumns Whether to include all columns or not */ - public ViewDefinition(String ksName, String viewName, UUID baseTableId, boolean includeAllColumns, CFMetaData metadata) + public ViewDefinition(String ksName, String viewName, UUID baseTableId, String baseTableName, boolean includeAllColumns, SelectStatement.RawStatement select, String whereClause, CFMetaData metadata) { this.ksName = ksName; this.viewName = viewName; this.baseTableId = baseTableId; + this.baseTableName = baseTableName; this.includeAllColumns = includeAllColumns; + this.select = select; + this.whereClause = whereClause; this.metadata = metadata; } @@ -63,7 +75,7 @@ public class ViewDefinition public ViewDefinition copy() { - return new ViewDefinition(ksName, viewName, baseTableId, includeAllColumns, metadata.copy()); + return new ViewDefinition(ksName, viewName, baseTableId, baseTableName, includeAllColumns, select, whereClause, metadata.copy()); } public CFMetaData baseTableMetadata() @@ -85,6 +97,7 @@ public class ViewDefinition && Objects.equals(viewName, other.viewName) && Objects.equals(baseTableId, other.baseTableId) && Objects.equals(includeAllColumns, other.includeAllColumns) + && Objects.equals(whereClause, other.whereClause) && Objects.equals(metadata, other.metadata); } @@ -96,6 +109,7 @@ public class ViewDefinition .append(viewName) .append(baseTableId) .append(includeAllColumns) + .append(whereClause) .append(metadata) .toHashCode(); } @@ -107,8 +121,58 @@ public class ViewDefinition .append("ksName", ksName) .append("viewName", viewName) .append("baseTableId", baseTableId) + .append("baseTableName", baseTableName) .append("includeAllColumns", includeAllColumns) + .append("whereClause", whereClause) .append("metadata", metadata) .toString(); } + + /** + * Replace the column {@param from} with {@param to} in this materialized view definition's partition, + * clustering, or included columns. + */ + public void renameColumn(ColumnIdentifier from, ColumnIdentifier to) + { + metadata.renameColumn(from, to); + + // convert whereClause to Relations, rename ids in Relations, then convert back to whereClause + List relations = whereClauseToRelations(whereClause); + ColumnIdentifier.Raw fromRaw = new ColumnIdentifier.Literal(from.toString(), true); + ColumnIdentifier.Raw toRaw = new ColumnIdentifier.Literal(to.toString(), true); + List newRelations = relations.stream() + .map(r -> r.renameIdentifier(fromRaw, toRaw)) + .collect(Collectors.toList()); + + this.whereClause = View.relationsToWhereClause(newRelations); + String rawSelect = View.buildSelectStatement(baseTableName, metadata.allColumns(), whereClause); + this.select = (SelectStatement.RawStatement) QueryProcessor.parseStatement(rawSelect); + } + + private static List whereClauseToRelations(String whereClause) + { + ErrorCollector errorCollector = new ErrorCollector(whereClause); + CharStream stream = new ANTLRStringStream(whereClause); + CqlLexer lexer = new CqlLexer(stream); + lexer.addErrorListener(errorCollector); + + TokenStream tokenStream = new CommonTokenStream(lexer); + CqlParser parser = new CqlParser(tokenStream); + parser.addErrorListener(errorCollector); + + try + { + List relations = parser.whereClause().build().relations; + + // The errorCollector has queued up any errors that the lexer and parser may have encountered + // along the way, if necessary, we turn the last error into exceptions here. + errorCollector.throwFirstSyntaxError(); + + return relations; + } + catch (RecognitionException | SyntaxException exc) + { + throw new RuntimeException("Unexpected error parsing materialized view's where clause while handling column rename: ", exc); + } + } } diff --git a/src/java/org/apache/cassandra/cql3/AbstractMarker.java b/src/java/org/apache/cassandra/cql3/AbstractMarker.java index d11b8e2273..d2bc022b0c 100644 --- a/src/java/org/apache/cassandra/cql3/AbstractMarker.java +++ b/src/java/org/apache/cassandra/cql3/AbstractMarker.java @@ -56,7 +56,7 @@ public abstract class AbstractMarker extends Term.NonTerminal /** * A parsed, but non prepared, bind marker. */ - public static class Raw implements Term.Raw + public static class Raw extends Term.Raw { protected final int bindIndex; @@ -85,7 +85,34 @@ public abstract class AbstractMarker extends Term.NonTerminal } @Override - public String toString() + public String getText() + { + return "?"; + } + } + + /** A MultiColumnRaw version of AbstractMarker.Raw */ + public static abstract class MultiColumnRaw extends Term.MultiColumnRaw + { + protected final int bindIndex; + + public MultiColumnRaw(int bindIndex) + { + this.bindIndex = bindIndex; + } + + public NonTerminal prepare(String keyspace, ColumnSpecification receiver) throws InvalidRequestException + { + throw new AssertionError("MultiColumnRaw..prepare() requires a list of receivers"); + } + + public AssignmentTestable.TestResult testAssignment(String keyspace, ColumnSpecification receiver) + { + return AssignmentTestable.TestResult.WEAKLY_ASSIGNABLE; + } + + @Override + public String getText() { return "?"; } diff --git a/src/java/org/apache/cassandra/cql3/ColumnIdentifier.java b/src/java/org/apache/cassandra/cql3/ColumnIdentifier.java index 4880c60750..eb16f93559 100644 --- a/src/java/org/apache/cassandra/cql3/ColumnIdentifier.java +++ b/src/java/org/apache/cassandra/cql3/ColumnIdentifier.java @@ -23,6 +23,7 @@ import java.util.List; import java.util.Locale; import java.nio.ByteBuffer; import java.util.concurrent.ConcurrentMap; +import java.util.regex.Pattern; import com.google.common.collect.MapMaker; @@ -55,6 +56,8 @@ public class ColumnIdentifier extends org.apache.cassandra.cql3.selection.Select public final long prefixComparison; private final boolean interned; + private static final Pattern UNQUOTED_IDENTIFIER = Pattern.compile("[a-z][a-z0-9_]*"); + private static final long EMPTY_SIZE = ObjectSizes.measure(new ColumnIdentifier(ByteBufferUtil.EMPTY_BYTE_BUFFER, "", false)); private static final ConcurrentMap internedInstances = new MapMaker().weakValues().makeMap(); @@ -150,6 +153,15 @@ public class ColumnIdentifier extends org.apache.cassandra.cql3.selection.Select return text; } + /** + * Returns a string representation of the identifier that is safe to use directly in CQL queries. + * In necessary, the string will be double-quoted, and any quotes inside the string will be escaped. + */ + public String toCQLString() + { + return maybeQuote(text); + } + public long unsharedHeapSize() { return EMPTY_SIZE @@ -198,6 +210,12 @@ public class ColumnIdentifier extends org.apache.cassandra.cql3.selection.Select { public ColumnIdentifier prepare(CFMetaData cfm); + + /** + * Returns a string representation of the identifier that is safe to use directly in CQL queries. + * In necessary, the string will be double-quoted, and any quotes inside the string will be escaped. + */ + public String toCQLString(); } public static class Literal implements Raw @@ -257,6 +275,11 @@ public class ColumnIdentifier extends org.apache.cassandra.cql3.selection.Select { return text; } + + public String toCQLString() + { + return maybeQuote(text); + } } public static class ColumnIdentifierValue implements Raw @@ -298,5 +321,17 @@ public class ColumnIdentifier extends org.apache.cassandra.cql3.selection.Select { return identifier.toString(); } + + public String toCQLString() + { + return maybeQuote(identifier.text); + } + } + + private static String maybeQuote(String text) + { + if (UNQUOTED_IDENTIFIER.matcher(text).matches()) + return text; + return "\"" + text.replace("\"", "\"\"") + "\""; } } diff --git a/src/java/org/apache/cassandra/cql3/Constants.java b/src/java/org/apache/cassandra/cql3/Constants.java index f10484d8be..425dd85c7f 100644 --- a/src/java/org/apache/cassandra/cql3/Constants.java +++ b/src/java/org/apache/cassandra/cql3/Constants.java @@ -46,7 +46,7 @@ public abstract class Constants public static final Value UNSET_VALUE = new Value(ByteBufferUtil.UNSET_BYTE_BUFFER); - public static final Term.Raw NULL_LITERAL = new Term.Raw() + private static class NullLiteral extends Term.Raw { public Term prepare(String keyspace, ColumnSpecification receiver) throws InvalidRequestException { @@ -63,12 +63,13 @@ public abstract class Constants : AssignmentTestable.TestResult.WEAKLY_ASSIGNABLE; } - @Override - public String toString() + public String getText() { - return "null"; + return "NULL"; } - }; + } + + public static final NullLiteral NULL_LITERAL = new NullLiteral(); public static final Term.Terminal NULL_VALUE = new Value(null) { @@ -86,7 +87,7 @@ public abstract class Constants } }; - public static class Literal implements Term.Raw + public static class Literal extends Term.Raw { private final Type type; private final String text; @@ -155,11 +156,6 @@ public abstract class Constants } } - public String getRawText() - { - return text; - } - public AssignmentTestable.TestResult testAssignment(String keyspace, ColumnSpecification receiver) { CQL3Type receiverType = receiver.type.asCQL3Type(); @@ -238,8 +234,12 @@ public abstract class Constants return AssignmentTestable.TestResult.NOT_ASSIGNABLE; } - @Override - public String toString() + public String getRawText() + { + return text; + } + + public String getText() { return type == Type.STRING ? String.format("'%s'", text) : text; } diff --git a/src/java/org/apache/cassandra/cql3/Cql.g b/src/java/org/apache/cassandra/cql3/Cql.g index cd52c1cedd..932ecd6bf2 100644 --- a/src/java/org/apache/cassandra/cql3/Cql.g +++ b/src/java/org/apache/cassandra/cql3/Cql.g @@ -762,19 +762,18 @@ createMaterializedViewStatement returns [CreateViewStatement expr] } : K_CREATE K_MATERIALIZED K_VIEW (K_IF K_NOT K_EXISTS { ifNotExists = true; })? cf=columnFamilyName K_AS K_SELECT sclause=selectClause K_FROM basecf=columnFamilyName - (K_WHERE wclause=mvWhereClause)? + (K_WHERE wclause=whereClause)? K_PRIMARY K_KEY ( '(' '(' k1=cident { partitionKeys.add(k1); } ( ',' kn=cident { partitionKeys.add(kn); } )* ')' ( ',' c1=cident { compositeKeys.add(c1); } )* ')' | '(' k1=cident { partitionKeys.add(k1); } ( ',' cn=cident { compositeKeys.add(cn); } )* ')' ) - { $expr = new CreateViewStatement(cf, basecf, sclause, wclause, partitionKeys, compositeKeys, ifNotExists); } + { + WhereClause where = wclause == null ? WhereClause.empty() : wclause.build(); + $expr = new CreateViewStatement(cf, basecf, sclause, where, partitionKeys, compositeKeys, ifNotExists); + } ( K_WITH cfamProperty[expr.properties] ( K_AND cfamProperty[expr.properties] )*)? ; -mvWhereClause returns [List expr] - : t1=cident { $expr = new ArrayList(); $expr.add(t1); } K_IS K_NOT K_NULL (K_AND tN=cident { $expr.add(tN); } K_IS K_NOT K_NULL)* - ; - /** * CREATE TRIGGER triggerName ON columnFamily USING 'triggerClass'; */ @@ -1423,6 +1422,7 @@ relationType returns [Operator op] relation[WhereClause.Builder clauses] : name=cident type=relationType t=term { $clauses.add(new SingleColumnRelation(name, type, t)); } + | name=cident K_IS K_NOT K_NULL { $clauses.add(new SingleColumnRelation(name, Operator.IS_NOT, Constants.NULL_LITERAL)); } | K_TOKEN l=tupleOfIdentifiers type=relationType t=term { $clauses.add(new TokenRelation(l, type, t)); } | name=cident K_IN marker=inMarker diff --git a/src/java/org/apache/cassandra/cql3/Json.java b/src/java/org/apache/cassandra/cql3/Json.java index 35c69ed458..df2d9ab180 100644 --- a/src/java/org/apache/cassandra/cql3/Json.java +++ b/src/java/org/apache/cassandra/cql3/Json.java @@ -143,9 +143,9 @@ public class Json this.columns = columns; } - public DelayedColumnValue getRawTermForColumn(ColumnDefinition def) + public RawDelayedColumnValue getRawTermForColumn(ColumnDefinition def) { - return new DelayedColumnValue(this, def); + return new RawDelayedColumnValue(this, def); } public void bind(QueryOptions options) throws InvalidRequestException @@ -173,7 +173,7 @@ public class Json * Note that this is intrinsically an already prepared term, but this still implements Term.Raw so that we can * easily use it to create raw operations. */ - private static class ColumnValue implements Term.Raw + private static class ColumnValue extends Term.Raw { private final Term term; @@ -193,19 +193,22 @@ public class Json { return TestResult.NOT_ASSIGNABLE; } + + public String getText() + { + return term.toString(); + } } /** - * A NonTerminal for a single column. - * - * As with {@code ColumnValue}, this is intrinsically a prepared term but implements Terms.Raw for convenience. + * A Raw term for a single column. Like ColumnValue, this is intrinsically already prepared. */ - private static class DelayedColumnValue extends Term.NonTerminal implements Term.Raw + private static class RawDelayedColumnValue extends Term.Raw { private final PreparedMarker marker; private final ColumnDefinition column; - public DelayedColumnValue(PreparedMarker prepared, ColumnDefinition column) + public RawDelayedColumnValue(PreparedMarker prepared, ColumnDefinition column) { this.marker = prepared; this.column = column; @@ -214,7 +217,7 @@ public class Json @Override public Term prepare(String keyspace, ColumnSpecification receiver) throws InvalidRequestException { - return this; + return new DelayedColumnValue(marker, column); } @Override @@ -223,6 +226,26 @@ public class Json return TestResult.WEAKLY_ASSIGNABLE; } + public String getText() + { + return marker.toString(); + } + } + + /** + * A NonTerminal for a single column. As with {@code ColumnValue}, this is intrinsically a prepared. + */ + private static class DelayedColumnValue extends Term.NonTerminal + { + private final PreparedMarker marker; + private final ColumnDefinition column; + + public DelayedColumnValue(PreparedMarker prepared, ColumnDefinition column) + { + this.marker = prepared; + this.column = column; + } + @Override public void collectMarkerSpecification(VariableSpecifications boundNames) { @@ -248,6 +271,7 @@ public class Json { return Collections.emptyList(); } + } /** diff --git a/src/java/org/apache/cassandra/cql3/Lists.java b/src/java/org/apache/cassandra/cql3/Lists.java index d9dac22dc9..830561e6c1 100644 --- a/src/java/org/apache/cassandra/cql3/Lists.java +++ b/src/java/org/apache/cassandra/cql3/Lists.java @@ -23,6 +23,7 @@ import java.nio.ByteBuffer; import java.util.ArrayList; import java.util.List; import java.util.concurrent.atomic.AtomicReference; +import java.util.stream.Collectors; import org.apache.cassandra.config.ColumnDefinition; import org.apache.cassandra.cql3.functions.Function; @@ -55,7 +56,7 @@ public abstract class Lists return new ColumnSpecification(column.ksName, column.cfName, new ColumnIdentifier("value(" + column.name + ")", true), ((ListType)column.type).getElementsType()); } - public static class Literal implements Term.Raw + public static class Literal extends Term.Raw { private final List elements; @@ -113,10 +114,9 @@ public abstract class Lists return AssignmentTestable.TestResult.testAll(keyspace, valueSpec, elements); } - @Override - public String toString() + public String getText() { - return elements.toString(); + return elements.stream().map(Term.Raw::getText).collect(Collectors.joining(", ", "[", "]")); } } diff --git a/src/java/org/apache/cassandra/cql3/Maps.java b/src/java/org/apache/cassandra/cql3/Maps.java index 0f0672f28f..d5df27956d 100644 --- a/src/java/org/apache/cassandra/cql3/Maps.java +++ b/src/java/org/apache/cassandra/cql3/Maps.java @@ -21,6 +21,7 @@ import static org.apache.cassandra.cql3.Constants.UNSET_VALUE; import java.nio.ByteBuffer; import java.util.*; +import java.util.stream.Collectors; import com.google.common.collect.Iterables; @@ -54,7 +55,7 @@ public abstract class Maps return new ColumnSpecification(column.ksName, column.cfName, new ColumnIdentifier("value(" + column.name + ")", true), ((MapType)column.type).getValuesType()); } - public static class Literal implements Term.Raw + public static class Literal extends Term.Raw { public final List> entries; @@ -129,18 +130,11 @@ public abstract class Maps return res; } - @Override - public String toString() + public String getText() { - StringBuilder sb = new StringBuilder(); - sb.append("{"); - for (int i = 0; i < entries.size(); i++) - { - if (i > 0) sb.append(", "); - sb.append(entries.get(i).left).append(":").append(entries.get(i).right); - } - sb.append("}"); - return sb.toString(); + return entries.stream() + .map(entry -> String.format("%s: %s", entry.left.getText(), entry.right.getText())) + .collect(Collectors.joining(", ", "{", "}")); } } diff --git a/src/java/org/apache/cassandra/cql3/MultiColumnRelation.java b/src/java/org/apache/cassandra/cql3/MultiColumnRelation.java index 7735c574af..143106d1bb 100644 --- a/src/java/org/apache/cassandra/cql3/MultiColumnRelation.java +++ b/src/java/org/apache/cassandra/cql3/MultiColumnRelation.java @@ -19,6 +19,7 @@ package org.apache.cassandra.cql3; import java.util.ArrayList; import java.util.List; +import java.util.stream.Collectors; import org.apache.cassandra.config.CFMetaData; import org.apache.cassandra.config.ColumnDefinition; @@ -114,11 +115,17 @@ public class MultiColumnRelation extends Relation * For non-IN relations, returns the Tuples.Literal or Tuples.Raw marker for a single tuple. * @return a Tuples.Literal for non-IN relations or Tuples.Raw marker for a single tuple. */ - private Term.MultiColumnRaw getValue() + public Term.MultiColumnRaw getValue() { return relationType == Operator.IN ? inMarker : valuesOrMarker; } + public List getInValues() + { + assert relationType == Operator.IN; + return inValues; + } + @Override public boolean isMultiColumn() { @@ -164,7 +171,15 @@ public class MultiColumnRelation extends Relation VariableSpecifications boundNames, boolean isKey) throws InvalidRequestException { - throw invalidRequest("%s cannot be used for Multi-column relations", operator()); + throw invalidRequest("%s cannot be used for multi-column relations", operator()); + } + + @Override + protected Restriction newIsNotRestriction(CFMetaData cfm, + VariableSpecifications boundNames) throws InvalidRequestException + { + // this is currently disallowed by the grammar + throw new AssertionError(String.format("%s cannot be used for multi-column relations", operator())); } @Override @@ -198,6 +213,15 @@ public class MultiColumnRelation extends Relation return names; } + public Relation renameIdentifier(ColumnIdentifier.Raw from, ColumnIdentifier.Raw to) + { + if (!entities.contains(from)) + return this; + + List newEntities = entities.stream().map(e -> e.equals(from) ? to : e).collect(Collectors.toList()); + return new MultiColumnRelation(newEntities, operator(), valuesOrMarker, inValues, inMarker); + } + @Override public String toString() { diff --git a/src/java/org/apache/cassandra/cql3/Operator.java b/src/java/org/apache/cassandra/cql3/Operator.java index 5ae988533f..7b28a30f72 100644 --- a/src/java/org/apache/cassandra/cql3/Operator.java +++ b/src/java/org/apache/cassandra/cql3/Operator.java @@ -21,8 +21,14 @@ import java.io.DataInput; import java.io.DataOutput; import java.io.IOException; import java.nio.ByteBuffer; +import java.util.List; +import java.util.Map; +import java.util.Set; import org.apache.cassandra.db.marshal.AbstractType; +import org.apache.cassandra.db.marshal.ListType; +import org.apache.cassandra.db.marshal.MapType; +import org.apache.cassandra.db.marshal.SetType; public enum Operator { @@ -87,6 +93,14 @@ public enum Operator { return "!="; } + }, + IS_NOT(9) + { + @Override + public String toString() + { + return "IS NOT"; + } }; /** @@ -114,6 +128,11 @@ public enum Operator output.writeInt(b); } + public int getValue() + { + return b; + } + /** * Deserializes a Operator instance from the specified input. * @@ -134,27 +153,48 @@ public enum Operator /** * Whether 2 values satisfy this operator (given the type they should be compared with). * - * @throws AssertionError for IN, CONTAINS and CONTAINS_KEY as this doesn't make sense for this function. + * @throws AssertionError for CONTAINS and CONTAINS_KEY as this doesn't support those operators yet */ public boolean isSatisfiedBy(AbstractType type, ByteBuffer leftOperand, ByteBuffer rightOperand) { - int comparison = type.compareForCQL(leftOperand, rightOperand); switch (this) { case EQ: - return comparison == 0; + return type.compareForCQL(leftOperand, rightOperand) == 0; case LT: - return comparison < 0; + return type.compareForCQL(leftOperand, rightOperand) < 0; case LTE: - return comparison <= 0; + return type.compareForCQL(leftOperand, rightOperand) <= 0; case GT: - return comparison > 0; + return type.compareForCQL(leftOperand, rightOperand) > 0; case GTE: - return comparison >= 0; + return type.compareForCQL(leftOperand, rightOperand) >= 0; case NEQ: - return comparison != 0; + return type.compareForCQL(leftOperand, rightOperand) != 0; + case IN: + List inValues = ((List) ListType.getInstance(type, false).getSerializer().deserialize(rightOperand)); + return inValues.contains(type.getSerializer().deserialize(leftOperand)); + case CONTAINS: + if (type instanceof ListType) + { + List list = (List) type.getSerializer().deserialize(leftOperand); + return list.contains(((ListType) type).getElementsType().getSerializer().deserialize(rightOperand)); + } + else if (type instanceof SetType) + { + Set set = (Set) type.getSerializer().deserialize(leftOperand); + return set.contains(((SetType) type).getElementsType().getSerializer().deserialize(rightOperand)); + } + else // MapType + { + Map map = (Map) type.getSerializer().deserialize(leftOperand); + return map.containsValue(((MapType) type).getValuesType().getSerializer().deserialize(rightOperand)); + } + case CONTAINS_KEY: + Map map = (Map) type.getSerializer().deserialize(leftOperand); + return map.containsKey(((MapType) type).getKeysType().getSerializer().deserialize(rightOperand)); default: - // we shouldn't get IN, CONTAINS, or CONTAINS KEY here + // we shouldn't get CONTAINS, CONTAINS KEY, or IS NOT here throw new AssertionError(); } } diff --git a/src/java/org/apache/cassandra/cql3/Relation.java b/src/java/org/apache/cassandra/cql3/Relation.java index 1337096b90..334464f067 100644 --- a/src/java/org/apache/cassandra/cql3/Relation.java +++ b/src/java/org/apache/cassandra/cql3/Relation.java @@ -38,6 +38,16 @@ public abstract class Relation { return relationType; } + /** + * Returns the raw value for this relation, or null if this is an IN relation. + */ + public abstract Term.Raw getValue(); + + /** + * Returns the list of raw IN values for this relation, or null if this is not an IN relation. + */ + public abstract List getInValues(); + /** * Checks if this relation apply to multiple columns. * @@ -132,6 +142,7 @@ public abstract class Relation { case IN: return newINRestriction(cfm, boundNames); case CONTAINS: return newContainsRestriction(cfm, boundNames, false); case CONTAINS_KEY: return newContainsRestriction(cfm, boundNames, true); + case IS_NOT: return newIsNotRestriction(cfm, boundNames); default: throw invalidRequest("Unsupported \"!=\" relation: %s", this); } } @@ -186,6 +197,9 @@ public abstract class Relation { VariableSpecifications boundNames, boolean isKey) throws InvalidRequestException; + protected abstract Restriction newIsNotRestriction(CFMetaData cfm, + VariableSpecifications boundNames) throws InvalidRequestException; + /** * Converts the specified Raw into a Term. * @param receivers the columns to which the values must be associated at @@ -246,4 +260,13 @@ public abstract class Relation { return def; } + + /** + * Renames an identifier in this Relation, if applicable. + * @param from the old identifier + * @param to the new identifier + * @return this object, if the old identifier is not in the set of entities that this relation covers; otherwise + * a new Relation with "from" replaced by "to" is returned. + */ + public abstract Relation renameIdentifier(ColumnIdentifier.Raw from, ColumnIdentifier.Raw to); } diff --git a/src/java/org/apache/cassandra/cql3/Sets.java b/src/java/org/apache/cassandra/cql3/Sets.java index 7ff38154ac..010abaaf78 100644 --- a/src/java/org/apache/cassandra/cql3/Sets.java +++ b/src/java/org/apache/cassandra/cql3/Sets.java @@ -21,6 +21,7 @@ import static org.apache.cassandra.cql3.Constants.UNSET_VALUE; import java.nio.ByteBuffer; import java.util.*; +import java.util.stream.Collectors; import com.google.common.base.Joiner; @@ -48,7 +49,7 @@ public abstract class Sets return new ColumnSpecification(column.ksName, column.cfName, new ColumnIdentifier("value(" + column.name + ")", true), ((SetType)column.type).getElementsType()); } - public static class Literal implements Term.Raw + public static class Literal extends Term.Raw { private final List elements; @@ -124,10 +125,9 @@ public abstract class Sets return AssignmentTestable.TestResult.testAll(keyspace, valueSpec, elements); } - @Override - public String toString() + public String getText() { - return "{" + Joiner.on(", ").join(elements) + "}"; + return elements.stream().map(Term.Raw::getText).collect(Collectors.joining(", ", "{", "}")); } } diff --git a/src/java/org/apache/cassandra/cql3/SingleColumnRelation.java b/src/java/org/apache/cassandra/cql3/SingleColumnRelation.java index 84e6274f10..e2c0b79cec 100644 --- a/src/java/org/apache/cassandra/cql3/SingleColumnRelation.java +++ b/src/java/org/apache/cassandra/cql3/SingleColumnRelation.java @@ -54,6 +54,9 @@ public final class SingleColumnRelation extends Relation this.relationType = type; this.value = value; this.inValues = inValues; + + if (type == Operator.IS_NOT) + assert value == Constants.NULL_LITERAL; } /** @@ -81,6 +84,16 @@ public final class SingleColumnRelation extends Relation this(entity, null, type, value); } + public Term.Raw getValue() + { + return value; + } + + public List getInValues() + { + return inValues; + } + public static SingleColumnRelation createInRelation(ColumnIdentifier.Raw entity, List inValues) { return new SingleColumnRelation(entity, null, Operator.IN, null, inValues); @@ -120,6 +133,13 @@ public final class SingleColumnRelation extends Relation } } + public Relation renameIdentifier(ColumnIdentifier.Raw from, ColumnIdentifier.Raw to) + { + return entity.equals(from) + ? new SingleColumnRelation(to, mapKey, operator(), value, inValues) + : this; + } + @Override public String toString() { @@ -185,6 +205,16 @@ public final class SingleColumnRelation extends Relation return new SingleColumnRestriction.ContainsRestriction(columnDef, term, isKey); } + @Override + protected Restriction newIsNotRestriction(CFMetaData cfm, + VariableSpecifications boundNames) throws InvalidRequestException + { + ColumnDefinition columnDef = toColumnDefinition(cfm, entity); + // currently enforced by the grammar + assert value == Constants.NULL_LITERAL : "Expected null literal for IS NOT relation: " + this.toString(); + return new SingleColumnRestriction.IsNotNullRestriction(columnDef); + } + /** * Returns the receivers for this relation. * @param columnDef the column definition diff --git a/src/java/org/apache/cassandra/cql3/Term.java b/src/java/org/apache/cassandra/cql3/Term.java index 6fa0c7621d..1f6bc62aaf 100644 --- a/src/java/org/apache/cassandra/cql3/Term.java +++ b/src/java/org/apache/cassandra/cql3/Term.java @@ -81,7 +81,7 @@ public interface Term * - a function call * - a marker */ - public interface Raw extends AssignmentTestable + public abstract class Raw implements AssignmentTestable { /** * This method validates this RawTerm is valid for provided column @@ -93,12 +93,23 @@ public interface Term * case this RawTerm describe a list index or a map key, etc... * @return the prepared term. */ - public Term prepare(String keyspace, ColumnSpecification receiver) throws InvalidRequestException; + public abstract Term prepare(String keyspace, ColumnSpecification receiver) throws InvalidRequestException; + + /** + * @return a String representation of the raw term that can be used when reconstructing a CQL query string. + */ + public abstract String getText(); + + @Override + public String toString() + { + return getText(); + } } - public interface MultiColumnRaw extends Raw + public abstract class MultiColumnRaw extends Term.Raw { - public Term prepare(String keyspace, List receiver) throws InvalidRequestException; + public abstract Term prepare(String keyspace, List receiver) throws InvalidRequestException; } /** diff --git a/src/java/org/apache/cassandra/cql3/TokenRelation.java b/src/java/org/apache/cassandra/cql3/TokenRelation.java index 6b487ef26a..2c13b199f9 100644 --- a/src/java/org/apache/cassandra/cql3/TokenRelation.java +++ b/src/java/org/apache/cassandra/cql3/TokenRelation.java @@ -20,6 +20,7 @@ package org.apache.cassandra.cql3; import java.util.ArrayList; import java.util.Collections; import java.util.List; +import java.util.stream.Collectors; import com.google.common.base.Joiner; @@ -63,6 +64,16 @@ public final class TokenRelation extends Relation return true; } + public Term.Raw getValue() + { + return value; + } + + public List getInValues() + { + return null; + } + @Override protected Restriction newEQRestriction(CFMetaData cfm, VariableSpecifications boundNames) throws InvalidRequestException { @@ -94,6 +105,12 @@ public final class TokenRelation extends Relation throw invalidRequest("%s cannot be used with the token function", operator()); } + @Override + protected Restriction newIsNotRestriction(CFMetaData cfm, VariableSpecifications boundNames) throws InvalidRequestException + { + throw invalidRequest("%s cannot be used with the token function", operator()); + } + @Override protected Term toTerm(List receivers, Raw raw, @@ -105,6 +122,15 @@ public final class TokenRelation extends Relation return term; } + public Relation renameIdentifier(ColumnIdentifier.Raw from, ColumnIdentifier.Raw to) + { + if (!entities.contains(from)) + return this; + + List newEntities = entities.stream().map(e -> e.equals(from) ? to : e).collect(Collectors.toList()); + return new TokenRelation(newEntities, operator(), value); + } + @Override public String toString() { diff --git a/src/java/org/apache/cassandra/cql3/Tuples.java b/src/java/org/apache/cassandra/cql3/Tuples.java index 933088f3ad..6c7df47918 100644 --- a/src/java/org/apache/cassandra/cql3/Tuples.java +++ b/src/java/org/apache/cassandra/cql3/Tuples.java @@ -21,6 +21,7 @@ import java.nio.ByteBuffer; import java.util.ArrayList; import java.util.Arrays; import java.util.List; +import java.util.stream.Collectors; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -53,7 +54,7 @@ public class Tuples * A raw, literal tuple. When prepared, this will become a Tuples.Value or Tuples.DelayedValue, depending * on whether the tuple holds NonTerminals. */ - public static class Literal implements Term.MultiColumnRaw + public static class Literal extends Term.MultiColumnRaw { private final List elements; @@ -133,10 +134,9 @@ public class Tuples } } - @Override - public String toString() + public String getText() { - return tupleToString(elements); + return elements.stream().map(Term.Raw::getText).collect(Collectors.joining(", ", "(", ")")); } } @@ -287,7 +287,7 @@ public class Tuples * For example, "SELECT ... WHERE (col1, col2) > ?". * } */ - public static class Raw extends AbstractMarker.Raw implements Term.MultiColumnRaw + public static class Raw extends AbstractMarker.MultiColumnRaw { public Raw(int bindIndex) { @@ -317,18 +317,12 @@ public class Tuples { return new Tuples.Marker(bindIndex, makeReceiver(receivers)); } - - @Override - public AbstractMarker prepare(String keyspace, ColumnSpecification receiver) - { - throw new AssertionError("Tuples.Raw.prepare() requires a list of receivers"); - } } /** * A raw marker for an IN list of tuples, like "SELECT ... WHERE (a, b, c) IN ?" */ - public static class INRaw extends AbstractMarker.Raw implements MultiColumnRaw + public static class INRaw extends AbstractMarker.MultiColumnRaw { public INRaw(int bindIndex) { @@ -362,12 +356,6 @@ public class Tuples { return new InMarker(bindIndex, makeInReceiver(receivers)); } - - @Override - public AbstractMarker prepare(String keyspace, ColumnSpecification receiver) - { - throw new AssertionError("Tuples.INRaw.prepare() requires a list of receivers"); - } } /** diff --git a/src/java/org/apache/cassandra/cql3/TypeCast.java b/src/java/org/apache/cassandra/cql3/TypeCast.java index 561a1585f9..890b34f76a 100644 --- a/src/java/org/apache/cassandra/cql3/TypeCast.java +++ b/src/java/org/apache/cassandra/cql3/TypeCast.java @@ -20,7 +20,7 @@ package org.apache.cassandra.cql3; import org.apache.cassandra.db.marshal.AbstractType; import org.apache.cassandra.exceptions.InvalidRequestException; -public class TypeCast implements Term.Raw +public class TypeCast extends Term.Raw { private final CQL3Type.Raw type; private final Term.Raw term; @@ -58,8 +58,7 @@ public class TypeCast implements Term.Raw return AssignmentTestable.TestResult.NOT_ASSIGNABLE; } - @Override - public String toString() + public String getText() { return "(" + type + ")" + term; } diff --git a/src/java/org/apache/cassandra/cql3/UserTypes.java b/src/java/org/apache/cassandra/cql3/UserTypes.java index 22c7987c11..0beff065f1 100644 --- a/src/java/org/apache/cassandra/cql3/UserTypes.java +++ b/src/java/org/apache/cassandra/cql3/UserTypes.java @@ -42,7 +42,7 @@ public abstract class UserTypes ut.fieldType(field)); } - public static class Literal implements Term.Raw + public static class Literal extends Term.Raw { public final Map entries; @@ -118,8 +118,7 @@ public abstract class UserTypes } } - @Override - public String toString() + public String getText() { StringBuilder sb = new StringBuilder(); sb.append("{"); @@ -127,7 +126,7 @@ public abstract class UserTypes while (iter.hasNext()) { Map.Entry entry = iter.next(); - sb.append(entry.getKey()).append(":").append(entry.getValue()); + sb.append(entry.getKey()).append(": ").append(entry.getValue().getText()); if (iter.hasNext()) sb.append(", "); } diff --git a/src/java/org/apache/cassandra/cql3/functions/FunctionCall.java b/src/java/org/apache/cassandra/cql3/functions/FunctionCall.java index b25d079e4f..1766a796bc 100644 --- a/src/java/org/apache/cassandra/cql3/functions/FunctionCall.java +++ b/src/java/org/apache/cassandra/cql3/functions/FunctionCall.java @@ -20,6 +20,7 @@ package org.apache.cassandra.cql3.functions; import java.nio.ByteBuffer; import java.util.ArrayList; import java.util.List; +import java.util.stream.Collectors; import com.google.common.collect.Iterables; @@ -110,7 +111,7 @@ public class FunctionCall extends Term.NonTerminal throw new AssertionError(); } - public static class Raw implements Term.Raw + public static class Raw extends Term.Raw { private FunctionName name; private final List terms; @@ -181,18 +182,9 @@ public class FunctionCall extends Term.NonTerminal } } - @Override - public String toString() + public String getText() { - StringBuilder sb = new StringBuilder(); - sb.append(name).append("("); - for (int i = 0; i < terms.size(); i++) - { - if (i > 0) - sb.append(", "); - sb.append(terms.get(i)); - } - return sb.append(")").toString(); + return name + terms.stream().map(Term.Raw::getText).collect(Collectors.joining(", ", "(", ")")); } } } diff --git a/src/java/org/apache/cassandra/cql3/restrictions/AbstractRestriction.java b/src/java/org/apache/cassandra/cql3/restrictions/AbstractRestriction.java index 0f56fd9a43..023c2acab1 100644 --- a/src/java/org/apache/cassandra/cql3/restrictions/AbstractRestriction.java +++ b/src/java/org/apache/cassandra/cql3/restrictions/AbstractRestriction.java @@ -17,17 +17,9 @@ */ package org.apache.cassandra.cql3.restrictions; -import java.nio.ByteBuffer; - -import org.apache.cassandra.cql3.ColumnSpecification; import org.apache.cassandra.cql3.QueryOptions; import org.apache.cassandra.cql3.statements.Bound; import org.apache.cassandra.db.MultiCBuilder; -import org.apache.cassandra.exceptions.InvalidRequestException; - -import static org.apache.cassandra.cql3.statements.RequestValidations.checkBindValueSet; -import static org.apache.cassandra.cql3.statements.RequestValidations.checkFalse; -import static org.apache.cassandra.cql3.statements.RequestValidations.checkNotNull; /** * Base class for Restrictions @@ -70,6 +62,12 @@ abstract class AbstractRestriction implements Restriction return false; } + @Override + public boolean isNotNull() + { + return false; + } + @Override public boolean hasBound(Bound b) { diff --git a/src/java/org/apache/cassandra/cql3/restrictions/ForwardingPrimaryKeyRestrictions.java b/src/java/org/apache/cassandra/cql3/restrictions/ForwardingPrimaryKeyRestrictions.java index f82cb11eda..18e7105d4d 100644 --- a/src/java/org/apache/cassandra/cql3/restrictions/ForwardingPrimaryKeyRestrictions.java +++ b/src/java/org/apache/cassandra/cql3/restrictions/ForwardingPrimaryKeyRestrictions.java @@ -166,6 +166,12 @@ abstract class ForwardingPrimaryKeyRestrictions implements PrimaryKeyRestriction return getDelegate().isContains(); } + @Override + public boolean isNotNull() + { + return getDelegate().isNotNull(); + } + @Override public boolean isMultiColumn() { diff --git a/src/java/org/apache/cassandra/cql3/restrictions/MultiColumnRestriction.java b/src/java/org/apache/cassandra/cql3/restrictions/MultiColumnRestriction.java index b60930e430..069a01b8bf 100644 --- a/src/java/org/apache/cassandra/cql3/restrictions/MultiColumnRestriction.java +++ b/src/java/org/apache/cassandra/cql3/restrictions/MultiColumnRestriction.java @@ -459,4 +459,59 @@ public abstract class MultiColumnRestriction extends AbstractRestriction return Collections.singletonList(terminal.get(options.getProtocolVersion())); } } + + public static class NotNullRestriction extends MultiColumnRestriction + { + public NotNullRestriction(List columnDefs) + { + super(columnDefs); + assert columnDefs.size() == 1; + } + + @Override + public Iterable getFunctions() + { + return Collections.emptyList(); + } + + @Override + public boolean isNotNull() + { + return true; + } + + @Override + public String toString() + { + return "IS NOT NULL"; + } + + @Override + public Restriction doMergeWith(Restriction otherRestriction) throws InvalidRequestException + { + throw invalidRequest("%s cannot be restricted by a relation if it includes an IS NOT NULL clause", + getColumnsInCommons(otherRestriction)); + } + + @Override + protected boolean isSupportedBy(Index index) + { + for(ColumnDefinition column : columnDefs) + if (index.supportsExpression(column, Operator.IS_NOT)) + return true; + return false; + } + + @Override + public MultiCBuilder appendTo(MultiCBuilder builder, QueryOptions options) + { + throw new UnsupportedOperationException("Cannot use IS NOT NULL restriction for slicing"); + } + + @Override + public final void addRowFilterTo(RowFilter filter, SecondaryIndexManager indexMananger, QueryOptions options) throws InvalidRequestException + { + throw new UnsupportedOperationException("Secondary indexes do not support IS NOT NULL restrictions"); + } + } } diff --git a/src/java/org/apache/cassandra/cql3/restrictions/Restriction.java b/src/java/org/apache/cassandra/cql3/restrictions/Restriction.java index a21087e41b..a84ebc4615 100644 --- a/src/java/org/apache/cassandra/cql3/restrictions/Restriction.java +++ b/src/java/org/apache/cassandra/cql3/restrictions/Restriction.java @@ -41,6 +41,7 @@ public interface Restriction public boolean isEQ(); public boolean isIN(); public boolean isContains(); + public boolean isNotNull(); public boolean isMultiColumn(); /** diff --git a/src/java/org/apache/cassandra/cql3/restrictions/SingleColumnRestriction.java b/src/java/org/apache/cassandra/cql3/restrictions/SingleColumnRestriction.java index 25146e5cd6..d85125361d 100644 --- a/src/java/org/apache/cassandra/cql3/restrictions/SingleColumnRestriction.java +++ b/src/java/org/apache/cassandra/cql3/restrictions/SingleColumnRestriction.java @@ -587,4 +587,62 @@ public abstract class SingleColumnRestriction extends AbstractRestriction super(columnDef); } } + + public static final class IsNotNullRestriction extends SingleColumnRestriction + { + public IsNotNullRestriction(ColumnDefinition columnDef) + { + super(columnDef); + } + + @Override + public Iterable getFunctions() + { + return Collections.emptyList(); + } + + @Override + public boolean isNotNull() + { + return true; + } + + @Override + MultiColumnRestriction toMultiColumnRestriction() + { + return new MultiColumnRestriction.NotNullRestriction(Collections.singletonList(columnDef)); + } + + @Override + public void addRowFilterTo(RowFilter filter, + SecondaryIndexManager indexManager, + QueryOptions options) + { + throw new UnsupportedOperationException("Secondary indexes do not support IS NOT NULL restrictions"); + } + + @Override + public MultiCBuilder appendTo(MultiCBuilder builder, QueryOptions options) + { + throw new UnsupportedOperationException("Cannot use IS NOT NULL restriction for slicing"); + } + + @Override + public String toString() + { + return "IS NOT NULL"; + } + + @Override + public Restriction doMergeWith(Restriction otherRestriction) throws InvalidRequestException + { + throw invalidRequest("%s cannot be restricted by a relation if it includes an IS NOT NULL", columnDef.name); + } + + @Override + protected boolean isSupportedBy(Index index) + { + return index.supportsExpression(columnDef, Operator.IS_NOT); + } + } } diff --git a/src/java/org/apache/cassandra/cql3/restrictions/StatementRestrictions.java b/src/java/org/apache/cassandra/cql3/restrictions/StatementRestrictions.java index b1c7aff441..1bd42183c2 100644 --- a/src/java/org/apache/cassandra/cql3/restrictions/StatementRestrictions.java +++ b/src/java/org/apache/cassandra/cql3/restrictions/StatementRestrictions.java @@ -77,6 +77,8 @@ public final class StatementRestrictions */ private RestrictionSet nonPrimaryKeyRestrictions; + private Set notNullColumns; + /** * The restrictions used to build the row filter */ @@ -111,6 +113,7 @@ public final class StatementRestrictions this.partitionKeyRestrictions = new PrimaryKeyRestrictionSet(cfm.getKeyValidatorAsClusteringComparator(), true); this.clusteringColumnsRestrictions = new PrimaryKeyRestrictionSet(cfm.comparator, false); this.nonPrimaryKeyRestrictions = new RestrictionSet(); + this.notNullColumns = new HashSet<>(); } public StatementRestrictions(StatementType type, @@ -119,7 +122,8 @@ public final class StatementRestrictions VariableSpecifications boundNames, boolean selectsOnlyStaticColumns, boolean selectACollection, - boolean useFiltering) + boolean useFiltering, + boolean forView) throws InvalidRequestException { this(type, cfm); @@ -133,7 +137,20 @@ public final class StatementRestrictions * in CQL so far) */ for (Relation relation : whereClause.relations) - addRestriction(relation.toRestriction(cfm, boundNames)); + { + if (relation.operator() == Operator.IS_NOT) + { + if (!forView) + throw new InvalidRequestException("Unsupported restriction: " + relation); + + for (ColumnDefinition def : relation.toRestriction(cfm, boundNames).getColumnDefs()) + this.notNullColumns.add(def); + } + else + { + addRestriction(relation.toRestriction(cfm, boundNames)); + } + } boolean hasQueriableClusteringColumnIndex = false; boolean hasQueriableIndex = false; @@ -180,7 +197,7 @@ public final class StatementRestrictions throw invalidRequest("Cannot restrict clustering columns when selecting only static columns"); } - processClusteringColumnsRestrictions(hasQueriableIndex, selectsOnlyStaticColumns, selectACollection); + processClusteringColumnsRestrictions(hasQueriableIndex, selectsOnlyStaticColumns, selectACollection, forView); // Covers indexes on the first clustering column (among others). if (isKeyRange && hasQueriableClusteringColumnIndex) @@ -244,18 +261,55 @@ public final class StatementRestrictions } /** - * Returns the non-PK column that are restricted. + * Returns the non-PK column that are restricted. If includeNotNullRestrictions is true, columns that are restricted + * by an IS NOT NULL restriction will be included, otherwise they will not be included (unless another restriction + * applies to them). */ - public Set nonPKRestrictedColumns() + public Set nonPKRestrictedColumns(boolean includeNotNullRestrictions) { Set columns = new HashSet<>(); for (Restrictions r : indexRestrictions.getRestrictions()) + { for (ColumnDefinition def : r.getColumnDefs()) if (!def.isPrimaryKeyColumn()) columns.add(def); + } + + if (includeNotNullRestrictions) + { + for (ColumnDefinition def : notNullColumns) + { + if (!def.isPrimaryKeyColumn()) + columns.add(def); + } + } + return columns; } + /** + * @return the set of columns that have an IS NOT NULL restriction on them + */ + public Set notNullColumns() + { + return notNullColumns; + } + + /** + * @return true if column is restricted by some restriction, false otherwise + */ + public boolean isRestricted(ColumnDefinition column) + { + if (notNullColumns.contains(column)) + return true; + else if (column.isPartitionKey()) + return partitionKeyRestrictions.getColumnDefs().contains(column); + else if (column.isClusteringColumn()) + return clusteringColumnsRestrictions.getColumnDefs().contains(column); + else + return nonPrimaryKeyRestrictions.getColumnDefs().contains(column); + } + /** * Checks if the restrictions on the partition key is an IN restriction. * @@ -370,7 +424,8 @@ public final class StatementRestrictions */ private void processClusteringColumnsRestrictions(boolean hasQueriableIndex, boolean selectsOnlyStaticColumns, - boolean selectACollection) + boolean selectACollection, + boolean forView) throws InvalidRequestException { checkFalse(!type.allowClusteringColumnSlices() && clusteringColumnsRestrictions.isSlice(), "Slice restrictions are not supported on the clustering columns in %s statements", type); @@ -401,7 +456,7 @@ public final class StatementRestrictions if (!clusteringColumn.equals(restrictedColumn)) { - checkTrue(hasQueriableIndex, + checkTrue(hasQueriableIndex || forView, "PRIMARY KEY column \"%s\" cannot be restricted as preceding column \"%s\" is not restricted", restrictedColumn.name, clusteringColumn.name); diff --git a/src/java/org/apache/cassandra/cql3/statements/AlterTableStatement.java b/src/java/org/apache/cassandra/cql3/statements/AlterTableStatement.java index 0d2011b334..c410f10144 100644 --- a/src/java/org/apache/cassandra/cql3/statements/AlterTableStatement.java +++ b/src/java/org/apache/cassandra/cql3/statements/AlterTableStatement.java @@ -341,7 +341,7 @@ public class AlterTableStatement extends SchemaAlteringStatement ViewDefinition viewCopy = view.copy(); ColumnIdentifier viewFrom = entry.getKey().prepare(viewCopy.metadata); ColumnIdentifier viewTo = entry.getValue().prepare(viewCopy.metadata); - viewCopy.metadata.renameColumn(viewFrom, viewTo); + viewCopy.renameColumn(viewFrom, viewTo); if (viewUpdates == null) viewUpdates = new ArrayList<>(); diff --git a/src/java/org/apache/cassandra/cql3/statements/CreateViewStatement.java b/src/java/org/apache/cassandra/cql3/statements/CreateViewStatement.java index 1a020ce16d..586b09b0ae 100644 --- a/src/java/org/apache/cassandra/cql3/statements/CreateViewStatement.java +++ b/src/java/org/apache/cassandra/cql3/statements/CreateViewStatement.java @@ -18,10 +18,8 @@ package org.apache.cassandra.cql3.statements; -import java.util.ArrayList; -import java.util.HashSet; -import java.util.List; -import java.util.Set; +import java.util.*; +import java.util.stream.Collectors; import com.google.common.collect.Iterables; @@ -30,8 +28,8 @@ import org.apache.cassandra.config.CFMetaData; import org.apache.cassandra.config.ColumnDefinition; import org.apache.cassandra.config.Schema; import org.apache.cassandra.config.ViewDefinition; -import org.apache.cassandra.cql3.CFName; -import org.apache.cassandra.cql3.ColumnIdentifier; +import org.apache.cassandra.cql3.*; +import org.apache.cassandra.cql3.restrictions.StatementRestrictions; import org.apache.cassandra.cql3.selection.RawSelector; import org.apache.cassandra.cql3.selection.Selectable; import org.apache.cassandra.db.marshal.AbstractType; @@ -52,7 +50,7 @@ public class CreateViewStatement extends SchemaAlteringStatement { private final CFName baseName; private final List selectClause; - private final List notNullWhereClause; + private final WhereClause whereClause; private final List partitionKeys; private final List clusteringKeys; public final CFProperties properties = new CFProperties(); @@ -61,7 +59,7 @@ public class CreateViewStatement extends SchemaAlteringStatement public CreateViewStatement(CFName viewName, CFName baseName, List selectClause, - List notNullWhereClause, + WhereClause whereClause, List partitionKeys, List clusteringKeys, boolean ifNotExists) @@ -69,7 +67,7 @@ public class CreateViewStatement extends SchemaAlteringStatement super(viewName); this.baseName = baseName; this.selectClause = selectClause; - this.notNullWhereClause = notNullWhereClause; + this.whereClause = whereClause; this.partitionKeys = partitionKeys; this.clusteringKeys = clusteringKeys; this.ifNotExists = ifNotExists; @@ -194,34 +192,50 @@ public class CreateViewStatement extends SchemaAlteringStatement throw new InvalidRequestException(String.format("Cannot use Static column '%s' in PRIMARY KEY of materialized view", identifier)); } + // build the select statement + Map orderings = Collections.emptyMap(); + SelectStatement.Parameters parameters = new SelectStatement.Parameters(orderings, false, true, false); + SelectStatement.RawStatement rawSelect = new SelectStatement.RawStatement(baseName, parameters, selectClause, whereClause, null); + + ClientState state = ClientState.forInternalCalls(); + state.setKeyspace(keyspace()); + + rawSelect.prepareKeyspace(state); + rawSelect.setBoundVariables(getBoundVariables()); + + ParsedStatement.Prepared prepared = rawSelect.prepare(true); + SelectStatement select = (SelectStatement) prepared.statement; + StatementRestrictions restrictions = select.getRestrictions(); + + if (!prepared.boundNames.isEmpty()) + throw new InvalidRequestException("Cannot use query parameters in CREATE MATERIALIZED VIEW statements"); + + if (!restrictions.nonPKRestrictedColumns(false).isEmpty()) + { + throw new InvalidRequestException(String.format( + "Non-primary key columns cannot be restricted in the SELECT statement used for materialized view " + + "creation (got restrictions on: %s)", + restrictions.nonPKRestrictedColumns(false).stream().map(def -> def.name.toString()).collect(Collectors.joining(", ")))); + } + + String whereClauseText = View.relationsToWhereClause(whereClause.relations); + Set basePrimaryKeyCols = new HashSet<>(); for (ColumnDefinition definition : Iterables.concat(cfm.partitionKeyColumns(), cfm.clusteringColumns())) basePrimaryKeyCols.add(definition.name); List targetClusteringColumns = new ArrayList<>(); List targetPartitionKeys = new ArrayList<>(); - Set notNullColumns = new HashSet<>(); - if (notNullWhereClause != null) - { - for (ColumnIdentifier.Raw raw : notNullWhereClause) - { - notNullColumns.add(raw.prepare(cfm)); - } - } // This is only used as an intermediate state; this is to catch whether multiple non-PK columns are used boolean hasNonPKColumn = false; for (ColumnIdentifier.Raw raw : partitionKeys) - { - hasNonPKColumn = getColumnIdentifier(cfm, basePrimaryKeyCols, hasNonPKColumn, raw, targetPartitionKeys, notNullColumns); - } + hasNonPKColumn = getColumnIdentifier(cfm, basePrimaryKeyCols, hasNonPKColumn, raw, targetPartitionKeys, restrictions); for (ColumnIdentifier.Raw raw : clusteringKeys) - { - hasNonPKColumn = getColumnIdentifier(cfm, basePrimaryKeyCols, hasNonPKColumn, raw, targetClusteringColumns, notNullColumns); - } + hasNonPKColumn = getColumnIdentifier(cfm, basePrimaryKeyCols, hasNonPKColumn, raw, targetClusteringColumns, restrictions); - // We need to include all of the primary key colums from the base table in order to make sure that we do not + // We need to include all of the primary key columns from the base table in order to make sure that we do not // overwrite values in the view. We cannot support "collapsing" the base table into a smaller number of rows in // the view because if we need to generate a tombstone, we have no way of knowing which value is currently being // used in the view and whether or not to generate a tombstone. In order to not surprise our users, we require @@ -269,7 +283,10 @@ public class CreateViewStatement extends SchemaAlteringStatement ViewDefinition definition = new ViewDefinition(keyspace(), columnFamily(), Schema.instance.getId(keyspace(), baseName.getColumnFamily()), + baseName.getColumnFamily(), included.isEmpty(), + rawSelect, + whereClauseText, viewCfm); try @@ -291,24 +308,21 @@ public class CreateViewStatement extends SchemaAlteringStatement boolean hasNonPKColumn, ColumnIdentifier.Raw raw, List columns, - Set allowedPKColumns) + StatementRestrictions restrictions) { ColumnIdentifier identifier = raw.prepare(cfm); + ColumnDefinition def = cfm.getColumnDefinition(identifier); boolean isPk = basePK.contains(identifier); if (!isPk && hasNonPKColumn) - { throw new InvalidRequestException(String.format("Cannot include more than one non-primary key column '%s' in materialized view partition key", identifier)); - } // We don't need to include the "IS NOT NULL" filter on a non-composite partition key // because we will never allow a single partition key to be NULL boolean isSinglePartitionKey = cfm.getColumnDefinition(identifier).isPartitionKey() && cfm.partitionKeyColumns().size() == 1; - if (!allowedPKColumns.remove(identifier) && !isSinglePartitionKey) - { + if (!isSinglePartitionKey && !restrictions.isRestricted(def)) throw new InvalidRequestException(String.format("Primary key column '%s' is required to be filtered by 'IS NOT NULL'", identifier)); - } columns.add(identifier); return !isPk; diff --git a/src/java/org/apache/cassandra/cql3/statements/IndexTarget.java b/src/java/org/apache/cassandra/cql3/statements/IndexTarget.java index 6210a864f8..8cdf2c881c 100644 --- a/src/java/org/apache/cassandra/cql3/statements/IndexTarget.java +++ b/src/java/org/apache/cassandra/cql3/statements/IndexTarget.java @@ -60,18 +60,10 @@ public class IndexTarget public String asCqlString(CFMetaData cfm) { - if (! cfm.getColumnDefinition(column).type.isCollection()) - return maybeEscapeQuotedName(column.toString()); + if (!cfm.getColumnDefinition(column).type.isCollection()) + return column.toCQLString(); - return String.format("%s(%s)", type.toString(), maybeEscapeQuotedName(column.toString())); - } - - // Quoted column names may themselves contain quotes, these need - // to be escaped with a preceding quote when written out as cql. - // Of course, the escaped name also needs to be wrapped in quotes. - private String maybeEscapeQuotedName(String name) - { - return quoteName ? '\"' + name.replace("\"", "\"\"") + '\"' : name; + return String.format("%s(%s)", type.toString(), column.toCQLString()); } public static class Raw diff --git a/src/java/org/apache/cassandra/cql3/statements/ModificationStatement.java b/src/java/org/apache/cassandra/cql3/statements/ModificationStatement.java index 54a4f28653..23a26d0db7 100644 --- a/src/java/org/apache/cassandra/cql3/statements/ModificationStatement.java +++ b/src/java/org/apache/cassandra/cql3/statements/ModificationStatement.java @@ -855,7 +855,7 @@ public abstract class ModificationStatement implements CQLStatement throw new InvalidRequestException(CUSTOM_EXPRESSIONS_NOT_ALLOWED); boolean applyOnlyToStaticColumns = appliesOnlyToStaticColumns(operations, conditions); - return new StatementRestrictions(type, cfm, where, boundNames, applyOnlyToStaticColumns, false, false); + return new StatementRestrictions(type, cfm, where, boundNames, applyOnlyToStaticColumns, false, false, false); } /** diff --git a/src/java/org/apache/cassandra/cql3/statements/ParsedStatement.java b/src/java/org/apache/cassandra/cql3/statements/ParsedStatement.java index 539a957dc7..4c3f8a9414 100644 --- a/src/java/org/apache/cassandra/cql3/statements/ParsedStatement.java +++ b/src/java/org/apache/cassandra/cql3/statements/ParsedStatement.java @@ -39,6 +39,11 @@ public abstract class ParsedStatement this.variables = new VariableSpecifications(boundNames); } + public void setBoundVariables(VariableSpecifications variables) + { + this.variables = variables; + } + public abstract Prepared prepare() throws RequestValidationException; public static class Prepared diff --git a/src/java/org/apache/cassandra/cql3/statements/SelectStatement.java b/src/java/org/apache/cassandra/cql3/statements/SelectStatement.java index 170bfdf066..7848556e6f 100644 --- a/src/java/org/apache/cassandra/cql3/statements/SelectStatement.java +++ b/src/java/org/apache/cassandra/cql3/statements/SelectStatement.java @@ -143,7 +143,7 @@ public class SelectStatement implements CQLStatement if (!def.isPrimaryKeyColumn()) builder.add(def); // as well as any restricted column (so we can actually apply the restriction) - builder.addAll(restrictions.nonPKRestrictedColumns()); + builder.addAll(restrictions.nonPKRestrictedColumns(true)); return builder.build(); } @@ -451,6 +451,17 @@ public class SelectStatement implements CQLStatement return new SinglePartitionReadCommand.Group(commands, limit); } + /** + * Returns a read command that can be used internally to filter individual rows for materialized views. + */ + public SinglePartitionReadCommand internalReadForView(DecoratedKey key, int nowInSec) + { + QueryOptions options = QueryOptions.forInternalCalls(Collections.emptyList()); + ClusteringIndexFilter filter = makeClusteringIndexFilter(options); + RowFilter rowFilter = getRowFilter(options); + return SinglePartitionReadCommand.create(cfm, nowInSec, queriedColumns, rowFilter, DataLimits.NONE, key, filter); + } + private ReadQuery getRangeCommand(QueryOptions options, DataLimits limit, int nowInSec) throws RequestValidationException { ClusteringIndexFilter clusteringIndexFilter = makeClusteringIndexFilter(options); @@ -737,6 +748,11 @@ public class SelectStatement implements CQLStatement } public ParsedStatement.Prepared prepare() throws InvalidRequestException + { + return prepare(false); + } + + public ParsedStatement.Prepared prepare(boolean forView) throws InvalidRequestException { CFMetaData cfm = ThriftValidation.validateColumnFamily(keyspace(), columnFamily()); VariableSpecifications boundNames = getBoundVariables(); @@ -745,7 +761,7 @@ public class SelectStatement implements CQLStatement ? Selection.wildcard(cfm) : Selection.fromSelectors(cfm, selectClause); - StatementRestrictions restrictions = prepareRestrictions(cfm, boundNames, selection); + StatementRestrictions restrictions = prepareRestrictions(cfm, boundNames, selection, forView); if (parameters.isDistinct) validateDistinctSelection(cfm, selection, restrictions); @@ -755,6 +771,7 @@ public class SelectStatement implements CQLStatement if (!parameters.orderings.isEmpty()) { + assert !forView; verifyOrderingIsAllowed(restrictions); orderingComparator = getOrderingComparator(cfm, selection, restrictions); isReversed = isReversed(cfm); @@ -787,7 +804,8 @@ public class SelectStatement implements CQLStatement */ private StatementRestrictions prepareRestrictions(CFMetaData cfm, VariableSpecifications boundNames, - Selection selection) throws InvalidRequestException + Selection selection, + boolean forView) throws InvalidRequestException { try { @@ -797,7 +815,8 @@ public class SelectStatement implements CQLStatement boundNames, selection.containsOnlyStaticColumns(), selection.containsACollection(), - parameters.allowFiltering); + parameters.allowFiltering, + forView); } catch (UnrecognizedEntityException e) { diff --git a/src/java/org/apache/cassandra/cql3/statements/UpdateStatement.java b/src/java/org/apache/cassandra/cql3/statements/UpdateStatement.java index f8435eb8b2..ce9aaee73a 100644 --- a/src/java/org/apache/cassandra/cql3/statements/UpdateStatement.java +++ b/src/java/org/apache/cassandra/cql3/statements/UpdateStatement.java @@ -186,6 +186,7 @@ public class UpdateStatement extends ModificationStatement boundNames, applyOnlyToStaticColumns, false, + false, false); return new UpdateStatement(StatementType.INSERT, @@ -254,6 +255,7 @@ public class UpdateStatement extends ModificationStatement boundNames, applyOnlyToStaticColumns, false, + false, false); return new UpdateStatement(StatementType.INSERT, diff --git a/src/java/org/apache/cassandra/db/PartitionRangeReadCommand.java b/src/java/org/apache/cassandra/db/PartitionRangeReadCommand.java index 4e96d8199a..f17f3e3591 100644 --- a/src/java/org/apache/cassandra/db/PartitionRangeReadCommand.java +++ b/src/java/org/apache/cassandra/db/PartitionRangeReadCommand.java @@ -139,15 +139,22 @@ public class PartitionRangeReadCommand extends ReadCommand return DatabaseDescriptor.getRangeRpcTimeout(); } - public boolean selects(DecoratedKey partitionKey, Clustering clustering) + public boolean selectsKey(DecoratedKey key) { - if (!dataRange().contains(partitionKey)) + if (!dataRange().contains(key)) return false; + return rowFilter().partitionKeyRestrictionsAreSatisfiedBy(key, metadata().getKeyValidator()); + } + + public boolean selectsClustering(DecoratedKey key, Clustering clustering) + { if (clustering == Clustering.STATIC_CLUSTERING) return !columnFilter().fetchedColumns().statics.isEmpty(); - return dataRange().clusteringIndexFilter(partitionKey).selects(clustering); + if (!dataRange().clusteringIndexFilter(key).selects(clustering)) + return false; + return rowFilter().clusteringKeyRestrictionsAreSatisfiedBy(clustering); } public PartitionIterator execute(ConsistencyLevel consistency, ClientState clientState) throws RequestExecutionException diff --git a/src/java/org/apache/cassandra/db/ReadCommand.java b/src/java/org/apache/cassandra/db/ReadCommand.java index 3d08f17faa..1c1da425d0 100644 --- a/src/java/org/apache/cassandra/db/ReadCommand.java +++ b/src/java/org/apache/cassandra/db/ReadCommand.java @@ -276,18 +276,6 @@ public abstract class ReadCommand implements ReadQuery */ public abstract ReadCommand copy(); - /** - * Whether the provided row, identified by its primary key components, is selected by - * this read command. - * - * @param partitionKey the partition key for the row to test. - * @param clustering the clustering for the row to test. - * - * @return whether the row of partition key {@code partitionKey} and clustering - * {@code clustering} is selected by this command. - */ - public abstract boolean selects(DecoratedKey partitionKey, Clustering clustering); - protected abstract UnfilteredPartitionIterator queryStorage(ColumnFamilyStore cfs, ReadOrderGroup orderGroup); protected abstract int oldestUnrepairedTombstone(); diff --git a/src/java/org/apache/cassandra/db/ReadQuery.java b/src/java/org/apache/cassandra/db/ReadQuery.java index 3abffd5672..d1f52727cc 100644 --- a/src/java/org/apache/cassandra/db/ReadQuery.java +++ b/src/java/org/apache/cassandra/db/ReadQuery.java @@ -33,7 +33,7 @@ import org.apache.cassandra.service.pager.PagingState; */ public interface ReadQuery { - public static final ReadQuery EMPTY = new ReadQuery() + ReadQuery EMPTY = new ReadQuery() { public ReadOrderGroup startOrderGroup() { @@ -67,6 +67,16 @@ public interface ReadQuery { return QueryPager.EMPTY; } + + public boolean selectsKey(DecoratedKey key) + { + return false; + } + + public boolean selectsClustering(DecoratedKey key, Clustering clustering) + { + return false; + } }; /** @@ -116,4 +126,16 @@ public interface ReadQuery * @return The limits for the query. */ public DataLimits limits(); + + /** + * @return true if the read query would select the given key, including checks against the row filter, if + * checkRowFilter is true + */ + public boolean selectsKey(DecoratedKey key); + + /** + * @return true if the read query would select the given clustering, including checks against the row filter, if + * checkRowFilter is true + */ + public boolean selectsClustering(DecoratedKey key, Clustering clustering); } diff --git a/src/java/org/apache/cassandra/db/SinglePartitionReadCommand.java b/src/java/org/apache/cassandra/db/SinglePartitionReadCommand.java index 49cf07ca95..a8e37b4068 100644 --- a/src/java/org/apache/cassandra/db/SinglePartitionReadCommand.java +++ b/src/java/org/apache/cassandra/db/SinglePartitionReadCommand.java @@ -21,6 +21,7 @@ import java.io.IOException; import java.nio.ByteBuffer; import java.util.*; +import com.google.common.collect.Iterables; import org.apache.cassandra.cache.IRowCacheEntry; import org.apache.cassandra.cache.RowCacheKey; import org.apache.cassandra.cache.RowCacheSentinel; @@ -190,15 +191,23 @@ public abstract class SinglePartitionReadCommand c.selectsKey(key)); + } + + public boolean selectsClustering(DecoratedKey key, Clustering clustering) + { + return Iterables.any(commands, c -> c.selectsClustering(key, clustering)); + } + @Override public String toString() { diff --git a/src/java/org/apache/cassandra/db/filter/RowFilter.java b/src/java/org/apache/cassandra/db/filter/RowFilter.java index b5968d506e..0ff30af348 100644 --- a/src/java/org/apache/cassandra/db/filter/RowFilter.java +++ b/src/java/org/apache/cassandra/db/filter/RowFilter.java @@ -114,6 +114,45 @@ public abstract class RowFilter implements Iterable */ public abstract UnfilteredPartitionIterator filter(UnfilteredPartitionIterator iter, int nowInSec); + /** + * Returns true if all of the expressions within this filter that apply to the partition key are satisfied by + * the given key, false otherwise. + */ + public boolean partitionKeyRestrictionsAreSatisfiedBy(DecoratedKey key, AbstractType keyValidator) + { + for (Expression e : expressions) + { + if (!e.column.isPartitionKey()) + continue; + + ByteBuffer value = keyValidator instanceof CompositeType + ? ((CompositeType) keyValidator).split(key.getKey())[e.column.position()] + : key.getKey(); + if (!e.operator().isSatisfiedBy(e.column.type, value, e.value)) + return false; + } + return true; + } + + /** + * Returns true if all of the expressions within this filter that apply to the clustering key are satisfied by + * the given Clustering, false otherwise. + */ + public boolean clusteringKeyRestrictionsAreSatisfiedBy(Clustering clustering) + { + for (Expression e : expressions) + { + if (!e.column.isClusteringColumn()) + continue; + + if (!e.operator().isSatisfiedBy(e.column.type, clustering.get(e.column.position()), e.value)) + { + return false; + } + } + return true; + } + /** * Returns this filter but without the provided expression. This method * *assumes* that the filter contains the provided expression. diff --git a/src/java/org/apache/cassandra/db/view/TemporalRow.java b/src/java/org/apache/cassandra/db/view/TemporalRow.java index 6eb9071467..46dc3fa33b 100644 --- a/src/java/org/apache/cassandra/db/view/TemporalRow.java +++ b/src/java/org/apache/cassandra/db/view/TemporalRow.java @@ -27,6 +27,7 @@ import java.util.Iterator; import java.util.List; import java.util.Map; +import com.google.common.base.MoreObjects; import com.google.common.collect.Iterables; import org.apache.cassandra.config.CFMetaData; @@ -94,6 +95,18 @@ public class TemporalRow this.isNew = isNew; } + @Override + public String toString() + { + return MoreObjects.toStringHelper(this) + .add("value", value == null ? "null" : ByteBufferUtil.bytesToHex(value)) + .add("timestamp", timestamp) + .add("ttl", ttl) + .add("localDeletionTime", localDeletionTime) + .add("isNew", isNew) + .toString(); + } + public TemporalCell reconcile(TemporalCell that) { int now = FBUtilities.nowInSeconds(); @@ -208,13 +221,13 @@ public class TemporalRow if (cell.isNew) { - assert newCell == null || newCell.equals(cell) : "Only one cell version can be marked New"; + assert newCell == null || newCell.equals(cell) : "Only one cell version can be marked New; newCell: " + newCell + ", cell: " + cell; newCell = cell; numSet = existingCell == null ? 1 : 2; } else { - assert existingCell == null || existingCell.equals(cell) : "Only one cell version can be marked Existing"; + assert existingCell == null || existingCell.equals(cell) : "Only one cell version can be marked Existing; existingCell: " + existingCell + ", cell: " + cell; existingCell = cell; numSet = newCell == null ? 1 : 2; } diff --git a/src/java/org/apache/cassandra/db/view/View.java b/src/java/org/apache/cassandra/db/view/View.java index 28ec489742..0a7f747c3d 100644 --- a/src/java/org/apache/cassandra/db/view/View.java +++ b/src/java/org/apache/cassandra/db/view/View.java @@ -18,53 +18,30 @@ package org.apache.cassandra.db.view; import java.nio.ByteBuffer; -import java.util.ArrayList; -import java.util.Collection; -import java.util.HashSet; -import java.util.Iterator; -import java.util.LinkedList; -import java.util.List; -import java.util.Set; -import java.util.UUID; +import java.util.*; +import java.util.stream.Collectors; import javax.annotation.Nullable; import com.google.common.collect.Iterables; -import org.apache.cassandra.config.CFMetaData; -import org.apache.cassandra.config.ColumnDefinition; -import org.apache.cassandra.config.ViewDefinition; -import org.apache.cassandra.config.Schema; +import org.apache.cassandra.cql3.*; +import org.apache.cassandra.cql3.statements.ParsedStatement; +import org.apache.cassandra.cql3.statements.SelectStatement; +import org.apache.cassandra.db.*; +import org.apache.cassandra.config.*; import org.apache.cassandra.cql3.ColumnIdentifier; -import org.apache.cassandra.db.AbstractReadCommandBuilder; import org.apache.cassandra.db.AbstractReadCommandBuilder.SinglePartitionSliceBuilder; -import org.apache.cassandra.db.CBuilder; -import org.apache.cassandra.db.Clustering; -import org.apache.cassandra.db.ColumnFamilyStore; -import org.apache.cassandra.db.DecoratedKey; -import org.apache.cassandra.db.DeletionInfo; -import org.apache.cassandra.db.DeletionTime; -import org.apache.cassandra.db.LivenessInfo; -import org.apache.cassandra.db.Mutation; -import org.apache.cassandra.db.RangeTombstone; -import org.apache.cassandra.db.ReadCommand; -import org.apache.cassandra.db.ReadOrderGroup; -import org.apache.cassandra.db.SinglePartitionReadCommand; -import org.apache.cassandra.db.Slice; import org.apache.cassandra.db.compaction.CompactionManager; import org.apache.cassandra.db.partitions.AbstractBTreePartition; import org.apache.cassandra.db.partitions.PartitionIterator; import org.apache.cassandra.db.partitions.PartitionUpdate; -import org.apache.cassandra.db.rows.BTreeRow; -import org.apache.cassandra.db.rows.Cell; -import org.apache.cassandra.db.rows.ColumnData; -import org.apache.cassandra.db.rows.ComplexColumnData; -import org.apache.cassandra.db.rows.Row; -import org.apache.cassandra.db.rows.RowIterator; +import org.apache.cassandra.db.rows.*; import org.apache.cassandra.schema.KeyspaceMetadata; +import org.apache.cassandra.service.ClientState; import org.apache.cassandra.service.pager.QueryPager; -import org.apache.cassandra.utils.FBUtilities; import org.apache.cassandra.transport.Server; +import org.apache.cassandra.utils.FBUtilities; /** * A View copies data from a base table into a view table which can be queried independently from the @@ -111,6 +88,13 @@ public class View private final boolean includeAllColumns; private ViewBuilder builder; + // Only the raw statement can be final, because the statement cannot always be prepared when the MV is initialized. + // For example, during startup, this view will be initialized as part of the Keyspace.open() work; preparing a statement + // also requires the keyspace to be open, so this results in double-initialization problems. + private final SelectStatement.RawStatement rawSelect; + private SelectStatement select; + private ReadQuery query; + public View(ViewDefinition definition, ColumnFamilyStore baseCfs) { @@ -120,6 +104,7 @@ public class View includeAllColumns = definition.includeAllColumns; viewHasAllPrimaryKeys = updateDefinition(definition); + this.rawSelect = definition.select; } public ViewDefinition getDefinition() @@ -205,9 +190,9 @@ public class View */ public boolean updateAffectsView(AbstractBTreePartition partition) { - // If we are including all of the columns, then any update will be included - if (includeAllColumns) - return true; + ReadQuery selectQuery = getReadQuery(); + if (!selectQuery.selectsKey(partition.partitionKey())) + return false; // If there are range tombstones, tombstones will also need to be generated for the view // This requires a query of the base rows and generating tombstones for all of those values @@ -217,7 +202,10 @@ public class View // Check each row for deletion or update for (Row row : partition) { - if (!row.deletion().isLive()) + if (!selectQuery.selectsClustering(partition.partitionKey(), row.clustering())) + continue; + + if (includeAllColumns || viewHasAllPrimaryKeys || !row.deletion().isLive()) return true; if (row.primaryKeyLivenessInfo().isLive(FBUtilities.nowInSeconds())) @@ -440,7 +428,7 @@ public class View if (!deletionInfo.getPartitionDeletion().isLive()) { - command = SinglePartitionReadCommand.fullPartitionRead(baseCfs.metadata, rowSet.nowInSec, dk); + command = getSelectStatement().internalReadForView(dk, rowSet.nowInSec); } else { @@ -459,11 +447,15 @@ public class View if (command == null) { + ReadQuery selectQuery = getReadQuery(); SinglePartitionSliceBuilder builder = null; for (Row row : partition) { if (!row.deletion().isLive()) { + if (!selectQuery.selectsClustering(rowSet.dk, row.clustering())) + continue; + if (builder == null) builder = new SinglePartitionSliceBuilder(baseCfs, rowSet.dk); builder.addSlice(Slice.make(row.clustering())); @@ -476,10 +468,10 @@ public class View if (command != null) { + ReadQuery selectQuery = getReadQuery(); + assert selectQuery.selectsKey(rowSet.dk); - //We may have already done this work for - //another MV update so check - + // We may have already done this work for another MV update so check if (!rowSet.hasTombstonedExisting()) { QueryPager pager = command.getPager(null, Server.CURRENT_VERSION); @@ -498,7 +490,8 @@ public class View while (rowIterator.hasNext()) { Row row = rowIterator.next(); - rowSet.addRow(row, false); + if (selectQuery.selectsClustering(rowSet.dk, row.clustering())) + rowSet.addRow(row, false); } } } @@ -609,6 +602,34 @@ public class View return rowSet; } + /** + * Returns the SelectStatement used to populate and filter this view. Internal users should access the select + * statement this way to ensure it has been prepared. + */ + public SelectStatement getSelectStatement() + { + if (select == null) + { + ClientState state = ClientState.forInternalCalls(); + state.setKeyspace(baseCfs.keyspace.getName()); + rawSelect.prepareKeyspace(state); + ParsedStatement.Prepared prepared = rawSelect.prepare(true); + select = (SelectStatement) prepared.statement; + } + + return select; + } + + /** + * Returns the ReadQuery used to filter this view. Internal users should access the query this way to ensure it + * has been prepared. + */ + public ReadQuery getReadQuery() + { + if (query == null) + query = getSelectStatement().getQuery(QueryOptions.forInternalCalls(Collections.emptyList()), FBUtilities.nowInSeconds()); + return query; + } /** * @param isBuilding If the view is currently being built, we do not query the values which are already stored, @@ -683,4 +704,55 @@ public class View final UUID baseId = Schema.instance.getId(keyspace, baseTable); return Iterables.filter(ksm.views, view -> view.baseTableId.equals(baseId)); } + + /** + * Builds the string text for a materialized view's SELECT statement. + */ + public static String buildSelectStatement(String cfName, Collection includedColumns, String whereClause) + { + StringBuilder rawSelect = new StringBuilder("SELECT "); + if (includedColumns == null || includedColumns.isEmpty()) + rawSelect.append("*"); + else + rawSelect.append(includedColumns.stream().map(id -> id.name.toCQLString()).collect(Collectors.joining(", "))); + rawSelect.append(" FROM \"").append(cfName).append("\" WHERE ") .append(whereClause).append(" ALLOW FILTERING"); + return rawSelect.toString(); + } + + public static String relationsToWhereClause(List whereClause) + { + List expressions = new ArrayList<>(whereClause.size()); + for (Relation rel : whereClause) + { + StringBuilder sb = new StringBuilder(); + + if (rel.isMultiColumn()) + { + sb.append(((MultiColumnRelation) rel).getEntities().stream() + .map(ColumnIdentifier.Raw::toCQLString) + .collect(Collectors.joining(", ", "(", ")"))); + } + else + { + sb.append(((SingleColumnRelation) rel).getEntity().toCQLString()); + } + + sb.append(" ").append(rel.operator()).append(" "); + + if (rel.isIN()) + { + sb.append(rel.getInValues().stream() + .map(Term.Raw::getText) + .collect(Collectors.joining(", ", "(", ")"))); + } + else + { + sb.append(rel.getValue().getText()); + } + + expressions.add(sb.toString()); + } + + return expressions.stream().collect(Collectors.joining(" AND ")); + } } diff --git a/src/java/org/apache/cassandra/db/view/ViewBuilder.java b/src/java/org/apache/cassandra/db/view/ViewBuilder.java index f0b01c7478..0a0fe087a0 100644 --- a/src/java/org/apache/cassandra/db/view/ViewBuilder.java +++ b/src/java/org/apache/cassandra/db/view/ViewBuilder.java @@ -29,12 +29,7 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.apache.cassandra.concurrent.ScheduledExecutors; -import org.apache.cassandra.db.ColumnFamilyStore; -import org.apache.cassandra.db.DecoratedKey; -import org.apache.cassandra.db.Mutation; -import org.apache.cassandra.db.ReadOrderGroup; -import org.apache.cassandra.db.SinglePartitionReadCommand; -import org.apache.cassandra.db.SystemKeyspace; +import org.apache.cassandra.db.*; import org.apache.cassandra.db.compaction.CompactionInfo; import org.apache.cassandra.db.compaction.CompactionManager; import org.apache.cassandra.db.compaction.OperationType; @@ -44,7 +39,6 @@ import org.apache.cassandra.db.partitions.PartitionIterator; import org.apache.cassandra.db.rows.RowIterator; import org.apache.cassandra.dht.Range; import org.apache.cassandra.dht.Token; -import org.apache.cassandra.exceptions.WriteTimeoutException; import org.apache.cassandra.io.sstable.ReducingKeyIterator; import org.apache.cassandra.io.sstable.format.SSTableReader; import org.apache.cassandra.service.StorageProxy; @@ -52,7 +46,6 @@ import org.apache.cassandra.service.StorageService; import org.apache.cassandra.service.pager.QueryPager; import org.apache.cassandra.transport.Server; import org.apache.cassandra.utils.FBUtilities; -import org.apache.cassandra.utils.NoSpamLogger; import org.apache.cassandra.utils.Pair; import org.apache.cassandra.utils.UUIDGen; import org.apache.cassandra.utils.concurrent.Refs; @@ -77,7 +70,11 @@ public class ViewBuilder extends CompactionInfo.Holder private void buildKey(DecoratedKey key) { - QueryPager pager = SinglePartitionReadCommand.fullPartitionRead(baseCfs.metadata, FBUtilities.nowInSeconds(), key).getPager(null, Server.CURRENT_VERSION); + ReadQuery selectQuery = view.getReadQuery(); + if (!selectQuery.selectsKey(key)) + return; + + QueryPager pager = view.getSelectStatement().internalReadForView(key, FBUtilities.nowInSeconds()).getPager(null, Server.CURRENT_VERSION); while (!pager.isExhausted()) { diff --git a/src/java/org/apache/cassandra/index/internal/composites/CompositesSearcher.java b/src/java/org/apache/cassandra/index/internal/composites/CompositesSearcher.java index 77867fcbd8..f1751f5007 100644 --- a/src/java/org/apache/cassandra/index/internal/composites/CompositesSearcher.java +++ b/src/java/org/apache/cassandra/index/internal/composites/CompositesSearcher.java @@ -51,7 +51,7 @@ public class CompositesSearcher extends CassandraIndexSearcher private boolean isMatchingEntry(DecoratedKey partitionKey, IndexEntry entry, ReadCommand command) { - return command.selects(partitionKey, entry.indexedEntryClustering); + return command.selectsKey(partitionKey) && command.selectsClustering(partitionKey, entry.indexedEntryClustering); } protected UnfilteredPartitionIterator queryDataFromIndex(final DecoratedKey indexKey, diff --git a/src/java/org/apache/cassandra/schema/SchemaKeyspace.java b/src/java/org/apache/cassandra/schema/SchemaKeyspace.java index fb97ca532b..5f27d82153 100644 --- a/src/java/org/apache/cassandra/schema/SchemaKeyspace.java +++ b/src/java/org/apache/cassandra/schema/SchemaKeyspace.java @@ -35,14 +35,14 @@ import org.slf4j.LoggerFactory; import org.apache.cassandra.config.*; import org.apache.cassandra.config.ColumnDefinition.ClusteringOrder; -import org.apache.cassandra.cql3.ColumnIdentifier; -import org.apache.cassandra.cql3.QueryProcessor; -import org.apache.cassandra.cql3.UntypedResultSet; +import org.apache.cassandra.cql3.*; import org.apache.cassandra.cql3.functions.*; +import org.apache.cassandra.cql3.statements.SelectStatement; import org.apache.cassandra.db.*; import org.apache.cassandra.db.marshal.*; import org.apache.cassandra.db.partitions.*; import org.apache.cassandra.db.rows.*; +import org.apache.cassandra.db.view.View; import org.apache.cassandra.exceptions.ConfigurationException; import org.apache.cassandra.exceptions.InvalidRequestException; import org.apache.cassandra.utils.ByteBufferUtil; @@ -75,6 +75,7 @@ public final class SchemaKeyspace public static final String AGGREGATES = "aggregates"; public static final String INDEXES = "indexes"; + public static final List ALL = ImmutableList.of(KEYSPACES, TABLES, COLUMNS, TRIGGERS, VIEWS, TYPES, FUNCTIONS, AGGREGATES, INDEXES); @@ -155,6 +156,7 @@ public final class SchemaKeyspace + "view_name text," + "base_table_id uuid," + "base_table_name text," + + "where_clause text," + "bloom_filter_fp_chance double," + "caching frozen>," + "comment text," @@ -1311,6 +1313,7 @@ public final class SchemaKeyspace builder.add("include_all_columns", view.includeAllColumns) .add("base_table_id", view.baseTableId) .add("base_table_name", view.baseTableMetadata().cfName) + .add("where_clause", view.whereClause) .add("id", table.cfId); addTableParamsToSchemaMutation(table.params, builder); @@ -1426,7 +1429,9 @@ public final class SchemaKeyspace String view = row.getString("view_name"); UUID id = row.getUUID("id"); UUID baseTableId = row.getUUID("base_table_id"); + String baseTableName = row.getString("base_table_name"); boolean includeAll = row.getBoolean("include_all_columns"); + String whereClause = row.getString("where_clause"); List columns = readSchemaPartitionForTableAndApply(COLUMNS, keyspace, view, SchemaKeyspace::createColumnsFromColumnsPartition); @@ -1447,7 +1452,10 @@ public final class SchemaKeyspace .params(createTableParamsFromRow(row)) .droppedColumns(droppedColumns); - return new ViewDefinition(keyspace, view, baseTableId, includeAll, cfm); + String rawSelect = View.buildSelectStatement(baseTableName, columns, whereClause); + SelectStatement.RawStatement rawStatement = (SelectStatement.RawStatement) QueryProcessor.parseStatement(rawSelect); + + return new ViewDefinition(keyspace, view, baseTableId, baseTableName, includeAll, rawStatement, whereClause, cfm); } /* diff --git a/test/unit/org/apache/cassandra/cql3/CQLTester.java b/test/unit/org/apache/cassandra/cql3/CQLTester.java index 61e4fc2324..e92563b2f9 100644 --- a/test/unit/org/apache/cassandra/cql3/CQLTester.java +++ b/test/unit/org/apache/cassandra/cql3/CQLTester.java @@ -29,6 +29,7 @@ import java.util.concurrent.CountDownLatch; import java.util.concurrent.ExecutionException; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicInteger; +import java.util.stream.Collectors; import com.datastax.driver.core.*; import com.datastax.driver.core.ResultSet; @@ -790,6 +791,83 @@ public abstract class CQLTester Assert.assertTrue(String.format("Got %s rows than expected. Expected %d but got %d", rows.length>i ? "less" : "more", rows.length, i), i == rows.length); } + /** + * Like assertRows(), but ignores the ordering of rows. + */ + public static void assertRowsIgnoringOrder(UntypedResultSet result, Object[]... rows) + { + if (result == null) + { + if (rows.length > 0) + Assert.fail(String.format("No rows returned by query but %d expected", rows.length)); + return; + } + + List meta = result.metadata(); + + Set> expectedRows = new HashSet<>(rows.length); + for (Object[] expected : rows) + { + Assert.assertEquals("Invalid number of (expected) values provided for row", expected.length, meta.size()); + List expectedRow = new ArrayList<>(meta.size()); + for (int j = 0; j < meta.size(); j++) + expectedRow.add(makeByteBuffer(expected[j], meta.get(j).type)); + expectedRows.add(expectedRow); + } + + Set> actualRows = new HashSet<>(result.size()); + for (UntypedResultSet.Row actual : result) + { + List actualRow = new ArrayList<>(meta.size()); + for (int j = 0; j < meta.size(); j++) + actualRow.add(actual.getBytes(meta.get(j).name.toString())); + actualRows.add(actualRow); + } + + com.google.common.collect.Sets.SetView> extra = com.google.common.collect.Sets.difference(actualRows, expectedRows); + com.google.common.collect.Sets.SetView> missing = com.google.common.collect.Sets.difference(expectedRows, actualRows); + if (!extra.isEmpty() || !missing.isEmpty()) + { + List extraRows = makeRowStrings(extra, meta); + List missingRows = makeRowStrings(missing, meta); + StringBuilder sb = new StringBuilder(); + if (!extra.isEmpty()) + { + sb.append("Got ").append(extra.size()).append(" extra row(s) "); + if (!missing.isEmpty()) + sb.append("and ").append(missing.size()).append(" missing row(s) "); + sb.append("in result. Extra rows:\n "); + sb.append(extraRows.stream().collect(Collectors.joining("\n "))); + if (!missing.isEmpty()) + sb.append("\nMissing Rows:\n ").append(missingRows.stream().collect(Collectors.joining("\n "))); + Assert.fail(sb.toString()); + } + + if (!missing.isEmpty()) + Assert.fail("Missing " + missing.size() + " row(s) in result: \n " + missingRows.stream().collect(Collectors.joining("\n "))); + } + + assert expectedRows.size() == actualRows.size(); + } + + private static List makeRowStrings(Iterable> rows, List meta) + { + List strings = new ArrayList<>(); + for (List row : rows) + { + StringBuilder sb = new StringBuilder("row("); + for (int j = 0; j < row.size(); j++) + { + ColumnSpecification column = meta.get(j); + sb.append(column.name.toString()).append("=").append(formatValue(row.get(j), column.type)); + if (j < (row.size() - 1)) + sb.append(", "); + } + strings.add(sb.append(")").toString()); + } + return strings; + } + protected void assertRowCount(UntypedResultSet result, int numExpectedRows) { if (result == null) diff --git a/test/unit/org/apache/cassandra/cql3/ViewFilteringTest.java b/test/unit/org/apache/cassandra/cql3/ViewFilteringTest.java new file mode 100644 index 0000000000..2d789c3004 --- /dev/null +++ b/test/unit/org/apache/cassandra/cql3/ViewFilteringTest.java @@ -0,0 +1,1292 @@ +/* + * 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.cql3; + +import java.util.*; + +import org.junit.After; +import org.junit.Before; +import org.junit.BeforeClass; +import org.junit.Test; + +import com.datastax.driver.core.exceptions.InvalidQueryException; +import junit.framework.Assert; + +import org.apache.cassandra.db.SystemKeyspace; + +public class ViewFilteringTest extends CQLTester +{ + int protocolVersion = 4; + private final List views = new ArrayList<>(); + + @BeforeClass + public static void startup() + { + requireNetwork(); + } + @Before + public void begin() + { + views.clear(); + } + + @After + public void end() throws Throwable + { + for (String viewName : views) + executeNet(protocolVersion, "DROP MATERIALIZED VIEW " + viewName); + } + + private void createView(String name, String query) throws Throwable + { + executeNet(protocolVersion, String.format(query, name)); + // If exception is thrown, the view will not be added to the list; since it shouldn't have been created, this is + // the desired behavior + views.add(name); + } + + private void dropView(String name) throws Throwable + { + executeNet(protocolVersion, "DROP MATERIALIZED VIEW " + name); + views.remove(name); + } + + @Test + public void testMVCreationSelectRestrictions() throws Throwable + { + createTable("CREATE TABLE %s (a int, b int, c int, d int, e int, PRIMARY KEY((a, b), c, d))"); + + execute("USE " + keyspace()); + executeNet(protocolVersion, "USE " + keyspace()); + + // IS NOT NULL is required on all PK statements that are not otherwise restricted + List badStatements = Arrays.asList( + "CREATE MATERIALIZED VIEW %s AS SELECT * FROM %%s WHERE b IS NOT NULL AND c IS NOT NULL AND d is NOT NULL PRIMARY KEY ((a, b), c, d)", + "CREATE MATERIALIZED VIEW %s AS SELECT * FROM %%s WHERE a IS NOT NULL AND c IS NOT NULL AND d is NOT NULL PRIMARY KEY ((a, b), c, d)", + "CREATE MATERIALIZED VIEW %s AS SELECT * FROM %%s WHERE a IS NOT NULL AND b IS NOT NULL AND d is NOT NULL PRIMARY KEY ((a, b), c, d)", + "CREATE MATERIALIZED VIEW %s AS SELECT * FROM %%s WHERE a IS NOT NULL AND b IS NOT NULL AND c is NOT NULL PRIMARY KEY ((a, b), c, d)", + "CREATE MATERIALIZED VIEW %s AS SELECT * FROM %%s WHERE a = ? AND b IS NOT NULL AND c is NOT NULL PRIMARY KEY ((a, b), c, d)", + "CREATE MATERIALIZED VIEW %s AS SELECT * FROM %%s WHERE a = blobAsInt(?) AND b IS NOT NULL AND c is NOT NULL PRIMARY KEY ((a, b), c, d)", + "CREATE MATERIALIZED VIEW %s AS SELECT * FROM %%s PRIMARY KEY (a, b, c, d)" + ); + + for (String badStatement : badStatements) + { + try + { + createView("mv1_test", badStatement); + Assert.fail("Create MV statement should have failed due to missing IS NOT NULL restriction: " + badStatement); + } + catch (InvalidQueryException exc) {} + } + + List goodStatements = Arrays.asList( + "CREATE MATERIALIZED VIEW %s AS SELECT * FROM %%s WHERE a = 1 AND b = 1 AND c IS NOT NULL AND d is NOT NULL PRIMARY KEY ((a, b), c, d)", + "CREATE MATERIALIZED VIEW %s AS SELECT * FROM %%s WHERE a IS NOT NULL AND b IS NOT NULL AND c = 1 AND d IS NOT NULL PRIMARY KEY ((a, b), c, d)", + "CREATE MATERIALIZED VIEW %s AS SELECT * FROM %%s WHERE a IS NOT NULL AND b IS NOT NULL AND c = 1 AND d = 1 PRIMARY KEY ((a, b), c, d)", + "CREATE MATERIALIZED VIEW %s AS SELECT * FROM %%s WHERE a = 1 AND b = 1 AND c = 1 AND d = 1 PRIMARY KEY ((a, b), c, d)", + "CREATE MATERIALIZED VIEW %s AS SELECT * FROM %%s WHERE a = 1 AND b = 1 AND c > 1 AND d IS NOT NULL PRIMARY KEY ((a, b), c, d)", + "CREATE MATERIALIZED VIEW %s AS SELECT * FROM %%s WHERE a = 1 AND b = 1 AND c = 1 AND d IN (1, 2, 3) PRIMARY KEY ((a, b), c, d)", + "CREATE MATERIALIZED VIEW %s AS SELECT * FROM %%s WHERE a = 1 AND b = 1 AND (c, d) = (1, 1) PRIMARY KEY ((a, b), c, d)", + "CREATE MATERIALIZED VIEW %s AS SELECT * FROM %%s WHERE a = 1 AND b = 1 AND (c, d) > (1, 1) PRIMARY KEY ((a, b), c, d)", + "CREATE MATERIALIZED VIEW %s AS SELECT * FROM %%s WHERE a = 1 AND b = 1 AND (c, d) IN ((1, 1), (2, 2)) PRIMARY KEY ((a, b), c, d)", + "CREATE MATERIALIZED VIEW %s AS SELECT * FROM %%s WHERE a = (int) 1 AND b = 1 AND c = 1 AND d = 1 PRIMARY KEY ((a, b), c, d)", + "CREATE MATERIALIZED VIEW %s AS SELECT * FROM %%s WHERE a = blobAsInt(intAsBlob(1)) AND b = 1 AND c = 1 AND d = 1 PRIMARY KEY ((a, b), c, d)" + ); + + for (int i = 0; i < goodStatements.size(); i++) + { + try + { + createView("mv" + i + "_test", goodStatements.get(i)); + } + catch (Exception e) + { + throw new RuntimeException("MV creation failed: " + goodStatements.get(i), e); + } + + try + { + executeNet(protocolVersion, "ALTER MATERIALIZED VIEW mv" + i + "_test WITH compaction = { 'class' : 'LeveledCompactionStrategy' }"); + } + catch (Exception e) + { + throw new RuntimeException("MV alter failed: " + goodStatements.get(i), e); + } + } + + try + { + createView("mv_foo", "CREATE MATERIALIZED VIEW %s AS SELECT * FROM %%s WHERE a = 1 AND b IS NOT NULL AND c IS NOT NULL AND d is NOT NULL PRIMARY KEY ((a, b), c, d)"); + Assert.fail("Partial partition key restriction should not be allowed"); + } + catch (InvalidQueryException exc) {} + } + + @Test + public void testCaseSensitivity() throws Throwable + { + createTable("CREATE TABLE %s (\"theKey\" int, \"theClustering\" int, \"the\"\"Value\" int, PRIMARY KEY (\"theKey\", \"theClustering\"))"); + + execute("USE " + keyspace()); + executeNet(protocolVersion, "USE " + keyspace()); + + execute("INSERT INTO %s (\"theKey\", \"theClustering\", \"the\"\"Value\") VALUES (?, ?, ?)", 0, 0, 0); + execute("INSERT INTO %s (\"theKey\", \"theClustering\", \"the\"\"Value\") VALUES (?, ?, ?)", 0, 1, 0); + execute("INSERT INTO %s (\"theKey\", \"theClustering\", \"the\"\"Value\") VALUES (?, ?, ?)", 1, 0, 0); + execute("INSERT INTO %s (\"theKey\", \"theClustering\", \"the\"\"Value\") VALUES (?, ?, ?)", 1, 1, 0); + + createView("mv_test", "CREATE MATERIALIZED VIEW %s AS SELECT * FROM %%s " + + "WHERE \"theKey\" = 1 AND \"theClustering\" = 1 AND \"the\"\"Value\" IS NOT NULL " + + "PRIMARY KEY (\"theKey\", \"theClustering\")"); + + while (!SystemKeyspace.isViewBuilt(keyspace(), "mv_test")) + Thread.sleep(10); + createView("mv_test2", "CREATE MATERIALIZED VIEW %s AS SELECT \"theKey\", \"theClustering\", \"the\"\"Value\" FROM %%s " + + "WHERE \"theKey\" = 1 AND \"theClustering\" = 1 AND \"the\"\"Value\" IS NOT NULL " + + "PRIMARY KEY (\"theKey\", \"theClustering\")"); + while (!SystemKeyspace.isViewBuilt(keyspace(), "mv_test2")) + Thread.sleep(10); + + for (String mvname : Arrays.asList("mv_test", "mv_test2")) + { + assertRowsIgnoringOrder(execute("SELECT \"theKey\", \"theClustering\", \"the\"\"Value\" FROM " + mvname), + row(1, 1, 0) + ); + } + + executeNet(protocolVersion, "ALTER TABLE %s RENAME \"theClustering\" TO \"Col\""); + + for (String mvname : Arrays.asList("mv_test", "mv_test2")) + { + assertRowsIgnoringOrder(execute("SELECT \"theKey\", \"Col\", \"the\"\"Value\" FROM " + mvname), + row(1, 1, 0) + ); + } + } + + @Test + public void testFilterWithFunction() throws Throwable + { + createTable("CREATE TABLE %s (a int, b int, c int, PRIMARY KEY (a, b))"); + + execute("USE " + keyspace()); + executeNet(protocolVersion, "USE " + keyspace()); + + execute("INSERT INTO %s (a, b, c) VALUES (?, ?, ?)", 0, 0, 0); + execute("INSERT INTO %s (a, b, c) VALUES (?, ?, ?)", 0, 1, 1); + execute("INSERT INTO %s (a, b, c) VALUES (?, ?, ?)", 1, 0, 2); + execute("INSERT INTO %s (a, b, c) VALUES (?, ?, ?)", 1, 1, 3); + + createView("mv_test", "CREATE MATERIALIZED VIEW %s AS SELECT * FROM %%s " + + "WHERE a = blobAsInt(intAsBlob(1)) AND b IS NOT NULL " + + "PRIMARY KEY (a, b)"); + + while (!SystemKeyspace.isViewBuilt(keyspace(), "mv_test")) + Thread.sleep(10); + + assertRows(execute("SELECT a, b, c FROM mv_test"), + row(1, 0, 2), + row(1, 1, 3) + ); + + executeNet(protocolVersion, "ALTER TABLE %s RENAME a TO foo"); + + assertRows(execute("SELECT foo, b, c FROM mv_test"), + row(1, 0, 2), + row(1, 1, 3) + ); + } + + @Test + public void testFilterWithTypecast() throws Throwable + { + createTable("CREATE TABLE %s (a int, b int, c int, PRIMARY KEY (a, b))"); + + execute("USE " + keyspace()); + executeNet(protocolVersion, "USE " + keyspace()); + + execute("INSERT INTO %s (a, b, c) VALUES (?, ?, ?)", 0, 0, 0); + execute("INSERT INTO %s (a, b, c) VALUES (?, ?, ?)", 0, 1, 1); + execute("INSERT INTO %s (a, b, c) VALUES (?, ?, ?)", 1, 0, 2); + execute("INSERT INTO %s (a, b, c) VALUES (?, ?, ?)", 1, 1, 3); + + createView("mv_test", "CREATE MATERIALIZED VIEW %s AS SELECT * FROM %%s " + + "WHERE a = (int) 1 AND b IS NOT NULL " + + "PRIMARY KEY (a, b)"); + + while (!SystemKeyspace.isViewBuilt(keyspace(), "mv_test")) + Thread.sleep(10); + + assertRows(execute("SELECT a, b, c FROM mv_test"), + row(1, 0, 2), + row(1, 1, 3) + ); + + executeNet(protocolVersion, "ALTER TABLE %s RENAME a TO foo"); + + assertRows(execute("SELECT foo, b, c FROM mv_test"), + row(1, 0, 2), + row(1, 1, 3) + ); + } + + @Test + public void testPartitionKeyRestrictions() throws Throwable + { + List mvPrimaryKeys = Arrays.asList("((a, b), c)", "((b, a), c)", "(a, b, c)", "(c, b, a)", "((c, a), b)"); + for (int i = 0; i < mvPrimaryKeys.size(); i++) + { + createTable("CREATE TABLE %s (a int, b int, c int, d int, PRIMARY KEY (a, b, c))"); + + execute("USE " + keyspace()); + executeNet(protocolVersion, "USE " + keyspace()); + + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 0, 0, 0, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 0, 1, 0, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 1, 0, 0, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 1, 0, 1, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 1, 1, 0, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 1, 1, 1, 0); + + logger.info("Testing MV primary key: {}", mvPrimaryKeys.get(i)); + + // only accept rows where a = 1 + createView("mv_test" + i, "CREATE MATERIALIZED VIEW %s AS SELECT * FROM %%s WHERE a = 1 AND b IS NOT NULL AND c IS NOT NULL PRIMARY KEY " + mvPrimaryKeys.get(i)); + + while (!SystemKeyspace.isViewBuilt(keyspace(), "mv_test" + i)) + Thread.sleep(10); + + assertRowsIgnoringOrder(execute("SELECT a, b, c, d FROM mv_test" + i), + row(1, 0, 0, 0), + row(1, 0, 1, 0), + row(1, 1, 0, 0), + row(1, 1, 1, 0) + ); + + // insert new rows that do not match the filter + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 2, 0, 0, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 2, 1, 0, 0); + assertRowsIgnoringOrder(execute("SELECT a, b, c, d FROM mv_test" + i), + row(1, 0, 0, 0), + row(1, 0, 1, 0), + row(1, 1, 0, 0), + row(1, 1, 1, 0) + ); + + // insert new row that does match the filter + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 1, 1, 2, 0); + assertRowsIgnoringOrder(execute("SELECT a, b, c, d FROM mv_test" + i), + row(1, 0, 0, 0), + row(1, 0, 1, 0), + row(1, 1, 0, 0), + row(1, 1, 1, 0), + row(1, 1, 2, 0) + ); + + // update rows that don't match the filter + execute("UPDATE %s SET d = ? WHERE a = ? AND b = ? AND c = ?", 1, 0, 0, 0); + execute("UPDATE %s SET d = ? WHERE a = ? AND b = ? AND c = ?", 1, 0, 1, 0); + assertRowsIgnoringOrder(execute("SELECT a, b, c, d FROM mv_test" + i), + row(1, 0, 0, 0), + row(1, 0, 1, 0), + row(1, 1, 0, 0), + row(1, 1, 1, 0), + row(1, 1, 2, 0) + ); + + // update a row that does match the filter + execute("UPDATE %s SET d = ? WHERE a = ? AND b = ? AND c = ?", 1, 1, 1, 0); + assertRowsIgnoringOrder(execute("SELECT a, b, c, d FROM mv_test" + i), + row(1, 0, 0, 0), + row(1, 0, 1, 0), + row(1, 1, 0, 1), + row(1, 1, 1, 0), + row(1, 1, 2, 0) + ); + + // delete rows that don't match the filter + execute("DELETE FROM %s WHERE a = ? AND b = ? AND c = ?", 0, 0, 0); + execute("DELETE FROM %s WHERE a = ? AND b = ? AND c = ?", 0, 1, 0); + execute("DELETE FROM %s WHERE a = ? AND b = ?", 0, 0); + assertRowsIgnoringOrder(execute("SELECT a, b, c, d FROM mv_test" + i), + row(1, 0, 0, 0), + row(1, 0, 1, 0), + row(1, 1, 0, 1), + row(1, 1, 1, 0), + row(1, 1, 2, 0) + ); + + // delete a row that does match the filter + execute("DELETE FROM %s WHERE a = ? AND b = ? AND c = ?", 1, 1, 0); + assertRowsIgnoringOrder(execute("SELECT a, b, c, d FROM mv_test" + i), + row(1, 0, 0, 0), + row(1, 0, 1, 0), + row(1, 1, 1, 0), + row(1, 1, 2, 0) + ); + + // delete a partition that matches the filter + execute("DELETE FROM %s WHERE a = ?", 1); + assertEmpty(execute("SELECT * FROM mv_test" + i)); + } + } + + @Test + public void testCompoundPartitionKeyRestrictions() throws Throwable + { + List mvPrimaryKeys = Arrays.asList("((a, b), c)", "((b, a), c)", "(a, b, c)", "(c, b, a)", "((c, a), b)"); + for (int i = 0; i < mvPrimaryKeys.size(); i++) + { + createTable("CREATE TABLE %s (a int, b int, c int, d int, PRIMARY KEY ((a, b), c))"); + + execute("USE " + keyspace()); + executeNet(protocolVersion, "USE " + keyspace()); + + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 0, 0, 0, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 0, 0, 1, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 0, 1, 0, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 0, 1, 1, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 1, 0, 0, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 1, 0, 1, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 1, 1, 0, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 1, 1, 1, 0); + + logger.info("Testing MV primary key: {}", mvPrimaryKeys.get(i)); + + // only accept rows where a = 1 and b = 1 + createView("mv_test" + i, "CREATE MATERIALIZED VIEW %s AS SELECT * FROM %%s WHERE a = 1 AND b = 1 AND c IS NOT NULL PRIMARY KEY " + mvPrimaryKeys.get(i)); + + while (!SystemKeyspace.isViewBuilt(keyspace(), "mv_test" + i)) + Thread.sleep(10); + + assertRowsIgnoringOrder(execute("SELECT a, b, c, d FROM mv_test" + i), + row(1, 1, 0, 0), + row(1, 1, 1, 0) + ); + + // insert new rows that do not match the filter + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 2, 0, 0, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 2, 1, 0, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 1, 2, 0, 0); + assertRowsIgnoringOrder(execute("SELECT a, b, c, d FROM mv_test" + i), + row(1, 1, 0, 0), + row(1, 1, 1, 0) + ); + + // insert new row that does match the filter + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 1, 1, 2, 0); + assertRowsIgnoringOrder(execute("SELECT a, b, c, d FROM mv_test" + i), + row(1, 1, 0, 0), + row(1, 1, 1, 0), + row(1, 1, 2, 0) + ); + + // update rows that don't match the filter + execute("UPDATE %s SET d = ? WHERE a = ? AND b = ? AND c = ?", 1, 0, 0, 0); + execute("UPDATE %s SET d = ? WHERE a = ? AND b = ? AND c = ?", 1, 1, 0, 0); + execute("UPDATE %s SET d = ? WHERE a = ? AND b = ? AND c = ?", 1, 0, 1, 0); + assertRowsIgnoringOrder(execute("SELECT a, b, c, d FROM mv_test" + i), + row(1, 1, 0, 0), + row(1, 1, 1, 0), + row(1, 1, 2, 0) + ); + + // update a row that does match the filter + execute("UPDATE %s SET d = ? WHERE a = ? AND b = ? AND c = ?", 1, 1, 1, 0); + assertRowsIgnoringOrder(execute("SELECT a, b, c, d FROM mv_test" + i), + row(1, 1, 0, 1), + row(1, 1, 1, 0), + row(1, 1, 2, 0) + ); + + // delete rows that don't match the filter + execute("DELETE FROM %s WHERE a = ? AND b = ? AND c = ?", 0, 0, 0); + execute("DELETE FROM %s WHERE a = ? AND b = ? AND c = ?", 1, 0, 0); + execute("DELETE FROM %s WHERE a = ? AND b = ? AND c = ?", 0, 1, 0); + execute("DELETE FROM %s WHERE a = ? AND b = ?", 0, 0); + assertRowsIgnoringOrder(execute("SELECT a, b, c, d FROM mv_test" + i), + row(1, 1, 0, 1), + row(1, 1, 1, 0), + row(1, 1, 2, 0) + ); + + // delete a row that does match the filter + execute("DELETE FROM %s WHERE a = ? AND b = ? AND c = ?", 1, 1, 0); + assertRowsIgnoringOrder(execute("SELECT a, b, c, d FROM mv_test" + i), + row(1, 1, 1, 0), + row(1, 1, 2, 0) + ); + + // delete a partition that matches the filter + execute("DELETE FROM %s WHERE a = ? AND b = ?", 1, 1); + assertEmpty(execute("SELECT * FROM mv_test" + i)); + } + } + + @Test + public void testCompoundPartitionKeyRestrictionsNotIncludeAll() throws Throwable + { + createTable("CREATE TABLE %s (a int, b int, c int, d int, PRIMARY KEY ((a, b), c))"); + execute("USE " + keyspace()); + executeNet(protocolVersion, "USE " + keyspace()); + + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 0, 0, 0, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 0, 0, 1, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 0, 1, 0, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 0, 1, 1, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 1, 0, 0, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 1, 0, 1, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 1, 1, 0, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 1, 1, 1, 0); + + // only accept rows where a = 1 and b = 1, don't include column d in the selection + createView("mv_test", "CREATE MATERIALIZED VIEW %s AS SELECT a, b, c FROM %%s WHERE a = 1 AND b = 1 AND c IS NOT NULL PRIMARY KEY ((a, b), c)"); + + while (!SystemKeyspace.isViewBuilt(keyspace(), "mv_test")) + Thread.sleep(10); + + assertRows(execute("SELECT * FROM mv_test"), + row(1, 1, 0), + row(1, 1, 1) + ); + + // insert new rows that do not match the filter + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 2, 0, 0, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 2, 1, 0, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 1, 2, 0, 0); + assertRows(execute("SELECT * FROM mv_test"), + row(1, 1, 0), + row(1, 1, 1) + ); + + // insert new row that does match the filter + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 1, 1, 2, 0); + assertRows(execute("SELECT * FROM mv_test"), + row(1, 1, 0), + row(1, 1, 1), + row(1, 1, 2) + ); + + // update rows that don't match the filter + execute("UPDATE %s SET d = ? WHERE a = ? AND b = ? AND c = ?", 1, 0, 0, 0); + execute("UPDATE %s SET d = ? WHERE a = ? AND b = ? AND c = ?", 1, 1, 0, 0); + execute("UPDATE %s SET d = ? WHERE a = ? AND b = ? AND c = ?", 1, 0, 1, 0); + assertRows(execute("SELECT * FROM mv_test"), + row(1, 1, 0), + row(1, 1, 1), + row(1, 1, 2) + ); + + // update a row that does match the filter + execute("UPDATE %s SET d = ? WHERE a = ? AND b = ? AND c = ?", 1, 1, 1, 0); + assertRows(execute("SELECT * FROM mv_test"), + row(1, 1, 0), + row(1, 1, 1), + row(1, 1, 2) + ); + + // delete rows that don't match the filter + execute("DELETE FROM %s WHERE a = ? AND b = ? AND c = ?", 0, 0, 0); + execute("DELETE FROM %s WHERE a = ? AND b = ? AND c = ?", 1, 0, 0); + execute("DELETE FROM %s WHERE a = ? AND b = ? AND c = ?", 0, 1, 0); + execute("DELETE FROM %s WHERE a = ? AND b = ?", 0, 0); + assertRows(execute("SELECT * FROM mv_test"), + row(1, 1, 0), + row(1, 1, 1), + row(1, 1, 2) + ); + + // delete a row that does match the filter + execute("DELETE FROM %s WHERE a = ? AND b = ? AND c = ?", 1, 1, 0); + assertRows(execute("SELECT * FROM mv_test"), + row(1, 1, 1), + row(1, 1, 2) + ); + + // delete a partition that matches the filter + execute("DELETE FROM %s WHERE a = ? AND b = ?", 1, 1); + assertEmpty(execute("SELECT * FROM mv_test")); + } + + @Test + public void testClusteringKeyEQRestrictions() throws Throwable + { + List mvPrimaryKeys = Arrays.asList("((a, b), c)", "((b, a), c)", "(a, b, c)", "(c, b, a)", "((c, a), b)"); + for (int i = 0; i < mvPrimaryKeys.size(); i++) + { + createTable("CREATE TABLE %s (a int, b int, c int, d int, PRIMARY KEY (a, b, c))"); + + execute("USE " + keyspace()); + executeNet(protocolVersion, "USE " + keyspace()); + + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 0, 0, 0, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 0, 0, 1, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 0, 1, 0, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 0, 1, 1, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 1, 0, 0, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 1, 0, 1, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 1, 1, 0, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 1, 1, 1, 0); + + logger.info("Testing MV primary key: {}", mvPrimaryKeys.get(i)); + + // only accept rows where b = 1 + createView("mv_test" + i, "CREATE MATERIALIZED VIEW %s AS SELECT * FROM %%s WHERE a IS NOT NULL AND b = 1 AND c IS NOT NULL PRIMARY KEY " + mvPrimaryKeys.get(i)); + + while (!SystemKeyspace.isViewBuilt(keyspace(), "mv_test" + i)) + Thread.sleep(10); + + assertRowsIgnoringOrder(execute("SELECT a, b, c, d FROM mv_test" + i), + row(0, 1, 0, 0), + row(0, 1, 1, 0), + row(1, 1, 0, 0), + row(1, 1, 1, 0) + ); + + // insert new rows that do not match the filter + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 2, 0, 0, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 2, 2, 0, 0); + assertRowsIgnoringOrder(execute("SELECT a, b, c, d FROM mv_test" + i), + row(0, 1, 0, 0), + row(0, 1, 1, 0), + row(1, 1, 0, 0), + row(1, 1, 1, 0) + ); + + // insert new row that does match the filter + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 1, 1, 2, 0); + assertRowsIgnoringOrder(execute("SELECT a, b, c, d FROM mv_test" + i), + row(0, 1, 0, 0), + row(0, 1, 1, 0), + row(1, 1, 0, 0), + row(1, 1, 1, 0), + row(1, 1, 2, 0) + ); + + // update rows that don't match the filter + execute("UPDATE %s SET d = ? WHERE a = ? AND b = ? AND c = ?", 1, 2, 0, 0); + execute("UPDATE %s SET d = ? WHERE a = ? AND b = ? AND c = ?", 1, 2, 2, 0); + assertRowsIgnoringOrder(execute("SELECT a, b, c, d FROM mv_test" + i), + row(0, 1, 0, 0), + row(0, 1, 1, 0), + row(1, 1, 0, 0), + row(1, 1, 1, 0), + row(1, 1, 2, 0) + ); + + // update a row that does match the filter + execute("UPDATE %s SET d = ? WHERE a = ? AND b = ? AND c = ?", 1, 1, 1, 0); + assertRowsIgnoringOrder(execute("SELECT a, b, c, d FROM mv_test" + i), + row(0, 1, 0, 0), + row(0, 1, 1, 0), + row(1, 1, 0, 1), + row(1, 1, 1, 0), + row(1, 1, 2, 0) + ); + + // delete rows that don't match the filter + execute("DELETE FROM %s WHERE a = ? AND b = ? AND c = ?", 2, 0, 0); + execute("DELETE FROM %s WHERE a = ? AND b = ? AND c = ?", 2, 2, 0); + execute("DELETE FROM %s WHERE a = ? AND b = ?", 0, 0); + assertRowsIgnoringOrder(execute("SELECT a, b, c, d FROM mv_test" + i), + row(0, 1, 0, 0), + row(0, 1, 1, 0), + row(1, 1, 0, 1), + row(1, 1, 1, 0), + row(1, 1, 2, 0) + ); + + // delete a row that does match the filter + execute("DELETE FROM %s WHERE a = ? AND b = ? AND c = ?", 1, 1, 0); + assertRowsIgnoringOrder(execute("SELECT a, b, c, d FROM mv_test" + i), + row(0, 1, 0, 0), + row(0, 1, 1, 0), + row(1, 1, 1, 0), + row(1, 1, 2, 0) + ); + + // delete a partition that matches the filter + execute("DELETE FROM %s WHERE a = ?", 1); + assertRowsIgnoringOrder(execute("SELECT a, b, c, d FROM mv_test" + i), + row(0, 1, 0, 0), + row(0, 1, 1, 0) + ); + + dropView("mv_test" + i); + dropTable("DROP TABLE %s"); + } + } + + @Test + public void testClusteringKeySliceRestrictions() throws Throwable + { + List mvPrimaryKeys = Arrays.asList("((a, b), c)", "((b, a), c)", "(a, b, c)", "(c, b, a)", "((c, a), b)"); + for (int i = 0; i < mvPrimaryKeys.size(); i++) + { + createTable("CREATE TABLE %s (a int, b int, c int, d int, PRIMARY KEY (a, b, c))"); + + execute("USE " + keyspace()); + executeNet(protocolVersion, "USE " + keyspace()); + + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 0, 0, 0, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 0, 0, 1, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 0, 1, 0, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 0, 1, 1, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 1, 0, 0, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 1, 0, 1, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 1, 1, 0, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 1, 1, 1, 0); + + logger.info("Testing MV primary key: {}", mvPrimaryKeys.get(i)); + + createView("mv_test" + i, "CREATE MATERIALIZED VIEW %s AS SELECT * FROM %%s WHERE a IS NOT NULL AND b >= 1 AND c IS NOT NULL PRIMARY KEY " + mvPrimaryKeys.get(i)); + + while (!SystemKeyspace.isViewBuilt(keyspace(), "mv_test" + i)) + Thread.sleep(10); + + assertRowsIgnoringOrder(execute("SELECT a, b, c, d FROM mv_test" + i), + row(0, 1, 0, 0), + row(0, 1, 1, 0), + row(1, 1, 0, 0), + row(1, 1, 1, 0) + ); + + // insert new rows that do not match the filter + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 2, -1, 0, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 2, 0, 0, 0); + assertRowsIgnoringOrder(execute("SELECT a, b, c, d FROM mv_test" + i), + row(0, 1, 0, 0), + row(0, 1, 1, 0), + row(1, 1, 0, 0), + row(1, 1, 1, 0) + ); + + // insert new row that does match the filter + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 1, 1, 2, 0); + assertRowsIgnoringOrder(execute("SELECT a, b, c, d FROM mv_test" + i), + row(0, 1, 0, 0), + row(0, 1, 1, 0), + row(1, 1, 0, 0), + row(1, 1, 1, 0), + row(1, 1, 2, 0) + ); + + // update rows that don't match the filter + execute("UPDATE %s SET d = ? WHERE a = ? AND b = ? AND c = ?", 1, 2, -1, 0); + execute("UPDATE %s SET d = ? WHERE a = ? AND b = ? AND c = ?", 1, 2, 0, 0); + assertRowsIgnoringOrder(execute("SELECT a, b, c, d FROM mv_test" + i), + row(0, 1, 0, 0), + row(0, 1, 1, 0), + row(1, 1, 0, 0), + row(1, 1, 1, 0), + row(1, 1, 2, 0) + ); + + // update a row that does match the filter + execute("UPDATE %s SET d = ? WHERE a = ? AND b = ? AND c = ?", 1, 1, 1, 0); + assertRowsIgnoringOrder(execute("SELECT a, b, c, d FROM mv_test" + i), + row(0, 1, 0, 0), + row(0, 1, 1, 0), + row(1, 1, 0, 1), + row(1, 1, 1, 0), + row(1, 1, 2, 0) + ); + + // delete rows that don't match the filter + execute("DELETE FROM %s WHERE a = ? AND b = ? AND c = ?", 2, -1, 0); + execute("DELETE FROM %s WHERE a = ? AND b = ? AND c = ?", 2, 0, 0); + execute("DELETE FROM %s WHERE a = ? AND b = ?", 0, 0); + assertRowsIgnoringOrder(execute("SELECT a, b, c, d FROM mv_test" + i), + row(0, 1, 0, 0), + row(0, 1, 1, 0), + row(1, 1, 0, 1), + row(1, 1, 1, 0), + row(1, 1, 2, 0) + ); + + // delete a row that does match the filter + execute("DELETE FROM %s WHERE a = ? AND b = ? AND c = ?", 1, 1, 0); + assertRowsIgnoringOrder(execute("SELECT a, b, c, d FROM mv_test" + i), + row(0, 1, 0, 0), + row(0, 1, 1, 0), + row(1, 1, 1, 0), + row(1, 1, 2, 0) + ); + + // delete a partition that matches the filter + execute("DELETE FROM %s WHERE a = ?", 1); + assertRowsIgnoringOrder(execute("SELECT a, b, c, d FROM mv_test" + i), + row(0, 1, 0, 0), + row(0, 1, 1, 0) + ); + + dropView("mv_test" + i); + dropTable("DROP TABLE %s"); + } + } + + @Test + public void testClusteringKeyINRestrictions() throws Throwable + { + List mvPrimaryKeys = Arrays.asList("((a, b), c)", "((b, a), c)", "(a, b, c)", "(c, b, a)", "((c, a), b)"); + for (int i = 0; i < mvPrimaryKeys.size(); i++) + { + createTable("CREATE TABLE %s (a int, b int, c int, d int, PRIMARY KEY (a, b, c))"); + + execute("USE " + keyspace()); + executeNet(protocolVersion, "USE " + keyspace()); + + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 0, 0, 0, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 0, 0, 1, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 0, 1, 0, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 0, 1, 1, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 1, 0, 0, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 1, 0, 1, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 1, 1, 0, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 1, 1, 1, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 1, 2, 1, 0); + + logger.info("Testing MV primary key: {}", mvPrimaryKeys.get(i)); + + // only accept rows where b = 1 + createView("mv_test" + i, "CREATE MATERIALIZED VIEW %s AS SELECT * FROM %%s WHERE a IS NOT NULL AND b IN (1, 2) AND c IS NOT NULL PRIMARY KEY " + mvPrimaryKeys.get(i)); + + while (!SystemKeyspace.isViewBuilt(keyspace(), "mv_test" + i)) + Thread.sleep(10); + + assertRowsIgnoringOrder(execute("SELECT a, b, c, d FROM mv_test" + i), + row(0, 1, 0, 0), + row(0, 1, 1, 0), + row(1, 1, 0, 0), + row(1, 1, 1, 0), + row(1, 2, 1, 0) + ); + + // insert new rows that do not match the filter + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 2, -1, 0, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 2, 0, 0, 0); + assertRowsIgnoringOrder(execute("SELECT a, b, c, d FROM mv_test" + i), + row(0, 1, 0, 0), + row(0, 1, 1, 0), + row(1, 1, 0, 0), + row(1, 1, 1, 0), + row(1, 2, 1, 0) + ); + + // insert new row that does match the filter + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 1, 1, 2, 0); + assertRowsIgnoringOrder(execute("SELECT a, b, c, d FROM mv_test" + i), + row(0, 1, 0, 0), + row(0, 1, 1, 0), + row(1, 1, 0, 0), + row(1, 1, 1, 0), + row(1, 1, 2, 0), + row(1, 2, 1, 0) + ); + + // update rows that don't match the filter + execute("UPDATE %s SET d = ? WHERE a = ? AND b = ? AND c = ?", 1, 2, -1, 0); + execute("UPDATE %s SET d = ? WHERE a = ? AND b = ? AND c = ?", 1, 2, 0, 0); + assertRowsIgnoringOrder(execute("SELECT a, b, c, d FROM mv_test" + i), + row(0, 1, 0, 0), + row(0, 1, 1, 0), + row(1, 1, 0, 0), + row(1, 1, 1, 0), + row(1, 1, 2, 0), + row(1, 2, 1, 0) + ); + + // update a row that does match the filter + execute("UPDATE %s SET d = ? WHERE a = ? AND b = ? AND c = ?", 1, 1, 1, 0); + assertRowsIgnoringOrder(execute("SELECT a, b, c, d FROM mv_test" + i), + row(0, 1, 0, 0), + row(0, 1, 1, 0), + row(1, 1, 0, 1), + row(1, 1, 1, 0), + row(1, 1, 2, 0), + row(1, 2, 1, 0) + ); + + // delete rows that don't match the filter + execute("DELETE FROM %s WHERE a = ? AND b = ? AND c = ?", 2, -1, 0); + execute("DELETE FROM %s WHERE a = ? AND b = ? AND c = ?", 2, 0, 0); + execute("DELETE FROM %s WHERE a = ? AND b = ?", 0, 0); + assertRowsIgnoringOrder(execute("SELECT a, b, c, d FROM mv_test" + i), + row(0, 1, 0, 0), + row(0, 1, 1, 0), + row(1, 1, 0, 1), + row(1, 1, 1, 0), + row(1, 1, 2, 0), + row(1, 2, 1, 0) + ); + + // delete a row that does match the filter + execute("DELETE FROM %s WHERE a = ? AND b = ? AND c = ?", 1, 1, 0); + assertRowsIgnoringOrder(execute("SELECT a, b, c, d FROM mv_test" + i), + row(0, 1, 0, 0), + row(0, 1, 1, 0), + row(1, 1, 1, 0), + row(1, 1, 2, 0), + row(1, 2, 1, 0) + ); + + // delete a partition that matches the filter + execute("DELETE FROM %s WHERE a = ?", 1); + assertRowsIgnoringOrder(execute("SELECT a, b, c, d FROM mv_test" + i), + row(0, 1, 0, 0), + row(0, 1, 1, 0) + ); + + dropView("mv_test" + i); + dropTable("DROP TABLE %s"); + } + } + + @Test + public void testClusteringKeyMultiColumnRestrictions() throws Throwable + { + List mvPrimaryKeys = Arrays.asList("((a, b), c)", "((b, a), c)", "(a, b, c)", "(c, b, a)", "((c, a), b)"); + for (int i = 0; i < mvPrimaryKeys.size(); i++) + { + createTable("CREATE TABLE %s (a int, b int, c int, d int, PRIMARY KEY (a, b, c))"); + + execute("USE " + keyspace()); + executeNet(protocolVersion, "USE " + keyspace()); + + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 0, 0, 0, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 0, 0, 1, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 0, 1, 0, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 0, 1, 1, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 1, 0, 0, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 1, 0, 1, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 1, 1, -1, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 1, 1, 0, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 1, 1, 1, 0); + + logger.info("Testing MV primary key: {}", mvPrimaryKeys.get(i)); + + // only accept rows where b = 1 + createView("mv_test" + i, "CREATE MATERIALIZED VIEW %s AS SELECT * FROM %%s WHERE a IS NOT NULL AND (b, c) >= (1, 0) PRIMARY KEY " + mvPrimaryKeys.get(i)); + + while (!SystemKeyspace.isViewBuilt(keyspace(), "mv_test" + i)) + Thread.sleep(10); + + assertRowsIgnoringOrder(execute("SELECT a, b, c, d FROM mv_test" + i), + row(0, 1, 0, 0), + row(0, 1, 1, 0), + row(1, 1, 0, 0), + row(1, 1, 1, 0) + ); + + // insert new rows that do not match the filter + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 2, -1, 0, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 2, 0, 0, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 2, 1, -1, 0); + assertRowsIgnoringOrder(execute("SELECT a, b, c, d FROM mv_test" + i), + row(0, 1, 0, 0), + row(0, 1, 1, 0), + row(1, 1, 0, 0), + row(1, 1, 1, 0) + ); + + // insert new row that does match the filter + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 1, 1, 2, 0); + assertRowsIgnoringOrder(execute("SELECT a, b, c, d FROM mv_test" + i), + row(0, 1, 0, 0), + row(0, 1, 1, 0), + row(1, 1, 0, 0), + row(1, 1, 1, 0), + row(1, 1, 2, 0) + ); + + // update rows that don't match the filter + execute("UPDATE %s SET d = ? WHERE a = ? AND b = ? AND c = ?", 1, 1, -1, 0); + execute("UPDATE %s SET d = ? WHERE a = ? AND b = ? AND c = ?", 1, 2, -1, 0); + execute("UPDATE %s SET d = ? WHERE a = ? AND b = ? AND c = ?", 1, 2, 0, 0); + assertRowsIgnoringOrder(execute("SELECT a, b, c, d FROM mv_test" + i), + row(0, 1, 0, 0), + row(0, 1, 1, 0), + row(1, 1, 0, 0), + row(1, 1, 1, 0), + row(1, 1, 2, 0) + ); + + // update a row that does match the filter + execute("UPDATE %s SET d = ? WHERE a = ? AND b = ? AND c = ?", 1, 1, 1, 0); + assertRowsIgnoringOrder(execute("SELECT a, b, c, d FROM mv_test" + i), + row(0, 1, 0, 0), + row(0, 1, 1, 0), + row(1, 1, 0, 1), + row(1, 1, 1, 0), + row(1, 1, 2, 0) + ); + + // delete rows that don't match the filter + execute("DELETE FROM %s WHERE a = ? AND b = ? AND c = ?", 1, 1, -1); + execute("DELETE FROM %s WHERE a = ? AND b = ? AND c = ?", 2, -1, 0); + execute("DELETE FROM %s WHERE a = ? AND b = ? AND c = ?", 2, 0, 0); + execute("DELETE FROM %s WHERE a = ? AND b = ?", 0, 0); + assertRowsIgnoringOrder(execute("SELECT a, b, c, d FROM mv_test" + i), + row(0, 1, 0, 0), + row(0, 1, 1, 0), + row(1, 1, 0, 1), + row(1, 1, 1, 0), + row(1, 1, 2, 0) + ); + + // delete a row that does match the filter + execute("DELETE FROM %s WHERE a = ? AND b = ? AND c = ?", 1, 1, 0); + assertRowsIgnoringOrder(execute("SELECT a, b, c, d FROM mv_test" + i), + row(0, 1, 0, 0), + row(0, 1, 1, 0), + row(1, 1, 1, 0), + row(1, 1, 2, 0) + ); + + // delete a partition that matches the filter + execute("DELETE FROM %s WHERE a = ?", 1); + assertRowsIgnoringOrder(execute("SELECT a, b, c, d FROM mv_test" + i), + row(0, 1, 0, 0), + row(0, 1, 1, 0) + ); + + dropView("mv_test" + i); + dropTable("DROP TABLE %s"); + } + } + + @Test + public void testClusteringKeyFilteringRestrictions() throws Throwable + { + List mvPrimaryKeys = Arrays.asList("((a, b), c)", "((b, a), c)", "(a, b, c)", "(c, b, a)", "((c, a), b)"); + for (int i = 0; i < mvPrimaryKeys.size(); i++) + { + createTable("CREATE TABLE %s (a int, b int, c int, d int, PRIMARY KEY (a, b, c))"); + + execute("USE " + keyspace()); + executeNet(protocolVersion, "USE " + keyspace()); + + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 0, 0, 0, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 0, 0, 1, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 0, 1, 0, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 0, 1, 1, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 1, 0, 0, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 1, 0, 1, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 1, 1, -1, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 1, 1, 0, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 1, 1, 1, 0); + + logger.info("Testing MV primary key: {}", mvPrimaryKeys.get(i)); + + // only accept rows where b = 1 + createView("mv_test" + i, "CREATE MATERIALIZED VIEW %s AS SELECT * FROM %%s WHERE a IS NOT NULL AND b IS NOT NULL AND c = 1 PRIMARY KEY " + mvPrimaryKeys.get(i)); + + while (!SystemKeyspace.isViewBuilt(keyspace(), "mv_test" + i)) + Thread.sleep(10); + + assertRowsIgnoringOrder(execute("SELECT a, b, c, d FROM mv_test" + i), + row(0, 0, 1, 0), + row(0, 1, 1, 0), + row(1, 0, 1, 0), + row(1, 1, 1, 0) + ); + + // insert new rows that do not match the filter + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 2, 0, 0, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 2, 1, -1, 0); + assertRowsIgnoringOrder(execute("SELECT a, b, c, d FROM mv_test" + i), + row(0, 0, 1, 0), + row(0, 1, 1, 0), + row(1, 0, 1, 0), + row(1, 1, 1, 0) + ); + + // insert new row that does match the filter + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 1, 2, 1, 0); + assertRowsIgnoringOrder(execute("SELECT a, b, c, d FROM mv_test" + i), + row(0, 0, 1, 0), + row(0, 1, 1, 0), + row(1, 0, 1, 0), + row(1, 1, 1, 0), + row(1, 2, 1, 0) + ); + + // update rows that don't match the filter + execute("UPDATE %s SET d = ? WHERE a = ? AND b = ? AND c = ?", 1, 1, -1, 0); + execute("UPDATE %s SET d = ? WHERE a = ? AND b = ? AND c = ?", 1, 2, 0, 0); + assertRowsIgnoringOrder(execute("SELECT a, b, c, d FROM mv_test" + i), + row(0, 0, 1, 0), + row(0, 1, 1, 0), + row(1, 0, 1, 0), + row(1, 1, 1, 0), + row(1, 2, 1, 0) + ); + + // update a row that does match the filter + execute("UPDATE %s SET d = ? WHERE a = ? AND b = ? AND c = ?", 2, 1, 1, 1); + assertRowsIgnoringOrder(execute("SELECT a, b, c, d FROM mv_test" + i), + row(0, 0, 1, 0), + row(0, 1, 1, 0), + row(1, 0, 1, 0), + row(1, 1, 1, 2), + row(1, 2, 1, 0) + ); + + // delete rows that don't match the filter + execute("DELETE FROM %s WHERE a = ? AND b = ? AND c = ?", 1, 1, -1); + execute("DELETE FROM %s WHERE a = ? AND b = ? AND c = ?", 2, -1, 0); + execute("DELETE FROM %s WHERE a = ? AND b = ? AND c = ?", 2, 0, 0); + execute("DELETE FROM %s WHERE a = ? AND b = ?", 0, -1); + assertRowsIgnoringOrder(execute("SELECT a, b, c, d FROM mv_test" + i), + row(0, 0, 1, 0), + row(0, 1, 1, 0), + row(1, 0, 1, 0), + row(1, 1, 1, 2), + row(1, 2, 1, 0) + ); + + // delete a row that does match the filter + execute("DELETE FROM %s WHERE a = ? AND b = ? AND c = ?", 1, 1, 1); + assertRowsIgnoringOrder(execute("SELECT a, b, c, d FROM mv_test" + i), + row(0, 0, 1, 0), + row(0, 1, 1, 0), + row(1, 0, 1, 0), + row(1, 2, 1, 0) + ); + + // delete a partition that matches the filter + execute("DELETE FROM %s WHERE a = ?", 1); + assertRowsIgnoringOrder(execute("SELECT a, b, c, d FROM mv_test" + i), + row(0, 0, 1, 0), + row(0, 1, 1, 0) + ); + + dropView("mv_test" + i); + dropTable("DROP TABLE %s"); + } + } + + @Test + public void testPartitionKeyAndClusteringKeyFilteringRestrictions() throws Throwable + { + List mvPrimaryKeys = Arrays.asList("((a, b), c)", "((b, a), c)", "(a, b, c)", "(c, b, a)", "((c, a), b)"); + for (int i = 0; i < mvPrimaryKeys.size(); i++) + { + createTable("CREATE TABLE %s (a int, b int, c int, d int, PRIMARY KEY (a, b, c))"); + + execute("USE " + keyspace()); + executeNet(protocolVersion, "USE " + keyspace()); + + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 0, 0, 0, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 0, 0, 1, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 0, 1, 0, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 0, 1, 1, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 1, 0, 0, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 1, 0, 1, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 1, 1, -1, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 1, 1, 0, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 1, 1, 1, 0); + + logger.info("Testing MV primary key: {}", mvPrimaryKeys.get(i)); + + // only accept rows where b = 1 + createView("mv_test" + i, "CREATE MATERIALIZED VIEW %s AS SELECT * FROM %%s WHERE a = 1 AND b IS NOT NULL AND c = 1 PRIMARY KEY " + mvPrimaryKeys.get(i)); + + while (!SystemKeyspace.isViewBuilt(keyspace(), "mv_test" + i)) + Thread.sleep(10); + + assertRowsIgnoringOrder(execute("SELECT a, b, c, d FROM mv_test" + i), + row(1, 0, 1, 0), + row(1, 1, 1, 0) + ); + + // insert new rows that do not match the filter + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 0, 0, 1, 0); + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 1, 1, 0, 0); + assertRowsIgnoringOrder(execute("SELECT a, b, c, d FROM mv_test" + i), + row(1, 0, 1, 0), + row(1, 1, 1, 0) + ); + + // insert new row that does match the filter + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?)", 1, 2, 1, 0); + assertRowsIgnoringOrder(execute("SELECT a, b, c, d FROM mv_test" + i), + row(1, 0, 1, 0), + row(1, 1, 1, 0), + row(1, 2, 1, 0) + ); + + // update rows that don't match the filter + execute("UPDATE %s SET d = ? WHERE a = ? AND b = ? AND c = ?", 1, 1, -1, 0); + execute("UPDATE %s SET d = ? WHERE a = ? AND b = ? AND c = ?", 0, 1, 1, 0); + assertRowsIgnoringOrder(execute("SELECT a, b, c, d FROM mv_test" + i), + row(1, 0, 1, 0), + row(1, 1, 1, 0), + row(1, 2, 1, 0) + ); + + // update a row that does match the filter + execute("UPDATE %s SET d = ? WHERE a = ? AND b = ? AND c = ?", 2, 1, 1, 1); + assertRowsIgnoringOrder(execute("SELECT a, b, c, d FROM mv_test" + i), + row(1, 0, 1, 0), + row(1, 1, 1, 2), + row(1, 2, 1, 0) + ); + + // delete rows that don't match the filter + execute("DELETE FROM %s WHERE a = ? AND b = ? AND c = ?", 1, 1, -1); + execute("DELETE FROM %s WHERE a = ? AND b = ? AND c = ?", 2, 0, 1); + execute("DELETE FROM %s WHERE a = ?", 0); + assertRowsIgnoringOrder(execute("SELECT a, b, c, d FROM mv_test" + i), + row(1, 0, 1, 0), + row(1, 1, 1, 2), + row(1, 2, 1, 0) + ); + + // delete a row that does match the filter + execute("DELETE FROM %s WHERE a = ? AND b = ? AND c = ?", 1, 1, 1); + assertRowsIgnoringOrder(execute("SELECT a, b, c, d FROM mv_test" + i), + row(1, 0, 1, 0), + row(1, 2, 1, 0) + ); + + // delete a partition that matches the filter + execute("DELETE FROM %s WHERE a = ?", 1); + assertEmpty(execute("SELECT a, b, c, d FROM mv_test" + i)); + + dropView("mv_test" + i); + dropTable("DROP TABLE %s"); + } + } + + @Test + public void testAllTypes() throws Throwable + { + String myType = createType("CREATE TYPE %s (a int, b uuid, c set)"); + String columnNames = "asciival, " + + "bigintval, " + + "blobval, " + + "booleanval, " + + "dateval, " + + "decimalval, " + + "doubleval, " + + "floatval, " + + "inetval, " + + "intval, " + + "textval, " + + "timeval, " + + "timestampval, " + + "timeuuidval, " + + "uuidval," + + "varcharval, " + + "varintval, " + + "frozenlistval, " + + "frozensetval, " + + "frozenmapval, " + + "tupleval, " + + "udtval"; + + createTable( + "CREATE TABLE %s (" + + "asciival ascii, " + + "bigintval bigint, " + + "blobval blob, " + + "booleanval boolean, " + + "dateval date, " + + "decimalval decimal, " + + "doubleval double, " + + "floatval float, " + + "inetval inet, " + + "intval int, " + + "textval text, " + + "timeval time, " + + "timestampval timestamp, " + + "timeuuidval timeuuid, " + + "uuidval uuid," + + "varcharval varchar, " + + "varintval varint, " + + "frozenlistval frozen>, " + + "frozensetval frozen>, " + + "frozenmapval frozen>," + + "tupleval frozen>," + + "udtval frozen<" + myType + ">, " + + "PRIMARY KEY (" + columnNames + "))"); + + execute("USE " + keyspace()); + executeNet(protocolVersion, "USE " + keyspace()); + + + createView( + "mv_test", + "CREATE MATERIALIZED VIEW %s AS SELECT * FROM %%s WHERE " + + "asciival = 'abc' AND " + + "bigintval = 123 AND " + + "blobval = 0xfeed AND " + + "booleanval = true AND " + + "dateval = '1987-03-23' AND " + + "decimalval = 123.123 AND " + + "doubleval = 123.123 AND " + + "floatval = 123.123 AND " + + "inetval = '127.0.0.1' AND " + + "intval = 123 AND " + + "textval = 'abc' AND " + + "timeval = '07:35:07.000111222' AND " + + "timestampval = 123123123 AND " + + "timeuuidval = 6BDDC89A-5644-11E4-97FC-56847AFE9799 AND " + + "uuidval = 6BDDC89A-5644-11E4-97FC-56847AFE9799 AND " + + "varcharval = 'abc' AND " + + "varintval = 123123123 AND " + + "frozenlistval = [1, 2, 3] AND " + + "frozensetval = {6BDDC89A-5644-11E4-97FC-56847AFE9799} AND " + + "frozenmapval = {'a': 1, 'b': 2} AND " + + "tupleval = (1, 'foobar', 6BDDC89A-5644-11E4-97FC-56847AFE9799) AND " + + "udtval = {a: 1, b: 6BDDC89A-5644-11E4-97FC-56847AFE9799, c: {'foo', 'bar'}} " + + "PRIMARY KEY (" + columnNames + ")"); + + execute("INSERT INTO %s (" + columnNames + ") VALUES (" + + "'abc'," + + "123," + + "0xfeed," + + "true," + + "'1987-03-23'," + + "123.123," + + "123.123," + + "123.123," + + "'127.0.0.1'," + + "123," + + "'abc'," + + "'07:35:07.000111222'," + + "123123123," + + "6BDDC89A-5644-11E4-97FC-56847AFE9799," + + "6BDDC89A-5644-11E4-97FC-56847AFE9799," + + "'abc'," + + "123123123," + + "[1, 2, 3]," + + "{6BDDC89A-5644-11E4-97FC-56847AFE9799}," + + "{'a': 1, 'b': 2}," + + "(1, 'foobar', 6BDDC89A-5644-11E4-97FC-56847AFE9799)," + + "{a: 1, b: 6BDDC89A-5644-11E4-97FC-56847AFE9799, c: {'foo', 'bar'}})"); + + assert !execute("SELECT * FROM mv_test").isEmpty(); + + executeNet(protocolVersion, "ALTER TABLE %s RENAME inetval TO foo"); + assert !execute("SELECT * FROM mv_test").isEmpty(); + } +} diff --git a/test/unit/org/apache/cassandra/cql3/ViewTest.java b/test/unit/org/apache/cassandra/cql3/ViewTest.java index 43f7747bf9..5d65115087 100644 --- a/test/unit/org/apache/cassandra/cql3/ViewTest.java +++ b/test/unit/org/apache/cassandra/cql3/ViewTest.java @@ -558,14 +558,14 @@ public class ViewTest extends CQLTester createView("mv_test2", "CREATE MATERIALIZED VIEW %s AS SELECT * FROM %%s WHERE textval2 IS NOT NULL AND k IS NOT NULL AND asciival IS NOT NULL AND bigintval IS NOT NULL AND textval1 IS NOT NULL PRIMARY KEY ((textval2, k), asciival, bigintval, textval1)"); while (!SystemKeyspace.isViewBuilt(keyspace(), "mv_test2")) - Thread.sleep(1000); + Thread.sleep(10); Assert.assertEquals(100, execute("select * from mv_test2").size()); createView("mv_test3", "CREATE MATERIALIZED VIEW %s AS SELECT * FROM %%s WHERE textval2 IS NOT NULL AND k IS NOT NULL AND asciival IS NOT NULL AND bigintval IS NOT NULL AND textval1 IS NOT NULL PRIMARY KEY ((textval2, k), bigintval, textval1, asciival)"); while (!SystemKeyspace.isViewBuilt(keyspace(), "mv_test3")) - Thread.sleep(1000); + Thread.sleep(10); Assert.assertEquals(100, execute("select * from mv_test3").size()); Assert.assertEquals(100, execute("select asciival from mv_test3 where textval2 = ? and k = ?", "baz", 0).size()); diff --git a/test/unit/org/apache/cassandra/cql3/validation/operations/SelectSingleColumnRelationTest.java b/test/unit/org/apache/cassandra/cql3/validation/operations/SelectSingleColumnRelationTest.java index f80e489f7e..fef34caad3 100644 --- a/test/unit/org/apache/cassandra/cql3/validation/operations/SelectSingleColumnRelationTest.java +++ b/test/unit/org/apache/cassandra/cql3/validation/operations/SelectSingleColumnRelationTest.java @@ -62,6 +62,10 @@ public class SelectSingleColumnRelationTest extends CQLTester "SELECT * FROM %s WHERE c = 0 AND b <= ?", set(0)); assertInvalidMessage("Collection column 'b' (set) cannot be restricted by a 'IN' relation", "SELECT * FROM %s WHERE c = 0 AND b IN (?)", set(0)); + assertInvalidMessage("Unsupported \"!=\" relation: b != 5", + "SELECT * FROM %s WHERE c = 0 AND b != 5"); + assertInvalidMessage("Unsupported restriction: b IS NOT NULL", + "SELECT * FROM %s WHERE c = 0 AND b IS NOT NULL"); } @Test