diff --git a/CHANGES.txt b/CHANGES.txt index 649a79f45f..d331967f82 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -115,6 +115,7 @@ * Fix cqlsh automatic protocol downgrade regression (CASSANDRA-13307) * Tracing payload not passed from QueryMessage to tracing session (CASSANDRA-12835) Merged from 3.0: + * Ensure consistent view of partition columns between coordinator and replica in ColumnFilter (CASSANDRA-13004) * Failed unregistering mbean during drop keyspace (CASSANDRA-13346) * nodetool scrub/cleanup/upgradesstables exit code is wrong (CASSANDRA-13542) * Fix the reported number of sstable data files accessed per read (CASSANDRA-13120) diff --git a/src/java/org/apache/cassandra/db/filter/ColumnFilter.java b/src/java/org/apache/cassandra/db/filter/ColumnFilter.java index b56870408a..dcd93e8c49 100644 --- a/src/java/org/apache/cassandra/db/filter/ColumnFilter.java +++ b/src/java/org/apache/cassandra/db/filter/ColumnFilter.java @@ -20,6 +20,7 @@ package org.apache.cassandra.db.filter; import java.io.IOException; import java.util.*; +import com.google.common.annotations.VisibleForTesting; import com.google.common.collect.Iterables; import com.google.common.collect.Iterators; import com.google.common.collect.SortedSetMultimap; @@ -68,12 +69,11 @@ public class ColumnFilter // True if _fetched_ includes all regular columns (and any static in _queried_), in which case metadata must not be // null. If false, then _fetched_ == _queried_ and we only store _queried_. - private final boolean fetchAllRegulars; - - private final TableMetadata metadata; // can be null if !isFetchAll + public final boolean fetchAllRegulars; + private final RegularAndStaticColumns fetched; private final RegularAndStaticColumns queried; // can be null if fetchAllRegulars, to represent a wildcard query (all - // static and regular columns are both _fetched_ and _queried_). + // static and regular columns are both _fetched_ and _queried_). private final SortedSetMultimap subSelections; // can be null private ColumnFilter(boolean fetchAllRegulars, @@ -84,7 +84,42 @@ public class ColumnFilter assert !fetchAllRegulars || metadata != null; assert fetchAllRegulars || queried != null; this.fetchAllRegulars = fetchAllRegulars; - this.metadata = metadata; + + if (fetchAllRegulars) + { + RegularAndStaticColumns all = metadata.regularAndStaticColumns(); + if (queried == null) + { + this.fetched = this.queried = all; + } + else + { + this.fetched = all.statics.isEmpty() + ? all + : new RegularAndStaticColumns(queried.statics, all.regulars); + this.queried = queried; + } + } + else + { + this.fetched = this.queried = queried; + } + + this.subSelections = subSelections; + } + + /** + * Used on replica for deserialisation + */ + private ColumnFilter(boolean fetchAllRegulars, + RegularAndStaticColumns fetched, + RegularAndStaticColumns queried, + SortedSetMultimap subSelections) + { + assert !fetchAllRegulars || fetched != null; + assert fetchAllRegulars || queried != null; + this.fetchAllRegulars = fetchAllRegulars; + this.fetched = fetchAllRegulars ? fetched : queried; this.queried = queried; this.subSelections = subSelections; } @@ -106,7 +141,7 @@ public class ColumnFilter */ public static ColumnFilter selection(RegularAndStaticColumns columns) { - return new ColumnFilter(false, null, columns, null); + return new ColumnFilter(false, (TableMetadata) null, columns, null); } /** @@ -125,15 +160,7 @@ public class ColumnFilter */ public RegularAndStaticColumns fetchedColumns() { - if (!fetchAllRegulars) - return queried; - - // We always fetch all regulars, but only fetch the statics in queried. Unless queried == null, in which - // case it's a wildcard and we fetch everything. - RegularAndStaticColumns all = metadata.regularAndStaticColumns(); - return queried == null || all.statics.isEmpty() - ? all - : new RegularAndStaticColumns(queried.statics, all.regulars); + return fetched; } /** @@ -143,8 +170,7 @@ public class ColumnFilter */ public RegularAndStaticColumns queriedColumns() { - assert queried != null || fetchAllRegulars; - return queried == null ? metadata.regularAndStaticColumns() : queried; + return queried; } /** @@ -431,6 +457,23 @@ public class ColumnFilter } } + @Override + public boolean equals(Object other) + { + if (other == this) + return true; + + if (!(other instanceof ColumnFilter)) + return false; + + ColumnFilter otherCf = (ColumnFilter) other; + + return otherCf.fetchAllRegulars == this.fetchAllRegulars && + Objects.equals(otherCf.fetched, this.fetched) && + Objects.equals(otherCf.queried, this.queried) && + Objects.equals(otherCf.subSelections, this.subSelections); + } + @Override public String toString() { @@ -476,8 +519,8 @@ public class ColumnFilter public static class Serializer { - private static final int FETCH_ALL_MASK = 0x01; - private static final int HAS_QUERIED_MASK = 0x02; + private static final int FETCH_ALL_MASK = 0x01; + private static final int HAS_QUERIED_MASK = 0x02; private static final int HAS_SUB_SELECTIONS_MASK = 0x04; private static int makeHeaderByte(ColumnFilter selection) @@ -487,7 +530,8 @@ public class ColumnFilter | (selection.subSelections != null ? HAS_SUB_SELECTIONS_MASK : 0); } - private static ColumnFilter maybeUpdateForBackwardCompatility(ColumnFilter selection, int version) + @VisibleForTesting + public static ColumnFilter maybeUpdateForBackwardCompatility(ColumnFilter selection, int version) { if (version > MessagingService.VERSION_30 || !selection.fetchAllRegulars || selection.queried == null) return selection; @@ -499,12 +543,11 @@ public class ColumnFilter // queried some columns that are actually only fetched, but it's fine during upgrade). // More concretely, we replace our filter by a non-fetch-all one that queries every columns that our // current filter fetches. - Columns allRegulars = selection.metadata.regularColumns(); Set queriedStatic = new HashSet<>(); Iterables.addAll(queriedStatic, Iterables.filter(selection.queried, ColumnMetadata::isStatic)); return new ColumnFilter(false, - null, - new RegularAndStaticColumns(Columns.from(queriedStatic), allRegulars), + (TableMetadata) null, + new RegularAndStaticColumns(Columns.from(queriedStatic), selection.fetched.regulars), selection.subSelections); } @@ -514,6 +557,12 @@ public class ColumnFilter out.writeByte(makeHeaderByte(selection)); + if (version >= MessagingService.VERSION_3014 && selection.fetchAllRegulars) + { + Columns.serializer.serialize(selection.fetched.statics, out); + Columns.serializer.serialize(selection.fetched.regulars, out); + } + if (selection.queried != null) { Columns.serializer.serialize(selection.queried.statics, out); @@ -535,7 +584,23 @@ public class ColumnFilter boolean hasQueried = (header & HAS_QUERIED_MASK) != 0; boolean hasSubSelections = (header & HAS_SUB_SELECTIONS_MASK) != 0; + RegularAndStaticColumns fetched = null; RegularAndStaticColumns queried = null; + + if (isFetchAll) + { + if (version >= MessagingService.VERSION_3014) + { + Columns statics = Columns.serializer.deserialize(in, metadata); + Columns regulars = Columns.serializer.deserialize(in, metadata); + fetched = new RegularAndStaticColumns(statics, regulars); + } + else + { + fetched = metadata.regularAndStaticColumns(); + } + } + if (hasQueried) { Columns statics = Columns.serializer.deserialize(in, metadata); @@ -564,7 +629,7 @@ public class ColumnFilter if (version <= MessagingService.VERSION_30 && isFetchAll && queried != null) queried = new RegularAndStaticColumns(metadata.staticColumns(), queried.regulars); - return new ColumnFilter(isFetchAll, isFetchAll ? metadata : null, queried, subSelections); + return new ColumnFilter(isFetchAll, fetched, queried, subSelections); } public long serializedSize(ColumnFilter selection, int version) @@ -573,6 +638,12 @@ public class ColumnFilter long size = 1; // header byte + if (version >= MessagingService.VERSION_3014 && selection.fetchAllRegulars) + { + size += Columns.serializer.serializedSize(selection.fetched.statics); + size += Columns.serializer.serializedSize(selection.fetched.regulars); + } + if (selection.queried != null) { size += Columns.serializer.serializedSize(selection.queried.statics); @@ -590,4 +661,4 @@ public class ColumnFilter return size; } } -} +} \ No newline at end of file diff --git a/src/java/org/apache/cassandra/net/MessagingService.java b/src/java/org/apache/cassandra/net/MessagingService.java index b3e7b614b4..41771e77b7 100644 --- a/src/java/org/apache/cassandra/net/MessagingService.java +++ b/src/java/org/apache/cassandra/net/MessagingService.java @@ -94,7 +94,8 @@ public final class MessagingService implements MessagingServiceMBean // 8 bits version, so don't waste versions public static final int VERSION_30 = 10; - public static final int VERSION_40 = 11; + public static final int VERSION_3014 = 11; + public static final int VERSION_40 = 12; public static final int current_version = VERSION_40; public static final String FAILURE_CALLBACK_PARAM = "CAL_BAC"; diff --git a/test/unit/org/apache/cassandra/db/filter/ColumnFilterTest.java b/test/unit/org/apache/cassandra/db/filter/ColumnFilterTest.java new file mode 100644 index 0000000000..eee2ad50a3 --- /dev/null +++ b/test/unit/org/apache/cassandra/db/filter/ColumnFilterTest.java @@ -0,0 +1,83 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.cassandra.db.filter; + +import org.junit.Test; + +import junit.framework.Assert; +import org.apache.cassandra.db.marshal.Int32Type; +import org.apache.cassandra.dht.Murmur3Partitioner; +import org.apache.cassandra.io.util.DataInputBuffer; +import org.apache.cassandra.io.util.DataInputPlus; +import org.apache.cassandra.io.util.DataOutputBuffer; +import org.apache.cassandra.net.MessagingService; +import org.apache.cassandra.schema.ColumnMetadata; +import org.apache.cassandra.schema.TableMetadata; +import org.apache.cassandra.utils.ByteBufferUtil; + +public class ColumnFilterTest +{ + final static ColumnFilter.Serializer serializer = new ColumnFilter.Serializer(); + + @Test + public void columnFilterSerialisationRoundTrip() throws Exception + { + TableMetadata metadata = TableMetadata.builder("ks", "table") + .partitioner(Murmur3Partitioner.instance) + .addPartitionKeyColumn("pk", Int32Type.instance) + .addClusteringColumn("ck", Int32Type.instance) + .addRegularColumn("v1", Int32Type.instance) + .addRegularColumn("v2", Int32Type.instance) + .addRegularColumn("v3", Int32Type.instance) + .build(); + + ColumnMetadata v1 = metadata.getColumn(ByteBufferUtil.bytes("v1")); + + ColumnFilter columnFilter; + + columnFilter = ColumnFilter.all(metadata); + testRoundTrip(columnFilter, ColumnFilter.Serializer.maybeUpdateForBackwardCompatility(columnFilter, MessagingService.VERSION_30), metadata, MessagingService.VERSION_30); + testRoundTrip(ColumnFilter.all(metadata), metadata, MessagingService.VERSION_3014); + testRoundTrip(ColumnFilter.all(metadata), metadata, MessagingService.VERSION_40); + + testRoundTrip(ColumnFilter.selection(metadata.regularAndStaticColumns().without(v1)), metadata, MessagingService.VERSION_30); + testRoundTrip(ColumnFilter.selection(metadata.regularAndStaticColumns().without(v1)), metadata, MessagingService.VERSION_3014); + testRoundTrip(ColumnFilter.selection(metadata.regularAndStaticColumns().without(v1)), metadata, MessagingService.VERSION_40); + + columnFilter = ColumnFilter.selection(metadata, metadata.regularAndStaticColumns().without(v1)); + testRoundTrip(columnFilter, ColumnFilter.Serializer.maybeUpdateForBackwardCompatility(columnFilter, MessagingService.VERSION_30), metadata, MessagingService.VERSION_30); + testRoundTrip(ColumnFilter.selection(metadata, metadata.regularAndStaticColumns().without(v1)), metadata, MessagingService.VERSION_3014); + testRoundTrip(ColumnFilter.selection(metadata, metadata.regularAndStaticColumns().without(v1)), metadata, MessagingService.VERSION_40); + } + + static void testRoundTrip(ColumnFilter columnFilter, TableMetadata metadata, int version) throws Exception + { + testRoundTrip(columnFilter, columnFilter, metadata, version); + } + + static void testRoundTrip(ColumnFilter columnFilter, ColumnFilter expected, TableMetadata metadata, int version) throws Exception + { + DataOutputBuffer output = new DataOutputBuffer(); + serializer.serialize(columnFilter, output, version); + Assert.assertEquals(serializer.serializedSize(columnFilter, version), output.position()); + DataInputPlus input = new DataInputBuffer(output.buffer(), false); + ColumnFilter deserialized = serializer.deserialize(input, version, metadata); + Assert.assertEquals(deserialized, expected); + } +} \ No newline at end of file