points = row.getList("points", BytesType.instance);
+
+ return PaxosRepairHistory.fromTupleBufferList(points);
}
/**
diff --git a/src/java/org/apache/cassandra/db/commitlog/CommitLogReplayer.java b/src/java/org/apache/cassandra/db/commitlog/CommitLogReplayer.java
index b59480e05a..cbed01885b 100644
--- a/src/java/org/apache/cassandra/db/commitlog/CommitLogReplayer.java
+++ b/src/java/org/apache/cassandra/db/commitlog/CommitLogReplayer.java
@@ -52,7 +52,6 @@ import org.apache.cassandra.schema.Schema;
import org.apache.cassandra.schema.SchemaConstants;
import org.apache.cassandra.schema.TableId;
import org.apache.cassandra.schema.TableMetadataRef;
-import org.apache.cassandra.service.StorageService;
import org.apache.cassandra.utils.FBUtilities;
import org.apache.cassandra.utils.WrappedRunnable;
diff --git a/src/java/org/apache/cassandra/db/compaction/AbstractCompactionStrategy.java b/src/java/org/apache/cassandra/db/compaction/AbstractCompactionStrategy.java
index 4aec19ff3c..5fe1df70bd 100644
--- a/src/java/org/apache/cassandra/db/compaction/AbstractCompactionStrategy.java
+++ b/src/java/org/apache/cassandra/db/compaction/AbstractCompactionStrategy.java
@@ -38,7 +38,6 @@ import org.apache.cassandra.db.lifecycle.LifecycleTransaction;
import org.apache.cassandra.dht.Range;
import org.apache.cassandra.dht.Token;
import org.apache.cassandra.exceptions.ConfigurationException;
-import org.apache.cassandra.io.sstable.Component;
import org.apache.cassandra.io.sstable.ISSTableScanner;
import org.apache.cassandra.io.sstable.metadata.MetadataCollector;
import org.apache.cassandra.io.sstable.metadata.StatsMetadata;
diff --git a/src/java/org/apache/cassandra/db/compaction/CompactionInterruptedException.java b/src/java/org/apache/cassandra/db/compaction/CompactionInterruptedException.java
index 809c10c8cf..b9174ec262 100644
--- a/src/java/org/apache/cassandra/db/compaction/CompactionInterruptedException.java
+++ b/src/java/org/apache/cassandra/db/compaction/CompactionInterruptedException.java
@@ -17,16 +17,14 @@
*/
package org.apache.cassandra.db.compaction;
+import org.apache.cassandra.utils.Shared;
+
+@Shared
public class CompactionInterruptedException extends RuntimeException
{
private static final long serialVersionUID = -8651427062512310398L;
- public CompactionInterruptedException(CompactionInfo info)
- {
- super("Compaction interrupted: " + info);
- }
-
- public CompactionInterruptedException(String info)
+ public CompactionInterruptedException(Object info)
{
super("Compaction interrupted: " + info);
}
diff --git a/src/java/org/apache/cassandra/db/compaction/CompactionIterator.java b/src/java/org/apache/cassandra/db/compaction/CompactionIterator.java
index fc4dc9c721..86ab5ffc6c 100644
--- a/src/java/org/apache/cassandra/db/compaction/CompactionIterator.java
+++ b/src/java/org/apache/cassandra/db/compaction/CompactionIterator.java
@@ -19,11 +19,17 @@ package org.apache.cassandra.db.compaction;
import java.util.*;
import java.util.function.LongPredicate;
+import java.util.concurrent.TimeUnit;
import com.google.common.collect.ImmutableSet;
import com.google.common.collect.Ordering;
+import org.apache.cassandra.config.DatabaseDescriptor;
+import org.apache.cassandra.dht.Token;
import org.apache.cassandra.io.sstable.format.SSTableReader;
+import org.apache.cassandra.schema.Schema;
+import org.apache.cassandra.schema.SchemaConstants;
+import org.apache.cassandra.schema.TableId;
import org.apache.cassandra.schema.TableMetadata;
import org.apache.cassandra.db.transform.DuplicateRowChecker;
@@ -37,8 +43,14 @@ import org.apache.cassandra.db.transform.Transformation;
import org.apache.cassandra.index.transactions.CompactionTransaction;
import org.apache.cassandra.io.sstable.ISSTableScanner;
import org.apache.cassandra.schema.CompactionParams.TombstoneOption;
+import org.apache.cassandra.service.paxos.PaxosRepairHistory;
+import org.apache.cassandra.service.paxos.uncommitted.PaxosRows;
import org.apache.cassandra.utils.TimeUUID;
+import static java.util.concurrent.TimeUnit.MICROSECONDS;
+import static org.apache.cassandra.config.Config.PaxosStatePurging.legacy;
+import static org.apache.cassandra.config.DatabaseDescriptor.paxosStatePurging;
+
/**
* Merge multiple iterators over the content of sstable into a "compacted" iterator.
*
@@ -110,7 +122,10 @@ public class CompactionIterator extends CompactionInfo.Holder implements Unfilte
? EmptyIterators.unfilteredPartition(controller.cfs.metadata())
: UnfilteredPartitionIterators.merge(scanners, listener());
merged = Transformation.apply(merged, new GarbageSkipper(controller));
- merged = Transformation.apply(merged, new Purger(controller, nowInSec));
+ Transformation purger = isPaxos(controller.cfs) && paxosStatePurging() != legacy
+ ? new PaxosPurger(nowInSec)
+ : new Purger(controller, nowInSec);
+ merged = Transformation.apply(merged, purger);
merged = DuplicateRowChecker.duringCompaction(merged, type);
compacted = Transformation.apply(merged, new AbortableUnfilteredPartitionTransformation(this));
}
@@ -556,6 +571,78 @@ public class CompactionIterator extends CompactionInfo.Holder implements Unfilte
}
}
+ private class PaxosPurger extends Transformation
+ {
+ private final long nowInSec;
+ private final long paxosPurgeGraceMicros = DatabaseDescriptor.getPaxosPurgeGrace(MICROSECONDS);
+ private final Map tableIdToHistory = new HashMap<>();
+ private Token currentToken;
+ private int compactedUnfiltered;
+
+ private PaxosPurger(long nowInSec)
+ {
+ this.nowInSec = nowInSec;
+ }
+
+ protected void onEmptyPartitionPostPurge(DecoratedKey key)
+ {
+ if (type == OperationType.COMPACTION)
+ controller.cfs.invalidateCachedPartition(key);
+ }
+
+ protected void updateProgress()
+ {
+ if ((++compactedUnfiltered) % UNFILTERED_TO_UPDATE_PROGRESS == 0)
+ updateBytesRead();
+ }
+
+ @Override
+ @SuppressWarnings("resource")
+ protected UnfilteredRowIterator applyToPartition(UnfilteredRowIterator partition)
+ {
+ currentToken = partition.partitionKey().getToken();
+ UnfilteredRowIterator purged = Transformation.apply(partition, this);
+ if (purged.isEmpty())
+ {
+ onEmptyPartitionPostPurge(purged.partitionKey());
+ purged.close();
+ return null;
+ }
+
+ return purged;
+ }
+
+ @Override
+ protected Row applyToRow(Row row)
+ {
+ updateProgress();
+ TableId tableId = PaxosRows.getTableId(row);
+
+ switch (paxosStatePurging())
+ {
+ default: throw new AssertionError();
+ case legacy:
+ case gc_grace:
+ {
+ TableMetadata metadata = Schema.instance.getTableMetadata(tableId);
+ return row.purgeDataOlderThan(TimeUnit.SECONDS.toMicros(nowInSec - (metadata == null ? (3 * 3600) : metadata.params.gcGraceSeconds)), false);
+ }
+ case repaired:
+ {
+ PaxosRepairHistory.Searcher history = tableIdToHistory.computeIfAbsent(tableId, find -> {
+ TableMetadata metadata = Schema.instance.getTableMetadata(find);
+ if (metadata == null)
+ return null;
+ return Keyspace.openAndGetStore(metadata).getPaxosRepairHistory().searcher();
+ });
+
+ return history == null ? row :
+ row.purgeDataOlderThan(history.ballotForToken(currentToken).unixMicros() - paxosPurgeGraceMicros, false);
+ }
+ }
+ }
+ }
+
private static class AbortableUnfilteredPartitionTransformation extends Transformation
{
private final AbortableUnfilteredRowTransformation abortableIter;
@@ -574,7 +661,7 @@ public class CompactionIterator extends CompactionInfo.Holder implements Unfilte
}
}
- private static class AbortableUnfilteredRowTransformation extends Transformation
+ private static class AbortableUnfilteredRowTransformation extends Transformation
{
private final CompactionIterator iter;
@@ -590,4 +677,9 @@ public class CompactionIterator extends CompactionInfo.Holder implements Unfilte
return row;
}
}
+
+ private static boolean isPaxos(ColumnFamilyStore cfs)
+ {
+ return cfs.name.equals(SystemKeyspace.PAXOS) && cfs.keyspace.getName().equals(SchemaConstants.SYSTEM_KEYSPACE_NAME);
+ }
}
diff --git a/src/java/org/apache/cassandra/db/marshal/ByteArrayAccessor.java b/src/java/org/apache/cassandra/db/marshal/ByteArrayAccessor.java
index bb26cba460..df24a627a4 100644
--- a/src/java/org/apache/cassandra/db/marshal/ByteArrayAccessor.java
+++ b/src/java/org/apache/cassandra/db/marshal/ByteArrayAccessor.java
@@ -29,6 +29,7 @@ import org.apache.cassandra.db.Digest;
import org.apache.cassandra.db.TypeSizes;
import org.apache.cassandra.io.util.DataInputPlus;
import org.apache.cassandra.io.util.DataOutputPlus;
+import org.apache.cassandra.service.paxos.Ballot;
import org.apache.cassandra.utils.TimeUUID;
import org.apache.cassandra.utils.ByteArrayUtil;
import org.apache.cassandra.utils.ByteBufferUtil;
@@ -241,6 +242,12 @@ public class ByteArrayAccessor implements ValueAccessor
return TimeUUID.fromBytes(getLong(value, 0), getLong(value, 8));
}
+ @Override
+ public Ballot toBallot(byte[] value)
+ {
+ return Ballot.deserialize(value);
+ }
+
@Override
public int putShort(byte[] dst, int offset, short value)
{
diff --git a/src/java/org/apache/cassandra/db/marshal/ByteBufferAccessor.java b/src/java/org/apache/cassandra/db/marshal/ByteBufferAccessor.java
index e428b7cd04..40a3bf4b34 100644
--- a/src/java/org/apache/cassandra/db/marshal/ByteBufferAccessor.java
+++ b/src/java/org/apache/cassandra/db/marshal/ByteBufferAccessor.java
@@ -28,6 +28,7 @@ import org.apache.cassandra.db.Digest;
import org.apache.cassandra.db.TypeSizes;
import org.apache.cassandra.io.util.DataInputPlus;
import org.apache.cassandra.io.util.DataOutputPlus;
+import org.apache.cassandra.service.paxos.Ballot;
import org.apache.cassandra.utils.TimeUUID;
import org.apache.cassandra.utils.ByteBufferUtil;
import org.apache.cassandra.utils.FastByteOperations;
@@ -245,6 +246,12 @@ public class ByteBufferAccessor implements ValueAccessor
return TimeUUID.fromBytes(value.getLong(value.position()), value.getLong(value.position() + 8));
}
+ @Override
+ public Ballot toBallot(ByteBuffer value)
+ {
+ return Ballot.deserialize(value);
+ }
+
@Override
public int putShort(ByteBuffer dst, int offset, short value)
{
diff --git a/src/java/org/apache/cassandra/db/marshal/TupleType.java b/src/java/org/apache/cassandra/db/marshal/TupleType.java
index 83fbb25d54..cc08487658 100644
--- a/src/java/org/apache/cassandra/db/marshal/TupleType.java
+++ b/src/java/org/apache/cassandra/db/marshal/TupleType.java
@@ -205,9 +205,17 @@ public class TupleType extends AbstractType
*/
public ByteBuffer[] split(ByteBuffer value)
{
- ByteBuffer[] components = new ByteBuffer[size()];
+ return split(value, size(), this);
+ }
+
+ /**
+ * Split a tuple value into its component values.
+ */
+ public static ByteBuffer[] split(ByteBuffer value, int numberOfElements, TupleType type)
+ {
+ ByteBuffer[] components = new ByteBuffer[numberOfElements];
ByteBuffer input = value.duplicate();
- for (int i = 0; i < size(); i++)
+ for (int i = 0; i < numberOfElements; i++)
{
if (!input.hasRemaining())
return Arrays.copyOfRange(components, 0, i);
@@ -226,7 +234,7 @@ public class TupleType extends AbstractType
{
throw new InvalidRequestException(String.format(
"Expected %s %s for %s column, but got more",
- size(), size() == 1 ? "value" : "values", this.asCQL3Type()));
+ numberOfElements, numberOfElements == 1 ? "value" : "values", type.asCQL3Type()));
}
return components;
diff --git a/src/java/org/apache/cassandra/db/marshal/ValueAccessor.java b/src/java/org/apache/cassandra/db/marshal/ValueAccessor.java
index 2f089a6c49..a51836e65a 100644
--- a/src/java/org/apache/cassandra/db/marshal/ValueAccessor.java
+++ b/src/java/org/apache/cassandra/db/marshal/ValueAccessor.java
@@ -37,6 +37,7 @@ import org.apache.cassandra.db.rows.CellPath;
import org.apache.cassandra.io.util.DataInputPlus;
import org.apache.cassandra.io.util.DataOutputPlus;
import org.apache.cassandra.schema.ColumnMetadata;
+import org.apache.cassandra.service.paxos.Ballot;
import org.apache.cassandra.utils.TimeUUID;
import static org.apache.cassandra.db.ClusteringPrefix.Kind.*;
@@ -326,6 +327,9 @@ public interface ValueAccessor
/** returns a TimeUUID from offset 0 */
TimeUUID toTimeUUID(V value);
+ /** returns a TimeUUID from offset 0 */
+ Ballot toBallot(V value);
+
/**
* writes the short value {@param value} to {@param dst} at offset {@param offset}
* @return the number of bytes written to {@param value}
diff --git a/src/java/org/apache/cassandra/db/partitions/AbstractBTreePartition.java b/src/java/org/apache/cassandra/db/partitions/AbstractBTreePartition.java
index 1d6603eec2..b51a6787a0 100644
--- a/src/java/org/apache/cassandra/db/partitions/AbstractBTreePartition.java
+++ b/src/java/org/apache/cassandra/db/partitions/AbstractBTreePartition.java
@@ -375,26 +375,52 @@ public abstract class AbstractBTreePartition implements Partition, Iterable
@Override
public String toString()
{
- StringBuilder sb = new StringBuilder();
+ return toString(true);
+ }
- sb.append(String.format("[%s] key=%s partition_deletion=%s columns=%s",
- metadata(),
- metadata().partitionKeyType.getString(partitionKey().getKey()),
- partitionLevelDeletion(),
- columns()));
+ public String toString(boolean includeFullDetails)
+ {
+ StringBuilder sb = new StringBuilder();
+ if (includeFullDetails)
+ {
+ sb.append(String.format("[%s.%s] key=%s partition_deletion=%s columns=%s",
+ metadata().keyspace,
+ metadata().name,
+ metadata().partitionKeyType.getString(partitionKey().getKey()),
+ partitionLevelDeletion(),
+ columns()));
+ }
+ else
+ {
+ sb.append("key=").append(metadata().partitionKeyType.getString(partitionKey().getKey()));
+ }
if (staticRow() != Rows.EMPTY_STATIC_ROW)
- sb.append("\n ").append(staticRow().toString(metadata(), true));
+ sb.append("\n ").append(staticRow().toString(metadata(), includeFullDetails));
try (UnfilteredRowIterator iter = unfilteredIterator())
{
while (iter.hasNext())
- sb.append("\n ").append(iter.next().toString(metadata(), true));
+ sb.append("\n ").append(iter.next().toString(metadata(), includeFullDetails));
}
-
return sb.toString();
}
+ @Override
+ public boolean equals(Object obj)
+ {
+ if (!(obj instanceof PartitionUpdate))
+ return false;
+
+ PartitionUpdate that = (PartitionUpdate) obj;
+ Holder a = this.holder(), b = that.holder();
+ return partitionKey.equals(that.partitionKey)
+ && metadata().id.equals(that.metadata().id)
+ && a.deletionInfo.equals(b.deletionInfo)
+ && a.staticRow.equals(b.staticRow)
+ && Iterators.elementsEqual(iterator(), that.iterator());
+ }
+
public int rowCount()
{
return BTree.size(holder().tree);
diff --git a/src/java/org/apache/cassandra/db/partitions/PartitionUpdate.java b/src/java/org/apache/cassandra/db/partitions/PartitionUpdate.java
index 9530e10d3a..f6fd259e88 100644
--- a/src/java/org/apache/cassandra/db/partitions/PartitionUpdate.java
+++ b/src/java/org/apache/cassandra/db/partitions/PartitionUpdate.java
@@ -17,13 +17,16 @@
*/
package org.apache.cassandra.db.partitions;
+import java.io.EOFException;
import java.io.IOException;
import java.nio.ByteBuffer;
import java.util.ArrayList;
import java.util.List;
import com.google.common.collect.Iterables;
+import com.google.common.collect.Iterators;
import com.google.common.collect.Lists;
+import com.google.common.primitives.Ints;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -38,6 +41,9 @@ import org.apache.cassandra.schema.TableId;
import org.apache.cassandra.schema.TableMetadata;
import org.apache.cassandra.utils.btree.BTree;
import org.apache.cassandra.utils.btree.UpdateFunction;
+import org.apache.cassandra.utils.vint.VIntCoding;
+
+import static org.apache.cassandra.db.rows.UnfilteredRowIteratorSerializer.IS_EMPTY;
/**
* Stores updates made on a partition.
@@ -322,23 +328,19 @@ public class PartitionUpdate extends AbstractBTreePartition
*/
public int dataSize()
{
- int size = 0;
+ return Ints.saturatedCast(BTree.accumulate(holder.tree, (row, value) -> row.dataSize() + value, 0L)
+ + holder.staticRow.dataSize() + holder.deletionInfo.dataSize());
+ }
- if (holder.staticRow != null)
- {
- for (ColumnData cd : holder.staticRow.columnData())
- {
- size += cd.dataSize();
- }
- }
-
- for (Row row : this)
- {
- size += row.clustering().dataSize();
- for (ColumnData cd : row)
- size += cd.dataSize();
- }
- return size;
+ /**
+ * The size of the data contained in this update.
+ *
+ * @return the size of the data contained in this update.
+ */
+ public long unsharedHeapSize()
+ {
+ return BTree.accumulate(holder.tree, (row, value) -> row.unsharedHeapSize() + value, 0L)
+ + holder.staticRow.unsharedHeapSize() + holder.deletionInfo.unsharedHeapSize();
}
public TableMetadata metadata()
@@ -667,6 +669,21 @@ public class PartitionUpdate extends AbstractBTreePartition
false);
}
+ public static boolean isEmpty(ByteBuffer in, DeserializationHelper.Flag flag, DecoratedKey key) throws IOException
+ {
+ int position = in.position();
+ position += 16; // CFMetaData.serializer.deserialize(in, version);
+ if (position >= in.limit())
+ throw new EOFException();
+ // DecoratedKey key = metadata.decorateKey(ByteBufferUtil.readWithVIntLength(in));
+ int keyLength = (int) VIntCoding.getUnsignedVInt(in, position);
+ position += keyLength + VIntCoding.computeUnsignedVIntSize(keyLength);
+ if (position >= in.limit())
+ throw new EOFException();
+ int flags = in.get(position) & 0xff;
+ return (flags & IS_EMPTY) != 0;
+ }
+
public long serializedSize(PartitionUpdate update, int version)
{
try (UnfilteredRowIterator iter = update.unfilteredIterator())
diff --git a/src/java/org/apache/cassandra/db/rows/AbstractCell.java b/src/java/org/apache/cassandra/db/rows/AbstractCell.java
index 5d4fc3c19e..21d2c10e98 100644
--- a/src/java/org/apache/cassandra/db/rows/AbstractCell.java
+++ b/src/java/org/apache/cassandra/db/rows/AbstractCell.java
@@ -98,6 +98,11 @@ public abstract class AbstractCell extends Cell
return this;
}
+ public Cell> purgeDataOlderThan(long timestamp)
+ {
+ return this.timestamp() < timestamp ? null : this;
+ }
+
public Cell> copy(AbstractAllocator allocator)
{
CellPath path = path();
diff --git a/src/java/org/apache/cassandra/db/rows/BTreeRow.java b/src/java/org/apache/cassandra/db/rows/BTreeRow.java
index 2d3ee83de5..64b494f953 100644
--- a/src/java/org/apache/cassandra/db/rows/BTreeRow.java
+++ b/src/java/org/apache/cassandra/db/rows/BTreeRow.java
@@ -474,6 +474,18 @@ public class BTreeRow extends AbstractRow
return transformAndFilter(newInfo, newDeletion, (cd) -> cd.purge(purger, nowInSec));
}
+ public Row purgeDataOlderThan(long timestamp, boolean enforceStrictLiveness)
+ {
+ LivenessInfo newInfo = primaryKeyLivenessInfo.timestamp() < timestamp ? LivenessInfo.EMPTY : primaryKeyLivenessInfo;
+ Deletion newDeletion = deletion.time().markedForDeleteAt() < timestamp ? Deletion.LIVE : deletion;
+
+ // when enforceStrictLiveness is set, a row is considered dead when it's PK liveness info is not present
+ if (enforceStrictLiveness && newDeletion.isLive() && newInfo.isEmpty())
+ return null;
+
+ return transformAndFilter(newInfo, newDeletion, cd -> cd.purgeDataOlderThan(timestamp));
+ }
+
private Row transformAndFilter(LivenessInfo info, Deletion deletion, Function function)
{
Object[] transformed = BTree.transformAndFilter(btree, function);
diff --git a/src/java/org/apache/cassandra/db/rows/BufferCell.java b/src/java/org/apache/cassandra/db/rows/BufferCell.java
index 85d28f8b46..171dfbb0d7 100644
--- a/src/java/org/apache/cassandra/db/rows/BufferCell.java
+++ b/src/java/org/apache/cassandra/db/rows/BufferCell.java
@@ -26,7 +26,6 @@ import org.apache.cassandra.schema.ColumnMetadata;
import org.apache.cassandra.db.marshal.ByteType;
import org.apache.cassandra.utils.ByteBufferUtil;
import org.apache.cassandra.utils.ObjectSizes;
-import org.apache.cassandra.utils.memory.AbstractAllocator;
import static java.lang.String.format;
@@ -134,14 +133,6 @@ public class BufferCell extends AbstractCell
return withUpdatedValue(ByteBufferUtil.EMPTY_BYTE_BUFFER);
}
- public Cell> copy(AbstractAllocator allocator)
- {
- if (!value.hasRemaining())
- return this;
-
- return new BufferCell(column, timestamp, ttl, localDeletionTime, allocator.clone(value), path == null ? null : path.copy(allocator));
- }
-
@Override
public long unsharedHeapSize()
{
diff --git a/src/java/org/apache/cassandra/db/rows/Cell.java b/src/java/org/apache/cassandra/db/rows/Cell.java
index 38d1c84880..2a19a986b3 100644
--- a/src/java/org/apache/cassandra/db/rows/Cell.java
+++ b/src/java/org/apache/cassandra/db/rows/Cell.java
@@ -167,6 +167,10 @@ public abstract class Cell extends ColumnData
// Overrides super type to provide a more precise return type.
public abstract Cell> purge(DeletionPurger purger, int nowInSec);
+ @Override
+ // Overrides super type to provide a more precise return type.
+ public abstract Cell> purgeDataOlderThan(long timestamp);
+
/**
* The serialization format for cell is:
* [ flags ][ timestamp ][ deletion time ][ ttl ][ path size ][ path ][ value size ][ value ]
diff --git a/src/java/org/apache/cassandra/db/rows/CellPath.java b/src/java/org/apache/cassandra/db/rows/CellPath.java
index cee24840b5..817a2d0314 100644
--- a/src/java/org/apache/cassandra/db/rows/CellPath.java
+++ b/src/java/org/apache/cassandra/db/rows/CellPath.java
@@ -65,6 +65,8 @@ public abstract class CellPath implements IMeasurableMemory
public abstract long unsharedHeapSizeExcludingData();
+ public abstract long unsharedHeapSize();
+
@Override
public final int hashCode()
{
diff --git a/src/java/org/apache/cassandra/db/rows/ColumnData.java b/src/java/org/apache/cassandra/db/rows/ColumnData.java
index 4146946616..8ac19cc496 100644
--- a/src/java/org/apache/cassandra/db/rows/ColumnData.java
+++ b/src/java/org/apache/cassandra/db/rows/ColumnData.java
@@ -58,6 +58,8 @@ public abstract class ColumnData implements IMeasurableMemory
public abstract long unsharedHeapSizeExcludingData();
+ public abstract long unsharedHeapSize();
+
/**
* Validate the column data.
*
@@ -95,6 +97,7 @@ public abstract class ColumnData implements IMeasurableMemory
public abstract ColumnData markCounterLocalToBeCleared();
public abstract ColumnData purge(DeletionPurger purger, int nowInSec);
+ public abstract ColumnData purgeDataOlderThan(long timestamp);
public abstract long maxTimestamp();
}
diff --git a/src/java/org/apache/cassandra/db/rows/ComplexColumnData.java b/src/java/org/apache/cassandra/db/rows/ComplexColumnData.java
index bf7714d9df..84f52e10bd 100644
--- a/src/java/org/apache/cassandra/db/rows/ComplexColumnData.java
+++ b/src/java/org/apache/cassandra/db/rows/ComplexColumnData.java
@@ -121,13 +121,10 @@ public class ComplexColumnData extends ColumnData implements Iterable>
return size;
}
- @Override
public long unsharedHeapSize()
{
- long heapSize = EMPTY_SIZE + ObjectSizes.sizeOfArray(cells);
- for (Cell> cell : this)
- heapSize += cell.unsharedHeapSize();
- return heapSize;
+ long heapSize = EMPTY_SIZE + ObjectSizes.sizeOfArray(cells) + complexDeletion.unsharedHeapSize();
+ return BTree.accumulate(cells, (cell, value) -> value + cell.unsharedHeapSize(), heapSize);
}
@Override
@@ -205,6 +202,12 @@ public class ComplexColumnData extends ColumnData implements Iterable>
return transformAndFilter(complexDeletion, (cell) -> filter.fetchedCellIsQueried(column, cell.path()) ? null : cell);
}
+ public ComplexColumnData purgeDataOlderThan(long timestamp)
+ {
+ DeletionTime newDeletion = complexDeletion.markedForDeleteAt() < timestamp ? DeletionTime.LIVE : complexDeletion;
+ return transformAndFilter(newDeletion, (cell) -> cell.purgeDataOlderThan(timestamp));
+ }
+
private ComplexColumnData transformAndFilter(DeletionTime newDeletion, Function super Cell>, ? extends Cell>> function)
{
Object[] transformed = BTree.transformAndFilter(cells, function);
diff --git a/src/java/org/apache/cassandra/db/rows/NativeCell.java b/src/java/org/apache/cassandra/db/rows/NativeCell.java
index 03dfc7090f..1267560aa2 100644
--- a/src/java/org/apache/cassandra/db/rows/NativeCell.java
+++ b/src/java/org/apache/cassandra/db/rows/NativeCell.java
@@ -177,5 +177,4 @@ public class NativeCell extends AbstractCell
{
return EMPTY_SIZE;
}
-
}
diff --git a/src/java/org/apache/cassandra/db/rows/Row.java b/src/java/org/apache/cassandra/db/rows/Row.java
index 85d27a44eb..7575f06552 100644
--- a/src/java/org/apache/cassandra/db/rows/Row.java
+++ b/src/java/org/apache/cassandra/db/rows/Row.java
@@ -116,7 +116,7 @@ public interface Row extends Unfiltered, Iterable, IMeasurableMemory
*
* @param nowInSec the current time to decide what is deleted and what isn't
* @param enforceStrictLiveness whether the row should be purged if there is no PK liveness info,
- * normally retrieved from {@link CFMetaData#enforceStrictLiveness()}
+ * normally retrieved from {@link TableMetadata#enforceStrictLiveness()}
* @return true if there is some live information
*/
public boolean hasLiveData(int nowInSec, boolean enforceStrictLiveness);
@@ -249,6 +249,11 @@ public interface Row extends Unfiltered, Iterable, IMeasurableMemory
*/
public Row withOnlyQueriedData(ColumnFilter filter);
+ /*
+ * Returns a copy of this row without any data with a timestamp older than the one provided
+ */
+ public Row purgeDataOlderThan(long timestamp, boolean enforceStrictLiveness);
+
/**
* Returns a copy of this row where all counter cells have they "local" shard marked for clearing.
*/
@@ -281,6 +286,7 @@ public interface Row extends Unfiltered, Iterable, IMeasurableMemory
public long unsharedHeapSizeExcludingData();
public String toString(TableMetadata metadata, boolean fullDetails);
+ public long unsharedHeapSize();
/**
* Apply a function to every column in a row
diff --git a/src/java/org/apache/cassandra/db/rows/UnfilteredRowIteratorSerializer.java b/src/java/org/apache/cassandra/db/rows/UnfilteredRowIteratorSerializer.java
index 938a3eed11..11541ee358 100644
--- a/src/java/org/apache/cassandra/db/rows/UnfilteredRowIteratorSerializer.java
+++ b/src/java/org/apache/cassandra/db/rows/UnfilteredRowIteratorSerializer.java
@@ -66,7 +66,7 @@ public class UnfilteredRowIteratorSerializer
{
protected static final Logger logger = LoggerFactory.getLogger(UnfilteredRowIteratorSerializer.class);
- private static final int IS_EMPTY = 0x01;
+ public static final int IS_EMPTY = 0x01;
private static final int IS_REVERSED = 0x02;
private static final int HAS_PARTITION_DELETION = 0x04;
private static final int HAS_STATIC_ROW = 0x08;
diff --git a/src/java/org/apache/cassandra/db/view/ViewUtils.java b/src/java/org/apache/cassandra/db/view/ViewUtils.java
index c248ddcde7..55a462c31b 100644
--- a/src/java/org/apache/cassandra/db/view/ViewUtils.java
+++ b/src/java/org/apache/cassandra/db/view/ViewUtils.java
@@ -23,7 +23,6 @@ import java.util.function.Predicate;
import com.google.common.collect.Iterables;
import org.apache.cassandra.config.DatabaseDescriptor;
-import org.apache.cassandra.db.Keyspace;
import org.apache.cassandra.dht.Token;
import org.apache.cassandra.locator.AbstractReplicationStrategy;
import org.apache.cassandra.locator.EndpointsForToken;
diff --git a/src/java/org/apache/cassandra/db/virtual/StreamingVirtualTable.java b/src/java/org/apache/cassandra/db/virtual/StreamingVirtualTable.java
index 1036f18426..f9fe0ef54d 100644
--- a/src/java/org/apache/cassandra/db/virtual/StreamingVirtualTable.java
+++ b/src/java/org/apache/cassandra/db/virtual/StreamingVirtualTable.java
@@ -23,10 +23,12 @@ import java.util.UUID;
import java.util.stream.Collectors;
import org.apache.cassandra.db.DecoratedKey;
+import org.apache.cassandra.db.marshal.TimeUUIDType;
import org.apache.cassandra.db.marshal.UUIDType;
import org.apache.cassandra.schema.TableMetadata;
import org.apache.cassandra.streaming.StreamManager;
import org.apache.cassandra.streaming.StreamingState;
+import org.apache.cassandra.utils.TimeUUID;
import static org.apache.cassandra.cql3.statements.schema.CreateTableStatement.parse;
@@ -35,7 +37,7 @@ public class StreamingVirtualTable extends AbstractVirtualTable
public StreamingVirtualTable(String keyspace)
{
super(parse("CREATE TABLE streaming (" +
- " id uuid,\n" +
+ " id timeuuid,\n" +
" follower boolean,\n" +
" operation text, \n" +
" peers frozen>,\n" +
@@ -76,7 +78,7 @@ public class StreamingVirtualTable extends AbstractVirtualTable
@Override
public DataSet data(DecoratedKey partitionKey)
{
- UUID id = UUIDType.instance.compose(partitionKey.getKey());
+ TimeUUID id = TimeUUIDType.instance.compose(partitionKey.getKey());
SimpleDataSet result = new SimpleDataSet(metadata());
StreamingState state = StreamManager.instance.getStreamingState(id);
if (state != null)
diff --git a/src/java/org/apache/cassandra/db/virtual/VirtualMutation.java b/src/java/org/apache/cassandra/db/virtual/VirtualMutation.java
index 09ac4a604b..3e26032383 100644
--- a/src/java/org/apache/cassandra/db/virtual/VirtualMutation.java
+++ b/src/java/org/apache/cassandra/db/virtual/VirtualMutation.java
@@ -19,6 +19,7 @@ package org.apache.cassandra.db.virtual;
import java.util.Collection;
import java.util.concurrent.TimeUnit;
+import java.util.function.Supplier;
import com.google.common.base.MoreObjects;
import com.google.common.collect.ImmutableMap;
@@ -26,6 +27,7 @@ import com.google.common.collect.ImmutableMap;
import org.apache.cassandra.config.DatabaseDescriptor;
import org.apache.cassandra.db.DecoratedKey;
import org.apache.cassandra.db.IMutation;
+import org.apache.cassandra.db.Mutation;
import org.apache.cassandra.db.partitions.PartitionUpdate;
import org.apache.cassandra.schema.TableId;
@@ -104,6 +106,12 @@ public final class VirtualMutation implements IMutation
return modifications.values();
}
+ @Override
+ public Supplier hintOnFailure()
+ {
+ return null;
+ }
+
@Override
public void validateIndexedColumns()
{
diff --git a/src/java/org/apache/cassandra/dht/Murmur3Partitioner.java b/src/java/org/apache/cassandra/dht/Murmur3Partitioner.java
index 2856f131f1..b0f8bf90e3 100644
--- a/src/java/org/apache/cassandra/dht/Murmur3Partitioner.java
+++ b/src/java/org/apache/cassandra/dht/Murmur3Partitioner.java
@@ -144,7 +144,7 @@ public class Murmur3Partitioner implements IPartitioner
{
static final long serialVersionUID = -5833580143318243006L;
- final long token;
+ public final long token;
public LongToken(long token)
{
@@ -204,11 +204,16 @@ public class Murmur3Partitioner implements IPartitioner
}
@Override
- public Token increaseSlightly()
+ public LongToken increaseSlightly()
{
return new LongToken(token + 1);
}
+ public LongToken decreaseSlightly()
+ {
+ return new LongToken(token - 1);
+ }
+
/**
* Reverses murmur3 to find a possible 16 byte key that generates a given token
*/
@@ -216,7 +221,7 @@ public class Murmur3Partitioner implements IPartitioner
public static ByteBuffer keyForToken(LongToken token)
{
ByteBuffer result = ByteBuffer.allocate(16);
- long[] inv = MurmurHash.inv_hash3_x64_128(new long[] {token.token, 0L});
+ long[] inv = MurmurHash.inv_hash3_x64_128(new long[]{ token.token, 0L });
result.putLong(inv[0]).putLong(inv[1]).position(0);
return result;
}
diff --git a/src/java/org/apache/cassandra/dht/Range.java b/src/java/org/apache/cassandra/dht/Range.java
index 5b2f3d9fbf..2c468990d2 100644
--- a/src/java/org/apache/cassandra/dht/Range.java
+++ b/src/java/org/apache/cassandra/dht/Range.java
@@ -483,7 +483,7 @@ public class Range> extends AbstractBounds implemen
* Given a list of unwrapped ranges sorted by left position, return an
* equivalent list of ranges but with no overlapping ranges.
*/
- private static > List> deoverlap(List> ranges)
+ public static > List> deoverlap(List> ranges)
{
if (ranges.isEmpty())
return ranges;
diff --git a/src/java/org/apache/cassandra/exceptions/CasWriteTimeoutException.java b/src/java/org/apache/cassandra/exceptions/CasWriteTimeoutException.java
index b1347645f1..32cc014da1 100644
--- a/src/java/org/apache/cassandra/exceptions/CasWriteTimeoutException.java
+++ b/src/java/org/apache/cassandra/exceptions/CasWriteTimeoutException.java
@@ -20,13 +20,14 @@ package org.apache.cassandra.exceptions;
import org.apache.cassandra.db.ConsistencyLevel;
import org.apache.cassandra.db.WriteType;
+
public class CasWriteTimeoutException extends WriteTimeoutException
{
public final int contentions;
public CasWriteTimeoutException(WriteType writeType, ConsistencyLevel consistency, int received, int blockFor, int contentions)
{
- super(writeType, consistency, received, blockFor, String.format("CAS operation timed out - encountered contentions: %d", contentions));
+ super(writeType, consistency, received, blockFor, String.format("CAS operation timed out: received %d of %d required responses after %d contention retries", received, blockFor, contentions));
this.contentions = contentions;
}
}
diff --git a/src/java/org/apache/cassandra/exceptions/RequestFailureException.java b/src/java/org/apache/cassandra/exceptions/RequestFailureException.java
index c75ff9c19d..a71534720b 100644
--- a/src/java/org/apache/cassandra/exceptions/RequestFailureException.java
+++ b/src/java/org/apache/cassandra/exceptions/RequestFailureException.java
@@ -35,7 +35,7 @@ public class RequestFailureException extends RequestExecutionException
this(code, buildErrorMessage(received, failureReasonByEndpoint), consistency, received, blockFor, failureReasonByEndpoint);
}
- protected RequestFailureException(ExceptionCode code, String msg, ConsistencyLevel consistency, int received, int blockFor, Map failureReasonByEndpoint)
+ public RequestFailureException(ExceptionCode code, String msg, ConsistencyLevel consistency, int received, int blockFor, Map failureReasonByEndpoint)
{
super(code, buildErrorMessage(msg, failureReasonByEndpoint));
this.consistency = consistency;
diff --git a/src/java/org/apache/cassandra/exceptions/RequestFailureReason.java b/src/java/org/apache/cassandra/exceptions/RequestFailureReason.java
index f205900bf8..3d3476a139 100644
--- a/src/java/org/apache/cassandra/exceptions/RequestFailureReason.java
+++ b/src/java/org/apache/cassandra/exceptions/RequestFailureReason.java
@@ -36,8 +36,9 @@ public enum RequestFailureReason
READ_TOO_MANY_TOMBSTONES (1),
TIMEOUT (2),
INCOMPATIBLE_SCHEMA (3),
- READ_SIZE (4);
-
+ READ_SIZE (4),
+ NODE_DOWN (5);
+
public static final Serializer serializer = new Serializer();
public final int code;
diff --git a/src/java/org/apache/cassandra/gms/EndpointState.java b/src/java/org/apache/cassandra/gms/EndpointState.java
index 2cc9c0debc..c60a4793bd 100644
--- a/src/java/org/apache/cassandra/gms/EndpointState.java
+++ b/src/java/org/apache/cassandra/gms/EndpointState.java
@@ -34,6 +34,7 @@ import org.apache.cassandra.io.IVersionedSerializer;
import org.apache.cassandra.io.util.DataInputPlus;
import org.apache.cassandra.io.util.DataOutputPlus;
import org.apache.cassandra.utils.CassandraVersion;
+import org.apache.cassandra.utils.NullableSerializer;
import static org.apache.cassandra.utils.Clock.Global.nanoTime;
@@ -41,13 +42,12 @@ import static org.apache.cassandra.utils.Clock.Global.nanoTime;
* This abstraction represents both the HeartBeatState and the ApplicationState in an EndpointState
* instance. Any state for a given endpoint can be retrieved from this instance.
*/
-
-
public class EndpointState
{
protected static final Logger logger = LoggerFactory.getLogger(EndpointState.class);
public final static IVersionedSerializer serializer = new EndpointStateSerializer();
+ public final static IVersionedSerializer nullableSerializer = NullableSerializer.wrap(serializer);
private volatile HeartBeatState hbState;
private final AtomicReference | | |