From 1b36740ebe66b8ed4c3d6cb64eb2419a9279dfbf Mon Sep 17 00:00:00 2001 From: Zhao Yang Date: Wed, 12 Jul 2017 17:49:38 +0800 Subject: [PATCH] Fix outstanding MV timestamp issues and add documentation about unsupported cases (see CASSANDRA-11500 for a summary of fixes) This patch introduces the following changes to fix MV timestamp issues: - Add strict liveness for view with non-key base column in pk - Deprecated shadowable tombstone and use expired livenessInfo instead - Include partition deletion for existing base row - Disallow dropping base column with MV Patch by Zhao Yang and Paulo Motta; reviewed by Paulo Motta for CASSANDRA-11500 --- NEWS.txt | 18 + doc/cql3/CQL.textile | 6 + .../apache/cassandra/config/CFMetaData.java | 13 + .../cassandra/cql3/UpdateParameters.java | 2 +- .../cql3/statements/AlterTableStatement.java | 18 +- .../org/apache/cassandra/db/LivenessInfo.java | 17 +- .../org/apache/cassandra/db/ReadCommand.java | 7 +- .../db/compaction/CompactionIterator.java | 7 +- .../apache/cassandra/db/filter/RowFilter.java | 4 +- .../db/partitions/PurgeFunction.java | 14 +- .../apache/cassandra/db/rows/BTreeRow.java | 6 +- .../org/apache/cassandra/db/rows/Row.java | 15 +- .../db/rows/UnfilteredSerializer.java | 5 + .../apache/cassandra/db/transform/Filter.java | 8 +- .../db/transform/FilteredPartitions.java | 4 +- .../cassandra/db/transform/FilteredRows.java | 2 +- .../apache/cassandra/db/view/TableViews.java | 18 +- .../org/apache/cassandra/db/view/View.java | 43 +- .../apache/cassandra/db/view/ViewManager.java | 5 + .../db/view/ViewUpdateGenerator.java | 163 +- .../cassandra/service/DataResolver.java | 4 +- .../org/apache/cassandra/cql3/CQLTester.java | 2 +- .../cassandra/cql3/ViewComplexTest.java | 1343 +++++++++++++++++ .../cassandra/cql3/ViewFilteringTest.java | 706 ++++----- .../org/apache/cassandra/cql3/ViewTest.java | 31 +- 25 files changed, 1973 insertions(+), 488 deletions(-) create mode 100644 test/unit/org/apache/cassandra/cql3/ViewComplexTest.java diff --git a/NEWS.txt b/NEWS.txt index bb5fdfe40e..7064c5d751 100644 --- a/NEWS.txt +++ b/NEWS.txt @@ -21,6 +21,24 @@ Upgrading - Nothing specific to this release, but please see previous upgrading sections, especially if you are upgrading from 2.2. +Materialized Views +------------------- + - Cassandra will no longer allow dropping columns on tables with Materialized Views. + - A change was made in the way the Materialized View timestamp is computed, which + may cause an old deletion to a base column which is view primary key (PK) column + to not be reflected in the view when repairing the base table post-upgrade. This + condition is only possible when a column deletion to an MV primary key (PK) column + not present in the base table PK (via UPDATE base SET view_pk_col = null or DELETE + view_pk_col FROM base) is missed before the upgrade and received by repair after the upgrade. + If such column deletions are done on a view PK column which is not a base PK, it's advisable + to run repair on the base table of all nodes prior to the upgrade. Alternatively it's possible + to fix potential inconsistencies by running repair on the views after upgrade or drop and + re-create the views. See CASSANDRA-11500 for more details. + - Removal of columns not selected in the Materialized View (via UPDATE base SET unselected_column + = null or DELETE unselected_column FROM base) may not be properly reflected in the view in some + situations so we advise against doing deletions on base columns not selected in views + until this is fixed on CASSANDRA-13826. + 3.0.14 ====== diff --git a/doc/cql3/CQL.textile b/doc/cql3/CQL.textile index 1efa6d467a..54888b8717 100644 --- a/doc/cql3/CQL.textile +++ b/doc/cql3/CQL.textile @@ -524,6 +524,12 @@ h4(#createMVWhere). @WHERE@ Clause The @@ is similar to the "where clause of a @SELECT@ statement":#selectWhere, with a few differences. First, the where clause must contain an expression that disallows @NULL@ values in columns in the view's primary key. If no other restriction is desired, this can be accomplished with an @IS NOT NULL@ expression. Second, only columns which are in the base table's primary key may be restricted with expressions other than @IS NOT NULL@. (Note that this second restriction may be lifted in the future.) +h4. MV Limitations + +__Note:__ +Removal of columns not selected in the Materialized View (via `UPDATE base SET unselected_column = null` or `DELETE unselected_column FROM base`) may shadow missed updates to other columns received by hints or repair. +For this reason, we advise against doing deletions on base columns not selected in views until this is fixed on CASSANDRA-13826. + h3(#alterMVStmt). ALTER MATERIALIZED VIEW __Syntax:__ diff --git a/src/java/org/apache/cassandra/config/CFMetaData.java b/src/java/org/apache/cassandra/config/CFMetaData.java index 44f3a96c45..1eb991a1fa 100644 --- a/src/java/org/apache/cassandra/config/CFMetaData.java +++ b/src/java/org/apache/cassandra/config/CFMetaData.java @@ -1078,6 +1078,19 @@ public final class CFMetaData return isView; } + /** + * A table with strict liveness filters/ignores rows without PK liveness info, + * effectively tying the row liveness to its primary key liveness. + * + * Currently this is only used by views with normal base column as PK column + * so updates to other columns do not make the row live when the base column + * is not live. See CASSANDRA-11500. + */ + public boolean enforceStrictLiveness() + { + return isView && Keyspace.open(ksName).viewManager.getByName(cfName).enforceStrictLiveness(); + } + public Serializers serializers() { return serializers; diff --git a/src/java/org/apache/cassandra/cql3/UpdateParameters.java b/src/java/org/apache/cassandra/cql3/UpdateParameters.java index 8ff5344170..d070f61127 100644 --- a/src/java/org/apache/cassandra/cql3/UpdateParameters.java +++ b/src/java/org/apache/cassandra/cql3/UpdateParameters.java @@ -233,6 +233,6 @@ public class UpdateParameters return pendingMutations; return Rows.merge(prefetchedRow, pendingMutations, nowInSec) - .purge(DeletionPurger.PURGE_ALL, nowInSec); + .purge(DeletionPurger.PURGE_ALL, nowInSec, metadata.enforceStrictLiveness()); } } diff --git a/src/java/org/apache/cassandra/cql3/statements/AlterTableStatement.java b/src/java/org/apache/cassandra/cql3/statements/AlterTableStatement.java index 5db4b9fdf8..befdd253c6 100644 --- a/src/java/org/apache/cassandra/cql3/statements/AlterTableStatement.java +++ b/src/java/org/apache/cassandra/cql3/statements/AlterTableStatement.java @@ -233,22 +233,10 @@ public class AlterTableStatement extends SchemaAlteringStatement .collect(Collectors.joining(",")))); } - // If a column is dropped which is included in a view, we don't allow the drop to take place. - boolean rejectAlter = false; - StringBuilder builder = new StringBuilder(); - for (ViewDefinition view : views) - { - if (!view.includes(columnName)) continue; - if (rejectAlter) - builder.append(','); - rejectAlter = true; - builder.append(view.viewName); - } - if (rejectAlter) - throw new InvalidRequestException(String.format("Cannot drop column %s, depended on by materialized views (%s.{%s})", + if (!Iterables.isEmpty(views)) + throw new InvalidRequestException(String.format("Cannot drop column %s on base table with materialized views.", columnName.toString(), - keyspace(), - builder.toString())); + keyspace())); break; case OPTS: if (attrs == null) diff --git a/src/java/org/apache/cassandra/db/LivenessInfo.java b/src/java/org/apache/cassandra/db/LivenessInfo.java index 411fb9a38c..ab61a2365a 100644 --- a/src/java/org/apache/cassandra/db/LivenessInfo.java +++ b/src/java/org/apache/cassandra/db/LivenessInfo.java @@ -176,13 +176,26 @@ public class LivenessInfo * Whether this liveness information supersedes another one (that is * whether is has a greater timestamp than the other or not). * - * @param other the {@code LivenessInfo} to compare this info to. + *
+ * + * If timestamps are the same, livenessInfo with greater TTL supersedes another. + * + * It also means, if timestamps are the same, ttl superseders no-ttl. + * + * This is the same rule as {@link Conflicts#resolveRegular} + * + * @param other + * the {@code LivenessInfo} to compare this info to. * * @return whether this {@code LivenessInfo} supersedes {@code other}. */ public boolean supersedes(LivenessInfo other) { - return timestamp > other.timestamp; + if (timestamp != other.timestamp) + return timestamp > other.timestamp; + if (isExpiring() == other.isExpiring()) + return localExpirationTime() > other.localExpirationTime(); + return isExpiring(); } /** diff --git a/src/java/org/apache/cassandra/db/ReadCommand.java b/src/java/org/apache/cassandra/db/ReadCommand.java index 6a21bb35a1..b73cddea6c 100644 --- a/src/java/org/apache/cassandra/db/ReadCommand.java +++ b/src/java/org/apache/cassandra/db/ReadCommand.java @@ -541,7 +541,12 @@ public abstract class ReadCommand implements ReadQuery { public WithoutPurgeableTombstones() { - super(isForThrift, nowInSec(), cfs.gcBefore(nowInSec()), oldestUnrepairedTombstone(), cfs.getCompactionStrategyManager().onlyPurgeRepairedTombstones()); + super(isForThrift, + nowInSec(), + cfs.gcBefore(nowInSec()), + oldestUnrepairedTombstone(), + cfs.getCompactionStrategyManager().onlyPurgeRepairedTombstones(), + cfs.metadata.enforceStrictLiveness()); } protected Predicate getPurgeEvaluator() diff --git a/src/java/org/apache/cassandra/db/compaction/CompactionIterator.java b/src/java/org/apache/cassandra/db/compaction/CompactionIterator.java index 9f0984fab6..bea365c347 100644 --- a/src/java/org/apache/cassandra/db/compaction/CompactionIterator.java +++ b/src/java/org/apache/cassandra/db/compaction/CompactionIterator.java @@ -266,7 +266,12 @@ public class CompactionIterator extends CompactionInfo.Holder implements Unfilte private Purger(boolean isForThrift, CompactionController controller, int nowInSec) { - super(isForThrift, nowInSec, controller.gcBefore, controller.compactingRepaired() ? Integer.MAX_VALUE : Integer.MIN_VALUE, controller.cfs.getCompactionStrategyManager().onlyPurgeRepairedTombstones()); + super(isForThrift, + nowInSec, + controller.gcBefore, + controller.compactingRepaired() ? Integer.MAX_VALUE : Integer.MIN_VALUE, + controller.cfs.getCompactionStrategyManager().onlyPurgeRepairedTombstones(), + controller.cfs.metadata.enforceStrictLiveness()); this.controller = controller; } diff --git a/src/java/org/apache/cassandra/db/filter/RowFilter.java b/src/java/org/apache/cassandra/db/filter/RowFilter.java index 8d11038898..9f3b8682b7 100644 --- a/src/java/org/apache/cassandra/db/filter/RowFilter.java +++ b/src/java/org/apache/cassandra/db/filter/RowFilter.java @@ -148,7 +148,7 @@ public abstract class RowFilter implements Iterable public boolean isSatisfiedBy(CFMetaData metadata, DecoratedKey partitionKey, Row row, int nowInSec) { // We purge all tombstones as the expressions isSatisfiedBy methods expects it - Row purged = row.purge(DeletionPurger.PURGE_ALL, nowInSec); + Row purged = row.purge(DeletionPurger.PURGE_ALL, nowInSec, metadata.enforceStrictLiveness()); if (purged == null) return expressions.isEmpty(); @@ -290,7 +290,7 @@ public abstract class RowFilter implements Iterable public Row applyToRow(Row row) { - Row purged = row.purge(DeletionPurger.PURGE_ALL, nowInSec); + Row purged = row.purge(DeletionPurger.PURGE_ALL, nowInSec, metadata.enforceStrictLiveness()); if (purged == null) return null; diff --git a/src/java/org/apache/cassandra/db/partitions/PurgeFunction.java b/src/java/org/apache/cassandra/db/partitions/PurgeFunction.java index 6679bdf793..5cc9145097 100644 --- a/src/java/org/apache/cassandra/db/partitions/PurgeFunction.java +++ b/src/java/org/apache/cassandra/db/partitions/PurgeFunction.java @@ -28,9 +28,16 @@ public abstract class PurgeFunction extends Transformation= oldestUnrepairedTombstone) && localDeletionTime < gcBefore && getPurgeEvaluator().test(timestamp); + this.enforceStrictLiveness = enforceStrictLiveness; } protected abstract Predicate getPurgeEvaluator(); @@ -84,14 +92,14 @@ public abstract class PurgeFunction extends Transformation cd.purge(purger, nowInSec)); } diff --git a/src/java/org/apache/cassandra/db/rows/Row.java b/src/java/org/apache/cassandra/db/rows/Row.java index 9ab1f09ea6..3bcc220880 100644 --- a/src/java/org/apache/cassandra/db/rows/Row.java +++ b/src/java/org/apache/cassandra/db/rows/Row.java @@ -200,10 +200,20 @@ public interface Row extends Unfiltered, Collection * * @param purger the {@code DeletionPurger} to use to decide what can be purged. * @param nowInSec the current time to decide what is deleted and what isn't (in the case of expired cells). + * @param enforceStrictLiveness whether the row should be purged if there is no PK liveness info, + * normally retrieved from {@link CFMetaData#enforceStrictLiveness()} + * + * When enforceStrictLiveness is set, rows with empty PK liveness info + * and no row deletion are purged. + * + * Currently this is only used by views with normal base column as PK column + * so updates to other base columns do not make the row live when the PK column + * is not live. See CASSANDRA-11500. + * * @return this row but without any deletion info purged by {@code purger}. If the purged row is empty, returns * {@code null}. */ - public Row purge(DeletionPurger purger, int nowInSec); + public Row purge(DeletionPurger purger, int nowInSec, boolean enforceStrictLiveness); /** * Returns a copy of this row where all counter cells have they "local" shard marked for clearing. @@ -215,7 +225,7 @@ public interface Row extends Unfiltered, Collection * timestamp by {@code newTimestamp - 1}. * * @param newTimestamp the timestamp to use for all live data in the returned row. - * @param a copy of this row with timestamp updated using {@code newTimestamp}. This can return {@code null} in the + * @return a copy of this row with timestamp updated using {@code newTimestamp}. This can return {@code null} in the * rare where the row only as a shadowable row deletion and the new timestamp supersedes it. * * @see Commit for why we need this. @@ -277,6 +287,7 @@ public interface Row extends Unfiltered, Collection return time.isLive() ? LIVE : new Deletion(time, false); } + @Deprecated public static Deletion shadowable(DeletionTime time) { return new Deletion(time, true); diff --git a/src/java/org/apache/cassandra/db/rows/UnfilteredSerializer.java b/src/java/org/apache/cassandra/db/rows/UnfilteredSerializer.java index bdc8388eca..c4684e1fa9 100644 --- a/src/java/org/apache/cassandra/db/rows/UnfilteredSerializer.java +++ b/src/java/org/apache/cassandra/db/rows/UnfilteredSerializer.java @@ -88,6 +88,11 @@ public class UnfilteredSerializer * Extended flags */ private final static int IS_STATIC = 0x01; // Whether the encoded row is a static. If there is no extended flag, the row is assumed not static. + /** + * A shadowable tombstone cannot replace a previous row deletion otherwise it could resurrect a + * previously deleted cell not updated by a subsequent update, SEE CASSANDRA-11500 + */ + @Deprecated private final static int HAS_SHADOWABLE_DELETION = 0x02; // Whether the row deletion is shadowable. If there is no extended flag (or no row deletion), the deletion is assumed not shadowable. public void serialize(Unfiltered unfiltered, SerializationHeader header, DataOutputPlus out, int version) diff --git a/src/java/org/apache/cassandra/db/transform/Filter.java b/src/java/org/apache/cassandra/db/transform/Filter.java index 747983f50e..48a1634bb1 100644 --- a/src/java/org/apache/cassandra/db/transform/Filter.java +++ b/src/java/org/apache/cassandra/db/transform/Filter.java @@ -26,10 +26,12 @@ import org.apache.cassandra.db.rows.*; public final class Filter extends Transformation { private final int nowInSec; + private final boolean enforceStrictLiveness; - public Filter(int nowInSec) + public Filter(int nowInSec, boolean enforceStrictLiveness) { this.nowInSec = nowInSec; + this.enforceStrictLiveness = enforceStrictLiveness; } @Override @@ -46,14 +48,14 @@ public final class Filter extends Transformation if (row.isEmpty()) return Rows.EMPTY_STATIC_ROW; - row = row.purge(DeletionPurger.PURGE_ALL, nowInSec); + row = row.purge(DeletionPurger.PURGE_ALL, nowInSec, enforceStrictLiveness); return row == null ? Rows.EMPTY_STATIC_ROW : row; } @Override protected Row applyToRow(Row row) { - return row.purge(DeletionPurger.PURGE_ALL, nowInSec); + return row.purge(DeletionPurger.PURGE_ALL, nowInSec, enforceStrictLiveness); } @Override diff --git a/src/java/org/apache/cassandra/db/transform/FilteredPartitions.java b/src/java/org/apache/cassandra/db/transform/FilteredPartitions.java index ad9446d6dc..b835a6b23f 100644 --- a/src/java/org/apache/cassandra/db/transform/FilteredPartitions.java +++ b/src/java/org/apache/cassandra/db/transform/FilteredPartitions.java @@ -52,7 +52,9 @@ public final class FilteredPartitions extends BasePartitions> implem */ public static RowIterator filter(UnfilteredRowIterator iterator, int nowInSecs) { - return new Filter(nowInSecs).applyToPartition(iterator); + return new Filter(nowInSecs, iterator.metadata().enforceStrictLiveness()).applyToPartition(iterator); } } diff --git a/src/java/org/apache/cassandra/db/view/TableViews.java b/src/java/org/apache/cassandra/db/view/TableViews.java index 1a3cbb1b3c..d2d4a4576f 100644 --- a/src/java/org/apache/cassandra/db/view/TableViews.java +++ b/src/java/org/apache/cassandra/db/view/TableViews.java @@ -159,6 +159,7 @@ public class TableViews extends AbstractCollection * but has simply some updated values. This will be empty for view building as we want to assume anything we'll pass * to {@code updates} is new. * @param nowInSec the current time in seconds. + * @param separateUpdates, if false, mutation is per partition. * @return the mutations to apply to the {@code views}. This can be empty. */ public Iterator> generateViewUpdates(Collection views, @@ -282,7 +283,10 @@ public class TableViews extends AbstractCollection continue; Row updateRow = (Row) update; - addToViewUpdateGenerators(emptyRow(updateRow.clustering(), DeletionTime.LIVE), updateRow, generators, nowInSec); + addToViewUpdateGenerators(emptyRow(updateRow.clustering(), existingsDeletion.currentDeletion()), + updateRow, + generators, + nowInSec); // If the updates have been filtered, then we won't have any mutations; we need to make sure that we // only return if the mutations are empty. Otherwise, we continue to search for an update which is @@ -321,7 +325,10 @@ public class TableViews extends AbstractCollection continue; Row updateRow = (Row) update; - addToViewUpdateGenerators(emptyRow(updateRow.clustering(), DeletionTime.LIVE), updateRow, generators, nowInSec); + addToViewUpdateGenerators(emptyRow(updateRow.clustering(), existingsDeletion.currentDeletion()), + updateRow, + generators, + nowInSec); } return Iterators.singletonIterator(buildMutations(baseTableMetadata, generators)); @@ -419,11 +426,12 @@ public class TableViews extends AbstractCollection ClusteringIndexFilter clusteringFilter = names == null ? new ClusteringIndexSliceFilter(sliceBuilder.build(), false) : new ClusteringIndexNamesFilter(names, false); + // since unselected columns also affect view liveness, we need to query all base columns if base and view have same key columns. // If we have more than one view, we should merge the queried columns by each views but to keep it simple we just // include everything. We could change that in the future. - ColumnFilter queriedColumns = views.size() == 1 - ? Iterables.getOnlyElement(views).getSelectStatement().queriedColumns() - : ColumnFilter.all(metadata); + ColumnFilter queriedColumns = views.size() == 1 && metadata.enforceStrictLiveness() + ? Iterables.getOnlyElement(views).getSelectStatement().queriedColumns() + : ColumnFilter.all(metadata); // Note that the views could have restrictions on regular columns, but even if that's the case we shouldn't apply those // when we read, because even if an existing row doesn't match the view filter, the update can change that in which // case we'll need to know the existing content. There is also no easy way to merge those RowFilter when we have multiple views. diff --git a/src/java/org/apache/cassandra/db/view/View.java b/src/java/org/apache/cassandra/db/view/View.java index e471349186..58e2a841c7 100644 --- a/src/java/org/apache/cassandra/db/view/View.java +++ b/src/java/org/apache/cassandra/db/view/View.java @@ -60,7 +60,6 @@ public class View public volatile List baseNonPKColumnsInViewPK; - 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. @@ -75,7 +74,6 @@ public class View { this.baseCfs = baseCfs; this.name = definition.viewName; - this.includeAllColumns = definition.includeAllColumns; this.rawSelect = definition.select; updateDefinition(definition); @@ -95,8 +93,6 @@ public class View public void updateDefinition(ViewDefinition definition) { this.definition = definition; - - CFMetaData viewCfm = definition.metadata; List nonPKDefPartOfViewPK = new ArrayList<>(); for (ColumnDefinition baseColumn : baseCfs.metadata.allColumns()) { @@ -148,21 +144,7 @@ public class View // neither included in the view, nor used by the view filter). if (!getReadQuery().selectsClustering(partitionKey, update.clustering())) return false; - - // We want to find if the update modify any of the columns that are part of the view (in which case the view is affected). - // But if the view include all the base table columns, or the update has either a row deletion or a row liveness (note - // that for the liveness, it would be more "precise" to check if it's live, but pushing an update that is already expired - // is dump so it's ok not to optimize for it and it saves us from having to pass nowInSec to the method), we know the view - // is affected right away. - if (includeAllColumns || !update.deletion().isLive() || !update.primaryKeyLivenessInfo().isEmpty()) - return true; - - for (ColumnData data : update) - { - if (definition.metadata.getColumnDefinition(data.column().name) != null) - return true; - } - return false; + return true; } /** @@ -293,4 +275,27 @@ public class View return expressions.stream().collect(Collectors.joining(" AND ")); } + + public boolean hasSamePrimaryKeyColumnsAsBaseTable() + { + return baseNonPKColumnsInViewPK.isEmpty(); + } + + /** + * When views contains a primary key column that is not part + * of the base table primary key, we use that column liveness + * info as the view PK, to ensure that whenever that column + * is not live in the base, the row is not live in the view. + * + * This is done to prevent cells other than the view PK from + * making the view row alive when the view PK column is not + * live in the base. So in this case we tie the row liveness, + * to the primary key liveness. + * + * See CASSANDRA-11500 for context. + */ + public boolean enforceStrictLiveness() + { + return !baseNonPKColumnsInViewPK.isEmpty(); + } } diff --git a/src/java/org/apache/cassandra/db/view/ViewManager.java b/src/java/org/apache/cassandra/db/view/ViewManager.java index 0a0fa7b4a7..d1cfd9e193 100644 --- a/src/java/org/apache/cassandra/db/view/ViewManager.java +++ b/src/java/org/apache/cassandra/db/view/ViewManager.java @@ -161,6 +161,11 @@ public class ViewManager SystemKeyspace.setViewRemoved(keyspace.getName(), view.name); } + public View getByName(String name) + { + return viewsByName.get(name); + } + public void buildAllViews() { for (View view : allViews()) diff --git a/src/java/org/apache/cassandra/db/view/ViewUpdateGenerator.java b/src/java/org/apache/cassandra/db/view/ViewUpdateGenerator.java index edb88d0263..0c8e078298 100644 --- a/src/java/org/apache/cassandra/db/view/ViewUpdateGenerator.java +++ b/src/java/org/apache/cassandra/db/view/ViewUpdateGenerator.java @@ -70,7 +70,7 @@ public class ViewUpdateGenerator UPDATE_EXISTING, // There was an entry and the update modifies it SWITCH_ENTRY // There was an entry and there is still one after update, // but they are not the same one. - }; + } /** * Creates a new {@code ViewUpdateBuilder}. @@ -121,14 +121,14 @@ public class ViewUpdateGenerator createEntry(mergedBaseRow); return; case DELETE_OLD: - deleteOldEntry(existingBaseRow); + deleteOldEntry(existingBaseRow, mergedBaseRow); return; case UPDATE_EXISTING: updateEntry(existingBaseRow, mergedBaseRow); return; case SWITCH_ENTRY: createEntry(mergedBaseRow); - deleteOldEntry(existingBaseRow); + deleteOldEntry(existingBaseRow, mergedBaseRow); return; } } @@ -180,6 +180,7 @@ public class ViewUpdateGenerator } assert view.baseNonPKColumnsInViewPK.size() <= 1 : "We currently only support one base non-PK column in the view PK"; + if (view.baseNonPKColumnsInViewPK.isEmpty()) { // The view entry is necessarily the same pre and post update. @@ -204,7 +205,9 @@ public class ViewUpdateGenerator if (!isLive(before)) return isLive(after) ? UpdateAction.NEW_ENTRY : UpdateAction.NONE; if (!isLive(after)) + { return UpdateAction.DELETE_OLD; + } return baseColumn.cellValueType().compare(before.value(), after.value()) == 0 ? UpdateAction.UPDATE_EXISTING @@ -269,7 +272,7 @@ public class ViewUpdateGenerator } if (!matchesViewFilter(mergedBaseRow)) { - deleteOldEntryInternal(existingBaseRow); + deleteOldEntryInternal(existingBaseRow, mergedBaseRow); return; } @@ -281,6 +284,12 @@ public class ViewUpdateGenerator currentViewEntryBuilder.addPrimaryKeyLivenessInfo(computeLivenessInfoForEntry(mergedBaseRow)); currentViewEntryBuilder.addRowDeletion(mergedBaseRow.deletion()); + addDifferentCells(existingBaseRow, mergedBaseRow); + submitUpdate(); + } + + private void addDifferentCells(Row existingBaseRow, Row mergedBaseRow) + { // We only add to the view update the cells from mergedBaseRow that differs from // existingBaseRow. For that and for speed we can just cell pointer equality: if the update // hasn't touched a cell, we know it will be the same object in existingBaseRow and @@ -362,8 +371,6 @@ public class ViewUpdateGenerator addCell(viewColumn, (Cell)mergedData); } } - - submitUpdate(); } /** @@ -371,20 +378,40 @@ public class ViewUpdateGenerator *

* This method checks that the base row does match the view filter before bothering. */ - private void deleteOldEntry(Row existingBaseRow) + private void deleteOldEntry(Row existingBaseRow, Row mergedBaseRow) { // Before deleting an old entry, make sure it was matching the view filter (otherwise there is nothing to delete) if (!matchesViewFilter(existingBaseRow)) return; - deleteOldEntryInternal(existingBaseRow); + deleteOldEntryInternal(existingBaseRow, mergedBaseRow); } - private void deleteOldEntryInternal(Row existingBaseRow) + private void deleteOldEntryInternal(Row existingBaseRow, Row mergedBaseRow) { startNewUpdate(existingBaseRow); - DeletionTime dt = new DeletionTime(computeTimestampForEntryDeletion(existingBaseRow), nowInSec); - currentViewEntryBuilder.addRowDeletion(Row.Deletion.shadowable(dt)); + long timestamp = computeTimestampForEntryDeletion(existingBaseRow, mergedBaseRow); + long rowDeletion = mergedBaseRow.deletion().time().markedForDeleteAt(); + assert timestamp >= rowDeletion; + + // If computed deletion timestamp greater than row deletion, it must be coming from + // 1. non-pk base column used in view pk, or + // 2. unselected base column + // any case, we need to use it as expired livenessInfo + // If computed deletion timestamp is from row deletion, we only need row deletion itself + if (timestamp > rowDeletion) + { + /** + * TODO: This is a hack and overload of LivenessInfo and we should probably modify + * the storage engine to properly support this, but on the meantime this + * should be fine because it only happens in some specific scenarios explained above. + */ + LivenessInfo info = LivenessInfo.create(timestamp, Integer.MAX_VALUE, nowInSec); + currentViewEntryBuilder.addPrimaryKeyLivenessInfo(info); + } + currentViewEntryBuilder.addRowDeletion(mergedBaseRow.deletion()); + + addDifferentCells(existingBaseRow, mergedBaseRow); submitUpdate(); } @@ -413,78 +440,92 @@ public class ViewUpdateGenerator private LivenessInfo computeLivenessInfoForEntry(Row baseRow) { - /* - * We need to compute both the timestamp and expiration. + /** + * There 3 cases: + * 1. No extra primary key in view and all base columns are selected in MV. all base row's components(livenessInfo, + * deletion, cells) are same as view row. Simply map base components to view row. + * 2. There is a base non-key column used in view pk. This base non-key column determines the liveness of view row. view's row level + * info should based on this column. + * 3. Most tricky case is no extra primary key in view and some base columns are not selected in MV. We cannot use 1 livenessInfo or + * row deletion to represent the liveness of unselected column properly, see CASSANDRA-11500. + * We could make some simplification: the unselected columns will be used only when it affects view row liveness. eg. if view row + * already exists and not expiring, there is no need to use unselected columns. + * Note: if the view row is removed due to unselected column removal(ttl or cell tombstone), we will have problem keeping view + * row alive with a smaller or equal timestamp than the max unselected column timestamp. * - * For the timestamp, it makes sense to use the bigger timestamp for all view PK columns. - * - * This is more complex for the expiration. We want to maintain consistency between the base and the view, so the - * entry should only exist as long as the base row exists _and_ has non-null values for all the columns that are part - * of the view PK. - * Which means we really have 2 cases: - * 1) either the columns for the base and view PKs are exactly the same: in that case, the view entry should live - * as long as the base row lives. This means the view entry should only expire once *everything* in the base row - * has expired. Which means the row TTL should be the max of any other TTL. - * 2) or there is a column that is not in the base PK but is in the view PK (we can only have one so far, we'll need - * to slightly adapt if we allow more later): in that case, as long as that column lives the entry does too, but - * as soon as it expires (or is deleted for that matter) the entry also should expire. So the expiration for the - * view is the one of that column, irregarding of any other expiration. - * To take an example of that case, if you have: - * CREATE TABLE t (a int, b int, c int, PRIMARY KEY (a, b)) - * CREATE MATERIALIZED VIEW mv AS SELECT * FROM t WHERE c IS NOT NULL AND a IS NOT NULL AND b IS NOT NULL PRIMARY KEY (c, a, b) - * INSERT INTO t(a, b) VALUES (0, 0) USING TTL 3; - * UPDATE t SET c = 0 WHERE a = 0 AND b = 0; - * then even after 3 seconds elapsed, the row will still exist (it just won't have a "row marker" anymore) and so - * the MV should still have a corresponding entry. */ assert view.baseNonPKColumnsInViewPK.size() <= 1; // This may change, but is currently an enforced limitation LivenessInfo baseLiveness = baseRow.primaryKeyLivenessInfo(); - if (view.baseNonPKColumnsInViewPK.isEmpty()) + if (view.hasSamePrimaryKeyColumnsAsBaseTable()) { - int ttl = baseLiveness.ttl(); - int expirationTime = baseLiveness.localExpirationTime(); + if (view.getDefinition().includeAllColumns) + return baseLiveness; + + long timestamp = baseLiveness.timestamp(); + boolean hasNonExpiringLiveCell = false; + Cell biggestExpirationCell = null; for (Cell cell : baseRow.cells()) { - if (cell.ttl() > ttl) + if (view.getViewColumn(cell.column()) != null) + continue; + if (!isLive(cell)) + continue; + timestamp = Math.max(timestamp, cell.maxTimestamp()); + if (!cell.isExpiring()) + hasNonExpiringLiveCell = true; + else { - ttl = cell.ttl(); - expirationTime = cell.localDeletionTime(); + if (biggestExpirationCell == null) + biggestExpirationCell = cell; + else if (cell.localDeletionTime() > biggestExpirationCell.localDeletionTime()) + biggestExpirationCell = cell; } } - return ttl == baseLiveness.ttl() - ? baseLiveness - : LivenessInfo.create(baseLiveness.timestamp(), ttl, expirationTime); + if (baseLiveness.isLive(nowInSec) && !baseLiveness.isExpiring()) + return LivenessInfo.create(viewMetadata, timestamp, nowInSec); + if (hasNonExpiringLiveCell) + return LivenessInfo.create(viewMetadata, timestamp, nowInSec); + if (biggestExpirationCell == null) + return baseLiveness; + if (biggestExpirationCell.localDeletionTime() > baseLiveness.localExpirationTime() + || !baseLiveness.isLive(nowInSec)) + return LivenessInfo.create(timestamp, + biggestExpirationCell.ttl(), + biggestExpirationCell.localDeletionTime()); + return baseLiveness; } - ColumnDefinition baseColumn = view.baseNonPKColumnsInViewPK.get(0); - Cell cell = baseRow.getCell(baseColumn); + Cell cell = baseRow.getCell(view.baseNonPKColumnsInViewPK.get(0)); assert isLive(cell) : "We shouldn't have got there if the base row had no associated entry"; - long timestamp = Math.max(baseLiveness.timestamp(), cell.timestamp()); - return LivenessInfo.create(timestamp, cell.ttl(), cell.localDeletionTime()); + return LivenessInfo.create(cell.timestamp(), cell.ttl(), cell.localDeletionTime()); } - private long computeTimestampForEntryDeletion(Row baseRow) + private long computeTimestampForEntryDeletion(Row existingBaseRow, Row mergedBaseRow) { - // We delete the old row with it's row entry timestamp using a shadowable deletion. - // We must make sure that the deletion deletes everything in the entry (or the entry will - // still show up), so we must use the bigger timestamp found in the existing row (for any - // column included in the view at least). - // TODO: We have a problem though: if the entry is "resurected" by a later update, we would - // need to ensure that the timestamp for then entry then is bigger than the tombstone - // we're just inserting, which is not currently guaranteed. - // This is a bug for a separate ticket though. - long timestamp = baseRow.primaryKeyLivenessInfo().timestamp(); - for (ColumnData data : baseRow) + DeletionTime deletion = mergedBaseRow.deletion().time(); + if (view.hasSamePrimaryKeyColumnsAsBaseTable()) { - if (!view.getDefinition().includes(data.column().name)) - continue; + long timestamp = Math.max(deletion.markedForDeleteAt(), existingBaseRow.primaryKeyLivenessInfo().timestamp()); + if (view.getDefinition().includeAllColumns) + return timestamp; - timestamp = Math.max(timestamp, data.maxTimestamp()); + for (Cell cell : existingBaseRow.cells()) + { + // selected column should not contribute to view deletion, itself is already included in view row + if (view.getViewColumn(cell.column()) != null) + continue; + // unselected column is used regardless live or dead, because we don't know if it was used for liveness. + timestamp = Math.max(timestamp, cell.maxTimestamp()); + } + return timestamp; } - return timestamp; + // has base non-pk column in view pk + Cell before = existingBaseRow.getCell(view.baseNonPKColumnsInViewPK.get(0)); + assert isLive(before) : "We shouldn't have got there if the base row had no associated entry"; + return deletion.deletes(before) ? deletion.markedForDeleteAt() : before.timestamp(); } private void addColumnData(ColumnDefinition viewColumn, ColumnData baseTableData) diff --git a/src/java/org/apache/cassandra/service/DataResolver.java b/src/java/org/apache/cassandra/service/DataResolver.java index 72c49509f8..c59d688226 100644 --- a/src/java/org/apache/cassandra/service/DataResolver.java +++ b/src/java/org/apache/cassandra/service/DataResolver.java @@ -86,7 +86,9 @@ public class DataResolver extends ResponseResolver DataLimits.Counter counter = command.limits().newCounter(command.nowInSec(), true, command.selectsFullPartition()); UnfilteredPartitionIterator merged = mergeWithShortReadProtection(iters, sources, counter); - FilteredPartitions filtered = FilteredPartitions.filter(merged, new Filter(command.nowInSec())); + FilteredPartitions filtered = FilteredPartitions.filter(merged, + new Filter(command.nowInSec(), + command.metadata().enforceStrictLiveness())); PartitionIterator counted = counter.applyTo(filtered); return command.isForThrift() diff --git a/test/unit/org/apache/cassandra/cql3/CQLTester.java b/test/unit/org/apache/cassandra/cql3/CQLTester.java index 40aec88846..3c0cefcb58 100644 --- a/test/unit/org/apache/cassandra/cql3/CQLTester.java +++ b/test/unit/org/apache/cassandra/cql3/CQLTester.java @@ -1001,7 +1001,7 @@ public abstract class CQLTester assert ignoreExtra || expectedRows.size() == actualRows.size(); } - private static List makeRowStrings(UntypedResultSet resultSet) + protected static List makeRowStrings(UntypedResultSet resultSet) { List> rows = new ArrayList<>(); for (UntypedResultSet.Row row : resultSet) diff --git a/test/unit/org/apache/cassandra/cql3/ViewComplexTest.java b/test/unit/org/apache/cassandra/cql3/ViewComplexTest.java new file mode 100644 index 0000000000..9e32620626 --- /dev/null +++ b/test/unit/org/apache/cassandra/cql3/ViewComplexTest.java @@ -0,0 +1,1343 @@ +package org.apache.cassandra.cql3; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertTrue; +import static org.junit.Assert.fail; + +import java.nio.ByteBuffer; +import java.util.ArrayList; +import java.util.Arrays; +import java.util.Comparator; +import java.util.HashMap; +import java.util.Iterator; +import java.util.List; +import java.util.Map; +import java.util.concurrent.TimeUnit; +import java.util.stream.Collectors; + +import org.apache.cassandra.concurrent.SEPExecutor; +import org.apache.cassandra.concurrent.Stage; +import org.apache.cassandra.concurrent.StageManager; +import org.apache.cassandra.db.ColumnFamilyStore; +import org.apache.cassandra.db.Keyspace; +import org.apache.cassandra.db.compaction.CompactionManager; +import org.apache.cassandra.utils.FBUtilities; +import org.junit.After; +import org.junit.Before; +import org.junit.BeforeClass; +import org.junit.Ignore; +import org.junit.Test; + +import com.google.common.base.Objects; + +public class ViewComplexTest 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 updateView(String query, Object... params) throws Throwable + { + updateViewWithFlush(query, false, params); + } + + private void updateViewWithFlush(String query, boolean flush, Object... params) throws Throwable + { + executeNet(protocolVersion, query, params); + while (!(((SEPExecutor) StageManager.getStage(Stage.VIEW_MUTATION)).getPendingTasks() == 0 + && ((SEPExecutor) StageManager.getStage(Stage.VIEW_MUTATION)).getActiveCount() == 0)) + { + Thread.sleep(1); + } + if (flush) + Keyspace.open(keyspace()).flush(); + } + + // for now, unselected column cannot be fully supported, SEE CASSANDRA-13826 + @Ignore + @Test + public void testPartialDeleteUnselectedColumn() throws Throwable + { + boolean flush = true; + execute("USE " + keyspace()); + executeNet(protocolVersion, "USE " + keyspace()); + createTable("CREATE TABLE %s (k int, c int, a int, b int, PRIMARY KEY (k, c))"); + createView("mv", + "CREATE MATERIALIZED VIEW %s AS SELECT k,c FROM %%s WHERE k IS NOT NULL AND c IS NOT NULL PRIMARY KEY (k,c)"); + Keyspace ks = Keyspace.open(keyspace()); + ks.getColumnFamilyStore("mv").disableAutoCompaction(); + + updateView("UPDATE %s USING TIMESTAMP 10 SET b=1 WHERE k=1 AND c=1"); + if (flush) + FBUtilities.waitOnFutures(ks.flush()); + assertRows(execute("SELECT * from %s"), row(1, 1, null, 1)); + assertRows(execute("SELECT * from mv"), row(1, 1)); + updateView("DELETE b FROM %s USING TIMESTAMP 11 WHERE k=1 AND c=1"); + if (flush) + FBUtilities.waitOnFutures(ks.flush()); + assertEmpty(execute("SELECT * from %s")); + assertEmpty(execute("SELECT * from mv")); + updateView("UPDATE %s USING TIMESTAMP 1 SET a=1 WHERE k=1 AND c=1"); + if (flush) + FBUtilities.waitOnFutures(ks.flush()); + assertRows(execute("SELECT * from %s"), row(1, 1, 1, null)); + assertRows(execute("SELECT * from mv"), row(1, 1)); + + execute("truncate %s;"); + + // removal generated by unselected column should not shadow PK update with smaller timestamp + updateViewWithFlush("UPDATE %s USING TIMESTAMP 18 SET a=1 WHERE k=1 AND c=1", flush); + assertRows(execute("SELECT * from %s"), row(1, 1, 1, null)); + assertRows(execute("SELECT * from mv"), row(1, 1)); + + updateViewWithFlush("UPDATE %s USING TIMESTAMP 20 SET a=null WHERE k=1 AND c=1", flush); + assertRows(execute("SELECT * from %s")); + assertRows(execute("SELECT * from mv")); + + updateViewWithFlush("INSERT INTO %s(k,c) VALUES(1,1) USING TIMESTAMP 15", flush); + assertRows(execute("SELECT * from %s"), row(1, 1, null, null)); + assertRows(execute("SELECT * from mv"), row(1, 1)); + } + + @Test + public void testPartialDeleteSelectedColumnWithFlush() throws Throwable + { + testPartialDeleteSelectedColumn(true); + } + + @Test + public void testPartialDeleteSelectedColumnWithoutFlush() throws Throwable + { + testPartialDeleteSelectedColumn(false); + } + + private void testPartialDeleteSelectedColumn(boolean flush) throws Throwable + { + execute("USE " + keyspace()); + executeNet(protocolVersion, "USE " + keyspace()); + createTable("CREATE TABLE %s (k int, c int, a int, b int, e int, f int, PRIMARY KEY (k, c))"); + createView("mv", + "CREATE MATERIALIZED VIEW %s AS SELECT a, b FROM %%s WHERE k IS NOT NULL AND c IS NOT NULL PRIMARY KEY (k,c)"); + Keyspace ks = Keyspace.open(keyspace()); + ks.getColumnFamilyStore("mv").disableAutoCompaction(); + + updateViewWithFlush("UPDATE %s USING TIMESTAMP 10 SET b=1 WHERE k=1 AND c=1", flush); + assertRows(execute("SELECT * from %s"), row(1, 1, null, 1, null, null)); + assertRows(execute("SELECT * from mv"), row(1, 1, null, 1)); + + updateViewWithFlush("DELETE b FROM %s USING TIMESTAMP 11 WHERE k=1 AND c=1", flush); + assertEmpty(execute("SELECT * from %s")); + assertEmpty(execute("SELECT * from mv")); + + updateViewWithFlush("UPDATE %s USING TIMESTAMP 1 SET a=1 WHERE k=1 AND c=1", flush); + assertRows(execute("SELECT * from %s"), row(1, 1, 1, null, null, null)); + assertRows(execute("SELECT * from mv"), row(1, 1, 1, null)); + + updateViewWithFlush("DELETE a FROM %s USING TIMESTAMP 1 WHERE k=1 AND c=1", flush); + assertEmpty(execute("SELECT * from %s")); + assertEmpty(execute("SELECT * from mv")); + + // view livenessInfo should not be affected by selected column ts or tb + updateViewWithFlush("INSERT INTO %s(k,c) VALUES(1,1) USING TIMESTAMP 0", flush); + assertRows(execute("SELECT * from %s"), row(1, 1, null, null, null, null)); + assertRows(execute("SELECT * from mv"), row(1, 1, null, null)); + + updateViewWithFlush("UPDATE %s USING TIMESTAMP 12 SET b=1 WHERE k=1 AND c=1", flush); + assertRows(execute("SELECT * from %s"), row(1, 1, null, 1, null, null)); + assertRows(execute("SELECT * from mv"), row(1, 1, null, 1)); + + updateViewWithFlush("DELETE b FROM %s USING TIMESTAMP 13 WHERE k=1 AND c=1", flush); + assertRows(execute("SELECT * from %s"), row(1, 1, null, null, null, null)); + assertRows(execute("SELECT * from mv"), row(1, 1, null, null)); + + updateViewWithFlush("DELETE FROM %s USING TIMESTAMP 14 WHERE k=1 AND c=1", flush); + assertEmpty(execute("SELECT * from %s")); + assertEmpty(execute("SELECT * from mv")); + + updateViewWithFlush("INSERT INTO %s(k,c) VALUES(1,1) USING TIMESTAMP 15", flush); + assertRows(execute("SELECT * from %s"), row(1, 1, null, null, null, null)); + assertRows(execute("SELECT * from mv"), row(1, 1, null, null)); + + updateViewWithFlush("UPDATE %s USING TTL 3 SET b=1 WHERE k=1 AND c=1", flush); + assertRows(execute("SELECT * from %s"), row(1, 1, null, 1, null, null)); + assertRows(execute("SELECT * from mv"), row(1, 1, null, 1)); + + TimeUnit.SECONDS.sleep(4); + + assertRows(execute("SELECT * from %s"), row(1, 1, null, null, null, null)); + assertRows(execute("SELECT * from mv"), row(1, 1, null, null)); + + updateViewWithFlush("DELETE FROM %s USING TIMESTAMP 15 WHERE k=1 AND c=1", flush); + assertEmpty(execute("SELECT * from %s")); + assertEmpty(execute("SELECT * from mv")); + + execute("truncate %s;"); + + // removal generated by unselected column should not shadow selected column with smaller timestamp + updateViewWithFlush("UPDATE %s USING TIMESTAMP 18 SET e=1 WHERE k=1 AND c=1", flush); + assertRows(execute("SELECT * from %s"), row(1, 1, null, null, 1, null)); + assertRows(execute("SELECT * from mv"), row(1, 1, null, null)); + + updateViewWithFlush("UPDATE %s USING TIMESTAMP 18 SET e=null WHERE k=1 AND c=1", flush); + assertRows(execute("SELECT * from %s")); + assertRows(execute("SELECT * from mv")); + + updateViewWithFlush("UPDATE %s USING TIMESTAMP 16 SET a=1 WHERE k=1 AND c=1", flush); + assertRows(execute("SELECT * from %s"), row(1, 1, 1, null, null, null)); + assertRows(execute("SELECT * from mv"), row(1, 1, 1, null)); + } + + @Test + public void testUpdateColumnInViewPKWithTTLWithFlush() throws Throwable + { + // CASSANDRA-13657 + testUpdateColumnInViewPKWithTTL(true); + } + + @Test + public void testUpdateColumnInViewPKWithTTLWithoutFlush() throws Throwable + { + // CASSANDRA-13657 + testUpdateColumnInViewPKWithTTL(false); + } + + private void testUpdateColumnInViewPKWithTTL(boolean flush) throws Throwable + { + // CASSANDRA-13657 if base column used in view pk is ttled, then view row is considered dead + createTable("create table %s (k int primary key, a int, b int)"); + + execute("USE " + keyspace()); + executeNet(protocolVersion, "USE " + keyspace()); + Keyspace ks = Keyspace.open(keyspace()); + + createView("mv", + "CREATE MATERIALIZED VIEW %s AS SELECT * FROM %%s WHERE k IS NOT NULL AND a IS NOT NULL PRIMARY KEY (a, k)"); + ks.getColumnFamilyStore("mv").disableAutoCompaction(); + + updateView("UPDATE %s SET a = 1 WHERE k = 1;"); + + if (flush) + FBUtilities.waitOnFutures(ks.flush()); + + assertRows(execute("SELECT * from %s"), row(1, 1, null)); + assertRows(execute("SELECT * from mv"), row(1, 1, null)); + + updateView("DELETE a FROM %s WHERE k = 1"); + + if (flush) + FBUtilities.waitOnFutures(ks.flush()); + + assertRows(execute("SELECT * from %s")); + assertEmpty(execute("SELECT * from mv")); + + updateView("INSERT INTO %s (k) VALUES (1);"); + + if (flush) + FBUtilities.waitOnFutures(ks.flush()); + + assertRows(execute("SELECT * from %s"), row(1, null, null)); + assertEmpty(execute("SELECT * from mv")); + + updateView("UPDATE %s USING TTL 5 SET a = 10 WHERE k = 1;"); + + if (flush) + FBUtilities.waitOnFutures(ks.flush()); + + assertRows(execute("SELECT * from %s"), row(1, 10, null)); + assertRows(execute("SELECT * from mv"), row(10, 1, null)); + + updateView("UPDATE %s SET b = 100 WHERE k = 1;"); + + if (flush) + FBUtilities.waitOnFutures(ks.flush()); + + assertRows(execute("SELECT * from %s"), row(1, 10, 100)); + assertRows(execute("SELECT * from mv"), row(10, 1, 100)); + + Thread.sleep(5000); + + // 'a' is TTL of 5 and removed. + assertRows(execute("SELECT * from %s"), row(1, null, 100)); + assertEmpty(execute("SELECT * from mv")); + assertEmpty(execute("SELECT * from mv WHERE k = ? AND a = ?", 1, 10)); + + updateView("DELETE b FROM %s WHERE k=1"); + + if (flush) + FBUtilities.waitOnFutures(ks.flush()); + + assertRows(execute("SELECT * from %s"), row(1, null, null)); + assertEmpty(execute("SELECT * from mv")); + + updateView("DELETE FROM %s WHERE k=1;"); + + if (flush) + FBUtilities.waitOnFutures(ks.flush()); + + assertEmpty(execute("SELECT * from %s")); + assertEmpty(execute("SELECT * from mv")); + } + + @Test + public void testUpdateColumnNotInViewWithFlush() throws Throwable + { + testUpdateColumnNotInView(true); + } + + @Test + public void testUpdateColumnNotInViewWithoutFlush() throws Throwable + { + // CASSANDRA-13127 + testUpdateColumnNotInView(false); + } + + private void testUpdateColumnNotInView(boolean flush) throws Throwable + { + // CASSANDRA-13127: if base column not selected in view are alive, then pk of view row should be alive + createTable("create table %s (p int, c int, v1 int, v2 int, primary key(p, c))"); + + execute("USE " + keyspace()); + executeNet(protocolVersion, "USE " + keyspace()); + Keyspace ks = Keyspace.open(keyspace()); + + createView("mv", + "CREATE MATERIALIZED VIEW %s AS SELECT p, c FROM %%s WHERE p IS NOT NULL AND c IS NOT NULL PRIMARY KEY (c, p);"); + ks.getColumnFamilyStore("mv").disableAutoCompaction(); + + updateView("UPDATE %s USING TIMESTAMP 0 SET v1 = 1 WHERE p = 0 AND c = 0"); + + if (flush) + FBUtilities.waitOnFutures(ks.flush()); + + assertRowsIgnoringOrder(execute("SELECT * from %s WHERE c = ? AND p = ?", 0, 0), row(0, 0, 1, null)); + assertRowsIgnoringOrder(execute("SELECT * from mv WHERE c = ? AND p = ?", 0, 0), row(0, 0)); + + updateView("DELETE v1 FROM %s USING TIMESTAMP 1 WHERE p = 0 AND c = 0"); + + if (flush) + FBUtilities.waitOnFutures(ks.flush()); + + assertEmpty(execute("SELECT * from %s WHERE c = ? AND p = ?", 0, 0)); + assertEmpty(execute("SELECT * from mv WHERE c = ? AND p = ?", 0, 0)); + + // shadowed by tombstone + updateView("UPDATE %s USING TIMESTAMP 1 SET v1 = 1 WHERE p = 0 AND c = 0"); + + if (flush) + FBUtilities.waitOnFutures(ks.flush()); + + assertEmpty(execute("SELECT * from %s WHERE c = ? AND p = ?", 0, 0)); + assertEmpty(execute("SELECT * from mv WHERE c = ? AND p = ?", 0, 0)); + + updateView("UPDATE %s USING TIMESTAMP 2 SET v2 = 1 WHERE p = 0 AND c = 0"); + + if (flush) + FBUtilities.waitOnFutures(ks.flush()); + + assertRowsIgnoringOrder(execute("SELECT * from %s WHERE c = ? AND p = ?", 0, 0), row(0, 0, null, 1)); + assertRowsIgnoringOrder(execute("SELECT * from mv WHERE c = ? AND p = ?", 0, 0), row(0, 0)); + + updateView("DELETE v1 FROM %s USING TIMESTAMP 3 WHERE p = 0 AND c = 0"); + + if (flush) + FBUtilities.waitOnFutures(ks.flush()); + + assertRowsIgnoringOrder(execute("SELECT * from %s WHERE c = ? AND p = ?", 0, 0), row(0, 0, null, 1)); + assertRowsIgnoringOrder(execute("SELECT * from mv WHERE c = ? AND p = ?", 0, 0), row(0, 0)); + + updateView("DELETE v2 FROM %s USING TIMESTAMP 4 WHERE p = 0 AND c = 0"); + + if (flush) + FBUtilities.waitOnFutures(ks.flush()); + + assertEmpty(execute("SELECT * from %s WHERE c = ? AND p = ?", 0, 0)); + assertEmpty(execute("SELECT * from mv WHERE c = ? AND p = ?", 0, 0)); + + updateView("UPDATE %s USING TTL 3 SET v2 = 1 WHERE p = 0 AND c = 0"); + + if (flush) + FBUtilities.waitOnFutures(ks.flush()); + + assertRowsIgnoringOrder(execute("SELECT * from %s WHERE c = ? AND p = ?", 0, 0), row(0, 0, null, 1)); + assertRowsIgnoringOrder(execute("SELECT * from mv WHERE c = ? AND p = ?", 0, 0), row(0, 0)); + + Thread.sleep(TimeUnit.SECONDS.toMillis(3)); + + assertRowsIgnoringOrder(execute("SELECT * from %s WHERE c = ? AND p = ?", 0, 0)); + assertRowsIgnoringOrder(execute("SELECT * from mv WHERE c = ? AND p = ?", 0, 0)); + + updateView("UPDATE %s SET v2 = 1 WHERE p = 0 AND c = 0"); + + if (flush) + FBUtilities.waitOnFutures(ks.flush()); + + assertRowsIgnoringOrder(execute("SELECT * from %s WHERE c = ? AND p = ?", 0, 0), row(0, 0, null, 1)); + assertRowsIgnoringOrder(execute("SELECT * from mv WHERE c = ? AND p = ?", 0, 0), row(0, 0)); + + assertInvalidMessage("Cannot drop column v2 on base table with materialized views", "ALTER TABLE %s DROP v2"); + // // drop unselected base column, unselected metadata should be removed, thus view row is dead + // updateView("ALTER TABLE %s DROP v2"); + // assertRowsIgnoringOrder(execute("SELECT * from %s WHERE c = ? AND p = ?", 0, 0)); + // assertRowsIgnoringOrder(execute("SELECT * from mv WHERE c = ? AND p = ?", 0, 0)); + // assertRowsIgnoringOrder(execute("SELECT * from %s")); + // assertRowsIgnoringOrder(execute("SELECT * from mv")); + } + + @Test + public void testPartialUpdateWithUnselectedCollectionsWithFlush() throws Throwable + { + testPartialUpdateWithUnselectedCollections(true); + } + + @Test + public void testPartialUpdateWithUnselectedCollectionsWithoutFlush() throws Throwable + { + testPartialUpdateWithUnselectedCollections(false); + } + + public void testPartialUpdateWithUnselectedCollections(boolean flush) throws Throwable + { + execute("USE " + keyspace()); + executeNet(protocolVersion, "USE " + keyspace()); + createTable("CREATE TABLE %s (k int, c int, a int, b int, l list, s set, m map, PRIMARY KEY (k, c))"); + createView("mv", + "CREATE MATERIALIZED VIEW %s AS SELECT a, b FROM %%s WHERE k IS NOT NULL AND c IS NOT NULL PRIMARY KEY (c, k)"); + Keyspace ks = Keyspace.open(keyspace()); + ks.getColumnFamilyStore("mv").disableAutoCompaction(); + + updateView("UPDATE %s SET l=l+[1,2,3] WHERE k = 1 AND c = 1"); + if (flush) + FBUtilities.waitOnFutures(ks.flush()); + assertRows(execute("SELECT * from mv"), row(1, 1, null, null)); + + updateView("UPDATE %s SET l=l-[1,2] WHERE k = 1 AND c = 1"); + if (flush) + FBUtilities.waitOnFutures(ks.flush()); + assertRows(execute("SELECT * from mv"), row(1, 1, null, null)); + + updateView("UPDATE %s SET b=3 WHERE k=1 AND c=1"); + if (flush) + FBUtilities.waitOnFutures(ks.flush()); + assertRows(execute("SELECT * from mv"), row(1, 1, null, 3)); + + updateView("UPDATE %s SET b=null, l=l-[3], s=s-{3} WHERE k = 1 AND c = 1"); + if (flush) + { + FBUtilities.waitOnFutures(ks.flush()); + ks.getColumnFamilyStore("mv").forceMajorCompaction(); + } + assertRowsIgnoringOrder(execute("SELECT k,c,a,b from %s")); + assertRowsIgnoringOrder(execute("SELECT * from mv")); + + updateView("UPDATE %s SET m=m+{3:3}, l=l-[1], s=s-{2} WHERE k = 1 AND c = 1"); + if (flush) + FBUtilities.waitOnFutures(ks.flush()); + assertRowsIgnoringOrder(execute("SELECT k,c,a,b from %s"), row(1, 1, null, null)); + assertRowsIgnoringOrder(execute("SELECT * from mv"), row(1, 1, null, null)); + + assertInvalidMessage("Cannot drop column m on base table with materialized views", "ALTER TABLE %s DROP m"); + // executeNet(protocolVersion, "ALTER TABLE %s DROP m"); + // ks.getColumnFamilyStore("mv").forceMajorCompaction(); + // assertRowsIgnoringOrder(execute("SELECT k,c,a,b from %s WHERE k = 1 AND c = 1")); + // assertRowsIgnoringOrder(execute("SELECT * from mv WHERE k = 1 AND c = 1")); + // assertRowsIgnoringOrder(execute("SELECT k,c,a,b from %s")); + // assertRowsIgnoringOrder(execute("SELECT * from mv")); + } + + @Test + public void testUnselectedColumnsTTLWithFlush() throws Throwable + { + // CASSANDRA-13127 + testUnselectedColumnsTTL(true); + } + + @Test + public void testUnselectedColumnsTTLWithoutFlush() throws Throwable + { + // CASSANDRA-13127 + testUnselectedColumnsTTL(false); + } + + private void testUnselectedColumnsTTL(boolean flush) throws Throwable + { + // CASSANDRA-13127 not ttled unselected column in base should keep view row alive + createTable("create table %s (p int, c int, v int, primary key(p, c))"); + + execute("USE " + keyspace()); + executeNet(protocolVersion, "USE " + keyspace()); + Keyspace ks = Keyspace.open(keyspace()); + + createView("mv", + "CREATE MATERIALIZED VIEW %s AS SELECT p, c FROM %%s WHERE p IS NOT NULL AND c IS NOT NULL PRIMARY KEY (c, p);"); + ks.getColumnFamilyStore("mv").disableAutoCompaction(); + + updateViewWithFlush("INSERT INTO %s (p, c) VALUES (0, 0) USING TTL 3;", flush); + + updateViewWithFlush("UPDATE %s USING TTL 1000 SET v = 0 WHERE p = 0 and c = 0;", flush); + + assertRowsIgnoringOrder(execute("SELECT * from mv WHERE c = ? AND p = ?", 0, 0), row(0, 0)); + + Thread.sleep(3000); + + UntypedResultSet.Row row = execute("SELECT v, ttl(v) from %s WHERE c = ? AND p = ?", 0, 0).one(); + assertTrue("row should have value of 0", row.getInt("v") == 0); + assertTrue("row should have ttl less than 1000", row.getInt("ttl(v)") < 1000); + assertRowsIgnoringOrder(execute("SELECT * from mv WHERE c = ? AND p = ?", 0, 0), row(0, 0)); + + updateViewWithFlush("DELETE FROM %s WHERE p = 0 and c = 0;", flush); + assertRowsIgnoringOrder(execute("SELECT * from mv WHERE c = ? AND p = ?", 0, 0)); + + updateViewWithFlush("INSERT INTO %s (p, c) VALUES (0, 0) ", flush); + assertRowsIgnoringOrder(execute("SELECT * from mv WHERE c = ? AND p = ?", 0, 0), row(0, 0)); + + // already have a live row, no need to apply the unselected cell ttl + updateViewWithFlush("UPDATE %s USING TTL 3 SET v = 0 WHERE p = 0 and c = 0;", flush); + assertRowsIgnoringOrder(execute("SELECT * from mv WHERE c = ? AND p = ?", 0, 0), row(0, 0)); + + updateViewWithFlush("INSERT INTO %s (p, c) VALUES (1, 1) USING TTL 3", flush); + assertRowsIgnoringOrder(execute("SELECT * from mv WHERE c = ? AND p = ?", 1, 1), row(1, 1)); + + Thread.sleep(4000); + + assertRowsIgnoringOrder(execute("SELECT * from mv WHERE c = ? AND p = ?", 0, 0), row(0, 0)); + assertRowsIgnoringOrder(execute("SELECT * from mv WHERE c = ? AND p = ?", 1, 1)); + + // unselected should keep view row alive + updateViewWithFlush("UPDATE %s SET v = 0 WHERE p = 1 and c = 1;", flush); + assertRowsIgnoringOrder(execute("SELECT * from mv WHERE c = ? AND p = ?", 1, 1), row(1, 1)); + + } + + @Test + public void testRangeDeletionWithFlush() throws Throwable + { + testRangeDeletion(true); + } + + @Test + public void testRangeDeletionWithoutFlush() throws Throwable + { + testRangeDeletion(false); + } + + public void testRangeDeletion(boolean flush) throws Throwable + { + // for partition range deletion, need to know that existing row is shadowed instead of not existed. + createTable("CREATE TABLE %s (a int, b int, c int, d int, PRIMARY KEY (a))"); + + execute("USE " + keyspace()); + executeNet(protocolVersion, "USE " + keyspace()); + + createView("mv_test1", + "CREATE MATERIALIZED VIEW %s AS SELECT * FROM %%s WHERE a IS NOT NULL AND b IS NOT NULL PRIMARY KEY (a, b)"); + + Keyspace ks = Keyspace.open(keyspace()); + ks.getColumnFamilyStore("mv_test1").disableAutoCompaction(); + + execute("INSERT INTO %s (a, b, c, d) VALUES (?, ?, ?, ?) using timestamp 0", 1, 1, 1, 1); + if (flush) + FBUtilities.waitOnFutures(ks.flush()); + + assertRowsIgnoringOrder(execute("SELECT * FROM mv_test1"), row(1, 1, 1, 1)); + + // remove view row + updateView("UPDATE %s using timestamp 1 set b = null WHERE a=1"); + if (flush) + FBUtilities.waitOnFutures(ks.flush()); + + assertRowsIgnoringOrder(execute("SELECT * FROM mv_test1")); + // remove base row, no view updated generated. + updateView("DELETE FROM %s using timestamp 2 where a=1"); + if (flush) + FBUtilities.waitOnFutures(ks.flush()); + + assertRowsIgnoringOrder(execute("SELECT * FROM mv_test1")); + + // restor view row with b,c column. d is still tombstone + updateView("UPDATE %s using timestamp 3 set b = 1,c = 1 where a=1"); // upsert + if (flush) + FBUtilities.waitOnFutures(ks.flush()); + + assertRowsIgnoringOrder(execute("SELECT * FROM mv_test1"), row(1, 1, 1, null)); + } + + @Test + public void testBaseTTLWithSameTimestampTest() throws Throwable + { + // CASSANDRA-13127 when liveness timestamp tie, greater localDeletionTime should win if both are expiring. + createTable("create table %s (p int, c int, v int, primary key(p, c))"); + + execute("USE " + keyspace()); + executeNet(protocolVersion, "USE " + keyspace()); + Keyspace ks = Keyspace.open(keyspace()); + + updateView("INSERT INTO %s (p, c, v) VALUES (0, 0, 0) using timestamp 1;"); + + FBUtilities.waitOnFutures(ks.flush()); + + updateView("INSERT INTO %s (p, c, v) VALUES (0, 0, 0) USING TTL 3 and timestamp 1;"); + + FBUtilities.waitOnFutures(ks.flush()); + + Thread.sleep(4000); + + assertEmpty(execute("SELECT * from %s WHERE c = ? AND p = ?", 0, 0)); + + // reversed order + execute("truncate %s;"); + + updateView("INSERT INTO %s (p, c, v) VALUES (0, 0, 0) USING TTL 3 and timestamp 1;"); + + FBUtilities.waitOnFutures(ks.flush()); + + updateView("INSERT INTO %s (p, c, v) VALUES (0, 0, 0) USING timestamp 1;"); + + FBUtilities.waitOnFutures(ks.flush()); + + Thread.sleep(4000); + + assertEmpty(execute("SELECT * from %s WHERE c = ? AND p = ?", 0, 0)); + + } + + @Test + public void testCommutativeRowDeletionFlush() throws Throwable + { + // CASSANDRA-13409 + testCommutativeRowDeletion(true); + } + + @Test + public void testCommutativeRowDeletionWithoutFlush() throws Throwable + { + // CASSANDRA-13409 + testCommutativeRowDeletion(false); + } + + private void testCommutativeRowDeletion(boolean flush) throws Throwable + { + // CASSANDRA-13409 new update should not resurrect previous deleted data in view + createTable("create table %s (p int primary key, v1 int, v2 int)"); + + execute("USE " + keyspace()); + executeNet(protocolVersion, "USE " + keyspace()); + Keyspace ks = Keyspace.open(keyspace()); + + createView("mv", + "create materialized view %s as select * from %%s where p is not null and v1 is not null primary key (v1, p);"); + ks.getColumnFamilyStore("mv").disableAutoCompaction(); + + // sstable-1, Set initial values TS=1 + updateView("Insert into %s (p, v1, v2) values (3, 1, 3) using timestamp 1;"); + + if (flush) + FBUtilities.waitOnFutures(ks.flush()); + + assertRowsIgnoringOrder(execute("SELECT v2, WRITETIME(v2) from mv WHERE v1 = ? AND p = ?", 1, 3), row(3, 1L)); + // sstable-2 + updateView("Delete from %s using timestamp 2 where p = 3;"); + + if (flush) + FBUtilities.waitOnFutures(ks.flush()); + + assertRowsIgnoringOrder(execute("SELECT v1, p, v2, WRITETIME(v2) from mv")); + // sstable-3 + updateView("Insert into %s (p, v1) values (3, 1) using timestamp 3;"); + + if (flush) + FBUtilities.waitOnFutures(ks.flush()); + + assertRowsIgnoringOrder(execute("SELECT v1, p, v2, WRITETIME(v2) from mv"), row(1, 3, null, null)); + // sstable-4 + updateView("UPdate %s using timestamp 4 set v1 = 2 where p = 3;"); + + if (flush) + FBUtilities.waitOnFutures(ks.flush()); + + assertRowsIgnoringOrder(execute("SELECT v1, p, v2, WRITETIME(v2) from mv"), row(2, 3, null, null)); + // sstable-5 + updateView("UPdate %s using timestamp 5 set v1 = 1 where p = 3;"); + + if (flush) + FBUtilities.waitOnFutures(ks.flush()); + + assertRowsIgnoringOrder(execute("SELECT v1, p, v2, WRITETIME(v2) from mv"), row(1, 3, null, null)); + + if (flush) + { + // compact sstable 2 and 4, 5; + ColumnFamilyStore cfs = ks.getColumnFamilyStore("mv"); + List sstables = cfs.getLiveSSTables() + .stream() + .sorted((s1, s2) -> s1.descriptor.generation - s2.descriptor.generation) + .map(s -> s.getFilename()) + .collect(Collectors.toList()); + String dataFiles = String.join(",", Arrays.asList(sstables.get(1), sstables.get(3), sstables.get(4))); + CompactionManager.instance.forceUserDefinedCompaction(dataFiles); + assertEquals(3, cfs.getLiveSSTables().size()); + } + // regular tombstone should be retained after compaction + assertRowsIgnoringOrder(execute("SELECT v1, p, v2, WRITETIME(v2) from mv"), row(1, 3, null, null)); + } + + @Test + public void testUnselectedColumnWithExpiredLivenessInfo() throws Throwable + { + boolean flush = true; + createTable("create table %s (k int, c int, a int, b int, PRIMARY KEY(k, c))"); + + execute("USE " + keyspace()); + executeNet(protocolVersion, "USE " + keyspace()); + Keyspace ks = Keyspace.open(keyspace()); + + createView("mv", + "create materialized view %s as select k,c,b from %%s where c is not null and k is not null primary key (c, k);"); + ks.getColumnFamilyStore("mv").disableAutoCompaction(); + + // sstable-1, Set initial values TS=1 + updateViewWithFlush("UPDATE %s SET a = 1 WHERE k = 1 AND c = 1;", flush); + + assertRowsIgnoringOrder(execute("SELECT * from %s WHERE k = 1 AND c = 1;"), + row(1, 1, 1, null)); + assertRowsIgnoringOrder(execute("SELECT k,c,b from mv WHERE k = 1 AND c = 1;"), + row(1, 1, null)); + + // sstable-2 + updateViewWithFlush("INSERT INTO %s(k,c) VALUES(1,1) USING TTL 5", flush); + + assertRowsIgnoringOrder(execute("SELECT * from %s WHERE k = 1 AND c = 1;"), + row(1, 1, 1, null)); + assertRowsIgnoringOrder(execute("SELECT k,c,b from mv WHERE k = 1 AND c = 1;"), + row(1, 1, null)); + + Thread.sleep(5001); + + assertRowsIgnoringOrder(execute("SELECT * from %s WHERE k = 1 AND c = 1;"), + row(1, 1, 1, null)); + assertRowsIgnoringOrder(execute("SELECT k,c,b from mv WHERE k = 1 AND c = 1;"), + row(1, 1, null)); + + // sstable-3 + updateViewWithFlush("Update %s set a = null where k = 1 AND c = 1;", flush); + + assertRowsIgnoringOrder(execute("SELECT * from %s WHERE k = 1 AND c = 1;")); + assertRowsIgnoringOrder(execute("SELECT k,c,b from mv WHERE k = 1 AND c = 1;")); + + // sstable-4 + updateViewWithFlush("Update %s USING TIMESTAMP 1 set b = 1 where k = 1 AND c = 1;", flush); + + assertRowsIgnoringOrder(execute("SELECT * from %s WHERE k = 1 AND c = 1;"), + row(1, 1, null, 1)); + assertRowsIgnoringOrder(execute("SELECT k,c,b from mv WHERE k = 1 AND c = 1;"), + row(1, 1, 1)); + } + + @Test + public void testUpdateWithColumnTimestampSmallerThanPkWithFlush() throws Throwable + { + testUpdateWithColumnTimestampSmallerThanPk(true); + } + + @Test + public void testUpdateWithColumnTimestampSmallerThanPkWithoutFlush() throws Throwable + { + testUpdateWithColumnTimestampSmallerThanPk(false); + } + + public void testUpdateWithColumnTimestampSmallerThanPk(boolean flush) throws Throwable + { + createTable("create table %s (p int primary key, v1 int, v2 int)"); + + execute("USE " + keyspace()); + executeNet(protocolVersion, "USE " + keyspace()); + Keyspace ks = Keyspace.open(keyspace()); + + createView("mv", + "create materialized view %s as select * from %%s where p is not null and v1 is not null primary key (v1, p);"); + ks.getColumnFamilyStore("mv").disableAutoCompaction(); + + // reset value + updateView("Insert into %s (p, v1, v2) values (3, 1, 3) using timestamp 6;"); + if (flush) + FBUtilities.waitOnFutures(ks.flush()); + assertRowsIgnoringOrder(execute("SELECT v1, p, v2, WRITETIME(v2) from mv"), row(1, 3, 3, 6L)); + // increase pk's timestamp to 20 + updateView("Insert into %s (p) values (3) using timestamp 20;"); + if (flush) + FBUtilities.waitOnFutures(ks.flush()); + assertRowsIgnoringOrder(execute("SELECT v1, p, v2, WRITETIME(v2) from mv"), row(1, 3, 3, 6L)); + // change v1's to 2 and remove existing view row with ts7 + updateView("UPdate %s using timestamp 7 set v1 = 2 where p = 3;"); + if (flush) + FBUtilities.waitOnFutures(ks.flush()); + assertRowsIgnoringOrder(execute("SELECT v1, p, v2, WRITETIME(v2) from mv"), row(2, 3, 3, 6L)); + // change v1's to 1 and remove existing view row with ts8 + updateView("UPdate %s using timestamp 8 set v1 = 1 where p = 3;"); + if (flush) + FBUtilities.waitOnFutures(ks.flush()); + assertRowsIgnoringOrder(execute("SELECT v1, p, v2, WRITETIME(v2) from mv"), row(1, 3, 3, 6L)); + } + + @Test + public void testUpdateWithColumnTimestampBiggerThanPkWithFlush() throws Throwable + { + // CASSANDRA-11500 + testUpdateWithColumnTimestampBiggerThanPk(true); + } + + @Test + public void testUpdateWithColumnTimestampBiggerThanPkWithoutFlush() throws Throwable + { + // CASSANDRA-11500 + testUpdateWithColumnTimestampBiggerThanPk(false); + } + + public void testUpdateWithColumnTimestampBiggerThanPk(boolean flush) throws Throwable + { + // CASSANDRA-11500 able to shadow old view row with column ts greater tahn pk's ts and re-insert the view row + createTable("CREATE TABLE %s (k int PRIMARY KEY, a int, b int);"); + + execute("USE " + keyspace()); + executeNet(protocolVersion, "USE " + keyspace()); + Keyspace ks = Keyspace.open(keyspace()); + + createView("mv", + "CREATE MATERIALIZED VIEW %s AS SELECT * FROM %%s WHERE k IS NOT NULL AND a IS NOT NULL PRIMARY KEY (k, a);"); + ks.getColumnFamilyStore("mv").disableAutoCompaction(); + updateView("DELETE FROM %s USING TIMESTAMP 0 WHERE k = 1;"); + if (flush) + FBUtilities.waitOnFutures(ks.flush()); + // sstable-1, Set initial values TS=1 + updateView("INSERT INTO %s(k, a, b) VALUES (1, 1, 1) USING TIMESTAMP 1;"); + if (flush) + FBUtilities.waitOnFutures(ks.flush()); + assertRowsIgnoringOrder(execute("SELECT k,a,b from mv"), row(1, 1, 1)); + updateView("UPDATE %s USING TIMESTAMP 10 SET b = 2 WHERE k = 1;"); + assertRowsIgnoringOrder(execute("SELECT k,a,b from mv"), row(1, 1, 2)); + if (flush) + FBUtilities.waitOnFutures(ks.flush()); + assertRowsIgnoringOrder(execute("SELECT k,a,b from mv"), row(1, 1, 2)); + updateView("UPDATE %s USING TIMESTAMP 2 SET a = 2 WHERE k = 1;"); + assertRowsIgnoringOrder(execute("SELECT k,a,b from mv"), row(1, 2, 2)); + if (flush) + FBUtilities.waitOnFutures(ks.flush()); + ks.getColumnFamilyStore("mv").forceMajorCompaction(); + assertRowsIgnoringOrder(execute("SELECT k,a,b from mv"), row(1, 2, 2)); + updateView("UPDATE %s USING TIMESTAMP 11 SET a = 1 WHERE k = 1;"); + if (flush) + FBUtilities.waitOnFutures(ks.flush()); + assertRowsIgnoringOrder(execute("SELECT k,a,b from mv"), row(1, 1, 2)); + assertRowsIgnoringOrder(execute("SELECT k,a,b from %s"), row(1, 1, 2)); + + // set non-key base column as tombstone, view row is removed with shadowable + updateView("UPDATE %s USING TIMESTAMP 12 SET a = null WHERE k = 1;"); + if (flush) + FBUtilities.waitOnFutures(ks.flush()); + assertRowsIgnoringOrder(execute("SELECT k,a,b from mv")); + assertRowsIgnoringOrder(execute("SELECT k,a,b from %s"), row(1, null, 2)); + + // column b should be alive + updateView("UPDATE %s USING TIMESTAMP 13 SET a = 1 WHERE k = 1;"); + if (flush) + FBUtilities.waitOnFutures(ks.flush()); + assertRowsIgnoringOrder(execute("SELECT k,a,b from mv"), row(1, 1, 2)); + assertRowsIgnoringOrder(execute("SELECT k,a,b from %s"), row(1, 1, 2)); + + assertInvalidMessage("Cannot drop column a on base table with materialized views", "ALTER TABLE %s DROP a"); + } + + @Test + public void testNonBaseColumnInViewPkWithFlush() throws Throwable + { + testNonBaseColumnInViewPk(true); + } + + @Test + public void testNonBaseColumnInViewPkWithoutFlush() throws Throwable + { + testNonBaseColumnInViewPk(true); + } + + public void testNonBaseColumnInViewPk(boolean flush) throws Throwable + { + createTable("create table %s (p1 int, p2 int, v1 int, v2 int, primary key (p1,p2))"); + + execute("USE " + keyspace()); + executeNet(protocolVersion, "USE " + keyspace()); + Keyspace ks = Keyspace.open(keyspace()); + + createView("mv", + "create materialized view %s as select * from %%s where p1 is not null and p2 is not null primary key (p2, p1)" + + " with gc_grace_seconds=5;"); + ColumnFamilyStore cfs = ks.getColumnFamilyStore("mv"); + cfs.disableAutoCompaction(); + + updateView("UPDATE %s USING TIMESTAMP 1 set v1 =1 where p1 = 1 AND p2 = 1;"); + if (flush) + FBUtilities.waitOnFutures(ks.flush()); + assertRowsIgnoringOrder(execute("SELECT p1, p2, v1, v2 from %s"), row(1, 1, 1, null)); + assertRowsIgnoringOrder(execute("SELECT p1, p2, v1, v2 from mv"), row(1, 1, 1, null)); + + updateView("UPDATE %s USING TIMESTAMP 2 set v1 = null, v2 = 1 where p1 = 1 AND p2 = 1;"); + if (flush) + FBUtilities.waitOnFutures(ks.flush()); + assertRowsIgnoringOrder(execute("SELECT p1, p2, v1, v2 from %s"), row(1, 1, null, 1)); + assertRowsIgnoringOrder(execute("SELECT p1, p2, v1, v2 from mv"), row(1, 1, null, 1)); + + updateView("UPDATE %s USING TIMESTAMP 2 set v2 = null where p1 = 1 AND p2 = 1;"); + if (flush) + FBUtilities.waitOnFutures(ks.flush()); + assertRowsIgnoringOrder(execute("SELECT p1, p2, v1, v2 from %s")); + assertRowsIgnoringOrder(execute("SELECT p1, p2, v1, v2 from mv")); + + updateView("INSERT INTO %s (p1,p2) VALUES(1,1) USING TIMESTAMP 3;"); + if (flush) + FBUtilities.waitOnFutures(ks.flush()); + assertRowsIgnoringOrder(execute("SELECT p1, p2, v1, v2 from %s"), row(1, 1, null, null)); + assertRowsIgnoringOrder(execute("SELECT p1, p2, v1, v2 from mv"), row(1, 1, null, null)); + + updateView("DELETE FROM %s USING TIMESTAMP 4 WHERE p1 =1 AND p2 = 1;"); + if (flush) + FBUtilities.waitOnFutures(ks.flush()); + assertRowsIgnoringOrder(execute("SELECT p1, p2, v1, v2 from %s")); + assertRowsIgnoringOrder(execute("SELECT p1, p2, v1, v2 from mv")); + + updateView("UPDATE %s USING TIMESTAMP 5 set v2 = 1 where p1 = 1 AND p2 = 1;"); + if (flush) + FBUtilities.waitOnFutures(ks.flush()); + assertRowsIgnoringOrder(execute("SELECT p1, p2, v1, v2 from %s"), row(1, 1, null, 1)); + assertRowsIgnoringOrder(execute("SELECT p1, p2, v1, v2 from mv"), row(1, 1, null, 1)); + } + + @Test + public void testStrictLivenessTombstone() throws Throwable + { + createTable("create table %s (p int primary key, v1 int, v2 int)"); + + execute("USE " + keyspace()); + executeNet(protocolVersion, "USE " + keyspace()); + Keyspace ks = Keyspace.open(keyspace()); + + createView("mv", + "create materialized view %s as select * from %%s where p is not null and v1 is not null primary key (v1, p)" + + " with gc_grace_seconds=5;"); + ColumnFamilyStore cfs = ks.getColumnFamilyStore("mv"); + cfs.disableAutoCompaction(); + + updateView("Insert into %s (p, v1, v2) values (1, 1, 1) ;"); + assertRowsIgnoringOrder(execute("SELECT p, v1, v2 from mv"), row(1, 1, 1)); + + updateView("Update %s set v1 = null WHERE p = 1"); + FBUtilities.waitOnFutures(ks.flush()); + assertRowsIgnoringOrder(execute("SELECT p, v1, v2 from mv")); + + cfs.forceMajorCompaction(); // before gc grace second, strict-liveness tombstoned dead row remains + assertEquals(1, cfs.getLiveSSTables().size()); + + Thread.sleep(6000); + assertEquals(1, cfs.getLiveSSTables().size()); // no auto compaction. + + cfs.forceMajorCompaction(); // after gc grace second, no data left + assertEquals(0, cfs.getLiveSSTables().size()); + + updateView("Update %s using ttl 5 set v1 = 1 WHERE p = 1"); + FBUtilities.waitOnFutures(ks.flush()); + assertRowsIgnoringOrder(execute("SELECT p, v1, v2 from mv"), row(1, 1, 1)); + + cfs.forceMajorCompaction(); // before ttl+gc_grace_second, strict-liveness ttled dead row remains + assertEquals(1, cfs.getLiveSSTables().size()); + assertRowsIgnoringOrder(execute("SELECT p, v1, v2 from mv"), row(1, 1, 1)); + + Thread.sleep(5500); // after expired, before gc_grace_second + cfs.forceMajorCompaction();// before ttl+gc_grace_second, strict-liveness ttled dead row remains + assertEquals(1, cfs.getLiveSSTables().size()); + assertRowsIgnoringOrder(execute("SELECT p, v1, v2 from mv")); + + Thread.sleep(5500); // after expired + gc_grace_second + assertEquals(1, cfs.getLiveSSTables().size()); // no auto compaction. + + cfs.forceMajorCompaction(); // after gc grace second, no data left + assertEquals(0, cfs.getLiveSSTables().size()); + } + + @Test + public void testCellTombstoneAndShadowableTombstonesWithFlush() throws Throwable + { + testCellTombstoneAndShadowableTombstones(true); + } + + @Test + public void testCellTombstoneAndShadowableTombstonesWithoutFlush() throws Throwable + { + testCellTombstoneAndShadowableTombstones(false); + } + + private void testCellTombstoneAndShadowableTombstones(boolean flush) throws Throwable + { + createTable("create table %s (p int primary key, v1 int, v2 int)"); + + execute("USE " + keyspace()); + executeNet(protocolVersion, "USE " + keyspace()); + Keyspace ks = Keyspace.open(keyspace()); + + createView("mv", + "create materialized view %s as select * from %%s where p is not null and v1 is not null primary key (v1, p);"); + ks.getColumnFamilyStore("mv").disableAutoCompaction(); + + // sstable 1, Set initial values TS=1 + updateView("Insert into %s (p, v1, v2) values (3, 1, 3) using timestamp 1;"); + + if (flush) + FBUtilities.waitOnFutures(ks.flush()); + + assertRowsIgnoringOrder(execute("SELECT v2, WRITETIME(v2) from mv WHERE v1 = ? AND p = ?", 1, 3), row(3, 1L)); + // sstable 2 + updateView("UPdate %s using timestamp 2 set v2 = null where p = 3"); + + if (flush) + FBUtilities.waitOnFutures(ks.flush()); + + assertRowsIgnoringOrder(execute("SELECT v2, WRITETIME(v2) from mv WHERE v1 = ? AND p = ?", 1, 3), + row(null, null)); + // sstable 3 + updateView("UPdate %s using timestamp 3 set v1 = 2 where p = 3"); + + if (flush) + FBUtilities.waitOnFutures(ks.flush()); + + assertRowsIgnoringOrder(execute("SELECT v1, p, v2, WRITETIME(v2) from mv"), row(2, 3, null, null)); + // sstable 4 + updateView("UPdate %s using timestamp 4 set v1 = 1 where p = 3"); + + if (flush) + FBUtilities.waitOnFutures(ks.flush()); + + assertRowsIgnoringOrder(execute("SELECT v1, p, v2, WRITETIME(v2) from mv"), row(1, 3, null, null)); + + if (flush) + { + // compact sstable 2 and 3; + ColumnFamilyStore cfs = ks.getColumnFamilyStore("mv"); + List sstables = cfs.getLiveSSTables() + .stream() + .sorted(Comparator.comparingInt(s -> s.descriptor.generation)) + .map(s -> s.getFilename()) + .collect(Collectors.toList()); + System.out.println("SSTables " + sstables); + String dataFiles = String.join(",", Arrays.asList(sstables.get(1), sstables.get(2))); + CompactionManager.instance.forceUserDefinedCompaction(dataFiles); + } + // cell-tombstone in sstable 4 is not compacted away, because the shadowable tombstone is shadowed by new row. + assertRowsIgnoringOrder(execute("SELECT v1, p, v2, WRITETIME(v2) from mv"), row(1, 3, null, null)); + } + + @Test + public void complexTimestampDeletionTestWithFlush() throws Throwable + { + complexTimestampWithbaseNonPKColumnsInViewPKDeletionTest(true); + complexTimestampWithbasePKColumnsInViewPKDeletionTest(true); + } + + @Test + public void complexTimestampDeletionTestWithoutFlush() throws Throwable + { + complexTimestampWithbaseNonPKColumnsInViewPKDeletionTest(false); + complexTimestampWithbasePKColumnsInViewPKDeletionTest(false); + } + + private void complexTimestampWithbasePKColumnsInViewPKDeletionTest(boolean flush) throws Throwable + { + createTable("create table %s (p1 int, p2 int, v1 int, v2 int, primary key(p1, p2))"); + + execute("USE " + keyspace()); + executeNet(protocolVersion, "USE " + keyspace()); + Keyspace ks = Keyspace.open(keyspace()); + + createView("mv2", + "create materialized view %s as select * from %%s where p1 is not null and p2 is not null primary key (p2, p1);"); + ks.getColumnFamilyStore("mv2").disableAutoCompaction(); + + // Set initial values TS=1 + updateView("Insert into %s (p1, p2, v1, v2) values (1, 2, 3, 4) using timestamp 1;"); + + if (flush) + FBUtilities.waitOnFutures(ks.flush()); + + assertRowsIgnoringOrder(execute("SELECT v1, v2, WRITETIME(v2) from mv2 WHERE p1 = ? AND p2 = ?", 1, 2), + row(3, 4, 1L)); + // remove row/mv TS=2 + updateView("Delete from %s using timestamp 2 where p1 = 1 and p2 = 2;"); + + if (flush) + FBUtilities.waitOnFutures(ks.flush()); + // view are empty + assertRowsIgnoringOrder(execute("SELECT * from mv2")); + // insert PK with TS=3 + updateView("Insert into %s (p1, p2) values (1, 2) using timestamp 3;"); + + if (flush) + FBUtilities.waitOnFutures(ks.flush()); + // deleted column in MV remained dead + assertRowsIgnoringOrder(execute("SELECT * from mv2"), row(2, 1, null, null)); + + ks.getColumnFamilyStore("mv2").forceMajorCompaction(); + assertRowsIgnoringOrder(execute("SELECT * from mv2"), row(2, 1, null, null)); + + // reset values + updateView("Insert into %s (p1, p2, v1, v2) values (1, 2, 3, 4) using timestamp 10;"); + if (flush) + FBUtilities.waitOnFutures(ks.flush()); + + assertRowsIgnoringOrder(execute("SELECT v1, v2, WRITETIME(v2) from mv2 WHERE p1 = ? AND p2 = ?", 1, 2), + row(3, 4, 10L)); + + updateView("UPDATE %s using timestamp 20 SET v2 = 5 WHERE p1 = 1 and p2 = 2"); + if (flush) + FBUtilities.waitOnFutures(ks.flush()); + + assertRowsIgnoringOrder(execute("SELECT v1, v2, WRITETIME(v2) from mv2 WHERE p1 = ? AND p2 = ?", 1, 2), + row(3, 5, 20L)); + + updateView("DELETE FROM %s using timestamp 10 WHERE p1 = 1 and p2 = 2"); + if (flush) + FBUtilities.waitOnFutures(ks.flush()); + + assertRowsIgnoringOrder(execute("SELECT v1, v2, WRITETIME(v2) from mv2 WHERE p1 = ? AND p2 = ?", 1, 2), + row(null, 5, 20L)); + } + + public void complexTimestampWithbaseNonPKColumnsInViewPKDeletionTest(boolean flush) throws Throwable + { + createTable("create table %s (p int primary key, v1 int, v2 int)"); + + execute("USE " + keyspace()); + executeNet(protocolVersion, "USE " + keyspace()); + Keyspace ks = Keyspace.open(keyspace()); + + createView("mv", + "create materialized view %s as select * from %%s where p is not null and v1 is not null primary key (v1, p);"); + ks.getColumnFamilyStore("mv").disableAutoCompaction(); + + // Set initial values TS=1 + updateView("Insert into %s (p, v1, v2) values (3, 1, 5) using timestamp 1;"); + + if (flush) + FBUtilities.waitOnFutures(ks.flush()); + + assertRowsIgnoringOrder(execute("SELECT v2, WRITETIME(v2) from mv WHERE v1 = ? AND p = ?", 1, 3), row(5, 1L)); + // remove row/mv TS=2 + updateView("Delete from %s using timestamp 2 where p = 3;"); + + if (flush) + FBUtilities.waitOnFutures(ks.flush()); + // view are empty + assertRowsIgnoringOrder(execute("SELECT * from mv")); + // insert PK with TS=3 + updateView("Insert into %s (p, v1) values (3, 1) using timestamp 3;"); + + if (flush) + FBUtilities.waitOnFutures(ks.flush()); + // deleted column in MV remained dead + assertRowsIgnoringOrder(execute("SELECT * from mv"), row(1, 3, null)); + + // insert values TS=2, it should be considered dead due to previous tombstone + updateView("Insert into %s (p, v1, v2) values (3, 1, 5) using timestamp 2;"); + + if (flush) + FBUtilities.waitOnFutures(ks.flush()); + // deleted column in MV remained dead + assertRowsIgnoringOrder(execute("SELECT * from mv"), row(1, 3, null)); + + // insert values TS=2, it should be considered dead due to previous tombstone + executeNet(protocolVersion, "UPDATE %s USING TIMESTAMP 3 SET v2 = ? WHERE p = ?", 4, 3); + + if (flush) + FBUtilities.waitOnFutures(ks.flush()); + + assertRows(execute("SELECT v1, p, v2, WRITETIME(v2) from mv"), row(1, 3, 4, 3L)); + + ks.getColumnFamilyStore("mv").forceMajorCompaction(); + assertRows(execute("SELECT v1, p, v2, WRITETIME(v2) from mv"), row(1, 3, 4, 3L)); + } + + @Test + public void testMVWithDifferentColumnsWithFlush() throws Throwable + { + testMVWithDifferentColumns(true); + } + + @Test + public void testMVWithDifferentColumnsWithoutFlush() throws Throwable + { + testMVWithDifferentColumns(false); + } + + private void testMVWithDifferentColumns(boolean flush) throws Throwable + { + createTable("CREATE TABLE %s (a int, b int, c int, d int, e int, f int, PRIMARY KEY(a, b))"); + + execute("USE " + keyspace()); + executeNet(protocolVersion, "USE " + keyspace()); + List viewNames = new ArrayList<>(); + List mvStatements = Arrays.asList( + // all selected + "CREATE MATERIALIZED VIEW %s AS SELECT * FROM %%s WHERE a IS NOT NULL AND b IS NOT NULL PRIMARY KEY (a,b)", + // unselected e,f + "CREATE MATERIALIZED VIEW %s AS SELECT c,d FROM %%s WHERE a IS NOT NULL AND b IS NOT NULL PRIMARY KEY (a,b)", + // no selected + "CREATE MATERIALIZED VIEW %s AS SELECT a,b FROM %%s WHERE a IS NOT NULL AND b IS NOT NULL PRIMARY KEY (a,b)", + // all selected, re-order keys + "CREATE MATERIALIZED VIEW %s AS SELECT * FROM %%s WHERE a IS NOT NULL AND b IS NOT NULL PRIMARY KEY (b,a)", + // unselected e,f, re-order keys + "CREATE MATERIALIZED VIEW %s AS SELECT a,b,c,d FROM %%s WHERE a IS NOT NULL AND b IS NOT NULL PRIMARY KEY (b,a)", + // no selected, re-order keys + "CREATE MATERIALIZED VIEW %s AS SELECT a,b FROM %%s WHERE a IS NOT NULL AND b IS NOT NULL PRIMARY KEY (b,a)"); + + Keyspace ks = Keyspace.open(keyspace()); + + for (int i = 0; i < mvStatements.size(); i++) + { + String name = "mv" + i; + viewNames.add(name); + createView(name, mvStatements.get(i)); + ks.getColumnFamilyStore(name).disableAutoCompaction(); + } + + // insert + updateViewWithFlush("INSERT INTO %s (a,b,c,d,e,f) VALUES(1,1,1,1,1,1) using timestamp 1", flush); + assertBaseViews(row(1, 1, 1, 1, 1, 1), viewNames); + + updateViewWithFlush("UPDATE %s using timestamp 2 SET c=0, d=0 WHERE a=1 AND b=1", flush); + assertBaseViews(row(1, 1, 0, 0, 1, 1), viewNames); + + updateViewWithFlush("UPDATE %s using timestamp 2 SET e=0, f=0 WHERE a=1 AND b=1", flush); + assertBaseViews(row(1, 1, 0, 0, 0, 0), viewNames); + + updateViewWithFlush("DELETE FROM %s using timestamp 2 WHERE a=1 AND b=1", flush); + assertBaseViews(null, viewNames); + + // partial update unselected, selected + updateViewWithFlush("UPDATE %s using timestamp 3 SET f=1 WHERE a=1 AND b=1", flush); + assertBaseViews(row(1, 1, null, null, null, 1), viewNames); + + updateViewWithFlush("UPDATE %s using timestamp 4 SET e = 1, f=null WHERE a=1 AND b=1", flush); + assertBaseViews(row(1, 1, null, null, 1, null), viewNames); + + updateViewWithFlush("UPDATE %s using timestamp 4 SET e = null WHERE a=1 AND b=1", flush); + assertBaseViews(null, viewNames); + + updateViewWithFlush("UPDATE %s using timestamp 5 SET c = 1 WHERE a=1 AND b=1", flush); + assertBaseViews(row(1, 1, 1, null, null, null), viewNames); + + updateViewWithFlush("UPDATE %s using timestamp 5 SET c = null WHERE a=1 AND b=1", flush); + assertBaseViews(null, viewNames); + + updateViewWithFlush("UPDATE %s using timestamp 6 SET d = 1 WHERE a=1 AND b=1", flush); + assertBaseViews(row(1, 1, null, 1, null, null), viewNames); + + updateViewWithFlush("UPDATE %s using timestamp 7 SET d = null WHERE a=1 AND b=1", flush); + assertBaseViews(null, viewNames); + + updateViewWithFlush("UPDATE %s using timestamp 8 SET f = 1 WHERE a=1 AND b=1", flush); + assertBaseViews(row(1, 1, null, null, null, 1), viewNames); + + updateViewWithFlush("UPDATE %s using timestamp 6 SET c = 1 WHERE a=1 AND b=1", flush); + assertBaseViews(row(1, 1, 1, null, null, 1), viewNames); + + // view row still alive due to c=1@6 + updateViewWithFlush("UPDATE %s using timestamp 8 SET f = null WHERE a=1 AND b=1", flush); + assertBaseViews(row(1, 1, 1, null, null, null), viewNames); + + updateViewWithFlush("UPDATE %s using timestamp 6 SET c = null WHERE a=1 AND b=1", flush); + assertBaseViews(null, viewNames); + } + + private void assertBaseViews(Object[] row, List viewNames) throws Throwable + { + UntypedResultSet result = execute("SELECT * FROM %s"); + if (row == null) + assertRowsIgnoringOrder(result); + else + assertRowsIgnoringOrder(result, row); + for (int i = 0; i < viewNames.size(); i++) + assertBaseView(result, execute(String.format("SELECT * FROM %s", viewNames.get(i))), viewNames.get(i)); + } + + private void assertBaseView(UntypedResultSet base, UntypedResultSet view, String mv) + { + List baseMeta = base.metadata(); + List viewMeta = view.metadata(); + + Iterator iter = base.iterator(); + Iterator viewIter = view.iterator(); + + List baseData = com.google.common.collect.Lists.newArrayList(iter); + List viewData = com.google.common.collect.Lists.newArrayList(viewIter); + + if (baseData.size() != viewData.size()) + fail(String.format("Mismatch number of rows in view %s: <%s>, in base <%s>", + mv, + makeRowStrings(view), + makeRowStrings(base))); + if (baseData.size() == 0) + return; + if (viewData.size() != 1) + fail(String.format("Expect only one row in view %s, but got <%s>", + mv, + makeRowStrings(view))); + + UntypedResultSet.Row row = baseData.get(0); + UntypedResultSet.Row viewRow = viewData.get(0); + + Map baseValues = new HashMap<>(); + for (int j = 0; j < baseMeta.size(); j++) + { + ColumnSpecification column = baseMeta.get(j); + ByteBuffer actualValue = row.getBytes(column.name.toString()); + baseValues.put(column.name.toString(), actualValue); + } + for (int j = 0; j < viewMeta.size(); j++) + { + ColumnSpecification column = viewMeta.get(j); + String name = column.name.toString(); + ByteBuffer viewValue = viewRow.getBytes(name); + if (!baseValues.containsKey(name)) + { + fail(String.format("Extra column: %s with value %s in view", name, column.type.compose(viewValue))); + } + else if (!Objects.equal(baseValues.get(name), viewValue)) + { + fail(String.format("Non equal column: %s, expected <%s> but got <%s>", + name, + column.type.compose(baseValues.get(name)), + column.type.compose(viewValue))); + } + } + } + +} diff --git a/test/unit/org/apache/cassandra/cql3/ViewFilteringTest.java b/test/unit/org/apache/cassandra/cql3/ViewFilteringTest.java index 245ceb78fc..fe618b6a51 100644 --- a/test/unit/org/apache/cassandra/cql3/ViewFilteringTest.java +++ b/test/unit/org/apache/cassandra/cql3/ViewFilteringTest.java @@ -77,13 +77,13 @@ public class ViewFilteringTest extends CQLTester // 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)" + "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) @@ -96,19 +96,19 @@ public class ViewFilteringTest extends CQLTester 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)" - ); + 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++) { @@ -153,21 +153,21 @@ public class ViewFilteringTest extends CQLTester 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\")"); + "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\")"); + "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) + row(1, 1, 0) ); } @@ -176,7 +176,7 @@ public class ViewFilteringTest extends CQLTester for (String mvname : Arrays.asList("mv_test", "mv_test2")) { assertRowsIgnoringOrder(execute("SELECT \"theKey\", \"Col\", \"the\"\"Value\" FROM " + mvname), - row(1, 1, 0) + row(1, 1, 0) ); } } @@ -195,22 +195,22 @@ public class ViewFilteringTest extends CQLTester 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)"); + "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) + 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) + row(1, 0, 2), + row(1, 1, 3) ); } @@ -228,22 +228,22 @@ public class ViewFilteringTest extends CQLTester 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)"); + "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) + 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) + row(1, 0, 2), + row(1, 1, 3) ); } @@ -274,51 +274,51 @@ public class ViewFilteringTest extends CQLTester 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) + 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) + 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) + 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) + 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) + 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 @@ -326,20 +326,20 @@ public class ViewFilteringTest extends CQLTester 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) + 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) + 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 @@ -377,8 +377,8 @@ public class ViewFilteringTest extends CQLTester Thread.sleep(10); assertRowsIgnoringOrder(execute("SELECT a, b, c, d FROM mv_test" + i), - row(1, 1, 0, 0), - row(1, 1, 1, 0) + row(1, 1, 0, 0), + row(1, 1, 1, 0) ); // insert new rows that do not match the filter @@ -386,16 +386,16 @@ public class ViewFilteringTest extends CQLTester 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) + 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) + row(1, 1, 0, 0), + row(1, 1, 1, 0), + row(1, 1, 2, 0) ); // update rows that don't match the filter @@ -403,17 +403,17 @@ public class ViewFilteringTest extends CQLTester 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) + 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) + row(1, 1, 0, 1), + row(1, 1, 1, 0), + row(1, 1, 2, 0) ); // delete rows that don't match the filter @@ -422,16 +422,16 @@ public class ViewFilteringTest extends CQLTester 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) + 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) + row(1, 1, 1, 0), + row(1, 1, 2, 0) ); // delete a partition that matches the filter @@ -463,8 +463,8 @@ public class ViewFilteringTest extends CQLTester Thread.sleep(10); assertRows(execute("SELECT * FROM mv_test"), - row(1, 1, 0), - row(1, 1, 1) + row(1, 1, 0), + row(1, 1, 1) ); // insert new rows that do not match the filter @@ -472,16 +472,16 @@ public class ViewFilteringTest extends CQLTester 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) + 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) + row(1, 1, 0), + row(1, 1, 1), + row(1, 1, 2) ); // update rows that don't match the filter @@ -489,17 +489,17 @@ public class ViewFilteringTest extends CQLTester 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) + 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) + row(1, 1, 0), + row(1, 1, 1), + row(1, 1, 2) ); // delete rows that don't match the filter @@ -508,16 +508,16 @@ public class ViewFilteringTest extends CQLTester 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) + 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) + row(1, 1, 1), + row(1, 1, 2) ); // delete a partition that matches the filter @@ -554,51 +554,51 @@ public class ViewFilteringTest extends CQLTester 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(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) + 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) + 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) + 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) + 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 @@ -606,27 +606,27 @@ public class ViewFilteringTest extends CQLTester 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) + 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) + 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) + row(0, 1, 0, 0), + row(0, 1, 1, 0) ); dropView("mv_test" + i); @@ -662,51 +662,51 @@ public class ViewFilteringTest extends CQLTester 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(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) + 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) + 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) + 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) + 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 @@ -714,27 +714,27 @@ public class ViewFilteringTest extends CQLTester 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(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) + 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) + row(0, 1, 0, 0), + row(0, 1, 1, 0) ); dropView("mv_test" + i); @@ -772,56 +772,56 @@ public class ViewFilteringTest extends CQLTester 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) + 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) + 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) + 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) + 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) + 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 @@ -829,29 +829,29 @@ public class ViewFilteringTest extends CQLTester 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) + 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) + 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) + row(0, 1, 0, 0), + row(0, 1, 1, 0) ); dropView("mv_test" + i); @@ -889,10 +889,10 @@ public class ViewFilteringTest extends CQLTester 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(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 @@ -900,20 +900,20 @@ public class ViewFilteringTest extends CQLTester 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) + 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) + 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 @@ -921,21 +921,21 @@ public class ViewFilteringTest extends CQLTester 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(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) + 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 @@ -944,27 +944,27 @@ public class ViewFilteringTest extends CQLTester 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(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) + 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) + row(0, 1, 0, 0), + row(0, 1, 1, 0) ); dropView("mv_test" + i); @@ -1002,51 +1002,51 @@ public class ViewFilteringTest extends CQLTester 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) + 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) + 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) + 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) + 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) + 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 @@ -1055,27 +1055,27 @@ public class ViewFilteringTest extends CQLTester 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) + 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) + 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) + row(0, 0, 1, 0), + row(0, 1, 1, 0) ); // insert a partition with one matching and one non-matching row using a batch (CASSANDRA-10614) @@ -1087,9 +1087,9 @@ public class ViewFilteringTest extends CQLTester 4, 4, 0, 0, 4, 4, 1, 1); assertRowsIgnoringOrder(execute("SELECT a, b, c, d FROM mv_test" + i), - row(0, 0, 1, 0), - row(0, 1, 1, 0), - row(4, 4, 1, 1) + row(0, 0, 1, 0), + row(0, 1, 1, 0), + row(4, 4, 1, 1) ); dropView("mv_test" + i); @@ -1127,41 +1127,41 @@ public class ViewFilteringTest extends CQLTester Thread.sleep(10); assertRowsIgnoringOrder(execute("SELECT a, b, c, d FROM mv_test" + i), - row(1, 0, 1, 0), - row(1, 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 (?, ?, ?, ?)", 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) + 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) + 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) + 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) + row(1, 0, 1, 0), + row(1, 1, 1, 2), + row(1, 2, 1, 0) ); // delete rows that don't match the filter @@ -1169,16 +1169,16 @@ public class ViewFilteringTest extends CQLTester 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) + 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) + row(1, 0, 1, 0), + row(1, 2, 1, 0) ); // delete a partition that matches the filter @@ -1218,61 +1218,61 @@ public class ViewFilteringTest extends CQLTester "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 + "))"); + "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 + ")"); + "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'," + diff --git a/test/unit/org/apache/cassandra/cql3/ViewTest.java b/test/unit/org/apache/cassandra/cql3/ViewTest.java index f8f8c9f3e6..61f4b4a8cd 100644 --- a/test/unit/org/apache/cassandra/cql3/ViewTest.java +++ b/test/unit/org/apache/cassandra/cql3/ViewTest.java @@ -18,12 +18,13 @@ package org.apache.cassandra.cql3; +import static org.junit.Assert.*; + import java.util.ArrayList; import java.util.HashSet; import java.util.List; import java.util.concurrent.Future; import java.util.concurrent.TimeUnit; - import com.google.common.util.concurrent.Uninterruptibles; import junit.framework.Assert; @@ -46,8 +47,6 @@ import org.apache.cassandra.db.SystemKeyspace; import org.apache.cassandra.db.compaction.CompactionManager; import org.apache.cassandra.utils.FBUtilities; -import static org.junit.Assert.assertTrue; - public class ViewTest extends CQLTester { int protocolVersion = 4; @@ -83,7 +82,7 @@ public class ViewTest extends CQLTester { executeNet(protocolVersion, query, params); while (!(((SEPExecutor) StageManager.getStage(Stage.VIEW_MUTATION)).getPendingTasks() == 0 - && ((SEPExecutor) StageManager.getStage(Stage.VIEW_MUTATION)).getActiveCount() == 0)) + && ((SEPExecutor) StageManager.getStage(Stage.VIEW_MUTATION)).getActiveCount() == 0)) { Thread.sleep(1); } @@ -342,9 +341,9 @@ public class ViewTest extends CQLTester public void testCreateMvWithTTL() throws Throwable { createTable("CREATE TABLE %s (" + - "k int PRIMARY KEY, " + - "c int, " + - "val int) WITH default_time_to_live = 60"); + "k int PRIMARY KEY, " + + "c int, " + + "val int) WITH default_time_to_live = 60"); execute("USE " + keyspace()); executeNet(protocolVersion, "USE " + keyspace()); @@ -421,8 +420,10 @@ public class ViewTest extends CQLTester if (flush) FBUtilities.waitOnFutures(ks.flush()); - //change c's value and TS=3, tombstones c=1 and adds c=0 record + // change c's value and TS=3, tombstones c=1 and adds c=0 record executeNet(protocolVersion, "UPDATE %s USING TIMESTAMP 3 SET c = ? WHERE a = ? and b = ? ", 0, 0, 0); + if (flush) + FBUtilities.waitOnFutures(ks.flush()); assertRows(execute("SELECT d from mv WHERE c = ? and a = ? and b = ?", 1, 0, 0)); if(flush) @@ -443,7 +444,7 @@ public class ViewTest extends CQLTester assertRows(execute("SELECT d,e from mv WHERE c = ? and a = ? and b = ?", 1, 0, 0), row(0, null)); - //Add e value @ TS=1 + //Add e value @ TS=1 executeNet(protocolVersion, "UPDATE %s USING TIMESTAMP 1 SET e = ? WHERE a = ? and b = ? ", 1, 0, 0); assertRows(execute("SELECT d,e from mv WHERE c = ? and a = ? and b = ?", 1, 0, 0), row(0, 1)); @@ -838,12 +839,12 @@ public class ViewTest extends CQLTester String table = KEYSPACE + "." + currentTable(); updateView("BEGIN BATCH " + "INSERT INTO " + table + " (a, b, c, d) VALUES (?, ?, ?, ?); " + // should be accepted - "UPDATE " + table + " SET d = ? WHERE a = ? AND b = ?; " + // should be ignored + "UPDATE " + table + " SET d = ? WHERE a = ? AND b = ?; " + // should be accepted "APPLY BATCH", 0, 0, 0, 0, 1, 0, 1); assertRows(execute("SELECT a, b, c from mv WHERE b = ?", 0), row(0, 0, 0)); - assertEmpty(execute("SELECT a, b, c from mv WHERE b = ?", 1)); + assertRows(execute("SELECT a, b, c from mv WHERE b = ?", 1), row(0, 1, null)); ColumnFamilyStore cfs = Keyspace.open(keyspace()).getColumnFamilyStore("mv"); cfs.forceBlockingFlush(); @@ -1005,9 +1006,9 @@ public class ViewTest extends CQLTester row(1, 3)); updateView(String.format("BEGIN UNLOGGED BATCH " + - "DELETE FROM %s WHERE a = 1 AND b > 1 AND b < 3;" + - "DELETE FROM %s WHERE a = 1;" + - "APPLY BATCH", currentTable(), currentTable())); + "DELETE FROM %s WHERE a = 1 AND b > 1 AND b < 3;" + + "DELETE FROM %s WHERE a = 1;" + + "APPLY BATCH", currentTable(), currentTable())); mvRows = executeNet(protocolVersion, "SELECT a, b FROM mv1"); assertRowsNet(protocolVersion, mvRows); @@ -1295,4 +1296,4 @@ public class ViewTest extends CQLTester assertRows(execute("SELECT k, toJson(listval) from mv"), row(0, "[[\"a\", \"1\"], [\"b\", \"2\"], [\"c\", \"3\"]]")); } -} +} \ No newline at end of file