mirror of https://github.com/apache/cassandra
Merge branch 'cassandra-3.11' into trunk
This commit is contained in:
commit
a07d327be8
|
|
@ -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)
|
||||
|
|
|
|||
|
|
@ -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<ColumnIdentifier, ColumnSubselection> 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<ColumnIdentifier, ColumnSubselection> 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<ColumnMetadata> 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;
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -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";
|
||||
|
|
|
|||
|
|
@ -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);
|
||||
}
|
||||
}
|
||||
Loading…
Reference in New Issue