this.asByteComparable = asByteComparable;
}
+ public static Kind fromOrdinal(int ordinal)
+ {
+ return VALUES[ordinal];
+ }
+
/**
* Compares the 2 provided kind.
*
@@ -476,7 +483,7 @@ public interface ClusteringPrefix extends IMeasurableMemory, Clusterable
public void skip(DataInputPlus in, int version, List> types) throws IOException
{
- Kind kind = Kind.values()[in.readByte()];
+ Kind kind = Kind.fromOrdinal(in.readByte());
// We shouldn't serialize static clusterings
assert kind != Kind.STATIC_CLUSTERING;
if (kind == Kind.CLUSTERING)
@@ -487,7 +494,7 @@ public interface ClusteringPrefix extends IMeasurableMemory, Clusterable
public ClusteringPrefix deserialize(DataInputPlus in, int version, List> types) throws IOException
{
- Kind kind = Kind.values()[in.readByte()];
+ Kind kind = Kind.fromOrdinal(in.readByte());
// We shouldn't serialize static clusterings
assert kind != Kind.STATIC_CLUSTERING;
if (kind == Kind.CLUSTERING)
@@ -657,7 +664,7 @@ public interface ClusteringPrefix extends IMeasurableMemory, Clusterable
throw new IOException("Corrupt flags value for clustering prefix (isStatic flag set): " + flags);
this.nextIsRow = UnfilteredSerializer.kind(flags) == Unfiltered.Kind.ROW;
- this.nextKind = nextIsRow ? Kind.CLUSTERING : ClusteringPrefix.Kind.values()[in.readByte()];
+ this.nextKind = nextIsRow ? Kind.CLUSTERING : Kind.fromOrdinal(in.readByte());
this.nextSize = nextIsRow ? comparator.size() : in.readUnsignedShort();
this.deserializedSize = 0;
diff --git a/src/java/org/apache/cassandra/db/DeletionTime.java b/src/java/org/apache/cassandra/db/DeletionTime.java
index 1d54a00fff..c220d1316b 100644
--- a/src/java/org/apache/cassandra/db/DeletionTime.java
+++ b/src/java/org/apache/cassandra/db/DeletionTime.java
@@ -178,7 +178,8 @@ public abstract class DeletionTime implements Comparable, IMeasura
public boolean deletes(Cell> cell)
{
- return deletes(cell.timestamp());
+ // check for LIVE first to avoid a potential cell megamorphic call
+ return markedForDeleteAt() != MARKED_FOR_DELETE_AT_LIVE && deletes(cell.timestamp());
}
public boolean deletes(long timestamp)
diff --git a/src/java/org/apache/cassandra/db/ReadCommand.java b/src/java/org/apache/cassandra/db/ReadCommand.java
index 53ad6b3d1b..d848e23b6b 100644
--- a/src/java/org/apache/cassandra/db/ReadCommand.java
+++ b/src/java/org/apache/cassandra/db/ReadCommand.java
@@ -198,6 +198,8 @@ public abstract class ReadCommand extends AbstractReadQuery
SINGLE_PARTITION (SinglePartitionReadCommand.selectionDeserializer, SinglePartitionReadCommand.accordSelectionDeserializer),
PARTITION_RANGE (PartitionRangeReadCommand.selectionDeserializer, ignore -> PartitionRangeReadCommand.selectionDeserializer);
+ private static final Kind[] VALUES = values();
+
private final SelectionDeserializer selectionDeserializer;
private final Function accordSelectionDeserializer;
@@ -206,6 +208,11 @@ public abstract class ReadCommand extends AbstractReadQuery
this.selectionDeserializer = selectionDeserializer;
this.accordSelectionDeserializer = accordSelectionDeserializer;
}
+
+ public static Kind fromOrdinal(int ordinal)
+ {
+ return VALUES[ordinal];
+ }
}
protected ReadCommand(Epoch serializedAtEpoch,
@@ -1450,7 +1457,7 @@ public abstract class ReadCommand extends AbstractReadQuery
public ReadCommand deserialize(DataInputPlus in, int version) throws IOException
{
- Kind kind = Kind.values()[in.readByte()];
+ Kind kind = Kind.fromOrdinal(in.readByte());
int flags = in.readByte();
// Shouldn't happen or it's a user error (see comment above) but
// better complain loudly than doing the wrong thing.
@@ -1488,7 +1495,7 @@ public abstract class ReadCommand extends AbstractReadQuery
public ReadCommand deserializeForAccord(Seekable key, TableMetadatas tables, DataInputPlus in, int version) throws IOException
{
- Kind kind = Kind.values()[in.readByte()];
+ Kind kind = Kind.fromOrdinal(in.readByte());
int flags = in.readByte();
if (isDigest(flags) || isForThrift(flags) || acceptsTransient(flags))
throw new IllegalStateException("Received an Accord command with a digest/thrift/transient flag set.");
diff --git a/src/java/org/apache/cassandra/db/StorageHook.java b/src/java/org/apache/cassandra/db/StorageHook.java
index f5fdec6a56..56938b441f 100644
--- a/src/java/org/apache/cassandra/db/StorageHook.java
+++ b/src/java/org/apache/cassandra/db/StorageHook.java
@@ -55,7 +55,7 @@ public interface StorageHook
String className = STORAGE_HOOK.getString();
if (className != null)
{
- return FBUtilities.construct(className, StorageHook.class.getSimpleName());
+ return FBUtilities.construct(className, StorageHook.class.getSimpleName(), StorageHook.class);
}
return new StorageHook()
@@ -89,4 +89,4 @@ public interface StorageHook
}
};
}
-}
\ No newline at end of file
+}
diff --git a/src/java/org/apache/cassandra/db/aggregation/AggregationSpecification.java b/src/java/org/apache/cassandra/db/aggregation/AggregationSpecification.java
index a4a1c57eca..2cd7d2a001 100644
--- a/src/java/org/apache/cassandra/db/aggregation/AggregationSpecification.java
+++ b/src/java/org/apache/cassandra/db/aggregation/AggregationSpecification.java
@@ -65,7 +65,14 @@ public abstract class AggregationSpecification
*/
public enum Kind
{
- AGGREGATE_EVERYTHING, AGGREGATE_BY_PK_PREFIX, AGGREGATE_BY_PK_PREFIX_WITH_SELECTOR
+ AGGREGATE_EVERYTHING, AGGREGATE_BY_PK_PREFIX, AGGREGATE_BY_PK_PREFIX_WITH_SELECTOR;
+
+ private static final Kind[] VALUES = values();
+
+ public static Kind fromOrdinal(int ordinal)
+ {
+ return VALUES[ordinal];
+ }
}
/**
@@ -253,7 +260,7 @@ public abstract class AggregationSpecification
public AggregationSpecification deserialize(DataInputPlus in, int version, TableMetadata metadata) throws IOException
{
- Kind kind = Kind.values()[in.readUnsignedByte()];
+ Kind kind = Kind.fromOrdinal(in.readUnsignedByte());
switch (kind)
{
case AGGREGATE_EVERYTHING:
diff --git a/src/java/org/apache/cassandra/db/filter/AbstractClusteringIndexFilter.java b/src/java/org/apache/cassandra/db/filter/AbstractClusteringIndexFilter.java
index d55799a484..f9fdb3f2e6 100644
--- a/src/java/org/apache/cassandra/db/filter/AbstractClusteringIndexFilter.java
+++ b/src/java/org/apache/cassandra/db/filter/AbstractClusteringIndexFilter.java
@@ -80,7 +80,7 @@ public abstract class AbstractClusteringIndexFilter implements ClusteringIndexFi
public ClusteringIndexFilter deserialize(DataInputPlus in, int version, TableMetadata metadata) throws IOException
{
- Kind kind = Kind.values()[in.readUnsignedByte()];
+ Kind kind = Kind.fromOrdinal(in.readUnsignedByte());
boolean reversed = in.readBoolean();
return kind.deserializer.deserialize(in, version, metadata, reversed);
diff --git a/src/java/org/apache/cassandra/db/filter/ClusteringIndexFilter.java b/src/java/org/apache/cassandra/db/filter/ClusteringIndexFilter.java
index 2ed144a0d1..c3286f60bb 100644
--- a/src/java/org/apache/cassandra/db/filter/ClusteringIndexFilter.java
+++ b/src/java/org/apache/cassandra/db/filter/ClusteringIndexFilter.java
@@ -46,12 +46,19 @@ public interface ClusteringIndexFilter
SLICE (ClusteringIndexSliceFilter.deserializer),
NAMES (ClusteringIndexNamesFilter.deserializer);
+ private static final Kind[] VALUES = values();
+
protected final InternalDeserializer deserializer;
private Kind(InternalDeserializer deserializer)
{
this.deserializer = deserializer;
}
+
+ public static Kind fromOrdinal(int ordinal)
+ {
+ return VALUES[ordinal];
+ }
}
static interface InternalDeserializer
diff --git a/src/java/org/apache/cassandra/db/filter/ColumnSubselection.java b/src/java/org/apache/cassandra/db/filter/ColumnSubselection.java
index b4f0346c76..f2baad98e2 100644
--- a/src/java/org/apache/cassandra/db/filter/ColumnSubselection.java
+++ b/src/java/org/apache/cassandra/db/filter/ColumnSubselection.java
@@ -44,7 +44,17 @@ public abstract class ColumnSubselection implements Comparable
// and this is why we have some UNUSEDX for values we don't use anymore
// (we could clean those on a major protocol update, but it's not worth
// the trouble for now)
- protected enum Kind { SIMPLE, MAP_ELEMENT, UNUSED1, CUSTOM, USER }
+ protected enum Kind
+ {
+ SIMPLE, MAP_ELEMENT, UNUSED1, CUSTOM, USER;
+
+ private static final Kind[] VALUES = values();
+
+ static Kind fromOrdinal(int ordinal)
+ {
+ return VALUES[ordinal];
+ }
+ }
protected abstract Kind kind();
protected final ColumnMetadata column;
@@ -685,7 +695,7 @@ public class RowFilter implements Iterable
public Expression deserialize(DataInputPlus in, int version, TableMetadata metadata) throws IOException
{
- Kind kind = Kind.values()[in.readByte()];
+ Kind kind = Kind.fromOrdinal(in.readByte());
// custom expressions (3.0+ only) do not contain a column or operator, only a value
if (kind == Kind.CUSTOM)
diff --git a/src/java/org/apache/cassandra/db/guardrails/GuardrailsConfigProvider.java b/src/java/org/apache/cassandra/db/guardrails/GuardrailsConfigProvider.java
index 990a07a1ce..ebae5d5394 100644
--- a/src/java/org/apache/cassandra/db/guardrails/GuardrailsConfigProvider.java
+++ b/src/java/org/apache/cassandra/db/guardrails/GuardrailsConfigProvider.java
@@ -65,7 +65,7 @@ public interface GuardrailsConfigProvider
*/
static GuardrailsConfigProvider build(String customImpl)
{
- return FBUtilities.construct(customImpl, "custom guardrails config provider");
+ return FBUtilities.construct(customImpl, "custom guardrails config provider", GuardrailsConfigProvider.class);
}
/**
diff --git a/src/java/org/apache/cassandra/db/guardrails/ValueGenerator.java b/src/java/org/apache/cassandra/db/guardrails/ValueGenerator.java
index c01e18f679..0a11734034 100644
--- a/src/java/org/apache/cassandra/db/guardrails/ValueGenerator.java
+++ b/src/java/org/apache/cassandra/db/guardrails/ValueGenerator.java
@@ -143,8 +143,11 @@ public abstract class ValueGenerator
try
{
+ Class extends ValueGenerator> rawGeneratorClass =
+ FBUtilities.classForNameWithoutInitialization(className, "generator", ValueGenerator.class);
+ @SuppressWarnings("unchecked")
Class extends ValueGenerator> generatorClass =
- FBUtilities.classForName(className, "generator");
+ (Class extends ValueGenerator>) rawGeneratorClass;
@SuppressWarnings("unchecked")
ValueGenerator generator = generatorClass.getConstructor(CustomGuardrailConfig.class)
@@ -165,4 +168,4 @@ public abstract class ValueGenerator
className, message), ex);
}
}
-}
\ No newline at end of file
+}
diff --git a/src/java/org/apache/cassandra/db/guardrails/ValueValidator.java b/src/java/org/apache/cassandra/db/guardrails/ValueValidator.java
index 3aa9533ff5..8e93e129e0 100644
--- a/src/java/org/apache/cassandra/db/guardrails/ValueValidator.java
+++ b/src/java/org/apache/cassandra/db/guardrails/ValueValidator.java
@@ -127,8 +127,11 @@ public abstract class ValueValidator
try
{
+ Class extends ValueValidator> rawValidatorClass =
+ FBUtilities.classForNameWithoutInitialization(className, "validator", ValueValidator.class);
+ @SuppressWarnings("unchecked")
Class extends ValueValidator> validatorClass =
- FBUtilities.classForName(className, "validator");
+ (Class extends ValueValidator>) rawValidatorClass;
@SuppressWarnings("unchecked")
ValueValidator validator = validatorClass.getConstructor(CustomGuardrailConfig.class)
diff --git a/src/java/org/apache/cassandra/db/marshal/AbstractTimeUUIDType.java b/src/java/org/apache/cassandra/db/marshal/AbstractTimeUUIDType.java
index ffa8fbc95b..fa1cd90273 100644
--- a/src/java/org/apache/cassandra/db/marshal/AbstractTimeUUIDType.java
+++ b/src/java/org/apache/cassandra/db/marshal/AbstractTimeUUIDType.java
@@ -39,7 +39,7 @@ public abstract class AbstractTimeUUIDType extends TemporalType
{
AbstractTimeUUIDType()
{
- super(ComparisonType.CUSTOM);
+ super(ComparisonType.CUSTOM, 16);
} // singleton
@Override
@@ -193,12 +193,6 @@ public abstract class AbstractTimeUUIDType extends TemporalType
return super.decomposeUntyped(value);
}
- @Override
- public int valueLengthIfFixed()
- {
- return 16;
- }
-
@Override
public long toTimeInMillis(ByteBuffer value)
{
diff --git a/src/java/org/apache/cassandra/db/marshal/AbstractType.java b/src/java/org/apache/cassandra/db/marshal/AbstractType.java
index 09545d486e..257a64f931 100644
--- a/src/java/org/apache/cassandra/db/marshal/AbstractType.java
+++ b/src/java/org/apache/cassandra/db/marshal/AbstractType.java
@@ -90,11 +90,18 @@ public abstract class AbstractType implements Comparator, Assignm
public final ComparisonType comparisonType;
public final boolean isByteOrderComparable;
public final ValueComparators comparatorSet;
+ private final int valueLengthIfFixed;
protected AbstractType(ComparisonType comparisonType)
+ {
+ this(comparisonType, VARIABLE_LENGTH);
+ }
+
+ protected AbstractType(ComparisonType comparisonType, int valueLengthIfFixed)
{
this.comparisonType = comparisonType;
this.isByteOrderComparable = comparisonType == ComparisonType.BYTE_ORDER;
+ this.valueLengthIfFixed = valueLengthIfFixed;
reverseComparator = (o1, o2) -> AbstractType.this.compare(o2, o1);
try
{
@@ -384,12 +391,12 @@ public abstract class AbstractType implements Comparator, Assignm
/**
* Similar to {@link #isValueCompatibleWith(AbstractType)}, but takes into account {@link Cell} encoding.
* In particular, this method doesn't consider two types serialization compatible if one of them has fixed
- * length (overrides {@link #valueLengthIfFixed()}, and the other one doesn't.
+ * length, and the other one doesn't.
*/
public boolean isSerializationCompatibleWith(AbstractType> previous)
{
return isValueCompatibleWith(previous)
- && valueLengthIfFixed() == previous.valueLengthIfFixed()
+ && valueLengthIfFixed == previous.valueLengthIfFixed
&& isMultiCell() == previous.isMultiCell();
}
@@ -498,7 +505,7 @@ public abstract class AbstractType implements Comparator, Assignm
*/
public int valueLengthIfFixed()
{
- return VARIABLE_LENGTH;
+ return valueLengthIfFixed;
}
/**
@@ -508,7 +515,7 @@ public abstract class AbstractType implements Comparator, Assignm
*/
public final boolean isValueLengthFixed()
{
- return valueLengthIfFixed() != VARIABLE_LENGTH;
+ return valueLengthIfFixed != VARIABLE_LENGTH;
}
/**
@@ -570,7 +577,7 @@ public abstract class AbstractType implements Comparator, Assignm
public void writeValue(V value, ValueAccessor accessor, DataOutputPlus out) throws IOException
{
assert !isNull(value, accessor) : "bytes should not be null for type " + this;
- int expectedValueLength = valueLengthIfFixed();
+ int expectedValueLength = valueLengthIfFixed;
if (expectedValueLength >= 0)
{
int actualValueLength = accessor.size(value);
@@ -589,7 +596,7 @@ public abstract class AbstractType implements Comparator, Assignm
public void writeValue(IndexedValueHolder valueHolder, int i, ValueAccessor accessor, DataOutputPlus out) throws IOException
{
assert !valueHolder.isNull(i) : "bytes should not be null for type " + this;
- int expectedValueLength = valueLengthIfFixed();
+ int expectedValueLength = valueLengthIfFixed;
if (expectedValueLength >= 0)
{
int actualValueLength = valueHolder.size(i);
@@ -613,7 +620,7 @@ public abstract class AbstractType implements Comparator, Assignm
public long writtenLength(V value, ValueAccessor accessor)
{
assert !accessor.isEmpty(value) : "bytes should not be empty for type " + this;
- return valueLengthIfFixed() >= 0
+ return valueLengthIfFixed >= 0
? accessor.size(value) // if the size is wrong, this will be detected in writeValue
: accessor.sizeWithVIntLength(value);
}
@@ -621,7 +628,7 @@ public abstract class AbstractType implements Comparator, Assignm
public long writtenLength(IndexedValueHolder valueHolder, int i, ValueAccessor accessor)
{
assert !valueHolder.isNull(i) : "bytes should not be null for type " + this;
- return valueLengthIfFixed() >= 0
+ return valueLengthIfFixed >= 0
? valueHolder.size(i) // if the size is wrong, this will be detected in writeValue
: accessor.sizeWithVIntLength(valueHolder, i);
}
@@ -643,7 +650,7 @@ public abstract class AbstractType implements Comparator, Assignm
public V read(ValueAccessor accessor, DataInputPlus in, int maxValueSize) throws IOException
{
- int length = valueLengthIfFixed();
+ int length = valueLengthIfFixed;
if (length >= 0)
return accessor.read(in, length);
@@ -664,7 +671,7 @@ public abstract class AbstractType implements Comparator, Assignm
public void skipValue(DataInputPlus in) throws IOException
{
- int length = valueLengthIfFixed();
+ int length = valueLengthIfFixed;
if (length >= 0)
in.skipBytesFully(length);
else
diff --git a/src/java/org/apache/cassandra/db/marshal/BooleanType.java b/src/java/org/apache/cassandra/db/marshal/BooleanType.java
index 171b239ad2..7f6238bbd6 100644
--- a/src/java/org/apache/cassandra/db/marshal/BooleanType.java
+++ b/src/java/org/apache/cassandra/db/marshal/BooleanType.java
@@ -36,7 +36,7 @@ public class BooleanType extends AbstractType
private static final ArgumentDeserializer ARGUMENT_DESERIALIZER = new DefaultArgumentDeserializer(instance);
private static final ByteBuffer MASKED_VALUE = instance.decompose(false);
- BooleanType() {super(ComparisonType.CUSTOM);} // singleton
+ BooleanType() {super(ComparisonType.CUSTOM, 1);} // singleton
@Override
public boolean allowsEmpty()
@@ -127,12 +127,6 @@ public class BooleanType extends AbstractType
return ARGUMENT_DESERIALIZER;
}
- @Override
- public int valueLengthIfFixed()
- {
- return 1;
- }
-
@Override
public ByteBuffer getMaskedValue()
{
diff --git a/src/java/org/apache/cassandra/db/marshal/DateType.java b/src/java/org/apache/cassandra/db/marshal/DateType.java
index 87a31ade55..6442c159a0 100644
--- a/src/java/org/apache/cassandra/db/marshal/DateType.java
+++ b/src/java/org/apache/cassandra/db/marshal/DateType.java
@@ -49,7 +49,7 @@ public class DateType extends AbstractType
private static final ArgumentDeserializer ARGUMENT_DESERIALIZER = new DefaultArgumentDeserializer(instance);
private static final ByteBuffer MASKED_VALUE = instance.decompose(new Date(0));
- DateType() {super(ComparisonType.BYTE_ORDER);} // singleton
+ DateType() {super(ComparisonType.BYTE_ORDER, 8);} // singleton
public boolean isEmptyValueMeaningless()
{
@@ -143,12 +143,6 @@ public class DateType extends AbstractType
return ARGUMENT_DESERIALIZER;
}
- @Override
- public int valueLengthIfFixed()
- {
- return 8;
- }
-
@Override
public ByteBuffer getMaskedValue()
{
diff --git a/src/java/org/apache/cassandra/db/marshal/DoubleType.java b/src/java/org/apache/cassandra/db/marshal/DoubleType.java
index 2d065e299d..2e423e91bf 100644
--- a/src/java/org/apache/cassandra/db/marshal/DoubleType.java
+++ b/src/java/org/apache/cassandra/db/marshal/DoubleType.java
@@ -40,7 +40,7 @@ public class DoubleType extends NumberType
private static final ByteBuffer MASKED_VALUE = instance.decompose(0d);
- DoubleType() {super(ComparisonType.CUSTOM);} // singleton
+ DoubleType() {super(ComparisonType.CUSTOM, 8);} // singleton
@Override
public boolean allowsEmpty()
@@ -145,12 +145,6 @@ public class DoubleType extends NumberType
};
}
- @Override
- public int valueLengthIfFixed()
- {
- return 8;
- }
-
@Override
public ByteBuffer add(Number left, Number right)
{
diff --git a/src/java/org/apache/cassandra/db/marshal/EmptyType.java b/src/java/org/apache/cassandra/db/marshal/EmptyType.java
index 23553111f7..98c82c87a1 100644
--- a/src/java/org/apache/cassandra/db/marshal/EmptyType.java
+++ b/src/java/org/apache/cassandra/db/marshal/EmptyType.java
@@ -71,7 +71,7 @@ public class EmptyType extends AbstractType
public static final EmptyType instance = new EmptyType();
- private EmptyType() {super(ComparisonType.CUSTOM);} // singleton
+ private EmptyType() {super(ComparisonType.CUSTOM, 0);} // singleton
@Override
public ByteSource asComparableBytes(ValueAccessor accessor, V data, ByteComparable.Version version)
@@ -137,12 +137,6 @@ public class EmptyType extends AbstractType
throw new UnsupportedOperationException();
}
- @Override
- public int valueLengthIfFixed()
- {
- return 0;
- }
-
@Override
public long writtenLength(V value, ValueAccessor accessor)
{
diff --git a/src/java/org/apache/cassandra/db/marshal/FloatType.java b/src/java/org/apache/cassandra/db/marshal/FloatType.java
index 9ccb7fca18..38bc5de766 100644
--- a/src/java/org/apache/cassandra/db/marshal/FloatType.java
+++ b/src/java/org/apache/cassandra/db/marshal/FloatType.java
@@ -41,7 +41,7 @@ public class FloatType extends NumberType
private static final ByteBuffer MASKED_VALUE = instance.decompose(0f);
- FloatType() {super(ComparisonType.CUSTOM);} // singleton
+ FloatType() {super(ComparisonType.CUSTOM, 4);} // singleton
@Override
public boolean allowsEmpty()
@@ -146,12 +146,6 @@ public class FloatType extends NumberType
};
}
- @Override
- public int valueLengthIfFixed()
- {
- return 4;
- }
-
@Override
public ByteBuffer add(Number left, Number right)
{
diff --git a/src/java/org/apache/cassandra/db/marshal/Int32Type.java b/src/java/org/apache/cassandra/db/marshal/Int32Type.java
index a68255edd3..fea3e1f366 100644
--- a/src/java/org/apache/cassandra/db/marshal/Int32Type.java
+++ b/src/java/org/apache/cassandra/db/marshal/Int32Type.java
@@ -43,7 +43,7 @@ public class Int32Type extends NumberType
Int32Type()
{
- super(ComparisonType.CUSTOM);
+ super(ComparisonType.CUSTOM, 4);
} // singleton
@Override
@@ -152,12 +152,6 @@ public class Int32Type extends NumberType
};
}
- @Override
- public int valueLengthIfFixed()
- {
- return 4;
- }
-
@Override
public ByteBuffer add(Number left, Number right)
{
diff --git a/src/java/org/apache/cassandra/db/marshal/LexicalUUIDType.java b/src/java/org/apache/cassandra/db/marshal/LexicalUUIDType.java
index fab451b291..4ee3024efb 100644
--- a/src/java/org/apache/cassandra/db/marshal/LexicalUUIDType.java
+++ b/src/java/org/apache/cassandra/db/marshal/LexicalUUIDType.java
@@ -44,7 +44,7 @@ public class LexicalUUIDType extends AbstractType
LexicalUUIDType()
{
- super(ComparisonType.CUSTOM);
+ super(ComparisonType.CUSTOM, 16);
} // singleton
@Override
@@ -148,12 +148,6 @@ public class LexicalUUIDType extends AbstractType
return ARGUMENT_DESERIALIZER;
}
- @Override
- public int valueLengthIfFixed()
- {
- return 16;
- }
-
@Override
public ByteBuffer getMaskedValue()
{
diff --git a/src/java/org/apache/cassandra/db/marshal/LongType.java b/src/java/org/apache/cassandra/db/marshal/LongType.java
index 400b6ae836..41ae32186a 100644
--- a/src/java/org/apache/cassandra/db/marshal/LongType.java
+++ b/src/java/org/apache/cassandra/db/marshal/LongType.java
@@ -41,7 +41,7 @@ public class LongType extends NumberType
private static final ByteBuffer MASKED_VALUE = instance.decompose(0L);
- LongType() {super(ComparisonType.CUSTOM);} // singleton
+ LongType() {super(ComparisonType.CUSTOM, 8);} // singleton
@Override
public boolean allowsEmpty()
@@ -170,12 +170,6 @@ public class LongType extends NumberType
};
}
- @Override
- public int valueLengthIfFixed()
- {
- return 8;
- }
-
@Override
public ByteBuffer add(Number left, Number right)
{
diff --git a/src/java/org/apache/cassandra/db/marshal/MultiElementType.java b/src/java/org/apache/cassandra/db/marshal/MultiElementType.java
index dc8f912e09..2bbede82ef 100644
--- a/src/java/org/apache/cassandra/db/marshal/MultiElementType.java
+++ b/src/java/org/apache/cassandra/db/marshal/MultiElementType.java
@@ -38,6 +38,11 @@ public abstract class MultiElementType extends AbstractType
super(comparisonType);
}
+ protected MultiElementType(ComparisonType comparisonType, int valueLengthIfFixed)
+ {
+ super(comparisonType, valueLengthIfFixed);
+ }
+
/**
* Returns the serialized representation of the value composed of the specified elements.
*
@@ -133,4 +138,3 @@ public abstract class MultiElementType extends AbstractType
throw new UnsupportedOperationException(this + " does not support retrieving elements by key or index");
}
}
-
diff --git a/src/java/org/apache/cassandra/db/marshal/NumberType.java b/src/java/org/apache/cassandra/db/marshal/NumberType.java
index 04c6e1dc1b..c007dac44d 100644
--- a/src/java/org/apache/cassandra/db/marshal/NumberType.java
+++ b/src/java/org/apache/cassandra/db/marshal/NumberType.java
@@ -34,6 +34,11 @@ public abstract class NumberType extends AbstractType
super(comparisonType);
}
+ protected NumberType(ComparisonType comparisonType, int valueLengthIfFixed)
+ {
+ super(comparisonType, valueLengthIfFixed);
+ }
+
/**
* Checks if this type support floating point numbers.
* @return {@code true} if this type support floating point numbers, {@code false} otherwise.
diff --git a/src/java/org/apache/cassandra/db/marshal/ReversedType.java b/src/java/org/apache/cassandra/db/marshal/ReversedType.java
index 5c81f0e558..85e4b44273 100644
--- a/src/java/org/apache/cassandra/db/marshal/ReversedType.java
+++ b/src/java/org/apache/cassandra/db/marshal/ReversedType.java
@@ -57,7 +57,7 @@ public class ReversedType extends AbstractType
private ReversedType(AbstractType baseType)
{
- super(ComparisonType.CUSTOM);
+ super(ComparisonType.CUSTOM, baseType.valueLengthIfFixed());
this.baseType = baseType;
}
@@ -181,12 +181,6 @@ public class ReversedType extends AbstractType
return getInstance(baseType.withUpdatedUserType(udt));
}
- @Override
- public int valueLengthIfFixed()
- {
- return baseType.valueLengthIfFixed();
- }
-
@Override
public boolean isReversed()
{
diff --git a/src/java/org/apache/cassandra/db/marshal/TemporalType.java b/src/java/org/apache/cassandra/db/marshal/TemporalType.java
index 5f5c4564d5..24f0eaa7a4 100644
--- a/src/java/org/apache/cassandra/db/marshal/TemporalType.java
+++ b/src/java/org/apache/cassandra/db/marshal/TemporalType.java
@@ -37,6 +37,11 @@ public abstract class TemporalType extends AbstractType
super(comparisonType);
}
+ protected TemporalType(ComparisonType comparisonType, int valueLengthIfFixed)
+ {
+ super(comparisonType, valueLengthIfFixed);
+ }
+
/**
* Returns the current temporal value.
* @return the current temporal value.
diff --git a/src/java/org/apache/cassandra/db/marshal/TimestampType.java b/src/java/org/apache/cassandra/db/marshal/TimestampType.java
index b408fdda53..68ac6c8919 100644
--- a/src/java/org/apache/cassandra/db/marshal/TimestampType.java
+++ b/src/java/org/apache/cassandra/db/marshal/TimestampType.java
@@ -53,7 +53,7 @@ public class TimestampType extends TemporalType
private static final ByteBuffer MASKED_VALUE = instance.decompose(new Date(0));
- private TimestampType() {super(ComparisonType.CUSTOM);} // singleton
+ private TimestampType() {super(ComparisonType.CUSTOM, 8);} // singleton
@Override
public boolean allowsEmpty()
@@ -166,12 +166,6 @@ public class TimestampType extends TemporalType
return TimestampSerializer.instance;
}
- @Override
- public int valueLengthIfFixed()
- {
- return 8;
- }
-
@Override
protected void validateDuration(Duration duration)
{
diff --git a/src/java/org/apache/cassandra/db/marshal/TypeParser.java b/src/java/org/apache/cassandra/db/marshal/TypeParser.java
index b2dc0202bc..d3a20fc7ca 100644
--- a/src/java/org/apache/cassandra/db/marshal/TypeParser.java
+++ b/src/java/org/apache/cassandra/db/marshal/TypeParser.java
@@ -449,8 +449,7 @@ public class TypeParser
private static AbstractType> getAbstractType(String compareWith) throws ConfigurationException
{
- String className = compareWith.contains(".") ? compareWith : "org.apache.cassandra.db.marshal." + compareWith;
- Class extends AbstractType>> typeClass = FBUtilities.>classForName(className, "abstract-type");
+ Class extends AbstractType>> typeClass = getAbstractTypeClass(compareWith);
try
{
Field field = typeClass.getDeclaredField("instance");
@@ -465,8 +464,7 @@ public class TypeParser
private static AbstractType> getAbstractType(String compareWith, TypeParser parser) throws SyntaxException, ConfigurationException
{
- String className = compareWith.contains(".") ? compareWith : "org.apache.cassandra.db.marshal." + compareWith;
- Class extends AbstractType>> typeClass = FBUtilities.>classForName(className, "abstract-type");
+ Class extends AbstractType>> typeClass = getAbstractTypeClass(compareWith);
if (PseudoUtf8Type.class.isAssignableFrom(typeClass))
{
if (StorageService.instance.isDaemonSetupCompleted())
@@ -491,6 +489,19 @@ public class TypeParser
}
}
+ private static Class extends AbstractType>> getAbstractTypeClass(String compareWith) throws ConfigurationException
+ {
+ String className = compareWith.contains(".") ? compareWith : "org.apache.cassandra.db.marshal." + compareWith;
+ // Defer class initialization until after confirming this is an AbstractType. The static instance field
+ // access or getInstance(TypeParser) invocation below performs the initialization for valid types.
+ @SuppressWarnings("unchecked")
+ Class extends AbstractType>> typeClass =
+ (Class extends AbstractType>>) FBUtilities.classForNameWithoutInitialization(className,
+ "abstract-type",
+ AbstractType.class);
+ return typeClass;
+ }
+
private static AbstractType> getRawAbstractType(Class extends AbstractType>> typeClass) throws ConfigurationException
{
try
diff --git a/src/java/org/apache/cassandra/db/marshal/UUIDType.java b/src/java/org/apache/cassandra/db/marshal/UUIDType.java
index 2b526f3cef..cf36f99dee 100644
--- a/src/java/org/apache/cassandra/db/marshal/UUIDType.java
+++ b/src/java/org/apache/cassandra/db/marshal/UUIDType.java
@@ -56,7 +56,7 @@ public class UUIDType extends AbstractType
UUIDType()
{
- super(ComparisonType.CUSTOM);
+ super(ComparisonType.CUSTOM, 16);
}
@Override
@@ -252,12 +252,6 @@ public class UUIDType extends AbstractType
return (uuid.get(6) & 0xf0) >> 4;
}
- @Override
- public int valueLengthIfFixed()
- {
- return 16;
- }
-
@Override
public ByteBuffer getMaskedValue()
{
diff --git a/src/java/org/apache/cassandra/db/marshal/VectorType.java b/src/java/org/apache/cassandra/db/marshal/VectorType.java
index 3cdf72270e..74123c126e 100644
--- a/src/java/org/apache/cassandra/db/marshal/VectorType.java
+++ b/src/java/org/apache/cassandra/db/marshal/VectorType.java
@@ -82,20 +82,16 @@ public final class VectorType extends MultiElementType>
public final AbstractType elementType;
public final int dimension;
private final TypeSerializer elementSerializer;
- private final int valueLengthIfFixed;
private final VectorSerializer serializer;
private VectorType(AbstractType elementType, int dimension)
{
- super(ComparisonType.CUSTOM);
+ super(ComparisonType.CUSTOM, valueLengthIfFixed(elementType, dimension));
if (dimension <= 0)
throw new InvalidRequestException(String.format("vectors may only have positive dimensions; given %d", dimension));
this.elementType = elementType;
this.dimension = dimension;
this.elementSerializer = elementType.getSerializer();
- this.valueLengthIfFixed = elementType.isValueLengthFixed() ?
- elementType.valueLengthIfFixed() * dimension :
- super.valueLengthIfFixed();
this.serializer = elementType.isValueLengthFixed() ?
new FixedLengthSerializer() :
new VariableLengthSerializer();
@@ -126,10 +122,10 @@ public final class VectorType extends MultiElementType>
return getSerializer().compareCustom(left, accessorL, right, accessorR);
}
- @Override
- public int valueLengthIfFixed()
+ private static int valueLengthIfFixed(AbstractType> elementType, int dimension)
{
- return valueLengthIfFixed;
+ int elementLength = elementType.valueLengthIfFixed();
+ return elementLength >= 0 ? elementLength * dimension : elementLength;
}
@Override
diff --git a/src/java/org/apache/cassandra/db/rows/AbstractCell.java b/src/java/org/apache/cassandra/db/rows/AbstractCell.java
index c398839742..ea6110cbf6 100644
--- a/src/java/org/apache/cassandra/db/rows/AbstractCell.java
+++ b/src/java/org/apache/cassandra/db/rows/AbstractCell.java
@@ -61,7 +61,18 @@ public abstract class AbstractCell extends Cell
public boolean isTombstone()
{
- return localDeletionTime() != NO_DELETION_TIME && ttl() == NO_TTL;
+ return isTombstone(localDeletionTime());
+ }
+
+ public long minDeletionTime()
+ {
+ long localDeletionTime = localDeletionTime();
+ return isTombstone(localDeletionTime) ? Long.MIN_VALUE : localDeletionTime;
+ }
+
+ private boolean isTombstone(long localDeletionTime)
+ {
+ return localDeletionTime != NO_DELETION_TIME && ttl() == NO_TTL;
}
public boolean isExpiring()
diff --git a/src/java/org/apache/cassandra/db/rows/ArrayCell.java b/src/java/org/apache/cassandra/db/rows/ArrayCell.java
index 77d44c3eb7..c74103f2e0 100644
--- a/src/java/org/apache/cassandra/db/rows/ArrayCell.java
+++ b/src/java/org/apache/cassandra/db/rows/ArrayCell.java
@@ -31,17 +31,12 @@ import org.apache.cassandra.utils.memory.ByteBufferCloner;
import static org.apache.cassandra.utils.ByteArrayUtil.EMPTY_BYTE_ARRAY;
-public class ArrayCell extends AbstractCell
+public class ArrayCell extends HeapAbstractCell
{
private static final long EMPTY_SIZE = ObjectSizes.measure(new ArrayCell(ColumnMetadata.regularColumn("", "", "", ByteType.instance, ColumnMetadata.NO_UNIQUE_ID), 0L, 0, 0, EMPTY_BYTE_ARRAY, null));
// Careful: Adding vars here has an impact on memtable size
- private final long timestamp;
- private final int ttl;
- private final int localDeletionTimeUnsignedInteger;
-
private final byte[] value;
- private final CellPath path;
// Please keep both int/long overloaded ctros public. Otherwise silent casts will mess timestamps when one is not
// available.
@@ -52,12 +47,8 @@ public class ArrayCell extends AbstractCell
public ArrayCell(ColumnMetadata column, long timestamp, int ttl, int localDeletionTimeUnsignedInteger, byte[] value, CellPath path)
{
- super(column);
- this.timestamp = timestamp;
- this.ttl = ttl;
- this.localDeletionTimeUnsignedInteger = localDeletionTimeUnsignedInteger;
+ super(column, timestamp, ttl, localDeletionTimeUnsignedInteger, path);
this.value = value;
- this.path = path;
}
public static ArrayCell live(ColumnMetadata column, long timestamp, byte[] value, CellPath path)
@@ -71,16 +62,6 @@ public class ArrayCell extends AbstractCell
return new ArrayCell(column, timestamp, ttl, ExpirationDateOverflowHandling.computeLocalExpirationTime(nowInSec, ttl), value, path);
}
- public long timestamp()
- {
- return timestamp;
- }
-
- public int ttl()
- {
- return ttl;
- }
-
public byte[] value()
{
return value;
@@ -91,10 +72,6 @@ public class ArrayCell extends AbstractCell
return ByteArrayAccessor.instance;
}
- public CellPath path()
- {
- return path;
- }
public Cell> withUpdatedColumn(ColumnMetadata newColumn)
{
@@ -144,9 +121,4 @@ public class ArrayCell extends AbstractCell
return EMPTY_SIZE + ObjectSizes.sizeOfArray(value) - value.length + (path == null ? 0 : path.unsharedHeapSizeExcludingData());
}
- @Override
- protected int localDeletionTimeAsUnsignedInt()
- {
- return localDeletionTimeUnsignedInteger;
- }
}
diff --git a/src/java/org/apache/cassandra/db/rows/BTreeRow.java b/src/java/org/apache/cassandra/db/rows/BTreeRow.java
index b8df8450c1..82cf43c6a7 100644
--- a/src/java/org/apache/cassandra/db/rows/BTreeRow.java
+++ b/src/java/org/apache/cassandra/db/rows/BTreeRow.java
@@ -172,7 +172,7 @@ public class BTreeRow extends AbstractRow
private static long minDeletionTime(Cell> cell)
{
- return cell.isTombstone() ? Long.MIN_VALUE : cell.localDeletionTime();
+ return cell.minDeletionTime();
}
private static long minDeletionTime(LivenessInfo info)
@@ -439,6 +439,11 @@ public class BTreeRow extends AbstractRow
return nowInSec >= minLocalDeletionTime;
}
+ public long minLocalDeletionTime()
+ {
+ return minLocalDeletionTime;
+ }
+
public boolean hasInvalidDeletions()
{
if (primaryKeyLivenessInfo().isExpiring() && (primaryKeyLivenessInfo().ttl() < 0 || primaryKeyLivenessInfo().localExpirationTime() < 0))
diff --git a/src/java/org/apache/cassandra/db/rows/BufferCell.java b/src/java/org/apache/cassandra/db/rows/BufferCell.java
index fa0826fa57..3f8de1d5fe 100644
--- a/src/java/org/apache/cassandra/db/rows/BufferCell.java
+++ b/src/java/org/apache/cassandra/db/rows/BufferCell.java
@@ -30,17 +30,12 @@ import org.apache.cassandra.utils.memory.ByteBufferCloner;
import static java.lang.String.format;
-public class BufferCell extends AbstractCell
+public class BufferCell extends HeapAbstractCell
{
private static final long EMPTY_SIZE = ObjectSizes.measure(new BufferCell(ColumnMetadata.regularColumn("", "", "", ByteType.instance, ColumnMetadata.NO_UNIQUE_ID), 0L, 0, 0, ByteBufferUtil.EMPTY_BYTE_BUFFER, null));
// Careful: Adding vars here has an impact on memtable size
- private final long timestamp;
- private final int ttl;
- private final int localDeletionTimeUnsignedInteger;
-
private final ByteBuffer value;
- private final CellPath path;
// Please keep both int/long overloaded ctros public. Otherwise silent casts will mess timestamps when one is not
// available.
@@ -51,14 +46,10 @@ public class BufferCell extends AbstractCell
public BufferCell(ColumnMetadata column, long timestamp, int ttl, int localDeletionTimeUnsignedInteger, ByteBuffer value, CellPath path)
{
- super(column);
+ super(column, timestamp, ttl, localDeletionTimeUnsignedInteger, path);
assert !column.isPrimaryKeyColumn();
assert column.isComplex() == (path != null) : format("Column %s.%s(%s: %s) isComplex: %b with cellpath: %s", column.ksName, column.cfName, column.name, column.type.toString(), column.isComplex(), path);
- this.timestamp = timestamp;
- this.ttl = ttl;
- this.localDeletionTimeUnsignedInteger = localDeletionTimeUnsignedInteger;
this.value = value;
- this.path = path;
}
public static BufferCell live(ColumnMetadata column, long timestamp, ByteBuffer value)
@@ -92,16 +83,6 @@ public class BufferCell extends AbstractCell
return new BufferCell(column, timestamp, NO_TTL, nowInSec, ByteBufferUtil.EMPTY_BYTE_BUFFER, path);
}
- public long timestamp()
- {
- return timestamp;
- }
-
- public int ttl()
- {
- return ttl;
- }
-
public ByteBuffer value()
{
return value;
@@ -112,10 +93,6 @@ public class BufferCell extends AbstractCell
return ByteBufferAccessor.instance;
}
- public CellPath path()
- {
- return path;
- }
public Cell> withUpdatedColumn(ColumnMetadata newColumn)
{
@@ -163,10 +140,4 @@ public class BufferCell extends AbstractCell
{
return EMPTY_SIZE + ObjectSizes.sizeOnHeapExcludingDataOf(value) + (path == null ? 0 : path.unsharedHeapSizeExcludingData());
}
-
- @Override
- protected int localDeletionTimeAsUnsignedInt()
- {
- return localDeletionTimeUnsignedInteger;
- }
}
diff --git a/src/java/org/apache/cassandra/db/rows/Cell.java b/src/java/org/apache/cassandra/db/rows/Cell.java
index 6dea396993..5a76d1bcd4 100644
--- a/src/java/org/apache/cassandra/db/rows/Cell.java
+++ b/src/java/org/apache/cassandra/db/rows/Cell.java
@@ -150,6 +150,8 @@ public abstract class Cell extends ColumnData
return deletionTimeUnsignedIntegerToLong(localDeletionTimeAsUnsignedInt());
}
+ public abstract long minDeletionTime();
+
/**
* Whether the cell is a tombstone or not.
*
diff --git a/src/java/org/apache/cassandra/db/rows/HeapAbstractCell.java b/src/java/org/apache/cassandra/db/rows/HeapAbstractCell.java
new file mode 100644
index 0000000000..d44186b2e9
--- /dev/null
+++ b/src/java/org/apache/cassandra/db/rows/HeapAbstractCell.java
@@ -0,0 +1,64 @@
+/*
+ * 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.rows;
+
+import org.apache.cassandra.schema.ColumnMetadata;
+
+public abstract class HeapAbstractCell extends AbstractCell
+{
+ // Careful: Adding vars here has an impact on memtable size
+ protected final long timestamp;
+ protected final int ttl;
+ protected final int localDeletionTimeUnsignedInteger;
+
+ protected final CellPath path;
+
+ protected HeapAbstractCell(ColumnMetadata column, long timestamp, int ttl, int localDeletionTimeUnsignedInteger, CellPath path)
+ {
+ super(column);
+ this.timestamp = timestamp;
+ this.ttl = ttl;
+ this.localDeletionTimeUnsignedInteger = localDeletionTimeUnsignedInteger;
+ this.path = path;
+ }
+
+ @Override
+ public long timestamp()
+ {
+ return timestamp;
+ }
+
+ @Override
+ public int ttl()
+ {
+ return ttl;
+ }
+
+ @Override
+ protected int localDeletionTimeAsUnsignedInt()
+ {
+ return localDeletionTimeUnsignedInteger;
+ }
+
+ @Override
+ public CellPath path()
+ {
+ return path;
+ }
+}
diff --git a/src/java/org/apache/cassandra/db/rows/Row.java b/src/java/org/apache/cassandra/db/rows/Row.java
index 9c281408ff..4137c42ca2 100644
--- a/src/java/org/apache/cassandra/db/rows/Row.java
+++ b/src/java/org/apache/cassandra/db/rows/Row.java
@@ -20,7 +20,6 @@ package org.apache.cassandra.db.rows;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.Collection;
-import java.util.Collections;
import java.util.Comparator;
import java.util.Iterator;
import java.util.List;
@@ -44,9 +43,10 @@ import org.apache.cassandra.schema.TableMetadata;
import org.apache.cassandra.service.paxos.Commit;
import org.apache.cassandra.utils.BiLongAccumulator;
import org.apache.cassandra.utils.BulkIterator;
+import org.apache.cassandra.utils.ComplexCellMergeIterator;
import org.apache.cassandra.utils.LongAccumulator;
-import org.apache.cassandra.utils.MergeIterator;
import org.apache.cassandra.utils.ObjectSizes;
+import org.apache.cassandra.utils.RowMergeIterator;
import org.apache.cassandra.utils.SearchIterator;
import org.apache.cassandra.utils.btree.BTree;
import org.apache.cassandra.utils.btree.UpdateFunction;
@@ -221,6 +221,15 @@ public interface Row extends Unfiltered, Iterable, IMeasurableMemory
*/
public boolean hasDeletion(long nowInSec);
+ /**
+ * The smallest local deletion time of all the data in this row (row deletion, primary key liveness, cells and
+ * complex deletions), or {@link Cell#MAX_DELETION_TIME} if the row has no deletion nor expiring data.
+ *
+ * Unlike {@link #hasDeletion(long)}, this value is independent of the current time. In particular, a value of
+ * {@link Cell#MAX_DELETION_TIME} guarantees the row carries neither tombstones nor expiring data.
+ */
+ public long minLocalDeletionTime();
+
/**
* An iterator to efficiently search data for a given column.
*
@@ -762,6 +771,11 @@ public interface Row extends Unfiltered, Iterable, IMeasurableMemory
LivenessInfo rowInfo = LivenessInfo.EMPTY;
Deletion rowDeletion = Deletion.LIVE;
+ int columnsCountEstimation = 0;
+ // Track the smallest local deletion time across all inputs: if none of them carries any deletion or
+ // expiring data (i.e. this stays at MAX_DELETION_TIME), the merged row can't either, so we can hand the
+ // value to BTreeRow.create() below and skip the full btree scan it would otherwise do to recompute it.
+ long minDeletionTime = Cell.MAX_DELETION_TIME;
for (Row row : rows)
{
if (row == null)
@@ -771,6 +785,11 @@ public interface Row extends Unfiltered, Iterable, IMeasurableMemory
rowInfo = row.primaryKeyLivenessInfo();
if (row.deletion().supersedes(rowDeletion))
rowDeletion = row.deletion();
+
+ minDeletionTime = Math.min(minDeletionTime, row.minLocalDeletionTime());
+
+ columnDataIterators.add(row.iterator());
+ columnsCountEstimation = Math.max(columnsCountEstimation, row.columnCount());
}
if (rowDeletion.isShadowedBy(rowInfo))
@@ -784,25 +803,12 @@ public interface Row extends Unfiltered, Iterable, IMeasurableMemory
if (activeDeletion.deletes(rowInfo))
rowInfo = LivenessInfo.EMPTY;
- int columnsCountEstimation = 0;
- for (Row row : rows)
- {
- if (row != null)
- {
- columnDataIterators.add(row.iterator());
- columnsCountEstimation = Math.max(columnsCountEstimation, row.columnCount());
- }
- else
- {
- columnDataIterators.add(Collections.emptyIterator());
- }
- }
// try to estimate and set a potential target capacity
if (dataBuffer.length < columnsCountEstimation)
dataBuffer = new ColumnData[columnsCountEstimation];
columnDataReducer.setActiveDeletion(activeDeletion);
- Iterator merged = MergeIterator.get(columnDataIterators, ColumnData.comparator, columnDataReducer);
+ Iterator merged = RowMergeIterator.get(columnDataIterators, ColumnData.comparator, columnDataReducer);
while (merged.hasNext())
{
ColumnData data = merged.next();
@@ -819,8 +825,12 @@ public interface Row extends Unfiltered, Iterable, IMeasurableMemory
try (BulkIterator it = BulkIterator.of(dataBuffer))
{
- return BTreeRow.create(clustering, rowInfo, rowDeletion,
- BTree.build(it, dataBufferSize, UpdateFunction.noOp()));
+ Object[] tree = BTree.build(it, dataBufferSize, UpdateFunction.noOp());
+ // If none of the merged rows had any deletion or expiring data, neither does the result, so we can
+ // pass the already-known min local deletion time and avoid rescanning the whole btree to recompute it.
+ return minDeletionTime == Cell.MAX_DELETION_TIME
+ ? BTreeRow.create(clustering, rowInfo, rowDeletion, tree, Cell.MAX_DELETION_TIME)
+ : BTreeRow.create(clustering, rowInfo, rowDeletion, tree);
}
}
@@ -841,10 +851,11 @@ public interface Row extends Unfiltered, Iterable, IMeasurableMemory
return rows;
}
- private static class ColumnDataReducer extends MergeIterator.Reducer
+ private static class ColumnDataReducer extends RowMergeIterator.Reducer
{
private ColumnMetadata column;
- private final List versions;
+ private final ColumnData[] versions;
+ private int versionsSize;
private DeletionTime activeDeletion;
@@ -854,7 +865,7 @@ public interface Row extends Unfiltered, Iterable, IMeasurableMemory
public ColumnDataReducer(int size, boolean hasComplex)
{
- this.versions = new ArrayList<>(size);
+ this.versions = new ColumnData[size];
this.complexBuilder = hasComplex ? ComplexColumnData.builder() : null;
this.complexCells = hasComplex ? new ArrayList<>(size) : null;
this.cellReducer = new CellReducer();
@@ -870,20 +881,24 @@ public interface Row extends Unfiltered, Iterable, IMeasurableMemory
if (useColumnMetadata(data.column()))
column = data.column();
- versions.add(data);
+ versions[versionsSize++] = data;
}
/**
- * Determines it the {@code ColumnMetadata} is the one that should be used.
- * @param dataColumn the {@code ColumnMetadata} to use.
- * @return {@code true} if the {@code ColumnMetadata} is the one that should be used, {@code false} otherwise.
+ * Determines whether {@code dataColumn} should replace the currently selected column metadata,
+ * i.e. whether no column has been selected yet or {@code dataColumn} is a newer version.
+ * @param dataColumn the candidate {@code ColumnMetadata} to evaluate.
+ * @return {@code true} if {@code dataColumn} should be used, {@code false} otherwise.
*/
private boolean useColumnMetadata(ColumnMetadata dataColumn)
{
- if (column == null)
+ ColumnMetadata currentColumn = column;
+ if (currentColumn == null)
return true;
+ if (currentColumn == dataColumn)
+ return false;
- return ColumnMetadataVersionComparator.INSTANCE.compare(column, dataColumn) < 0;
+ return ColumnMetadataVersionComparator.INSTANCE.compare(currentColumn, dataColumn) < 0;
}
protected ColumnData getReduced()
@@ -891,9 +906,9 @@ public interface Row extends Unfiltered, Iterable, IMeasurableMemory
if (column.isSimple())
{
Cell> merged = null;
- for (int i=0, isize=versions.size(); i cell = (Cell>) versions.get(i);
+ Cell> cell = (Cell>) versions[i];
if (!activeDeletion.deletes(cell))
merged = merged == null ? cell : Cells.reconcile(merged, cell);
}
@@ -904,9 +919,9 @@ public interface Row extends Unfiltered, Iterable, IMeasurableMemory
complexBuilder.newColumn(column);
complexCells.clear();
DeletionTime complexDeletion = DeletionTime.LIVE;
- for (int i=0, isize=versions.size(); i, IMeasurableMemory
cellReducer.setActiveDeletion(activeDeletion);
}
- Iterator> cells = MergeIterator.get(complexCells, Cell.comparator, cellReducer);
+ Iterator> cells = ComplexCellMergeIterator.get(complexCells, Cell.comparator, cellReducer);
while (cells.hasNext())
{
Cell> merged = cells.next();
@@ -937,11 +952,12 @@ public interface Row extends Unfiltered, Iterable, IMeasurableMemory
protected void onKeyChange()
{
column = null;
- versions.clear();
+ Arrays.fill(versions, 0, versionsSize, null);
+ versionsSize = 0;
}
}
- private static class CellReducer extends MergeIterator.Reducer, Cell>>
+ private static class CellReducer extends ComplexCellMergeIterator.Reducer| , Cell>>
{
private DeletionTime activeDeletion;
private Cell> merged;
diff --git a/src/java/org/apache/cassandra/db/rows/UnfilteredRowIterators.java b/src/java/org/apache/cassandra/db/rows/UnfilteredRowIterators.java
index bc4bfd15ba..a306841892 100644
--- a/src/java/org/apache/cassandra/db/rows/UnfilteredRowIterators.java
+++ b/src/java/org/apache/cassandra/db/rows/UnfilteredRowIterators.java
@@ -38,7 +38,7 @@ import org.apache.cassandra.schema.TableMetadata;
import org.apache.cassandra.serializers.MarshalException;
import org.apache.cassandra.utils.FBUtilities;
import org.apache.cassandra.utils.IMergeIterator;
-import org.apache.cassandra.utils.MergeIterator;
+import org.apache.cassandra.utils.UnfilteredMergeIterator;
/**
* Static methods to work with atom iterators.
@@ -415,7 +415,7 @@ public abstract class UnfilteredRowIterators
reversed,
EncodingStats.merge(iterators, UnfilteredRowIterator::stats));
- this.mergeIterator = MergeIterator.get(iterators,
+ this.mergeIterator = UnfilteredMergeIterator.get(iterators,
reversed ? metadata.comparator.reversed() : metadata.comparator,
new MergeReducer(iterators.size(), reversed, listener));
this.listener = listener;
@@ -540,7 +540,7 @@ public abstract class UnfilteredRowIterators
listener.close();
}
- private class MergeReducer extends MergeIterator.Reducer
+ private class MergeReducer extends UnfilteredMergeIterator.Reducer
{
private final MergeListener listener;
diff --git a/src/java/org/apache/cassandra/diag/DiagnosticEventPersistence.java b/src/java/org/apache/cassandra/diag/DiagnosticEventPersistence.java
index 81820382f3..b1e8df8783 100644
--- a/src/java/org/apache/cassandra/diag/DiagnosticEventPersistence.java
+++ b/src/java/org/apache/cassandra/diag/DiagnosticEventPersistence.java
@@ -128,6 +128,7 @@ public final class DiagnosticEventPersistence
LastEventIdBroadcaster.instance().setLastEventId(event.getClass().getName(), store.getLastEventId());
}
+ @SuppressWarnings("unchecked")
private Class getEventClass(String eventClazz) throws ClassNotFoundException, InvalidClassException
{
// get class by eventClazz argument name
@@ -135,12 +136,12 @@ public final class DiagnosticEventPersistence
if (!eventClazz.startsWith("org.apache.cassandra."))
throw new RuntimeException("Not a Cassandra event class: " + eventClazz);
- Class clazz = (Class) Class.forName(eventClazz);
+ Class> clazz = Class.forName(eventClazz, false, DiagnosticEventPersistence.class.getClassLoader());
if (!(DiagnosticEvent.class.isAssignableFrom(clazz)))
throw new InvalidClassException("Event class must be of type DiagnosticEvent");
- return clazz;
+ return (Class) clazz.asSubclass(DiagnosticEvent.class);
}
private DiagnosticEventStore getStore(Class cls)
diff --git a/src/java/org/apache/cassandra/index/SecondaryIndexManager.java b/src/java/org/apache/cassandra/index/SecondaryIndexManager.java
index 2846a8ede5..aeb4556ff4 100644
--- a/src/java/org/apache/cassandra/index/SecondaryIndexManager.java
+++ b/src/java/org/apache/cassandra/index/SecondaryIndexManager.java
@@ -921,6 +921,11 @@ public class SecondaryIndexManager implements IndexRegistry, INotificationConsum
return indexes.get(indexName);
}
+ static Class extends Index> loadIndexClass(String className)
+ {
+ return FBUtilities.classForNameWithoutInitialization(className, "Index", Index.class);
+ }
+
private Index createInstance(IndexMetadata indexDef)
{
Index newIndex;
@@ -933,7 +938,7 @@ public class SecondaryIndexManager implements IndexRegistry, INotificationConsum
try
{
- Class extends Index> indexClass = FBUtilities.classForName(className, "Index");
+ Class extends Index> indexClass = loadIndexClass(className);
Constructor extends Index> ctor = indexClass.getConstructor(ColumnFamilyStore.class, IndexMetadata.class);
newIndex = ctor.newInstance(baseCfs, indexDef);
}
diff --git a/src/java/org/apache/cassandra/index/sai/disk/io/IndexInputReader.java b/src/java/org/apache/cassandra/index/sai/disk/io/IndexInputReader.java
index 0c93f3c1fe..ed7dcb488c 100644
--- a/src/java/org/apache/cassandra/index/sai/disk/io/IndexInputReader.java
+++ b/src/java/org/apache/cassandra/index/sai/disk/io/IndexInputReader.java
@@ -20,6 +20,8 @@ package org.apache.cassandra.index.sai.disk.io;
import java.io.IOException;
+import javax.annotation.concurrent.NotThreadSafe;
+
import org.apache.lucene.store.DataInput;
import org.apache.lucene.store.IndexInput;
@@ -30,10 +32,13 @@ import org.apache.cassandra.io.util.RandomAccessReader;
* This is a wrapper over a Cassandra {@link RandomAccessReader} that provides an {@link IndexInput}
* interface for Lucene classes that need {@link IndexInput}. This is an optimisation because the
* Lucene {@link DataInput} reads bytes one at a time whereas the {@link RandomAccessReader} is
- * optimised to read multibyte objects faster.
+ * optimized to read multibyte objects faster.
*/
+@NotThreadSafe
public class IndexInputReader extends IndexInput
{
+ public static final Runnable NO_OP_ON_CLOSE = () -> {};
+
/**
* the byte order of `input`'s native readX operations doesn't matter,
* because we only use `readFully` and `readByte` methods. IndexInput calls these
@@ -42,27 +47,47 @@ public class IndexInputReader extends IndexInput
private final RandomAccessReader input;
private final Runnable doOnClose;
- private IndexInputReader(RandomAccessReader input, Runnable doOnClose)
+ /** Absolute offset in the underlying file that this input's position 0 refers to. */
+ private final long offset;
+
+ /** Bounded length of this input, in bytes. */
+ private final long length;
+
+ private IndexInputReader(RandomAccessReader input, Runnable doOnClose, long offset, long length)
{
super(input.getPath());
this.input = input;
this.doOnClose = doOnClose;
+ this.offset = offset;
+ this.length = length;
}
public static IndexInputReader create(RandomAccessReader input)
{
- return new IndexInputReader(input, () -> {});
+ // Top-level inputs own the underlying reader; folding its close into doOnClose lets us
+ // avoid a separate ownership flag on the class.
+ return new IndexInputReader(input, input::close, 0L, input.length());
}
public static IndexInputReader create(RandomAccessReader input, Runnable doOnClose)
{
- return new IndexInputReader(input, doOnClose);
+ Runnable close = () -> {
+ try
+ {
+ input.close();
+ }
+ finally
+ {
+ doOnClose.run();
+ }
+ };
+ return new IndexInputReader(input, close, 0L, input.length());
}
public static IndexInputReader create(FileHandle handle)
{
RandomAccessReader reader = handle.createReader();
- return new IndexInputReader(reader, () -> {});
+ return new IndexInputReader(reader, reader::close, 0L, reader.length());
}
@Override
@@ -80,37 +105,42 @@ public class IndexInputReader extends IndexInput
@Override
public void close()
{
- try
- {
- input.close();
- }
- finally
- {
- doOnClose.run();
- }
+ doOnClose.run();
}
@Override
public long getFilePointer()
{
- return input.getFilePointer();
+ return input.getFilePointer() - offset;
}
@Override
public void seek(long position)
{
- input.seek(position);
+ if (position > length)
+ throw new IllegalArgumentException("Cannot seek to position " + position + " past length of " + length);
+
+ input.seek(offset + position);
}
@Override
public long length()
{
- return input.length();
+ return length;
}
@Override
public IndexInput slice(String sliceDescription, long offset, long length)
{
- throw new UnsupportedOperationException("Slice operations are not supported");
+ if (offset < 0 || length < 0 || offset + length > this.length)
+ throw new IllegalArgumentException("Invalid slice: offset=" + offset + ", length=" + length + ", parent length=" + this.length + " for " + sliceDescription);
+
+ // Slices share the underlying reader with their parent; the no-op close keeps the parent's lifecycle intact.
+ IndexInputReader slice = new IndexInputReader(input, NO_OP_ON_CLOSE, this.offset + offset, length);
+
+ // Seek to the beginning of the slice...
+ slice.seek(0);
+
+ return slice;
}
}
diff --git a/src/java/org/apache/cassandra/index/sai/disk/v1/V1OnDiskFormat.java b/src/java/org/apache/cassandra/index/sai/disk/v1/V1OnDiskFormat.java
index ba8b13ca4c..5068e039f4 100644
--- a/src/java/org/apache/cassandra/index/sai/disk/v1/V1OnDiskFormat.java
+++ b/src/java/org/apache/cassandra/index/sai/disk/v1/V1OnDiskFormat.java
@@ -21,11 +21,14 @@ package org.apache.cassandra.index.sai.disk.v1;
import java.io.IOException;
import java.io.UncheckedIOException;
import java.util.EnumSet;
+import java.util.List;
import java.util.Set;
import com.codahale.metrics.Gauge;
import com.google.common.annotations.VisibleForTesting;
+import org.apache.lucene.codecs.CodecUtil;
+import org.apache.lucene.index.CorruptIndexException;
import org.apache.lucene.store.IndexInput;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -44,6 +47,7 @@ import org.apache.cassandra.index.sai.disk.format.IndexComponent;
import org.apache.cassandra.index.sai.disk.format.IndexDescriptor;
import org.apache.cassandra.index.sai.disk.format.OnDiskFormat;
import org.apache.cassandra.index.sai.disk.v1.segment.SegmentBuilder;
+import org.apache.cassandra.index.sai.disk.v1.segment.SegmentMetadata;
import org.apache.cassandra.index.sai.metrics.AbstractMetrics;
import org.apache.cassandra.index.sai.utils.IndexIdentifier;
import org.apache.cassandra.index.sai.utils.IndexTermType;
@@ -96,6 +100,18 @@ public class V1OnDiskFormat implements OnDiskFormat
IndexComponent.TERMS_DATA,
IndexComponent.POSTING_LISTS);
+ /**
+ * Per-column components whose files are written in append mode with one SAI codec footer
+ * per segment (see {@link org.apache.cassandra.index.sai.disk.v1.bbtree.NumericIndexWriter},
+ * {@link org.apache.cassandra.index.sai.disk.v1.trie.TrieTermsDictionaryWriter},
+ * {@link org.apache.cassandra.index.sai.disk.v1.postings.PostingsWriter}, and
+ * {@link org.apache.cassandra.index.sai.disk.v1.vector.OnHeapGraph}).
+ */
+ private static final Set SEGMENTED_COMPONENTS = EnumSet.of(IndexComponent.BALANCED_TREE,
+ IndexComponent.POSTING_LISTS,
+ IndexComponent.TERMS_DATA,
+ IndexComponent.COMPRESSED_VECTORS);
+
/**
* Global limit on heap consumed by all index segment building that occurs outside the context of Memtable flush.
*
@@ -219,15 +235,90 @@ public class V1OnDiskFormat implements OnDiskFormat
}
}
+ if (isEmptyIndex)
+ return;
+
+ // Safely read the segment metadata so we can validate per-segment checksums below...
+ List segments = null;
+ if (checksum)
+ {
+ validateIndexComponent(indexDescriptor, indexIdentifier, IndexComponent.META, true);
+ try
+ {
+ segments = SegmentMetadata.load(MetadataSource.loadColumnMetadata(indexDescriptor, indexIdentifier), indexDescriptor.primaryKeyFactory);
+ }
+ catch (IOException e)
+ {
+ rethrowIOException(e);
+ }
+ }
+
for (IndexComponent indexComponent : perColumnIndexComponents(indexTermType))
{
- if (!isEmptyIndex && isNotBuildCompletionMarker(indexComponent))
+ if (isNotBuildCompletionMarker(indexComponent))
{
- validateIndexComponent(indexDescriptor, indexIdentifier, indexComponent, checksum);
+ // META was validated up-front in CHECKSUM mode; don't validate it twice.
+ if (checksum && indexComponent == IndexComponent.META)
+ continue;
+
+ if (checksum && SEGMENTED_COMPONENTS.contains(indexComponent))
+ {
+ assert segments != null : "No segment metadata available!";
+ validateSegmentedIndexComponent(indexDescriptor, indexIdentifier, indexComponent, segments, indexTermType.isVector());
+ }
+ else
+ validateIndexComponent(indexDescriptor, indexIdentifier, indexComponent, checksum);
}
}
}
+ private static void validateSegmentedIndexComponent(IndexDescriptor indexDescriptor,
+ IndexIdentifier indexIdentifier,
+ IndexComponent indexComponent,
+ List segments,
+ boolean payloadOnlyMetadata)
+ {
+ try (IndexInput input = indexDescriptor.openPerIndexInput(indexComponent, indexIdentifier))
+ {
+ long fileLength = input.length();
+ long frameStart = 0;
+
+ for (SegmentMetadata segment : segments)
+ {
+ SegmentMetadata.ComponentMetadata cm = segment.componentMetadatas.get(indexComponent);
+
+ // Non-vector writers record offsets as the codec-framed segment starts (before
+ // the header) and length as the full framed length (through the footer). The vector
+ // writer (OnHeapGraph#writeData) instead records the offset as the payload start (after
+ // the header) and length as just the payload length, because vector readers seek
+ // directly at the payload. Segments are written contiguously in append mode, so we can
+ // recover the vector-path frame extent by walking segment ends and adding the trailing
+ // 16-byte codec footer.
+ long frameEnd = payloadOnlyMetadata ? cm.offset + cm.length + CodecUtil.footerLength() : cm.offset + cm.length;
+
+ if (frameEnd > fileLength || frameEnd < frameStart)
+ throw new CorruptIndexException(String.format("Segment frame [%d, %d) is inconsistent with component file length %d",
+ frameStart, frameEnd, fileLength),
+ indexComponent.name + '@' + frameStart);
+
+ IndexInput slice = input.slice(indexComponent.name + '@' + frameStart, frameStart, frameEnd - frameStart);
+ SAICodecUtils.validateChecksum(slice);
+ frameStart = frameEnd;
+ }
+
+ if (frameStart != fileLength)
+ throw new CorruptIndexException(String.format("Component file length %d does not match combined frame length of all segments %d",
+ fileLength, frameStart),
+ indexComponent.name);
+ }
+ catch (Exception e)
+ {
+ logger.warn(indexDescriptor.logMessage("Segmented checksum validation failed for index component {} on SSTable {}"),
+ indexComponent, indexDescriptor.sstableDescriptor);
+ rethrowIOException(e);
+ }
+ }
+
private static void validateIndexComponent(IndexDescriptor indexDescriptor,
IndexIdentifier indexContext,
IndexComponent indexComponent,
@@ -245,9 +336,7 @@ public class V1OnDiskFormat implements OnDiskFormat
catch (Exception e)
{
logger.warn(indexDescriptor.logMessage("{} failed for index component {} on SSTable {}"),
- checksum ? "Checksum validation" : "Validation",
- indexComponent,
- indexDescriptor.sstableDescriptor);
+ checksum ? "Checksum validation" : "Validation", indexComponent, indexDescriptor.sstableDescriptor);
rethrowIOException(e);
}
}
diff --git a/src/java/org/apache/cassandra/index/sai/utils/CellWithSource.java b/src/java/org/apache/cassandra/index/sai/utils/CellWithSource.java
index ce697f361a..b0ee8284b4 100644
--- a/src/java/org/apache/cassandra/index/sai/utils/CellWithSource.java
+++ b/src/java/org/apache/cassandra/index/sai/utils/CellWithSource.java
@@ -103,6 +103,12 @@ public class CellWithSource extends Cell
return cell.localDeletionTime();
}
+ @Override
+ public long minDeletionTime()
+ {
+ return cell.minDeletionTime();
+ }
+
@Override
public boolean isTombstone()
{
diff --git a/src/java/org/apache/cassandra/index/sai/utils/RowWithSource.java b/src/java/org/apache/cassandra/index/sai/utils/RowWithSource.java
index 33f5f07deb..6b13a92236 100644
--- a/src/java/org/apache/cassandra/index/sai/utils/RowWithSource.java
+++ b/src/java/org/apache/cassandra/index/sai/utils/RowWithSource.java
@@ -213,6 +213,12 @@ public class RowWithSource implements Row
return row.hasDeletion(nowInSec);
}
+ @Override
+ public long minLocalDeletionTime()
+ {
+ return row.minLocalDeletionTime();
+ }
+
@Override
public SearchIterator searchIterator()
{
diff --git a/src/java/org/apache/cassandra/index/sasi/conf/IndexMode.java b/src/java/org/apache/cassandra/index/sasi/conf/IndexMode.java
index c827b2b820..339a2cae5d 100644
--- a/src/java/org/apache/cassandra/index/sasi/conf/IndexMode.java
+++ b/src/java/org/apache/cassandra/index/sasi/conf/IndexMode.java
@@ -37,6 +37,7 @@ import org.apache.cassandra.index.sasi.disk.OnDiskIndexBuilder.Mode;
import org.apache.cassandra.index.sasi.plan.Expression.Op;
import org.apache.cassandra.schema.ColumnMetadata;
import org.apache.cassandra.schema.IndexMetadata;
+import org.apache.cassandra.utils.FBUtilities;
public class IndexMode
{
@@ -60,10 +61,10 @@ public class IndexMode
public final Mode mode;
public final boolean isAnalyzed, isLiteral;
- public final Class analyzerClass;
+ public final Class extends AbstractAnalyzer> analyzerClass;
public final long maxCompactionFlushMemoryInBytes;
- private IndexMode(Mode mode, boolean isLiteral, boolean isAnalyzed, Class analyzerClass, long maxMemBytes)
+ private IndexMode(Mode mode, boolean isLiteral, boolean isAnalyzed, Class extends AbstractAnalyzer> analyzerClass, long maxMemBytes)
{
this.mode = mode;
this.isLiteral = isLiteral;
@@ -81,7 +82,7 @@ public class IndexMode
if (isAnalyzed)
{
if (analyzerClass != null)
- analyzer = (AbstractAnalyzer) analyzerClass.newInstance();
+ analyzer = analyzerClass.newInstance();
else if (TOKENIZABLE_TYPES.contains(validator))
analyzer = new StandardAnalyzer();
}
@@ -99,21 +100,14 @@ public class IndexMode
// validate that a valid analyzer class was provided if specified
if (indexOptions.containsKey(INDEX_ANALYZER_CLASS_OPTION))
{
- Class> analyzerClass;
- try
- {
- analyzerClass = Class.forName(indexOptions.get(INDEX_ANALYZER_CLASS_OPTION));
- }
- catch (ClassNotFoundException e)
- {
- throw new ConfigurationException(String.format("Invalid analyzer class option specified [%s]",
- indexOptions.get(INDEX_ANALYZER_CLASS_OPTION)));
- }
+ Class extends AbstractAnalyzer> analyzerClass = FBUtilities.classForNameWithoutInitialization(indexOptions.get(INDEX_ANALYZER_CLASS_OPTION),
+ "analyzer",
+ AbstractAnalyzer.class);
AbstractAnalyzer analyzer;
try
{
- analyzer = (AbstractAnalyzer) analyzerClass.newInstance();
+ analyzer = analyzerClass.newInstance();
analyzer.validate(indexOptions, cd);
}
catch (InstantiationException | IllegalAccessException e)
@@ -148,25 +142,30 @@ public class IndexMode
}
boolean isAnalyzed = false;
- Class analyzerClass = null;
- try
+ Class extends AbstractAnalyzer> analyzerClass = null;
+ if (indexOptions.get(INDEX_ANALYZER_CLASS_OPTION) != null)
{
- if (indexOptions.get(INDEX_ANALYZER_CLASS_OPTION) != null)
+ try
{
- analyzerClass = Class.forName(indexOptions.get(INDEX_ANALYZER_CLASS_OPTION));
+ analyzerClass = FBUtilities.classForNameWithoutInitialization(indexOptions.get(INDEX_ANALYZER_CLASS_OPTION),
+ "analyzer",
+ AbstractAnalyzer.class);
isAnalyzed = indexOptions.get(INDEX_ANALYZED_OPTION) == null
- ? true : Boolean.parseBoolean(indexOptions.get(INDEX_ANALYZED_OPTION));
+ ? true : Boolean.parseBoolean(indexOptions.get(INDEX_ANALYZED_OPTION));
}
- else if (indexOptions.get(INDEX_ANALYZED_OPTION) != null)
+ catch (ConfigurationException e)
{
- isAnalyzed = Boolean.parseBoolean(indexOptions.get(INDEX_ANALYZED_OPTION));
+ if (!(e.getCause() instanceof ClassNotFoundException))
+ throw e;
+
+ // Should not happen as we already validated we could instantiate an instance in validateAnalyzer().
+ logger.error("Failed to find specified analyzer class [{}]. Falling back to default analyzer",
+ indexOptions.get(INDEX_ANALYZER_CLASS_OPTION));
}
}
- catch (ClassNotFoundException e)
+ else if (indexOptions.get(INDEX_ANALYZED_OPTION) != null)
{
- // should not happen as we already validated we could instantiate an instance in validateAnalyzer()
- logger.error("Failed to find specified analyzer class [{}]. Falling back to default analyzer",
- indexOptions.get(INDEX_ANALYZER_CLASS_OPTION));
+ isAnalyzed = Boolean.parseBoolean(indexOptions.get(INDEX_ANALYZED_OPTION));
}
boolean isLiteral = false;
diff --git a/src/java/org/apache/cassandra/io/sstable/ClusteringDescriptor.java b/src/java/org/apache/cassandra/io/sstable/ClusteringDescriptor.java
index 27bd3cea89..45015df281 100644
--- a/src/java/org/apache/cassandra/io/sstable/ClusteringDescriptor.java
+++ b/src/java/org/apache/cassandra/io/sstable/ClusteringDescriptor.java
@@ -61,7 +61,7 @@ public class ClusteringDescriptor extends ResizableByteBuffer
protected void loadClustering(RandomAccessReader dataReader, byte clusteringKind, int clusteringColumnsBound) throws IOException
{
- set(ClusteringPrefix.Kind.values()[clusteringKind], clusteringKind, clusteringColumnsBound);
+ set(ClusteringPrefix.Kind.fromOrdinal(clusteringKind), clusteringKind, clusteringColumnsBound);
if (clusteringKind != STATIC_CLUSTERING_KIND)
readUnfilteredClustering(dataReader, clusteringTypes, this.clusteringColumnsBound, this);
else
@@ -103,7 +103,7 @@ public class ClusteringDescriptor extends ResizableByteBuffer
}
private void set(byte clusteringKindEncoded, int clusteringColumnsBound) {
- set(ClusteringPrefix.Kind.values()[clusteringKindEncoded], clusteringKindEncoded, clusteringColumnsBound);
+ set(ClusteringPrefix.Kind.fromOrdinal(clusteringKindEncoded), clusteringKindEncoded, clusteringColumnsBound);
}
private void set(ClusteringPrefix.Kind clusteringKind, byte clusteringKindEncoded, int clusteringColumnsBound)
diff --git a/src/java/org/apache/cassandra/io/sstable/SSTableCursorWriter.java b/src/java/org/apache/cassandra/io/sstable/SSTableCursorWriter.java
index 102e109ab3..bc0b830f0f 100644
--- a/src/java/org/apache/cassandra/io/sstable/SSTableCursorWriter.java
+++ b/src/java/org/apache/cassandra/io/sstable/SSTableCursorWriter.java
@@ -513,7 +513,7 @@ public class SSTableCursorWriter implements AutoCloseable
public void writeRangeTombstone(UnfilteredDescriptor rangeTombstone, boolean updateClusteringMetadata) throws IOException
{
int tombstoneKind = rangeTombstone.clusteringKindEncoded();
- ClusteringPrefix.Kind kind = ClusteringPrefix.Kind.values()[tombstoneKind];
+ ClusteringPrefix.Kind kind = ClusteringPrefix.Kind.fromOrdinal(tombstoneKind);
long unfilteredStartPosition = getPosition();
/** See: {@link org.apache.cassandra.db.rows.UnfilteredSerializer#serialize */
dataWriter.writeByte((byte)IS_MARKER);
diff --git a/src/java/org/apache/cassandra/locator/AbstractReplicationStrategy.java b/src/java/org/apache/cassandra/locator/AbstractReplicationStrategy.java
index 4dd68b3dfd..4be3f0e605 100644
--- a/src/java/org/apache/cassandra/locator/AbstractReplicationStrategy.java
+++ b/src/java/org/apache/cassandra/locator/AbstractReplicationStrategy.java
@@ -340,11 +340,11 @@ public abstract class AbstractReplicationStrategy
if ("org.apache.cassandra.locator.OldNetworkTopologyStrategy".equals(className)) // see CASSANDRA-16301
throw new ConfigurationException("The support for the OldNetworkTopologyStrategy has been removed in C* version 4.0. The keyspace strategy should be switch to NetworkTopologyStrategy");
- Class strategyClass = FBUtilities.classForName(className, "replication strategy");
- if (!AbstractReplicationStrategy.class.isAssignableFrom(strategyClass))
- {
- throw new ConfigurationException(String.format("Specified replication strategy class (%s) is not derived from AbstractReplicationStrategy", className));
- }
+ @SuppressWarnings("unchecked")
+ Class strategyClass =
+ (Class) FBUtilities.classForNameWithoutInitialization(className,
+ "replication strategy",
+ AbstractReplicationStrategy.class);
return strategyClass;
}
diff --git a/src/java/org/apache/cassandra/metrics/ThreadLocalMetrics.java b/src/java/org/apache/cassandra/metrics/ThreadLocalMetrics.java
index feb96c02bd..d7b429ada6 100644
--- a/src/java/org/apache/cassandra/metrics/ThreadLocalMetrics.java
+++ b/src/java/org/apache/cassandra/metrics/ThreadLocalMetrics.java
@@ -26,6 +26,8 @@ import java.util.List;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.CopyOnWriteArrayList;
+import java.util.concurrent.RejectedExecutionException;
+import java.util.concurrent.ScheduledFuture;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;
import java.util.concurrent.atomic.AtomicInteger;
@@ -41,6 +43,7 @@ import org.apache.cassandra.concurrent.ScheduledExecutors;
import org.apache.cassandra.concurrent.Shutdownable;
import io.netty.util.concurrent.FastThreadLocal;
+import io.netty.util.concurrent.FastThreadLocalThread;
import static com.google.common.collect.ImmutableList.of;
import static org.apache.cassandra.concurrent.ExecutorFactory.Global.executorFactory;
@@ -64,10 +67,8 @@ public class ThreadLocalMetrics
static final AtomicInteger idGenerator = new AtomicInteger();
- private static final Object freeMetricIdSetGuard = new Object();
-
@VisibleForTesting
- static final BitSet freeMetricIdSet = new BitSet();
+ static final FreeMetricIdSetTracker freeMetricIdSetTracker = new FreeMetricIdSetTracker();
private static final List allThreadLocalMetrics = new CopyOnWriteArrayList<>();
@@ -101,7 +102,13 @@ public class ThreadLocalMetrics
{
ThreadLocalMetrics result = new ThreadLocalMetrics();
allThreadLocalMetrics.add(result);
- destroyWhenUnreachable(Thread.currentThread(), result::release);
+
+ Thread thread = Thread.currentThread();
+ // use phantom references ony if needed
+ // CassandraThread is FastThreadLocalThread too
+ if (!(thread instanceof FastThreadLocalThread))
+ destroyWhenUnreachable(thread, result::release);
+
return result;
}
@@ -317,13 +324,7 @@ public class ThreadLocalMetrics
static int allocateMetricId()
{
- int metricId;
- synchronized (freeMetricIdSetGuard)
- {
- metricId = freeMetricIdSet.nextSetBit(0);
- if (metricId >= 0)
- freeMetricIdSet.clear(metricId);
- }
+ int metricId = freeMetricIdSetTracker.getFreeMetricId();
if (metricId < 0)
metricId = idGenerator.getAndIncrement();
@@ -374,26 +375,92 @@ public class ThreadLocalMetrics
lock.unlock();
}
- // there's no an obvious happens-before relation between currentCounterValues[metricId] = 0 write we just did
- // and an initial read of the entry by a thread which updates the reused metric
- // as a workaround we introduce a delay in recyling to provide the write visibility in practice
- // even if it is not formally guaranteed by the JMM
- ScheduledExecutors.scheduledTasks.schedule(() -> {
- synchronized (freeMetricIdSetGuard)
+ freeMetricIdSetTracker.markAsFree(metricId);
+ }
+
+ @VisibleForTesting
+ static class FreeMetricIdSetTracker
+ {
+ private final BitSet freeMetricIdSet = new BitSet();
+
+ private final BitSet tickDelayedToFreeMetricIdSet = new BitSet();
+ private final BitSet tockDelayedToFreeMetricIdSet = new BitSet();
+
+ private BitSet delayedToFreeMetricIdSet = tickDelayedToFreeMetricIdSet;
+
+ private ScheduledFuture> cleanupTask;
+
+ @VisibleForTesting
+ synchronized void triggerRecycling()
+ {
+ cleanupTask = null;
+ BitSet toProcess = otherSet(delayedToFreeMetricIdSet);
+ freeMetricIdSet.or(toProcess);
+ toProcess.clear();
+ if (!delayedToFreeMetricIdSet.isEmpty())
+ scheduleCleanupTask();
+ delayedToFreeMetricIdSet = toProcess;
+ }
+
+ private BitSet otherSet(BitSet set)
+ {
+ return set == tickDelayedToFreeMetricIdSet ? tockDelayedToFreeMetricIdSet : tickDelayedToFreeMetricIdSet;
+ }
+
+ public synchronized int getFreeMetricId()
+ {
+ int metricId = freeMetricIdSet.nextSetBit(0);
+ if (metricId >= 0)
+ freeMetricIdSet.clear(metricId);
+ return metricId;
+ }
+
+ public synchronized void markAsFree(int metricId)
+ {
+ // there's no an obvious happens-before relation between currentCounterValues[metricId] = 0 write we just did
+ // and an initial read of the entry by a thread which updates the reused metric
+ // as a workaround we introduce a delay in recyling to provide the write visibility in practice
+ // even if it is not formally guaranteed by the JMM
+ delayedToFreeMetricIdSet.set(metricId);
+ scheduleCleanupTask();
+ }
+
+ // must be called while holding this monitor (from a synchronized method)
+ @VisibleForTesting
+ protected void scheduleCleanupTask()
+ {
+ try
{
- freeMetricIdSet.set(metricId);
+ if (cleanupTask == null)
+ cleanupTask = ScheduledExecutors.scheduledTasks.schedule(this::triggerRecycling, 5, TimeUnit.SECONDS);
}
- }, 5, TimeUnit.SECONDS);
+ catch (RejectedExecutionException e)
+ {
+ // ignore theoretically possible rejections during a shutdown
+ }
+ }
+
+ public synchronized int getFreeMetricSetCardinality()
+ {
+ return freeMetricIdSet.cardinality();
+ }
+
+ @Override
+ public synchronized String toString()
+ {
+ return "FreeMetricIdSetTracker{" +
+ "freeMetricIdSet=" + freeMetricIdSet +
+ ", tickDelayedToFreeMetricIdSet=" + tickDelayedToFreeMetricIdSet +
+ ", tockDelayedToFreeMetricIdSet=" + tockDelayedToFreeMetricIdSet +
+ ", delayedToFreeMetricIdSet=" + (delayedToFreeMetricIdSet == tickDelayedToFreeMetricIdSet ? "tick" : "tock") +
+ '}';
+ }
}
@VisibleForTesting
static int getAllocatedMetricsCount()
{
- int freeCount;
- synchronized (freeMetricIdSetGuard)
- {
- freeCount = freeMetricIdSet.cardinality();
- }
+ int freeCount = freeMetricIdSetTracker.getFreeMetricSetCardinality();
return idGenerator.get() - freeCount;
}
diff --git a/src/java/org/apache/cassandra/net/AbstractMessageHandler.java b/src/java/org/apache/cassandra/net/AbstractMessageHandler.java
index 0ef03990fd..abbd5c77b9 100644
--- a/src/java/org/apache/cassandra/net/AbstractMessageHandler.java
+++ b/src/java/org/apache/cassandra/net/AbstractMessageHandler.java
@@ -31,7 +31,6 @@ import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.apache.cassandra.concurrent.ManyToOneConcurrentLinkedQueue;
-import org.apache.cassandra.metrics.ClientMetrics;
import org.apache.cassandra.net.FrameDecoder.CorruptFrame;
import org.apache.cassandra.net.FrameDecoder.Frame;
import org.apache.cassandra.net.FrameDecoder.FrameProcessor;
@@ -314,7 +313,7 @@ public abstract class AbstractMessageHandler extends ChannelInboundHandlerAdapte
decoder.reactivate();
if (decoder.isActive())
- ClientMetrics.instance.unpauseConnection();
+ onConnectionUnpaused();
}
}
catch (Throwable t)
@@ -323,6 +322,10 @@ public abstract class AbstractMessageHandler extends ChannelInboundHandlerAdapte
}
}
+ protected void onConnectionUnpaused()
+ {
+ }
+
protected abstract void fatalExceptionCaught(Throwable t);
// return true if the handler should be reactivated - if no new hurdles were encountered,
diff --git a/src/java/org/apache/cassandra/repair/autorepair/AutoRepairConfig.java b/src/java/org/apache/cassandra/repair/autorepair/AutoRepairConfig.java
index dd5c0d6ed9..2c3ba9dc20 100644
--- a/src/java/org/apache/cassandra/repair/autorepair/AutoRepairConfig.java
+++ b/src/java/org/apache/cassandra/repair/autorepair/AutoRepairConfig.java
@@ -395,7 +395,10 @@ public class AutoRepairConfig implements Serializable
className = parameterizedClass.class_name.contains(".") ?
parameterizedClass.class_name :
"org.apache.cassandra.repair.autorepair." + parameterizedClass.class_name;
- tokenRangeSplitterClass = FBUtilities.classForName(className, "token_range_splitter");
+ tokenRangeSplitterClass =
+ FBUtilities.classForNameWithoutInitialization(className,
+ "token_range_splitter",
+ IAutoRepairTokenRangeSplitter.class);
}
else
{
diff --git a/src/java/org/apache/cassandra/schema/CompactionParams.java b/src/java/org/apache/cassandra/schema/CompactionParams.java
index 9217a36376..1cfdd31458 100644
--- a/src/java/org/apache/cassandra/schema/CompactionParams.java
+++ b/src/java/org/apache/cassandra/schema/CompactionParams.java
@@ -311,15 +311,7 @@ public final class CompactionParams
String className = name.contains(".")
? name
: "org.apache.cassandra.db.compaction." + name;
- Class strategyClass = FBUtilities.classForName(className, "compaction strategy");
-
- if (!AbstractCompactionStrategy.class.isAssignableFrom(strategyClass))
- {
- throw new ConfigurationException(format("Compaction strategy class %s is not derived from AbstractReplicationStrategy",
- className));
- }
-
- return strategyClass;
+ return FBUtilities.classForNameWithoutInitialization(className, "compaction strategy", AbstractCompactionStrategy.class);
}
/*
diff --git a/src/java/org/apache/cassandra/schema/CompressionParams.java b/src/java/org/apache/cassandra/schema/CompressionParams.java
index a438e34eb1..078f58d8f7 100644
--- a/src/java/org/apache/cassandra/schema/CompressionParams.java
+++ b/src/java/org/apache/cassandra/schema/CompressionParams.java
@@ -46,6 +46,7 @@ import org.apache.cassandra.io.compress.ZstdDictionaryCompressor;
import org.apache.cassandra.io.util.DataInputPlus;
import org.apache.cassandra.io.util.DataOutputPlus;
import org.apache.cassandra.net.MessagingService;
+import org.apache.cassandra.utils.FBUtilities;
import static java.lang.String.format;
@@ -302,7 +303,7 @@ public final class CompressionParams
return maxCompressedLength;
}
- private static Class> parseCompressorClass(String className) throws ConfigurationException
+ private static Class extends ICompressor> parseCompressorClass(String className) throws ConfigurationException
{
if (className == null || className.isEmpty())
return null;
@@ -310,15 +311,17 @@ public final class CompressionParams
className = className.contains(".") ? className : "org.apache.cassandra.io.compress." + className;
try
{
- return Class.forName(className);
+ return FBUtilities.classForNameWithoutInitialization(className, "compression", ICompressor.class);
}
- catch (Exception e)
+ catch (ConfigurationException e)
{
- throw new ConfigurationException("Could not create Compression for type " + className, e);
+ if (e.getCause() instanceof ClassNotFoundException || e.getCause() instanceof NoClassDefFoundError)
+ throw new ConfigurationException("Could not create Compression for type " + className, e);
+ throw e;
}
}
- private static ICompressor createCompressor(Class> compressorClass, Map compressionOptions) throws ConfigurationException
+ private static ICompressor createCompressor(Class extends ICompressor> compressorClass, Map compressionOptions) throws ConfigurationException
{
if (compressorClass == null)
{
diff --git a/src/java/org/apache/cassandra/schema/IndexMetadata.java b/src/java/org/apache/cassandra/schema/IndexMetadata.java
index 15a0670e3a..b471260ccc 100644
--- a/src/java/org/apache/cassandra/schema/IndexMetadata.java
+++ b/src/java/org/apache/cassandra/schema/IndexMetadata.java
@@ -148,9 +148,7 @@ public final class IndexMetadata
// Get the fully qualified class name:
String className = getIndexClassName();
- Class indexerClass = FBUtilities.classForName(className, "custom indexer");
- if (!Index.class.isAssignableFrom(indexerClass))
- throw new ConfigurationException(String.format("Specified Indexer class (%s) does not implement the Indexer interface", className));
+ Class extends Index> indexerClass = FBUtilities.classForNameWithoutInitialization(className, "custom indexer", Index.class);
validateCustomIndexOptions(table, indexerClass, options);
}
}
diff --git a/src/java/org/apache/cassandra/schema/MemtableParams.java b/src/java/org/apache/cassandra/schema/MemtableParams.java
index 7d88f6518e..c2111b9b0b 100644
--- a/src/java/org/apache/cassandra/schema/MemtableParams.java
+++ b/src/java/org/apache/cassandra/schema/MemtableParams.java
@@ -228,19 +228,25 @@ public final class MemtableParams
try
{
Memtable.Factory factory;
- Class> clazz = Class.forName(className);
+ Class> clazz = Class.forName(className, false, MemtableParams.class.getClassLoader());
final Map parametersCopy = options.parameters != null
? new HashMap<>(options.parameters)
: new HashMap<>();
try
{
Method factoryMethod = clazz.getDeclaredMethod("factory", Map.class);
+ if (!Memtable.Factory.class.isAssignableFrom(factoryMethod.getReturnType()))
+ throw new ClassCastException("Memtable factory method on " + className +
+ " must return " + Memtable.Factory.class.getName());
factory = (Memtable.Factory) factoryMethod.invoke(null, parametersCopy);
}
catch (NoSuchMethodException e)
{
// continue with FACTORY field
Field factoryField = clazz.getDeclaredField("FACTORY");
+ if (!Memtable.Factory.class.isAssignableFrom(factoryField.getType()))
+ throw new ClassCastException("Memtable FACTORY field on " + className +
+ " must be of type " + Memtable.Factory.class.getName());
factory = (Memtable.Factory) factoryField.get(null);
}
if (!parametersCopy.isEmpty())
diff --git a/src/java/org/apache/cassandra/schema/ReplicationParams.java b/src/java/org/apache/cassandra/schema/ReplicationParams.java
index c1e2643897..63cad56057 100644
--- a/src/java/org/apache/cassandra/schema/ReplicationParams.java
+++ b/src/java/org/apache/cassandra/schema/ReplicationParams.java
@@ -286,7 +286,7 @@ public final class ReplicationParams
Map options = new HashMap<>(size);
for (int i = 0; i < size; i++)
options.put(in.readUTF(), in.readUTF());
- return new ReplicationParams(FBUtilities.classForName(klassName, "ReplicationStrategy"), options);
+ return new ReplicationParams(FBUtilities.classForNameWithoutInitialization(klassName, "ReplicationStrategy", AbstractReplicationStrategy.class), options);
}
public long serializedSize(ReplicationParams t, Version version)
@@ -322,7 +322,7 @@ public final class ReplicationParams
Map options = new HashMap<>(size);
for (int i=0; i keyProviderClass = (Class)Class.forName(options.key_provider.class_name);
- Constructor ctor = keyProviderClass.getConstructor(TransparentDataEncryptionOptions.class);
- keyProvider = (KeyProvider)ctor.newInstance(options);
+ Class extends KeyProvider> keyProviderClass =
+ FBUtilities.classForNameWithoutInitialization(options.key_provider.class_name, "key provider", KeyProvider.class);
+ Constructor extends KeyProvider> ctor = keyProviderClass.getConstructor(TransparentDataEncryptionOptions.class);
+ keyProvider = ctor.newInstance(options);
}
catch (Exception e)
{
diff --git a/src/java/org/apache/cassandra/service/CacheService.java b/src/java/org/apache/cassandra/service/CacheService.java
index 8240c2880f..c755c06347 100644
--- a/src/java/org/apache/cassandra/service/CacheService.java
+++ b/src/java/org/apache/cassandra/service/CacheService.java
@@ -147,9 +147,11 @@ public class CacheService implements CacheServiceMBean
? DatabaseDescriptor.getRowCacheClassName() : "org.apache.cassandra.cache.NopCacheProvider";
try
{
- Class> cacheProviderClass =
- (Class>) Class.forName(cacheProviderClassName);
- cacheProvider = cacheProviderClass.newInstance();
+ Class extends CacheProvider> cacheProviderClass =
+ FBUtilities.classForNameWithoutInitialization(cacheProviderClassName, "row cache provider", CacheProvider.class);
+ @SuppressWarnings("unchecked")
+ CacheProvider typedCacheProvider = cacheProviderClass.newInstance();
+ cacheProvider = typedCacheProvider;
}
catch (Exception e)
{
diff --git a/src/java/org/apache/cassandra/service/ClientState.java b/src/java/org/apache/cassandra/service/ClientState.java
index c119d5f4c3..a17688d116 100644
--- a/src/java/org/apache/cassandra/service/ClientState.java
+++ b/src/java/org/apache/cassandra/service/ClientState.java
@@ -124,7 +124,7 @@ public class ClientState
{
try
{
- handler = FBUtilities.construct(customHandlerClass, "QueryHandler");
+ handler = FBUtilities.construct(customHandlerClass, "QueryHandler", QueryHandler.class);
logger.info("Using {} as a query handler for native protocol queries (as requested by the {} system property)",
customHandlerClass, CUSTOM_QUERY_HANDLER_CLASS.getKey());
}
diff --git a/src/java/org/apache/cassandra/service/DiskErrorsHandlerService.java b/src/java/org/apache/cassandra/service/DiskErrorsHandlerService.java
index 036559a738..7556e368d4 100644
--- a/src/java/org/apache/cassandra/service/DiskErrorsHandlerService.java
+++ b/src/java/org/apache/cassandra/service/DiskErrorsHandlerService.java
@@ -78,7 +78,7 @@ public class DiskErrorsHandlerService
String fsErrorHandlerClass = CassandraRelevantProperties.CUSTOM_DISK_ERROR_HANDLER.getString();
DiskErrorsHandler fsErrorHandler = fsErrorHandlerClass == null
? new DefaultDiskErrorsHandler()
- : FBUtilities.construct(fsErrorHandlerClass, "disk error handler");
+ : FBUtilities.construct(fsErrorHandlerClass, "disk error handler", DiskErrorsHandler.class);
DiskErrorsHandlerService.set(fsErrorHandler);
}
}
diff --git a/src/java/org/apache/cassandra/service/StorageService.java b/src/java/org/apache/cassandra/service/StorageService.java
index 397e50fa0d..b4c8343558 100644
--- a/src/java/org/apache/cassandra/service/StorageService.java
+++ b/src/java/org/apache/cassandra/service/StorageService.java
@@ -3403,6 +3403,11 @@ public class StorageService extends NotificationBroadcasterSupport implements IE
*/
@Deprecated(since = "4.0")
public List getNaturalEndpoints(String keyspaceName, String cf, String key)
+ {
+ return getNaturalReplicas(keyspaceName, cf, key);
+ }
+
+ public List getNaturalReplicas(String keyspaceName, String cf, String key)
{
EndpointsForToken replicas = getNaturalReplicasForToken(keyspaceName, cf, key);
List inetList = new ArrayList<>(replicas.size());
@@ -3410,11 +3415,16 @@ public class StorageService extends NotificationBroadcasterSupport implements IE
return inetList;
}
- public List getNaturalEndpointsWithPort(String keyspaceName, String cf, String key)
+ public List getNaturalReplicasWithPort(String keyspaceName, String cf, String key)
{
return Replicas.stringify(getNaturalReplicasForToken(keyspaceName, cf, key), true);
}
+ public List getNaturalEndpointsWithPort(String keyspaceName, String cf, String key)
+ {
+ return getNaturalReplicasWithPort(keyspaceName, cf, key);
+ }
+
/** @deprecated See CASSANDRA-7544 */
@Deprecated(since = "4.0")
public List getNaturalEndpoints(String keyspaceName, ByteBuffer key)
diff --git a/src/java/org/apache/cassandra/service/StorageServiceMBean.java b/src/java/org/apache/cassandra/service/StorageServiceMBean.java
index e6f7a18c8e..abc8689031 100644
--- a/src/java/org/apache/cassandra/service/StorageServiceMBean.java
+++ b/src/java/org/apache/cassandra/service/StorageServiceMBean.java
@@ -262,12 +262,28 @@ public interface StorageServiceMBean extends NotificationEmitter
* @param key - key for which we need to find the endpoint return value -
* the endpoint responsible for this key
* @deprecated See CASSANDRA-7544
+ * @link getNaturalReplicas
+ * @link getNaturalReplicasWithPort
*/
@Deprecated(since = "4.0") public List getNaturalEndpoints(String keyspaceName, String cf, String key);
- public List getNaturalEndpointsWithPort(String keyspaceName, String cf, String key);
+ /** @deprecated See CASSANDRA-17665 */
+ @Deprecated(since = "7.0") public List getNaturalEndpointsWithPort(String keyspaceName, String cf, String key);
/** @deprecated See CASSANDRA-7544 */
@Deprecated(since = "4.0") public List getNaturalEndpoints(String keyspaceName, ByteBuffer key);
- public List getNaturalEndpointsWithPort(String keysapceName, ByteBuffer key);
+ /** @deprecated See CASSANDRA-17665 */
+ @Deprecated(since = "7.0") public List getNaturalEndpointsWithPort(String keysapceName, ByteBuffer key);
+
+ /**
+ * This method returns the N replicas that are responsible for storing the
+ * specified key i.e for replication.
+ *
+ * @param keyspaceName keyspace name
+ * @param cf Column family name
+ * @param key - key for which we need to find the replica return value -
+ * the replica responsible for this key
+ */
+ public List getNaturalReplicas(String keyspaceName, String cf, String key);
+ public List getNaturalReplicasWithPort(String keyspaceName, String cf, String key);
/**
* @deprecated use {@link #takeSnapshot(String tag, Map options, String... entities)} instead. See CASSANDRA-10907
diff --git a/src/java/org/apache/cassandra/service/accord/AccordService.java b/src/java/org/apache/cassandra/service/accord/AccordService.java
index 340760370c..111e207474 100644
--- a/src/java/org/apache/cassandra/service/accord/AccordService.java
+++ b/src/java/org/apache/cassandra/service/accord/AccordService.java
@@ -458,7 +458,9 @@ public class AccordService implements IAccordService, Shutdownable
{
Invariants.require(localId != null, "static localId must be set before instantiating AccordService");
logger.info("Starting accord with nodeId {}", localId);
- AccordAgent agent = FBUtilities.construct(CassandraRelevantProperties.ACCORD_AGENT_CLASS.getString(AccordAgent.class.getName()), "AccordAgent");
+ AccordAgent agent = FBUtilities.construct(CassandraRelevantProperties.ACCORD_AGENT_CLASS.getString(AccordAgent.class.getName()),
+ "AccordAgent",
+ AccordAgent.class);
agent.setup(localId);
AccordTimeService time = new AccordTimeService();
this.scheduler = new AccordScheduler();
diff --git a/src/java/org/apache/cassandra/streaming/StreamHook.java b/src/java/org/apache/cassandra/streaming/StreamHook.java
index 84db420f04..df1e8ee76e 100644
--- a/src/java/org/apache/cassandra/streaming/StreamHook.java
+++ b/src/java/org/apache/cassandra/streaming/StreamHook.java
@@ -37,7 +37,7 @@ public interface StreamHook
String className = STREAM_HOOK.getString();
if (className != null)
{
- return FBUtilities.construct(className, StreamHook.class.getSimpleName());
+ return FBUtilities.construct(className, StreamHook.class.getSimpleName(), StreamHook.class);
}
else
{
diff --git a/src/java/org/apache/cassandra/tcm/extensions/ExtensionKey.java b/src/java/org/apache/cassandra/tcm/extensions/ExtensionKey.java
index 976e327951..2165fceddd 100644
--- a/src/java/org/apache/cassandra/tcm/extensions/ExtensionKey.java
+++ b/src/java/org/apache/cassandra/tcm/extensions/ExtensionKey.java
@@ -41,7 +41,7 @@ public class ExtensionKey> extends MetadataKey
public K newValue()
{
- return valueType.cast(FBUtilities.construct(valueType.getName(), "extension value"));
+ return FBUtilities.construct(valueType.getName(), "extension value", valueType);
}
public static final class Serializer implements MetadataSerializer>
@@ -58,7 +58,9 @@ public class ExtensionKey> extends MetadataKey
{
String id = in.readUTF();
String valType = in.readUTF();
- return new ExtensionKey(id, FBUtilities.classForName(valType, "value type"));
+ Class extends ExtensionValue> valueType =
+ FBUtilities.classForNameWithoutInitialization(valType, "value type", ExtensionValue.class);
+ return new ExtensionKey(id, valueType);
}
@Override
@@ -68,4 +70,3 @@ public class ExtensionKey> extends MetadataKey
}
}
}
-
diff --git a/src/java/org/apache/cassandra/tcm/sequences/SingleNodeSequences.java b/src/java/org/apache/cassandra/tcm/sequences/SingleNodeSequences.java
index 313700aaaf..82a1059b6e 100644
--- a/src/java/org/apache/cassandra/tcm/sequences/SingleNodeSequences.java
+++ b/src/java/org/apache/cassandra/tcm/sequences/SingleNodeSequences.java
@@ -92,6 +92,7 @@ public interface SingleNodeSequences
else if (InProgressSequences.isLeave(inProgress))
{
logger.info("Resuming decommission @ {} (current epoch = {}): {}", inProgress.latestModification, metadata.epoch, inProgress.status());
+ StorageService.instance.clearTransientMode();
}
else
{
diff --git a/src/java/org/apache/cassandra/tools/NodeProbe.java b/src/java/org/apache/cassandra/tools/NodeProbe.java
index 353324e0ea..c60e0427e5 100644
--- a/src/java/org/apache/cassandra/tools/NodeProbe.java
+++ b/src/java/org/apache/cassandra/tools/NodeProbe.java
@@ -1200,14 +1200,14 @@ public class NodeProbe implements AutoCloseable
ssProxy.setHintedHandoffThrottleInKB(throttleInKB);
}
- public List getEndpointsWithPort(String keyspace, String cf, String key)
+ public List getReplicasWithPort(String keyspace, String cf, String key)
{
- return ssProxy.getNaturalEndpointsWithPort(keyspace, cf, key);
+ return ssProxy.getNaturalReplicasWithPort(keyspace, cf, key);
}
- public List getEndpoints(String keyspace, String cf, String key)
+ public List getReplicas(String keyspace, String cf, String key)
{
- return ssProxy.getNaturalEndpoints(keyspace, cf, key);
+ return ssProxy.getNaturalReplicas(keyspace, cf, key);
}
public List getSSTables(String keyspace, String cf, String key, boolean hexFormat)
diff --git a/src/java/org/apache/cassandra/tools/nodetool/GetEndpoints.java b/src/java/org/apache/cassandra/tools/nodetool/GetEndpoints.java
index 3f30680951..6f0b6ea972 100644
--- a/src/java/org/apache/cassandra/tools/nodetool/GetEndpoints.java
+++ b/src/java/org/apache/cassandra/tools/nodetool/GetEndpoints.java
@@ -17,62 +17,10 @@
*/
package org.apache.cassandra.tools.nodetool;
-import java.net.InetAddress;
-import java.util.ArrayList;
-import java.util.List;
-
-import org.apache.cassandra.tools.NodeProbe;
-import org.apache.cassandra.tools.nodetool.layout.CassandraUsage;
-
import picocli.CommandLine.Command;
-import picocli.CommandLine.Mixin;
-import picocli.CommandLine.Parameters;
-import static com.google.common.base.Preconditions.checkArgument;
-import static org.apache.cassandra.tools.nodetool.CommandUtils.concatArgs;
-
-@Command(name = "getendpoints", description = "Print the end points that owns the key")
-public class GetEndpoints extends AbstractCommand
+@Deprecated(since = "7.0") // this is alias to getreplicas
+@Command(name = "getendpoints", description = "Print the end points that owns the key, deprecated, use getreplicas instead")
+public class GetEndpoints extends GetReplicas
{
- @CassandraUsage(usage = " ", description = "The keyspace, the table, and the partition key for which we need to find the endpoint")
- private List args = new ArrayList<>();
-
- @Parameters(index = "0", arity = "0..1", description = "The keyspace for which we need to find the endpoint")
- private String keyspace;
-
- @Parameters(index = "1", arity = "0..1", description = "The table for which we need to find the endpoint")
- private String table;
-
- @Parameters(index = "2", arity = "0..1", description = "The partition key for which we need to find the endpoint")
- private String key;
-
- @Mixin
- private PrintPortMixin printPortMixin = new PrintPortMixin();
-
- @Override
- public void execute(NodeProbe probe)
- {
- args = concatArgs(keyspace, table, key);
-
- checkArgument(args.size() == 3, "getendpoints requires keyspace, table and partition key arguments");
- String ks = args.get(0);
- String table = args.get(1);
- String key = args.get(2);
-
- if (printPortMixin.printPort)
- {
- for (String endpoint : probe.getEndpointsWithPort(ks, table, key))
- {
- probe.output().out.println(endpoint);
- }
- }
- else
- {
- List endpoints = probe.getEndpoints(ks, table, key);
- for (InetAddress endpoint : endpoints)
- {
- probe.output().out.println(endpoint.getHostAddress());
- }
- }
- }
}
diff --git a/src/java/org/apache/cassandra/tools/nodetool/GetReplicas.java b/src/java/org/apache/cassandra/tools/nodetool/GetReplicas.java
new file mode 100644
index 0000000000..5c91d20a7d
--- /dev/null
+++ b/src/java/org/apache/cassandra/tools/nodetool/GetReplicas.java
@@ -0,0 +1,79 @@
+/*
+ * 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.tools.nodetool;
+
+import java.net.InetAddress;
+import java.util.ArrayList;
+import java.util.List;
+
+import org.apache.cassandra.tools.NodeProbe;
+import org.apache.cassandra.tools.nodetool.layout.CassandraUsage;
+
+import picocli.CommandLine.Command;
+import picocli.CommandLine.Mixin;
+import picocli.CommandLine.Parameters;
+
+import static com.google.common.base.Preconditions.checkArgument;
+import static org.apache.cassandra.tools.nodetool.CommandUtils.concatArgs;
+
+@Command(name = "getreplicas", description = "Print the replicas that own the key")
+public class GetReplicas extends AbstractCommand
+{
+ @CassandraUsage(usage = " ", description = "The keyspace, the table, and the partition key for which we need to find the replica (e.g., pk1:pk2:pk3 for compound keys)")
+ private List args = new ArrayList<>();
+
+ @Parameters(index = "0", arity = "0..1", description = "The keyspace for which we need to find the replica")
+ private String keyspace;
+
+ @Parameters(index = "1", arity = "0..1", description = "The table for which we need to find the replica")
+ private String table;
+
+ @Parameters(index = "2", arity = "0..1", description = "The partition key for which we need to find the replica (e.g., pk1:pk2:pk3 for compound keys)")
+ private String key;
+
+ @Mixin
+ private PrintPortMixin printPortMixin = new PrintPortMixin();
+
+ @Override
+ public void execute(NodeProbe probe)
+ {
+ args = concatArgs(keyspace, table, key);
+
+ checkArgument(args.size() == 3, "requires keyspace, table and partition key arguments");
+ String ks = args.get(0);
+ String table = args.get(1);
+ String key = args.get(2);
+
+ if (printPortMixin.printPort)
+ {
+ for (String replica : probe.getReplicasWithPort(ks, table, key))
+ {
+ probe.output().out.println(replica);
+ }
+ }
+ else
+ {
+ List replicas = probe.getReplicas(ks, table, key);
+ for (InetAddress replica : replicas)
+ {
+ probe.output().out.println(replica.getHostAddress());
+ }
+ }
+ }
+}
diff --git a/src/java/org/apache/cassandra/tools/nodetool/InvalidatePermissionsCache.java b/src/java/org/apache/cassandra/tools/nodetool/InvalidatePermissionsCache.java
index e1fcc6b26a..7f1372b33d 100644
--- a/src/java/org/apache/cassandra/tools/nodetool/InvalidatePermissionsCache.java
+++ b/src/java/org/apache/cassandra/tools/nodetool/InvalidatePermissionsCache.java
@@ -38,7 +38,7 @@ import static com.google.common.base.Preconditions.checkArgument;
@Command(name = "invalidatepermissionscache", description = "Invalidate the permissions cache")
public class InvalidatePermissionsCache extends AbstractCommand
{
- @Parameters(paramLabel = "role", description = "A role for which permissions to specified resources need to be invalidated", arity = "0..1", index = "0")
+ @Parameters(paramLabel = "role_name", description = "A role for which permissions to specified resources need to be invalidated", arity = "0..1", index = "0")
private String roleName;
// Data Resources
diff --git a/src/java/org/apache/cassandra/tools/nodetool/NodetoolCommand.java b/src/java/org/apache/cassandra/tools/nodetool/NodetoolCommand.java
index af1e835d91..fd2ccde672 100644
--- a/src/java/org/apache/cassandra/tools/nodetool/NodetoolCommand.java
+++ b/src/java/org/apache/cassandra/tools/nodetool/NodetoolCommand.java
@@ -120,6 +120,7 @@ import static org.apache.cassandra.tools.nodetool.Help.printTopCommandUsage;
GetInterDCStreamThroughput.class,
GetLoggingLevels.class,
GetMaxHintWindow.class,
+ GetReplicas.class,
GetSSTables.class,
GetSeeds.class,
GetSnapshotThrottle.class,
diff --git a/src/java/org/apache/cassandra/tools/nodetool/Sjk.java b/src/java/org/apache/cassandra/tools/nodetool/Sjk.java
index c117ed638b..f9d19eabd2 100644
--- a/src/java/org/apache/cassandra/tools/nodetool/Sjk.java
+++ b/src/java/org/apache/cassandra/tools/nodetool/Sjk.java
@@ -393,6 +393,7 @@ public class Sjk extends AbstractCommand
List> result = new ArrayList<>();
try
{
+ ClassLoader cl = Thread.currentThread().getContextClassLoader();
String path = packageName.replace('.', '/');
for (String f : findFiles(path))
{
@@ -400,7 +401,7 @@ public class Sjk extends AbstractCommand
{
f = f.substring(0, f.length() - ".class".length());
f = f.replace('/', '.');
- result.add(Class.forName(f));
+ result.add(Class.forName(f, false, cl));
}
}
return result;
diff --git a/src/java/org/apache/cassandra/tracing/Tracing.java b/src/java/org/apache/cassandra/tracing/Tracing.java
index 4ef2e500cb..8d621cb668 100644
--- a/src/java/org/apache/cassandra/tracing/Tracing.java
+++ b/src/java/org/apache/cassandra/tracing/Tracing.java
@@ -117,7 +117,7 @@ public abstract class Tracing extends ExecutorLocals.Impl
{
try
{
- tracing = FBUtilities.construct(customTracingClass, "Tracing");
+ tracing = FBUtilities.construct(customTracingClass, "Tracing", Tracing.class);
logger.info("Using the {} class to trace queries (as requested by the {} system property)",
customTracingClass, CUSTOM_TRACING_CLASS.getKey());
}
diff --git a/src/java/org/apache/cassandra/transport/CBUtil.java b/src/java/org/apache/cassandra/transport/CBUtil.java
index a456b6e2a3..260350cce3 100644
--- a/src/java/org/apache/cassandra/transport/CBUtil.java
+++ b/src/java/org/apache/cassandra/transport/CBUtil.java
@@ -472,6 +472,9 @@ public abstract class CBUtil
int length = cb.readInt();
if (length < 0)
return null;
+ if (length > cb.readableBytes())
+ throw new ProtocolException(String.format("Cannot read value of length %d, only %d bytes remaining in the message",
+ length, cb.readableBytes()));
ByteBuffer buffer = cb.nioBuffer(cb.readerIndex(), length);
cb.skipBytes(length);
@@ -738,6 +741,9 @@ public abstract class CBUtil
private static byte[] readRawBytes(ByteBuf cb, int length)
{
+ if (length > cb.readableBytes())
+ throw new ProtocolException(String.format("Cannot read value of length %d, only %d bytes remaining in the message",
+ length, cb.readableBytes()));
byte[] bytes = new byte[length];
cb.readBytes(bytes);
return bytes;
diff --git a/src/java/org/apache/cassandra/transport/CQLMessageHandler.java b/src/java/org/apache/cassandra/transport/CQLMessageHandler.java
index 2486497c2f..13cb7463dd 100644
--- a/src/java/org/apache/cassandra/transport/CQLMessageHandler.java
+++ b/src/java/org/apache/cassandra/transport/CQLMessageHandler.java
@@ -106,7 +106,7 @@ public class CQLMessageHandler extends AbstractMessageHandler
interface MessageConsumer
{
- void dispatch(Channel channel, M message, Dispatcher.FlushItemConverter toFlushItem, Overload backpressure);
+ void dispatch(Channel channel, M message, Dispatcher.FlushItemConverter toFlushItem, P param, Overload backpressure);
boolean hasQueueCapacity();
}
@@ -159,6 +159,12 @@ public class CQLMessageHandler extends AbstractMessageHandler
return super.process(frame);
}
+ @Override
+ protected void onConnectionUnpaused()
+ {
+ ClientMetrics.instance.unpauseConnection();
+ }
+
/**
* Checks limits on bytes in flight and the request rate limiter (if enabled), then takes one of three actions:
*
@@ -389,7 +395,7 @@ public class CQLMessageHandler extends AbstractMessageHandler
try
{
message = messageDecoder.decode(channel, request);
- dispatcher.dispatch(channel, message, this::toFlushItem, backpressure);
+ dispatcher.dispatch(channel, message, CQLMessageHandler::toFlushItem, this, backpressure);
// sucessfully delivered a CQL message to the execution
// stage, so reset the counter of consecutive errors
@@ -485,7 +491,8 @@ public class CQLMessageHandler extends AbstractMessageHandler
// The Dispatcher will call this to obtain the FlushItem to enqueue with its Flusher once
// a dispatched request has been processed.
- Envelope responseFrame = response.encode(request.getSource().header.version);
+ Envelope.Header header = request.getSource().header;
+ Envelope responseFrame = response.encode(header.version, header.streamId);
int responseSize = envelopeSize(responseFrame.header);
ClientMessageSizeMetrics.bytesSent.inc(responseSize);
ClientMessageSizeMetrics.bytesSentPerResponse.update(responseSize);
@@ -500,9 +507,9 @@ public class CQLMessageHandler extends AbstractMessageHandler
@Override
public void cleanup(Flusher.FlushItem flushItem)
{
- release(flushItem.request.header);
- flushItem.request.release();
- flushItem.response.release();
+ release(flushItem.requestEnvelope.header);
+ flushItem.requestEnvelope.release();
+ flushItem.responseEnvelope.release();
}
private void release(Envelope.Header header)
@@ -524,8 +531,9 @@ public class CQLMessageHandler extends AbstractMessageHandler
if (!extracted.isSuccess())
{
// Hard fail on any decoding error as we can't trust the subsequent frames of
- // the large message
- handleError(ProtocolException.toFatalException(extracted.error()));
+ // the large message. The stream id is a best-effort value read before extraction
+ // failed, so route it back where possible rather than defaulting.
+ handleError(ProtocolException.toFatalException(extracted.error()), extracted.streamId());
return false;
}
@@ -542,7 +550,7 @@ public class CQLMessageHandler extends AbstractMessageHandler
// not make sense to continue processing subsequent frames
handleError(ProtocolException.toFatalException(new OversizedAuthMessageException(
MULTI_FRAME_AUTH_ERROR_MESSAGE_PREFIX +
- "type = " + header.type + ", size = " + header.bodySizeInBytes)));
+ "type = " + header.type + ", size = " + header.bodySizeInBytes)), header.streamId);
ClientMetrics.instance.markRequestDiscarded();
return false;
}
diff --git a/src/java/org/apache/cassandra/transport/Dispatcher.java b/src/java/org/apache/cassandra/transport/Dispatcher.java
index 3f6468c218..2e1b36a4f9 100644
--- a/src/java/org/apache/cassandra/transport/Dispatcher.java
+++ b/src/java/org/apache/cassandra/transport/Dispatcher.java
@@ -93,10 +93,9 @@ public class Dispatcher implements CQLMessageHandler.MessageConsumer
{
- FlushItem> toFlushItem(Channel channel, Message.Request request, Message.Response response);
+ FlushItem> toFlushItem(P param, Channel channel, Message.Request request, Message.Response response);
}
public Dispatcher(boolean useLegacyFlusher)
@@ -105,18 +104,17 @@ public class Dispatcher implements CQLMessageHandler.MessageConsumer void dispatch(Channel channel, Message.Request request, FlushItemConverter forFlusher, P param, Overload backpressure)
{
if (!request.connection().getTracker().isRunning())
{
// We can not respond with a custom, transport, or server exceptions since, given current implementation of clients,
// they will defunct the connection. Without a protocol version bump that introduces an "I am going away message",
// we have to stick to an existing error code.
- Message.Response response = ErrorMessage.fromException(new OverloadedException("Server is shutting down"));
- response.setStreamId(request.getStreamId());
+ Message.Response response = ErrorMessage.fromTransportException(new OverloadedException("Server is shutting down"));
response.setWarnings(ClientWarn.instance.getWarnings());
response.attach(request.connection);
- FlushItem> toFlush = forFlusher.toFlushItem(channel, request, response);
+ FlushItem> toFlush = forFlusher.toFlushItem(param, channel, request, response);
flush(toFlush);
return;
}
@@ -128,7 +126,7 @@ public class Dispatcher implements CQLMessageHandler.MessageConsumer(channel, request, forFlusher, param, backpressure));
ClientMetrics.instance.markRequestDispatched();
}
@@ -293,20 +291,22 @@ public class Dispatcher implements CQLMessageHandler.MessageConsumer implements DebuggableTask.RunnableDebuggableTask
{
private final Channel channel;
private final Message.Request request;
- private final FlushItemConverter forFlusher;
+ private final FlushItemConverter forFlusher;
+ private final P flusherParam;
private final Overload backpressure;
private volatile long startTimeNanos;
- public RequestProcessor(Channel channel, Message.Request request, FlushItemConverter forFlusher, Overload backpressure)
+ public RequestProcessor(Channel channel, Message.Request request, FlushItemConverter forFlusher, P flusherParam, Overload backpressure)
{
this.channel = channel;
this.request = request;
this.forFlusher = forFlusher;
+ this.flusherParam = flusherParam;
this.backpressure = backpressure;
}
@@ -314,7 +314,7 @@ public class Dispatcher implements CQLMessageHandler.MessageConsumer DatabaseDescriptor.getNativeTransportTimeout(TimeUnit.NANOSECONDS))
{
ClientMetrics.instance.markTimedOutBeforeProcessing();
- return ErrorMessage.fromException(new OverloadedException("Query timed out before it could start"));
+ return ErrorMessage.fromTransportException(new OverloadedException("Query timed out before it could start"));
}
if (connection.getVersion().isGreaterOrEqualTo(ProtocolVersion.V4))
@@ -434,8 +434,6 @@ public class Dispatcher implements CQLMessageHandler.MessageConsumer handler = ExceptionHandlers.getUnexpectedExceptionHandler(channel, true);
- ErrorMessage error = ErrorMessage.fromException(t, handler);
- error.setStreamId(request.getStreamId());
- error.setWarnings(ClientWarn.instance.getWarnings());
- return error;
+ response = ErrorMessage.fromExceptionNoStreamId(t, handler);
}
finally
{
+ if (response != null)
+ response.setWarnings(ClientWarn.instance.getWarnings());
CoordinatorWarnings.reset();
CoordinatorWriteWarnings.reset();
ClientWarn.instance.resetWarnings();
}
+ return response;
}
/**
* Note: this method is not expected to execute on the netty event loop.
*/
- void processRequest(Channel channel, Message.Request request, FlushItemConverter forFlusher, Overload backpressure, RequestTime requestTime)
+ | | | |