Fully reconciled SSTable promotion

patch by Abe Ratnofsky; reviewed by Ariel Weisberg, Blake Eggleston for CASSANDRA-20381
This commit is contained in:
Abe Ratnofsky 2025-05-28 11:17:34 -04:00 committed by Blake Eggleston
parent 51e63f0ff3
commit 65caf3f58b
65 changed files with 1177 additions and 653 deletions

View File

@ -136,6 +136,7 @@ import org.apache.cassandra.metrics.TopPartitionTracker;
import org.apache.cassandra.repair.TableRepairManager;
import org.apache.cassandra.repair.consistent.admin.CleanupSummary;
import org.apache.cassandra.repair.consistent.admin.PendingStat;
import org.apache.cassandra.replication.ImmutableCoordinatorLogOffsets;
import org.apache.cassandra.replication.MutationId;
import org.apache.cassandra.schema.ColumnMetadata;
import org.apache.cassandra.schema.CompactionParams;
@ -675,19 +676,19 @@ public class ColumnFamilyStore implements ColumnFamilyStoreMBean, Memtable.Owner
return memtableFactory.streamFromMemtable();
}
public SSTableMultiWriter createSSTableMultiWriter(Descriptor descriptor, long keyCount, long repairedAt, TimeUUID pendingRepair, boolean isTransient, CoordinatorLogBoundaries coordinatorLogBoundaries, SerializationHeader header, ILifecycleTransaction txn)
public SSTableMultiWriter createSSTableMultiWriter(Descriptor descriptor, long keyCount, long repairedAt, TimeUUID pendingRepair, boolean isTransient, ImmutableCoordinatorLogOffsets coordinatorLogOffsets, SerializationHeader header, ILifecycleTransaction txn)
{
return createSSTableMultiWriter(descriptor, keyCount, repairedAt, pendingRepair, isTransient, coordinatorLogBoundaries, null, 0, header, txn);
return createSSTableMultiWriter(descriptor, keyCount, repairedAt, pendingRepair, isTransient, coordinatorLogOffsets, null, 0, header, txn);
}
public SSTableMultiWriter createSSTableMultiWriter(Descriptor descriptor, long keyCount, long repairedAt, TimeUUID pendingRepair, boolean isTransient, CoordinatorLogBoundaries coordinatorLogBoundaries, IntervalSet<CommitLogPosition> commitLogPositions, SerializationHeader header, ILifecycleTransaction txn)
public SSTableMultiWriter createSSTableMultiWriter(Descriptor descriptor, long keyCount, long repairedAt, TimeUUID pendingRepair, boolean isTransient, ImmutableCoordinatorLogOffsets coordinatorLogOffsets, IntervalSet<CommitLogPosition> commitLogPositions, SerializationHeader header, ILifecycleTransaction txn)
{
return createSSTableMultiWriter(descriptor, keyCount, repairedAt, pendingRepair, isTransient, coordinatorLogBoundaries, commitLogPositions, 0, header, txn);
return createSSTableMultiWriter(descriptor, keyCount, repairedAt, pendingRepair, isTransient, coordinatorLogOffsets, commitLogPositions, 0, header, txn);
}
public SSTableMultiWriter createSSTableMultiWriter(Descriptor descriptor, long keyCount, long repairedAt, TimeUUID pendingRepair, boolean isTransient, CoordinatorLogBoundaries coordinatorLogBoundaries, IntervalSet<CommitLogPosition> commitLogPositions, int sstableLevel, SerializationHeader header, ILifecycleTransaction txn)
public SSTableMultiWriter createSSTableMultiWriter(Descriptor descriptor, long keyCount, long repairedAt, TimeUUID pendingRepair, boolean isTransient, ImmutableCoordinatorLogOffsets coordinatorLogOffsets, IntervalSet<CommitLogPosition> commitLogPositions, int sstableLevel, SerializationHeader header, ILifecycleTransaction txn)
{
return getCompactionStrategyManager().createSSTableMultiWriter(descriptor, keyCount, repairedAt, pendingRepair, isTransient, coordinatorLogBoundaries, commitLogPositions, sstableLevel, header, indexManager.listIndexGroups(), txn);
return getCompactionStrategyManager().createSSTableMultiWriter(descriptor, keyCount, repairedAt, pendingRepair, isTransient, coordinatorLogOffsets, commitLogPositions, sstableLevel, header, indexManager.listIndexGroups(), txn);
}
public boolean supportsEarlyOpen()
@ -2420,14 +2421,14 @@ public class ColumnFamilyStore implements ColumnFamilyStoreMBean, Memtable.Owner
return null;
List<Memtable.FlushablePartitionSet<?>> dataSets = new ArrayList<>(ranges.size());
CoordinatorLogBoundariesBuilder boundaries = new CoordinatorLogBoundariesBuilder();
ImmutableCoordinatorLogOffsets.Builder logOffsetsBuilder = new ImmutableCoordinatorLogOffsets.Builder();
IntervalSet.Builder<CommitLogPosition> commitLogIntervals = new IntervalSet.Builder();
long keys = 0;
for (Range<PartitionPosition> range : ranges)
{
Memtable.FlushablePartitionSet<?> dataSet = current.getFlushSet(range.left, range.right);
dataSets.add(dataSet);
boundaries.addAll(dataSet.coordinatorLogBoundaries());
logOffsetsBuilder.addAll(dataSet.coordinatorLogOffsets());
commitLogIntervals.add(dataSet.commitLogLowerBound(), dataSet.commitLogUpperBound());
keys += dataSet.partitionCount();
}
@ -2441,7 +2442,7 @@ public class ColumnFamilyStore implements ColumnFamilyStoreMBean, Memtable.Owner
0,
repairSessionID,
false,
boundaries.build(),
logOffsetsBuilder.build(),
commitLogIntervals.build(),
new SerializationHeader(true,
firstDataSet.metadata(),

View File

@ -1,88 +0,0 @@
/*
* 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;
import java.io.IOException;
import org.apache.cassandra.io.IVersionedSerializer;
import org.apache.cassandra.io.sstable.metadata.StatsMetadata;
import org.apache.cassandra.io.util.DataInputPlus;
import org.apache.cassandra.io.util.DataOutputPlus;
import org.apache.cassandra.net.MessagingService;
import org.apache.cassandra.replication.CoordinatorLogId;
import org.apache.cassandra.replication.MutationId;
import org.apache.cassandra.utils.vint.VIntCoding;
/**
* Max mutation ID present in this SSTable for each coordinator log, to determine whether an SSTable is reconciled or
* not. Once max mutation IDs are reconciled, next compaction can safely mark this SSTabled as repaired. Note that peers
* may have reconciled all mutations included in an SSTable, but {@link StatsMetadata#repairedAt} is dependent on
* compaction timing, so "nodetool repair --validate" may report temporary disagreements on the repaired set.
* <p>
* A reference to this class should be treated as immutable. Do not cast to {@link CoordinatorLogBoundariesMap}.
* Iterable over {@link CoordinatorLogId}.
*/
public interface CoordinatorLogBoundaries extends Iterable<Long>
{
int maxOffset(long logId);
MutationId max(long logId);
int size();
IVersionedSerializer<CoordinatorLogBoundaries> serializer = new IVersionedSerializer<>()
{
@Override
public void serialize(CoordinatorLogBoundaries boundaries, DataOutputPlus out, int version) throws IOException
{
if (version < MessagingService.VERSION_52)
return;
out.writeUnsignedVInt32(boundaries.size());
for (long logId : boundaries)
MutationId.serializer.serialize(boundaries.max(logId), out, version);
}
@Override
public CoordinatorLogBoundaries deserialize(DataInputPlus in, int version) throws IOException
{
if (version < MessagingService.VERSION_52)
return CoordinatorLogBoundaries.NONE;
int size = in.readUnsignedVInt32();
CoordinatorLogBoundariesMap boundaries = new CoordinatorLogBoundariesMap(size);
for (int i = 0; i < size; i++)
{
MutationId mutationId = MutationId.serializer.deserialize(in, version);
boundaries.add(mutationId);
}
return boundaries;
}
@Override
public long serializedSize(CoordinatorLogBoundaries boundaries, int version)
{
if (version < MessagingService.VERSION_52)
return 0;
long size = 0;
size += VIntCoding.computeUnsignedVIntSize(boundaries.size());
for (long logId : boundaries)
size += MutationId.serializer.serializedSize(boundaries.max(logId), version);
return size;
}
};
CoordinatorLogBoundaries NONE = new CoordinatorLogBoundariesMap();
}

View File

@ -1,92 +0,0 @@
/*
* 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;
import java.util.Iterator;
import javax.annotation.concurrent.NotThreadSafe;
import com.google.common.collect.Iterators;
import org.agrona.collections.Long2ObjectHashMap;
import org.apache.cassandra.replication.MutationId;
import org.apache.cassandra.replication.ShortMutationId;
@NotThreadSafe
public class CoordinatorLogBoundariesBuilder
{
private static final MutationId NONE = MutationId.none();
private static final int NONE_OFFSET = NONE.offset();
private final Long2ObjectHashMap<MutationId> ids;
public CoordinatorLogBoundariesBuilder()
{
this.ids = new Long2ObjectHashMap<>();
}
public CoordinatorLogBoundariesBuilder add(MutationId mutationId)
{
if (mutationId.isNone())
return this;
long logId = mutationId.logId();
MutationId existing = ids.get(logId);
if (existing == null || ShortMutationId.comparator.compare(existing, mutationId) < 0)
ids.put(mutationId.logId(), mutationId);
return this;
}
public CoordinatorLogBoundariesBuilder addAll(CoordinatorLogBoundaries boundaries)
{
for (long logId : boundaries)
add(boundaries.max(logId));
return this;
}
public CoordinatorLogBoundaries build()
{
return new CoordinatorLogBoundaries()
{
@Override
public int maxOffset(long logId)
{
MutationId id = ids.get(logId);
return id == null ? NONE_OFFSET : id.offset();
}
@Override
public MutationId max(long logId)
{
return ids.getOrDefault(logId, CoordinatorLogBoundariesBuilder.NONE);
}
@Override
public int size()
{
return ids.size();
}
@Override
public Iterator<Long> iterator()
{
return Iterators.unmodifiableIterator(ids.keySet().iterator());
}
};
}
}

View File

@ -1,80 +0,0 @@
/*
* 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;
import java.util.Iterator;
import javax.annotation.concurrent.ThreadSafe;
import com.google.common.collect.Iterators;
import org.apache.cassandra.replication.MutationId;
import org.apache.cassandra.replication.ShortMutationId;
import org.jctools.maps.NonBlockingHashMapLong;
/**
* A replica can only receive writes from another replica it shares ranges with, and tracked writes are executed by
* coordinators, so this should contain up to (2*RF - 1) keys.
* Consider wrapping value in AtomicReference to avoid false sharing, see: https://trishagee.com/2011/07/22/dissecting_the_disruptor_why_its_so_fast_part_two__magic_cache_line_padding/
*/
@ThreadSafe
class CoordinatorLogBoundariesMap extends NonBlockingHashMapLong<MutationId> implements MutableCoordinatorLogBoundaries
{
private static final MutationId NONE = MutationId.none();
private static final int NONE_OFFSET = NONE.offset();
protected CoordinatorLogBoundariesMap(int size)
{
super(size);
}
protected CoordinatorLogBoundariesMap()
{
super();
}
public void add(MutationId mutationId)
{
long logId = mutationId.logId();
merge(logId, mutationId, (existing, updating) -> {
if (ShortMutationId.comparator.compare(existing, updating) < 0)
return updating;
return existing;
});
}
@Override
public int maxOffset(long logId)
{
MutationId id = get(logId);
return id == null ? NONE_OFFSET : id.offset();
}
@Override
public MutationId max(long logId)
{
return getOrDefault(logId, NONE);
}
@Override
public Iterator<Long> iterator()
{
return Iterators.unmodifiableIterator(keySet().iterator());
}
}

View File

@ -51,6 +51,7 @@ import org.apache.cassandra.io.sstable.SimpleSSTableMultiWriter;
import org.apache.cassandra.io.sstable.format.SSTableReader;
import org.apache.cassandra.io.sstable.metadata.MetadataCollector;
import org.apache.cassandra.io.sstable.metadata.StatsMetadata;
import org.apache.cassandra.replication.ImmutableCoordinatorLogOffsets;
import org.apache.cassandra.schema.CompactionParams;
import org.apache.cassandra.utils.TimeUUID;
@ -560,7 +561,7 @@ public abstract class AbstractCompactionStrategy
long repairedAt,
TimeUUID pendingRepair,
boolean isTransient,
CoordinatorLogBoundaries coordinatorLogBoundaries,
ImmutableCoordinatorLogOffsets coordinatorLogOffsets,
IntervalSet<CommitLogPosition> commitLogPositions,
int sstableLevel,
SerializationHeader header,
@ -572,7 +573,7 @@ public abstract class AbstractCompactionStrategy
repairedAt,
pendingRepair,
isTransient,
coordinatorLogBoundaries,
coordinatorLogOffsets,
cfs.metadata,
commitLogPositions,
sstableLevel,

View File

@ -28,7 +28,6 @@ import java.util.function.Supplier;
import com.google.common.base.Preconditions;
import org.apache.cassandra.db.ColumnFamilyStore;
import org.apache.cassandra.db.CoordinatorLogBoundaries;
import org.apache.cassandra.db.SerializationHeader;
import org.apache.cassandra.db.commitlog.CommitLogPosition;
import org.apache.cassandra.db.commitlog.IntervalSet;
@ -40,6 +39,7 @@ import org.apache.cassandra.io.sstable.Descriptor;
import org.apache.cassandra.io.sstable.ISSTableScanner;
import org.apache.cassandra.io.sstable.SSTableMultiWriter;
import org.apache.cassandra.io.sstable.format.SSTableReader;
import org.apache.cassandra.replication.ImmutableCoordinatorLogOffsets;
import org.apache.cassandra.schema.CompactionParams;
import org.apache.cassandra.utils.TimeUUID;
@ -197,7 +197,7 @@ public abstract class AbstractStrategyHolder
long repairedAt,
TimeUUID pendingRepair,
boolean isTransient,
CoordinatorLogBoundaries coordinatorLogBoundaries,
ImmutableCoordinatorLogOffsets coordinatorLogOffsets,
IntervalSet<CommitLogPosition> commitLogPositions,
int sstableLevel,
SerializationHeader header,

View File

@ -70,8 +70,6 @@ import org.apache.cassandra.concurrent.ExecutorFactory;
import org.apache.cassandra.concurrent.WrappedExecutorPlus;
import org.apache.cassandra.config.DatabaseDescriptor;
import org.apache.cassandra.db.ColumnFamilyStore;
import org.apache.cassandra.db.CoordinatorLogBoundaries;
import org.apache.cassandra.db.CoordinatorLogBoundariesBuilder;
import org.apache.cassandra.db.DecoratedKey;
import org.apache.cassandra.db.Directories;
import org.apache.cassandra.db.DiskBoundaries;
@ -114,6 +112,7 @@ import org.apache.cassandra.locator.Replica;
import org.apache.cassandra.metrics.CompactionMetrics;
import org.apache.cassandra.metrics.TableMetrics;
import org.apache.cassandra.repair.NoSuchRepairSessionException;
import org.apache.cassandra.replication.ImmutableCoordinatorLogOffsets;
import org.apache.cassandra.schema.CompactionParams.TombstoneOption;
import org.apache.cassandra.schema.Schema;
import org.apache.cassandra.schema.TableMetadata;
@ -1652,8 +1651,8 @@ public class CompactionManager implements CompactionManagerMBean, ICompactionMan
CompactionIterator ci = new CompactionIterator(OperationType.CLEANUP, Collections.singletonList(scanner), controller, nowInSec, nextTimeUUID(), active, null))
{
StatsMetadata metadata = sstable.getSSTableMetadata();
// TODO(aratnofsky): filter coordinatorLogBoundaries to exclude any CoordinatorLogIds we're no longer responsible for, after ownership change
writer.switchWriter(createWriter(cfs, compactionFileLocation, expectedBloomFilterSize, metadata.repairedAt, metadata.pendingRepair, metadata.isTransient, metadata.coordinatorLogBoundaries, sstable, txn));
// TODO(aratnofsky): filter coordinatorLogOffsets to exclude any CoordinatorLogIds we're no longer responsible for, after ownership change
writer.switchWriter(createWriter(cfs, compactionFileLocation, expectedBloomFilterSize, metadata.repairedAt, metadata.pendingRepair, metadata.isTransient, metadata.coordinatorLogOffsets, sstable, txn));
long lastBytesScanned = 0;
while (ci.hasNext())
@ -1818,7 +1817,7 @@ public class CompactionManager implements CompactionManagerMBean, ICompactionMan
long repairedAt,
TimeUUID pendingRepair,
boolean isTransient,
CoordinatorLogBoundaries coordinatorLogBoundaries,
ImmutableCoordinatorLogOffsets coordinatorLogOffsets,
SSTableReader sstable,
LifecycleTransaction txn)
{
@ -1830,7 +1829,7 @@ public class CompactionManager implements CompactionManagerMBean, ICompactionMan
.setRepairedAt(repairedAt)
.setPendingRepair(pendingRepair)
.setTransientSSTable(isTransient)
.setCoordinatorLogBoundaries(coordinatorLogBoundaries)
.setCoordinatorLogOffsets(coordinatorLogOffsets)
.setTableMetadataRef(cfs.metadata)
.setMetadataCollector(new MetadataCollector(cfs.metadata().comparator).sstableLevel(sstable.getSSTableLevel()))
.setSerializationHeader(sstable.header)
@ -1851,10 +1850,10 @@ public class CompactionManager implements CompactionManagerMBean, ICompactionMan
{
FileUtils.createDirectory(compactionFileLocation);
int minLevel = Integer.MAX_VALUE;
CoordinatorLogBoundariesBuilder boundaries = new CoordinatorLogBoundariesBuilder();
ImmutableCoordinatorLogOffsets.Builder logOffsetsBuilder = new ImmutableCoordinatorLogOffsets.Builder();
for (SSTableReader sstable : sstables)
{
boundaries.addAll(sstable.getCoordinatorLogBoundaries());
logOffsetsBuilder.addAll(sstable.getCoordinatorLogOffsets());
// if all sstables have the same level, we can compact them together without creating overlap during anticompaction
// note that we only anticompact from unrepaired sstables, which is not leveled, but we still keep original level
@ -1873,7 +1872,7 @@ public class CompactionManager implements CompactionManagerMBean, ICompactionMan
.setKeyCount(expectedBloomFilterSize)
.setRepairedAt(repairedAt)
.setPendingRepair(pendingRepair)
.setCoordinatorLogBoundaries(boundaries.build())
.setCoordinatorLogOffsets(logOffsetsBuilder.build())
.setTransientSSTable(isTransient)
.setTableMetadataRef(cfs.metadata)
.setMetadataCollector(new MetadataCollector(sstables, cfs.metadata().comparator).sstableLevel(minLevel))

View File

@ -26,7 +26,6 @@ import com.google.common.base.Preconditions;
import com.google.common.collect.Iterables;
import org.apache.cassandra.db.ColumnFamilyStore;
import org.apache.cassandra.db.CoordinatorLogBoundaries;
import org.apache.cassandra.db.SerializationHeader;
import org.apache.cassandra.db.commitlog.CommitLogPosition;
import org.apache.cassandra.db.commitlog.IntervalSet;
@ -38,6 +37,7 @@ import org.apache.cassandra.io.sstable.Descriptor;
import org.apache.cassandra.io.sstable.ISSTableScanner;
import org.apache.cassandra.io.sstable.SSTableMultiWriter;
import org.apache.cassandra.io.sstable.format.SSTableReader;
import org.apache.cassandra.replication.ImmutableCoordinatorLogOffsets;
import org.apache.cassandra.schema.CompactionParams;
import org.apache.cassandra.service.ActiveRepairService;
import org.apache.cassandra.utils.TimeUUID;
@ -227,7 +227,7 @@ public class CompactionStrategyHolder extends AbstractStrategyHolder
long repairedAt,
TimeUUID pendingRepair,
boolean isTransient,
CoordinatorLogBoundaries coordinatorLogBoundaries,
ImmutableCoordinatorLogOffsets coordinatorLogOffsets,
IntervalSet<CommitLogPosition> commitLogPositions,
int sstableLevel,
SerializationHeader header,
@ -253,7 +253,7 @@ public class CompactionStrategyHolder extends AbstractStrategyHolder
repairedAt,
pendingRepair,
isTransient,
coordinatorLogBoundaries,
coordinatorLogOffsets,
commitLogPositions,
sstableLevel,
header,

View File

@ -50,7 +50,6 @@ import org.slf4j.LoggerFactory;
import org.apache.cassandra.config.DatabaseDescriptor;
import org.apache.cassandra.db.ColumnFamilyStore;
import org.apache.cassandra.db.CoordinatorLogBoundaries;
import org.apache.cassandra.db.Directories;
import org.apache.cassandra.db.DiskBoundaries;
import org.apache.cassandra.db.SerializationHeader;
@ -82,6 +81,7 @@ import org.apache.cassandra.notifications.SSTableListChangedNotification;
import org.apache.cassandra.notifications.SSTableMetadataChanged;
import org.apache.cassandra.notifications.SSTableRepairStatusChanged;
import org.apache.cassandra.repair.consistent.admin.CleanupSummary;
import org.apache.cassandra.replication.ImmutableCoordinatorLogOffsets;
import org.apache.cassandra.schema.CompactionParams;
import org.apache.cassandra.service.ActiveRepairService;
import org.apache.cassandra.utils.TimeUUID;
@ -1412,7 +1412,7 @@ public class CompactionStrategyManager implements INotificationConsumer
long repairedAt,
TimeUUID pendingRepair,
boolean isTransient,
CoordinatorLogBoundaries coordinatorLogBoundaries,
ImmutableCoordinatorLogOffsets coordinatorLogOffsets,
IntervalSet<CommitLogPosition> commitLogPositions,
int sstableLevel,
SerializationHeader header,
@ -1429,7 +1429,7 @@ public class CompactionStrategyManager implements INotificationConsumer
repairedAt,
pendingRepair,
isTransient,
coordinatorLogBoundaries,
coordinatorLogOffsets,
commitLogPositions,
sstableLevel,
header,

View File

@ -37,8 +37,6 @@ import org.slf4j.LoggerFactory;
import org.apache.cassandra.config.DatabaseDescriptor;
import org.apache.cassandra.db.ColumnFamilyStore;
import org.apache.cassandra.db.CoordinatorLogBoundaries;
import org.apache.cassandra.db.CoordinatorLogBoundariesBuilder;
import org.apache.cassandra.db.Directories;
import org.apache.cassandra.db.SystemKeyspace;
import org.apache.cassandra.db.WriteContext;
@ -57,6 +55,7 @@ import org.apache.cassandra.io.sstable.ISSTableScanner;
import org.apache.cassandra.io.sstable.format.SSTableReader;
import org.apache.cassandra.io.sstable.metadata.MetadataCollector;
import org.apache.cassandra.io.util.File;
import org.apache.cassandra.replication.ImmutableCoordinatorLogOffsets;
import org.apache.cassandra.service.ActiveRepairService;
import org.apache.cassandra.service.snapshot.SnapshotManager;
import org.apache.cassandra.service.snapshot.SnapshotOptions;
@ -460,11 +459,11 @@ public class CompactionTask extends AbstractCompactionTask
return isTransient;
}
public static CoordinatorLogBoundaries getCoordinatorLogBoundaries(Set<SSTableReader> sstables)
public static ImmutableCoordinatorLogOffsets getCoordinatorLogOffsets(Set<SSTableReader> sstables)
{
CoordinatorLogBoundariesBuilder builder = new CoordinatorLogBoundariesBuilder();
ImmutableCoordinatorLogOffsets.Builder builder = new ImmutableCoordinatorLogOffsets.Builder();
for (SSTableReader sstable : sstables)
builder.addAll(sstable.getCoordinatorLogBoundaries());
builder.addAll(sstable.getCoordinatorLogOffsets());
return builder.build();
}

View File

@ -27,7 +27,6 @@ import com.google.common.base.Preconditions;
import com.google.common.collect.Iterables;
import org.apache.cassandra.db.ColumnFamilyStore;
import org.apache.cassandra.db.CoordinatorLogBoundaries;
import org.apache.cassandra.db.SerializationHeader;
import org.apache.cassandra.db.commitlog.CommitLogPosition;
import org.apache.cassandra.db.commitlog.IntervalSet;
@ -39,6 +38,7 @@ import org.apache.cassandra.io.sstable.Descriptor;
import org.apache.cassandra.io.sstable.ISSTableScanner;
import org.apache.cassandra.io.sstable.SSTableMultiWriter;
import org.apache.cassandra.io.sstable.format.SSTableReader;
import org.apache.cassandra.replication.ImmutableCoordinatorLogOffsets;
import org.apache.cassandra.schema.CompactionParams;
import org.apache.cassandra.service.ActiveRepairService;
import org.apache.cassandra.utils.TimeUUID;
@ -247,7 +247,7 @@ public class PendingRepairHolder extends AbstractStrategyHolder
long repairedAt,
TimeUUID pendingRepair,
boolean isTransient,
CoordinatorLogBoundaries coordinatorLogBoundaries,
ImmutableCoordinatorLogOffsets coordinatorLogOffsets,
IntervalSet<CommitLogPosition> commitLogPositions,
int sstableLevel,
SerializationHeader header,
@ -265,7 +265,7 @@ public class PendingRepairHolder extends AbstractStrategyHolder
repairedAt,
pendingRepair,
isTransient,
coordinatorLogBoundaries,
coordinatorLogOffsets,
commitLogPositions,
sstableLevel,
header,

View File

@ -55,6 +55,7 @@ import org.apache.cassandra.index.Index;
import org.apache.cassandra.io.sstable.Descriptor;
import org.apache.cassandra.io.sstable.SSTableMultiWriter;
import org.apache.cassandra.io.sstable.format.SSTableReader;
import org.apache.cassandra.replication.ImmutableCoordinatorLogOffsets;
import org.apache.cassandra.schema.TableMetadata;
import org.apache.cassandra.tcm.ClusterMetadata;
import org.apache.cassandra.utils.Clock;
@ -300,7 +301,7 @@ public class UnifiedCompactionStrategy extends AbstractCompactionStrategy
long repairedAt,
TimeUUID pendingRepair,
boolean isTransient,
CoordinatorLogBoundaries coordinatorLogBoundaries,
ImmutableCoordinatorLogOffsets coordinatorLogOffsets,
IntervalSet<CommitLogPosition> commitLogPositions,
int sstableLevel,
SerializationHeader header,
@ -318,7 +319,7 @@ public class UnifiedCompactionStrategy extends AbstractCompactionStrategy
repairedAt,
pendingRepair,
isTransient,
coordinatorLogBoundaries,
coordinatorLogOffsets,
commitLogPositions,
header,
indexGroups,

View File

@ -80,7 +80,7 @@ public class Upgrader
.setRepairedAt(metadata.repairedAt)
.setPendingRepair(metadata.pendingRepair)
.setTransientSSTable(metadata.isTransient)
.setCoordinatorLogBoundaries(metadata.coordinatorLogBoundaries)
.setCoordinatorLogOffsets(metadata.coordinatorLogOffsets)
.setTableMetadataRef(cfs.metadata)
.setMetadataCollector(sstableMetadataCollector)
.setSerializationHeader(SerializationHeader.make(cfs.metadata(), Sets.newHashSet(sstable)))

View File

@ -26,7 +26,6 @@ import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.apache.cassandra.db.ColumnFamilyStore;
import org.apache.cassandra.db.CoordinatorLogBoundaries;
import org.apache.cassandra.db.DecoratedKey;
import org.apache.cassandra.db.SerializationHeader;
import org.apache.cassandra.db.commitlog.CommitLogPosition;
@ -40,6 +39,7 @@ import org.apache.cassandra.io.sstable.SSTableMultiWriter;
import org.apache.cassandra.io.sstable.format.SSTableReader;
import org.apache.cassandra.io.sstable.format.SSTableWriter;
import org.apache.cassandra.io.sstable.metadata.MetadataCollector;
import org.apache.cassandra.replication.ImmutableCoordinatorLogOffsets;
import org.apache.cassandra.schema.TableId;
import org.apache.cassandra.utils.FBUtilities;
import org.apache.cassandra.utils.TimeUUID;
@ -62,7 +62,7 @@ public class ShardedMultiWriter implements SSTableMultiWriter
private final long repairedAt;
private final TimeUUID pendingRepair;
private final boolean isTransient;
private final CoordinatorLogBoundaries coordinatorLogBoundaries;
private final ImmutableCoordinatorLogOffsets coordinatorLogOffsets;
private final IntervalSet<CommitLogPosition> commitLogPositions;
private final SerializationHeader header;
private final Collection<Index.Group> indexGroups;
@ -77,7 +77,7 @@ public class ShardedMultiWriter implements SSTableMultiWriter
long repairedAt,
TimeUUID pendingRepair,
boolean isTransient,
CoordinatorLogBoundaries coordinatorLogBoundaries,
ImmutableCoordinatorLogOffsets coordinatorLogOffsets,
IntervalSet<CommitLogPosition> commitLogPositions,
SerializationHeader header,
Collection<Index.Group> indexGroups,
@ -90,7 +90,7 @@ public class ShardedMultiWriter implements SSTableMultiWriter
this.repairedAt = repairedAt;
this.pendingRepair = pendingRepair;
this.isTransient = isTransient;
this.coordinatorLogBoundaries = coordinatorLogBoundaries;
this.coordinatorLogOffsets = coordinatorLogOffsets;
this.commitLogPositions = commitLogPositions;
this.header = header;
this.indexGroups = indexGroups;
@ -116,7 +116,7 @@ public class ShardedMultiWriter implements SSTableMultiWriter
.setKeyCount(forSplittingKeysBy(boundaries.count()))
.setRepairedAt(repairedAt)
.setPendingRepair(pendingRepair)
.setCoordinatorLogBoundaries(coordinatorLogBoundaries)
.setCoordinatorLogOffsets(coordinatorLogOffsets)
.setTransientSSTable(isTransient)
.setTableMetadataRef(cfs.metadata)
.setMetadataCollector(metadataCollector)

View File

@ -27,7 +27,6 @@ import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.apache.cassandra.db.ColumnFamilyStore;
import org.apache.cassandra.db.CoordinatorLogBoundaries;
import org.apache.cassandra.db.DecoratedKey;
import org.apache.cassandra.db.Directories;
import org.apache.cassandra.db.DiskBoundaries;
@ -41,6 +40,7 @@ import org.apache.cassandra.io.sstable.SSTableRewriter;
import org.apache.cassandra.io.sstable.format.SSTableReader;
import org.apache.cassandra.io.sstable.format.SSTableWriter;
import org.apache.cassandra.io.sstable.metadata.MetadataCollector;
import org.apache.cassandra.replication.ImmutableCoordinatorLogOffsets;
import org.apache.cassandra.utils.FBUtilities;
import org.apache.cassandra.utils.TimeUUID;
import org.apache.cassandra.utils.concurrent.Transactional;
@ -62,7 +62,7 @@ public abstract class CompactionAwareWriter extends Transactional.AbstractTransa
protected final long minRepairedAt;
protected final TimeUUID pendingRepair;
protected final boolean isTransient;
protected final CoordinatorLogBoundaries coordinatorLogBoundaries;
protected final ImmutableCoordinatorLogOffsets coordinatorLogOffsets;
protected final SSTableRewriter sstableWriter;
protected final ILifecycleTransaction txn;
@ -98,7 +98,7 @@ public abstract class CompactionAwareWriter extends Transactional.AbstractTransa
minRepairedAt = CompactionTask.getMinRepairedAt(nonExpiredSSTables);
pendingRepair = CompactionTask.getPendingRepair(nonExpiredSSTables);
isTransient = CompactionTask.getIsTransient(nonExpiredSSTables);
coordinatorLogBoundaries = CompactionTask.getCoordinatorLogBoundaries(nonExpiredSSTables);
coordinatorLogOffsets = CompactionTask.getCoordinatorLogOffsets(nonExpiredSSTables);
DiskBoundaries db = cfs.getDiskBoundaries();
diskBoundaries = db.positions;
locations = db.directories;
@ -331,7 +331,7 @@ public abstract class CompactionAwareWriter extends Transactional.AbstractTransa
.setTransientSSTable(isTransient)
.setRepairedAt(minRepairedAt)
.setPendingRepair(pendingRepair)
.setCoordinatorLogBoundaries(coordinatorLogBoundaries)
.setCoordinatorLogOffsets(coordinatorLogOffsets)
.setSecondaryIndexGroups(cfs.indexManager.listIndexGroups())
.addDefaultComponents(cfs.indexManager.listIndexGroups())
.setCompressionDictionaryManager(cfs.compressionDictionaryManager());

View File

@ -244,7 +244,7 @@ public class Flushing
ActiveRepairService.UNREPAIRED_SSTABLE,
ActiveRepairService.NO_PENDING_REPAIR,
false,
flushSet.coordinatorLogBoundaries(),
flushSet.coordinatorLogOffsets(),
new IntervalSet<>(flushSet.commitLogLowerBound(),
flushSet.commitLogUpperBound()),
new SerializationHeader(true,

View File

@ -26,7 +26,6 @@ import javax.annotation.concurrent.NotThreadSafe;
import org.apache.cassandra.db.CellSourceIdentifier;
import org.apache.cassandra.db.ColumnFamilyStore;
import org.apache.cassandra.db.CoordinatorLogBoundaries;
import org.apache.cassandra.db.PartitionPosition;
import org.apache.cassandra.db.RegularAndStaticColumns;
import org.apache.cassandra.db.commitlog.CommitLogPosition;
@ -38,6 +37,7 @@ import org.apache.cassandra.db.rows.UnfilteredSource;
import org.apache.cassandra.dht.Token;
import org.apache.cassandra.index.transactions.UpdateTransaction;
import org.apache.cassandra.io.sstable.format.SSTableWriter;
import org.apache.cassandra.replication.ImmutableCoordinatorLogOffsets;
import org.apache.cassandra.replication.MutationId;
import org.apache.cassandra.schema.TableMetadata;
import org.apache.cassandra.schema.TableMetadataRef;
@ -336,8 +336,8 @@ public interface Memtable extends Comparable<Memtable>, UnfilteredSource, CellSo
/** Statistics required for writing an sstable efficiently */
EncodingStats encodingStats();
/** The boundaries in coordinator logs for all included tracked mutations */
CoordinatorLogBoundaries coordinatorLogBoundaries();
/** The offsets in coordinator logs for all included tracked mutations */
ImmutableCoordinatorLogOffsets coordinatorLogOffsets();
default TableMetadata metadata()
{

View File

@ -25,18 +25,16 @@ import java.util.concurrent.ConcurrentSkipListMap;
import java.util.concurrent.atomic.AtomicLong;
import java.util.concurrent.atomic.AtomicReference;
import com.google.common.annotations.VisibleForTesting;
import com.google.common.collect.Iterators;
import org.github.jamm.Unmetered;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.apache.cassandra.db.CoordinatorLogBoundaries;
import org.apache.cassandra.db.CoordinatorLogBoundariesBuilder;
import org.apache.cassandra.db.DataRange;
import org.apache.cassandra.db.DecoratedKey;
import org.apache.cassandra.db.MutableCoordinatorLogBoundaries;
import org.apache.cassandra.replication.ImmutableCoordinatorLogOffsets;
import org.apache.cassandra.replication.MutableCoordinatorLogOffsets;
import org.apache.cassandra.db.PartitionPosition;
import org.apache.cassandra.db.RegularAndStaticColumns;
import org.apache.cassandra.db.Slices;
@ -93,19 +91,21 @@ public class ShardedSkipListMemtable extends AbstractShardedMemtable
ShardedSkipListMemtable(AtomicReference<CommitLogPosition> commitLogLowerBound,
TableMetadataRef metadataRef,
Owner owner,
Integer shardCountOption)
Integer shardCountOption,
boolean locking)
{
super(commitLogLowerBound, metadataRef, owner, shardCountOption);
this.shards = generatePartitionShards(boundaries.shardCount(), allocator, metadataRef);
this.shards = generatePartitionShards(boundaries.shardCount(), allocator, metadataRef, locking);
}
private static MemtableShard[] generatePartitionShards(int splits,
MemtableAllocator allocator,
TableMetadataRef metadata)
TableMetadataRef metadata,
boolean locking)
{
MemtableShard[] partitionMapContainer = new MemtableShard[splits];
for (int i = 0; i < splits; i++)
partitionMapContainer[i] = new MemtableShard(metadata, allocator);
partitionMapContainer[i] = new MemtableShard(metadata, allocator, locking);
return partitionMapContainer;
}
@ -301,11 +301,11 @@ public class ShardedSkipListMemtable extends AbstractShardedMemtable
int partitionCount = keyCount;
Iterator<AtomicBTreePartition> toFlush = getPartitionIterator(from, true, to, false);
CoordinatorLogBoundaries flushableBoundaries;
ImmutableCoordinatorLogOffsets flushableBoundaries;
{
CoordinatorLogBoundariesBuilder builder = new CoordinatorLogBoundariesBuilder();
ImmutableCoordinatorLogOffsets.Builder builder = new ImmutableCoordinatorLogOffsets.Builder();
for (MemtableShard shard : shards)
builder.addAll(shard.coordinatorLogBoundaries);
builder.addAll(shard.coordinatorLogOffsets);
flushableBoundaries = builder.build();
}
@ -350,7 +350,7 @@ public class ShardedSkipListMemtable extends AbstractShardedMemtable
}
@Override
public CoordinatorLogBoundaries coordinatorLogBoundaries()
public ImmutableCoordinatorLogOffsets coordinatorLogOffsets()
{
return flushableBoundaries;
}
@ -379,20 +379,21 @@ public class ShardedSkipListMemtable extends AbstractShardedMemtable
private final ColumnsCollector columnsCollector;
private final StatsCollector statsCollector;
private final MutableCoordinatorLogBoundaries coordinatorLogBoundaries = MutableCoordinatorLogBoundaries.create();
private final MutableCoordinatorLogOffsets coordinatorLogOffsets;
@Unmetered // total pool size should not be included in memtable's deep size
private final MemtableAllocator allocator;
private final TableMetadataRef metadata;
@VisibleForTesting
MemtableShard(TableMetadataRef metadata, MemtableAllocator allocator)
private MemtableShard(TableMetadataRef metadata, MemtableAllocator allocator, boolean locking)
{
this.columnsCollector = new ColumnsCollector(metadata.get().regularAndStaticColumns());
this.statsCollector = new StatsCollector();
this.allocator = allocator;
this.metadata = metadata;
this.coordinatorLogOffsets = MutableCoordinatorLogOffsets.create(locking);
}
public long put(MutationId mutationId, DecoratedKey key, PartitionUpdate update, UpdateTransaction indexer, OpOrder.Group opGroup, boolean assumeMissing)
@ -424,7 +425,7 @@ public class ShardedSkipListMemtable extends AbstractShardedMemtable
liveDataSize.addAndGet(initialSize + updater.dataSize);
columnsCollector.update(update.columns());
statsCollector.update(update.stats());
coordinatorLogBoundaries.add(mutationId);
coordinatorLogOffsets.add(mutationId);
currentOperations.addAndGet(update.operationCount());
return updater.colUpdateTimeDelta;
}
@ -525,7 +526,7 @@ public class ShardedSkipListMemtable extends AbstractShardedMemtable
{
Locking(AtomicReference<CommitLogPosition> commitLogLowerBound, TableMetadataRef metadataRef, Owner owner, Integer shardCountOption)
{
super(commitLogLowerBound, metadataRef, owner, shardCountOption);
super(commitLogLowerBound, metadataRef, owner, shardCountOption, true);
}
/**
@ -571,7 +572,7 @@ public class ShardedSkipListMemtable extends AbstractShardedMemtable
{
return isLocking
? new Locking(commitLogLowerBound, metadataRef, owner, shardCount)
: new ShardedSkipListMemtable(commitLogLowerBound, metadataRef, owner, shardCount);
: new ShardedSkipListMemtable(commitLogLowerBound, metadataRef, owner, shardCount, false);
}
public boolean equals(Object o)

View File

@ -32,11 +32,10 @@ import org.slf4j.LoggerFactory;
import org.apache.cassandra.config.DatabaseDescriptor;
import org.apache.cassandra.db.BufferDecoratedKey;
import org.apache.cassandra.db.CoordinatorLogBoundaries;
import org.apache.cassandra.db.CoordinatorLogBoundariesBuilder;
import org.apache.cassandra.db.DataRange;
import org.apache.cassandra.db.DecoratedKey;
import org.apache.cassandra.db.MutableCoordinatorLogBoundaries;
import org.apache.cassandra.replication.ImmutableCoordinatorLogOffsets;
import org.apache.cassandra.replication.MutableCoordinatorLogOffsets;
import org.apache.cassandra.db.PartitionPosition;
import org.apache.cassandra.db.Slices;
import org.apache.cassandra.db.commitlog.CommitLogPosition;
@ -90,7 +89,7 @@ public class SkipListMemtable extends AbstractAllocatorMemtable
// actually only store DecoratedKey.
private final ConcurrentNavigableMap<PartitionPosition, AtomicBTreePartition> partitions = new ConcurrentSkipListMap<>();
private final MutableCoordinatorLogBoundaries coordinatorLogBoundaries = MutableCoordinatorLogBoundaries.create();
private final MutableCoordinatorLogOffsets coordinatorLogOffsets = MutableCoordinatorLogOffsets.create(false);
private final AtomicLong liveDataSize = new AtomicLong(0);
@ -149,7 +148,7 @@ public class SkipListMemtable extends AbstractAllocatorMemtable
liveDataSize.addAndGet(initialSize + updater.dataSize);
columnsCollector.update(update.columns());
statsCollector.update(update.stats());
coordinatorLogBoundaries.add(mutationId);
coordinatorLogOffsets.add(mutationId);
currentOperations.addAndGet(update.operationCount());
return updater.colUpdateTimeDelta;
}
@ -295,8 +294,8 @@ public class SkipListMemtable extends AbstractAllocatorMemtable
}
final long partitionKeysSize = keysSize;
final long partitionCount = keyCount;
CoordinatorLogBoundaries flushableBoundaries = new CoordinatorLogBoundariesBuilder()
.addAll(coordinatorLogBoundaries)
ImmutableCoordinatorLogOffsets flushableBoundaries = new ImmutableCoordinatorLogOffsets.Builder()
.addAll(coordinatorLogOffsets)
.build();
return new AbstractFlushablePartitionSet<AtomicBTreePartition>()
@ -346,7 +345,7 @@ public class SkipListMemtable extends AbstractAllocatorMemtable
}
@Override
public CoordinatorLogBoundaries coordinatorLogBoundaries()
public ImmutableCoordinatorLogOffsets coordinatorLogOffsets()
{
return flushableBoundaries;
}

View File

@ -38,12 +38,11 @@ import org.apache.cassandra.config.DatabaseDescriptor;
import org.apache.cassandra.db.BufferDecoratedKey;
import org.apache.cassandra.db.Clustering;
import org.apache.cassandra.db.ColumnFamilyStore;
import org.apache.cassandra.db.CoordinatorLogBoundaries;
import org.apache.cassandra.db.CoordinatorLogBoundariesBuilder;
import org.apache.cassandra.db.DataRange;
import org.apache.cassandra.db.DecoratedKey;
import org.apache.cassandra.db.DeletionInfo;
import org.apache.cassandra.db.MutableCoordinatorLogBoundaries;
import org.apache.cassandra.replication.ImmutableCoordinatorLogOffsets;
import org.apache.cassandra.replication.MutableCoordinatorLogOffsets;
import org.apache.cassandra.db.PartitionPosition;
import org.apache.cassandra.db.RegularAndStaticColumns;
import org.apache.cassandra.db.Slices;
@ -451,11 +450,11 @@ public class TrieMemtable extends AbstractShardedMemtable
partitionKeySize = keySize;
partitionCount = keyCount;
CoordinatorLogBoundaries flushableBoundaries;
ImmutableCoordinatorLogOffsets flushableBoundaries;
{
CoordinatorLogBoundariesBuilder builder = new CoordinatorLogBoundariesBuilder();
ImmutableCoordinatorLogOffsets.Builder builder = new ImmutableCoordinatorLogOffsets.Builder();
for (MemtableShard shard : shards)
builder.addAll(shard.coordinatorLogBoundaries);
builder.addAll(shard.coordinatorLogOffsets);
flushableBoundaries = builder.build();
}
@ -503,7 +502,7 @@ public class TrieMemtable extends AbstractShardedMemtable
}
@Override
public CoordinatorLogBoundaries coordinatorLogBoundaries()
public ImmutableCoordinatorLogOffsets coordinatorLogOffsets()
{
return flushableBoundaries;
}
@ -548,7 +547,7 @@ public class TrieMemtable extends AbstractShardedMemtable
private final ColumnsCollector columnsCollector;
private final StatsCollector statsCollector;
private final MutableCoordinatorLogBoundaries coordinatorLogBoundaries = MutableCoordinatorLogBoundaries.create();
private final MutableCoordinatorLogOffsets coordinatorLogOffsets = MutableCoordinatorLogOffsets.create(true);
@Unmetered // total pool size should not be included in memtable's deep size
private final MemtableAllocator allocator;
@ -605,7 +604,7 @@ public class TrieMemtable extends AbstractShardedMemtable
columnsCollector.update(update.columns());
statsCollector.update(update.stats());
coordinatorLogBoundaries.add(mutationId);
coordinatorLogOffsets.add(mutationId);
}
}
finally

View File

@ -74,7 +74,7 @@ public class CassandraCompressedStreamReader extends CassandraStreamReader
try (CompressedInputStream cis = new CompressedInputStream(inputPlus, compressionInfo, ChecksumType.CRC32, cfs::getCrcCheckChance))
{
TrackedDataInputPlus in = new TrackedDataInputPlus(cis);
writer = createWriter(cfs, totalSize, repairedAt, pendingRepair, coordinatorLogBoundaries, inputVersion.format);
writer = createWriter(cfs, totalSize, repairedAt, pendingRepair, coordinatorLogOffsets, inputVersion.format);
deserializer = new StreamDeserializer(cfs.metadata(), in, inputVersion, getHeader(cfs.metadata()), session, writer);
String filename = writer.getFilename();
String sectionName = filename + '-' + fileSeqNum;

View File

@ -26,10 +26,10 @@ import com.google.common.annotations.VisibleForTesting;
import com.google.common.base.Preconditions;
import org.apache.cassandra.config.DatabaseDescriptor;
import org.apache.cassandra.db.CoordinatorLogBoundaries;
import org.apache.cassandra.dht.Range;
import org.apache.cassandra.dht.Token;
import org.apache.cassandra.io.sstable.format.SSTableReader;
import org.apache.cassandra.replication.ImmutableCoordinatorLogOffsets;
import org.apache.cassandra.schema.TableId;
import org.apache.cassandra.streaming.OutgoingStream;
import org.apache.cassandra.streaming.StreamOperation;
@ -153,9 +153,9 @@ public class CassandraOutgoingFile implements OutgoingStream
}
@Override
public CoordinatorLogBoundaries getCoordinatorLogBoundaries()
public ImmutableCoordinatorLogOffsets getCoordinatorLogOffsets()
{
return ref.get().getCoordinatorLogBoundaries();
return ref.get().getCoordinatorLogOffsets();
}
@Override

View File

@ -31,7 +31,6 @@ import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.apache.cassandra.db.ColumnFamilyStore;
import org.apache.cassandra.db.CoordinatorLogBoundaries;
import org.apache.cassandra.db.DecoratedKey;
import org.apache.cassandra.db.DeletionTime;
import org.apache.cassandra.db.Directories;
@ -58,6 +57,7 @@ import org.apache.cassandra.io.sstable.format.Version;
import org.apache.cassandra.io.util.DataInputPlus;
import org.apache.cassandra.io.util.TrackedDataInputPlus;
import org.apache.cassandra.metrics.StorageMetrics;
import org.apache.cassandra.replication.ImmutableCoordinatorLogOffsets;
import org.apache.cassandra.schema.TableId;
import org.apache.cassandra.schema.TableMetadata;
import org.apache.cassandra.service.StorageService;
@ -88,7 +88,7 @@ public class CassandraStreamReader implements IStreamReader
protected final Version inputVersion;
protected final long repairedAt;
protected final TimeUUID pendingRepair;
protected final CoordinatorLogBoundaries coordinatorLogBoundaries;
protected final ImmutableCoordinatorLogOffsets coordinatorLogOffsets;
protected final int sstableLevel;
protected final SerializationHeader.Component header;
protected final int fileSeqNum;
@ -108,7 +108,7 @@ public class CassandraStreamReader implements IStreamReader
this.inputVersion = streamHeader.version;
this.repairedAt = header.repairedAt;
this.pendingRepair = header.pendingRepair;
this.coordinatorLogBoundaries = header.coordinatorLogBoundaries;
this.coordinatorLogOffsets = header.coordinatorLogOffsets;
this.sstableLevel = streamHeader.sstableLevel;
this.header = streamHeader.serializationHeader;
this.fileSeqNum = header.sequenceNumber;
@ -138,7 +138,7 @@ public class CassandraStreamReader implements IStreamReader
try (StreamCompressionInputStream streamCompressionInputStream = new StreamCompressionInputStream(inputPlus, current_version))
{
TrackedDataInputPlus in = new TrackedDataInputPlus(streamCompressionInputStream);
writer = createWriter(cfs, totalSize, repairedAt, pendingRepair, coordinatorLogBoundaries, inputVersion.format);
writer = createWriter(cfs, totalSize, repairedAt, pendingRepair, coordinatorLogOffsets, inputVersion.format);
deserializer = getDeserializer(cfs.metadata(), in, inputVersion, session, writer);
String sequenceName = writer.getFilename() + '-' + fileSeqNum;
long lastBytesRead = 0;
@ -179,7 +179,7 @@ public class CassandraStreamReader implements IStreamReader
{
return header != null? header.toHeader(metadata) : null; //pre-3.0 sstable have no SerializationHeader
}
protected SSTableTxnSingleStreamWriter createWriter(ColumnFamilyStore cfs, long totalSize, long repairedAt, TimeUUID pendingRepair, CoordinatorLogBoundaries coordinatorLogBoundaries, SSTableFormat<?, ?> format) throws IOException
protected SSTableTxnSingleStreamWriter createWriter(ColumnFamilyStore cfs, long totalSize, long repairedAt, TimeUUID pendingRepair, ImmutableCoordinatorLogOffsets coordinatorLogOffsets, SSTableFormat<?, ?> format) throws IOException
{
Directories.DataDirectory localDir = cfs.getDirectories().getWriteableLocation(totalSize);
if (localDir == null)
@ -188,7 +188,7 @@ public class CassandraStreamReader implements IStreamReader
StreamReceiver streamReceiver = session.getAggregator(tableId);
Preconditions.checkState(streamReceiver instanceof CassandraStreamReceiver);
ILifecycleTransaction txn = createTxn();
RangeAwareSSTableWriter writer = new RangeAwareSSTableWriter(cfs, estimatedKeys, repairedAt, pendingRepair, false, coordinatorLogBoundaries, format, sstableLevel, totalSize, txn, getHeader(cfs.metadata()));
RangeAwareSSTableWriter writer = new RangeAwareSSTableWriter(cfs, estimatedKeys, repairedAt, pendingRepair, false, coordinatorLogOffsets, format, sstableLevel, totalSize, txn, getHeader(cfs.metadata()));
return new SSTableTxnSingleStreamWriter(txn, writer);
}

View File

@ -32,7 +32,6 @@ import java.util.stream.Stream;
import com.google.common.annotations.VisibleForTesting;
import org.apache.cassandra.config.DatabaseDescriptor;
import org.apache.cassandra.db.CoordinatorLogBoundaries;
import org.apache.cassandra.db.DecoratedKey;
import org.apache.cassandra.db.RegularAndStaticColumns;
import org.apache.cassandra.db.SerializationHeader;
@ -43,6 +42,7 @@ import org.apache.cassandra.index.Index;
import org.apache.cassandra.io.sstable.format.SSTableFormat;
import org.apache.cassandra.io.sstable.format.SSTableReader;
import org.apache.cassandra.io.util.File;
import org.apache.cassandra.replication.ImmutableCoordinatorLogOffsets;
import org.apache.cassandra.schema.TableMetadataRef;
import org.apache.cassandra.service.ActiveRepairService;
@ -130,7 +130,7 @@ public abstract class AbstractSSTableSimpleWriter implements Closeable
SerializationHeader header = new SerializationHeader(true, metadata.get(), columns, EncodingStats.NO_STATS);
if (makeRangeAware)
return SSTableTxnWriter.createRangeAware(metadata, 0, ActiveRepairService.UNREPAIRED_SSTABLE, ActiveRepairService.NO_PENDING_REPAIR, false, CoordinatorLogBoundaries.NONE, format, header);
return SSTableTxnWriter.createRangeAware(metadata, 0, ActiveRepairService.UNREPAIRED_SSTABLE, ActiveRepairService.NO_PENDING_REPAIR, false, ImmutableCoordinatorLogOffsets.NONE, format, header);
SSTable.Owner effectiveOwner;
@ -152,7 +152,7 @@ public abstract class AbstractSSTableSimpleWriter implements Closeable
ActiveRepairService.UNREPAIRED_SSTABLE,
ActiveRepairService.NO_PENDING_REPAIR,
false,
CoordinatorLogBoundaries.NONE,
ImmutableCoordinatorLogOffsets.NONE,
header,
indexGroups,
effectiveOwner);

View File

@ -23,7 +23,6 @@ import java.util.Collection;
import java.util.List;
import org.apache.cassandra.db.ColumnFamilyStore;
import org.apache.cassandra.db.CoordinatorLogBoundaries;
import org.apache.cassandra.db.DecoratedKey;
import org.apache.cassandra.db.Directories;
import org.apache.cassandra.db.DiskBoundaries;
@ -33,6 +32,7 @@ import org.apache.cassandra.db.lifecycle.ILifecycleTransaction;
import org.apache.cassandra.db.rows.UnfilteredRowIterator;
import org.apache.cassandra.io.sstable.format.SSTableFormat;
import org.apache.cassandra.io.sstable.format.SSTableReader;
import org.apache.cassandra.replication.ImmutableCoordinatorLogOffsets;
import org.apache.cassandra.schema.TableId;
import org.apache.cassandra.utils.FBUtilities;
import org.apache.cassandra.utils.TimeUUID;
@ -46,7 +46,7 @@ public class RangeAwareSSTableWriter implements SSTableMultiWriter
private final long repairedAt;
private final TimeUUID pendingRepair;
private final boolean isTransient;
private final CoordinatorLogBoundaries coordinatorLogBoundaries;
private final ImmutableCoordinatorLogOffsets coordinatorLogOffsets;
private final SSTableFormat<?, ?> format;
private final SerializationHeader header;
private final ILifecycleTransaction txn;
@ -55,7 +55,7 @@ public class RangeAwareSSTableWriter implements SSTableMultiWriter
private final List<SSTableMultiWriter> finishedWriters = new ArrayList<>();
private SSTableMultiWriter currentWriter = null;
public RangeAwareSSTableWriter(ColumnFamilyStore cfs, long estimatedKeys, long repairedAt, TimeUUID pendingRepair, boolean isTransient, CoordinatorLogBoundaries coordinatorLogBoundaries, SSTableFormat<?, ?> format, int sstableLevel, long totalSize, ILifecycleTransaction txn, SerializationHeader header) throws IOException
public RangeAwareSSTableWriter(ColumnFamilyStore cfs, long estimatedKeys, long repairedAt, TimeUUID pendingRepair, boolean isTransient, ImmutableCoordinatorLogOffsets coordinatorLogOffsets, SSTableFormat<?, ?> format, int sstableLevel, long totalSize, ILifecycleTransaction txn, SerializationHeader header) throws IOException
{
DiskBoundaries db = cfs.getDiskBoundaries();
directories = db.directories;
@ -65,7 +65,7 @@ public class RangeAwareSSTableWriter implements SSTableMultiWriter
this.repairedAt = repairedAt;
this.pendingRepair = pendingRepair;
this.isTransient = isTransient;
this.coordinatorLogBoundaries = coordinatorLogBoundaries;
this.coordinatorLogOffsets = coordinatorLogOffsets;
this.format = format;
this.txn = txn;
this.header = header;
@ -77,7 +77,7 @@ public class RangeAwareSSTableWriter implements SSTableMultiWriter
throw new IOException(String.format("Insufficient disk space to store %s",
FBUtilities.prettyPrintMemory(totalSize)));
Descriptor desc = cfs.newSSTableDescriptor(cfs.getDirectories().getLocationForDisk(localDir), format);
currentWriter = cfs.createSSTableMultiWriter(desc, estimatedKeys, repairedAt, pendingRepair, isTransient, coordinatorLogBoundaries, null, sstableLevel, header, txn);
currentWriter = cfs.createSSTableMultiWriter(desc, estimatedKeys, repairedAt, pendingRepair, isTransient, coordinatorLogOffsets, null, sstableLevel, header, txn);
}
}
@ -99,7 +99,7 @@ public class RangeAwareSSTableWriter implements SSTableMultiWriter
finishedWriters.add(currentWriter);
Descriptor desc = cfs.newSSTableDescriptor(cfs.getDirectories().getLocationForDisk(directories.get(currentIndex)), format);
currentWriter = cfs.createSSTableMultiWriter(desc, estimatedKeys, repairedAt, pendingRepair, isTransient, coordinatorLogBoundaries, null, sstableLevel, header, txn);
currentWriter = cfs.createSSTableMultiWriter(desc, estimatedKeys, repairedAt, pendingRepair, isTransient, coordinatorLogOffsets, null, sstableLevel, header, txn);
}
}

View File

@ -22,7 +22,6 @@ import java.io.IOException;
import java.util.Collection;
import org.apache.cassandra.db.ColumnFamilyStore;
import org.apache.cassandra.db.CoordinatorLogBoundaries;
import org.apache.cassandra.db.Keyspace;
import org.apache.cassandra.db.SerializationHeader;
import org.apache.cassandra.db.compaction.OperationType;
@ -31,6 +30,7 @@ import org.apache.cassandra.db.rows.UnfilteredRowIterator;
import org.apache.cassandra.index.Index;
import org.apache.cassandra.io.sstable.format.SSTableFormat;
import org.apache.cassandra.io.sstable.format.SSTableReader;
import org.apache.cassandra.replication.ImmutableCoordinatorLogOffsets;
import org.apache.cassandra.schema.TableMetadataRef;
import org.apache.cassandra.utils.TimeUUID;
import org.apache.cassandra.utils.concurrent.Transactional;
@ -114,10 +114,10 @@ public class SSTableTxnWriter extends Transactional.AbstractTransactional implem
}
@SuppressWarnings({"resource", "RedundantSuppression"}) // log and writer closed during doPostCleanup
public static SSTableTxnWriter create(ColumnFamilyStore cfs, Descriptor descriptor, long keyCount, long repairedAt, TimeUUID pendingRepair, boolean isTransient, CoordinatorLogBoundaries coordinatorLogBoundaries, SerializationHeader header)
public static SSTableTxnWriter create(ColumnFamilyStore cfs, Descriptor descriptor, long keyCount, long repairedAt, TimeUUID pendingRepair, boolean isTransient, ImmutableCoordinatorLogOffsets coordinatorLogOffsets, SerializationHeader header)
{
LifecycleTransaction txn = LifecycleTransaction.offline(OperationType.WRITE);
SSTableMultiWriter writer = cfs.createSSTableMultiWriter(descriptor, keyCount, repairedAt, pendingRepair, isTransient, coordinatorLogBoundaries, header, txn);
SSTableMultiWriter writer = cfs.createSSTableMultiWriter(descriptor, keyCount, repairedAt, pendingRepair, isTransient, coordinatorLogOffsets, header, txn);
return new SSTableTxnWriter(txn, writer);
}
@ -126,7 +126,7 @@ public class SSTableTxnWriter extends Transactional.AbstractTransactional implem
long repairedAt,
TimeUUID pendingRepair,
boolean isTransient,
CoordinatorLogBoundaries coordinatorLogBoundaries,
ImmutableCoordinatorLogOffsets coordinatorLogOffsets,
SSTableFormat<?, ?> type,
SerializationHeader header)
{
@ -136,7 +136,7 @@ public class SSTableTxnWriter extends Transactional.AbstractTransactional implem
SSTableMultiWriter writer;
try
{
writer = new RangeAwareSSTableWriter(cfs, keyCount, repairedAt, pendingRepair, isTransient, coordinatorLogBoundaries, type, 0, 0, txn, header);
writer = new RangeAwareSSTableWriter(cfs, keyCount, repairedAt, pendingRepair, isTransient, coordinatorLogOffsets, type, 0, 0, txn, header);
}
catch (IOException e)
{
@ -155,14 +155,14 @@ public class SSTableTxnWriter extends Transactional.AbstractTransactional implem
long repairedAt,
TimeUUID pendingRepair,
boolean isTransient,
CoordinatorLogBoundaries coordinatorLogBoundaries,
ImmutableCoordinatorLogOffsets coordinatorLogOffsets,
SerializationHeader header,
Collection<Index.Group> indexGroups,
SSTable.Owner owner)
{
// if the column family store does not exist, we create a new default SSTableMultiWriter to use:
LifecycleTransaction txn = LifecycleTransaction.offline(OperationType.WRITE);
SSTableMultiWriter writer = SimpleSSTableMultiWriter.create(descriptor, keyCount, repairedAt, pendingRepair, isTransient, coordinatorLogBoundaries, metadata, null, 0, header, indexGroups, txn, owner);
SSTableMultiWriter writer = SimpleSSTableMultiWriter.create(descriptor, keyCount, repairedAt, pendingRepair, isTransient, coordinatorLogOffsets, metadata, null, 0, header, indexGroups, txn, owner);
return new SSTableTxnWriter(txn, writer);
}
}

View File

@ -20,7 +20,6 @@ package org.apache.cassandra.io.sstable;
import java.util.Collection;
import java.util.Collections;
import org.apache.cassandra.db.CoordinatorLogBoundaries;
import org.apache.cassandra.db.SerializationHeader;
import org.apache.cassandra.db.commitlog.CommitLogPosition;
import org.apache.cassandra.db.commitlog.IntervalSet;
@ -31,6 +30,7 @@ import org.apache.cassandra.index.Index;
import org.apache.cassandra.io.sstable.format.SSTableReader;
import org.apache.cassandra.io.sstable.format.SSTableWriter;
import org.apache.cassandra.io.sstable.metadata.MetadataCollector;
import org.apache.cassandra.replication.ImmutableCoordinatorLogOffsets;
import org.apache.cassandra.schema.TableId;
import org.apache.cassandra.schema.TableMetadataRef;
import org.apache.cassandra.utils.TimeUUID;
@ -118,7 +118,7 @@ public class SimpleSSTableMultiWriter implements SSTableMultiWriter
long repairedAt,
TimeUUID pendingRepair,
boolean isTransient,
CoordinatorLogBoundaries coordinatorLogBoundaries,
ImmutableCoordinatorLogOffsets coordinatorLogOffsets,
TableMetadataRef metadata,
IntervalSet<CommitLogPosition> commitLogPositions,
int sstableLevel,
@ -138,7 +138,7 @@ public class SimpleSSTableMultiWriter implements SSTableMultiWriter
.setRepairedAt(repairedAt)
.setPendingRepair(pendingRepair)
.setTransientSSTable(isTransient)
.setCoordinatorLogBoundaries(coordinatorLogBoundaries)
.setCoordinatorLogOffsets(coordinatorLogOffsets)
.setTableMetadataRef(metadata)
.setMetadataCollector(metadataCollector)
.setSerializationHeader(header)

View File

@ -58,7 +58,6 @@ import org.apache.cassandra.concurrent.ScheduledExecutors;
import org.apache.cassandra.config.CassandraRelevantProperties;
import org.apache.cassandra.config.DatabaseDescriptor;
import org.apache.cassandra.db.ColumnFamilyStore;
import org.apache.cassandra.db.CoordinatorLogBoundaries;
import org.apache.cassandra.db.DecoratedKey;
import org.apache.cassandra.db.PartitionPosition;
import org.apache.cassandra.db.SerializationHeader;
@ -99,6 +98,7 @@ import org.apache.cassandra.io.util.FileUtils;
import org.apache.cassandra.io.util.FileUtils.DuplicateHardlinkException;
import org.apache.cassandra.io.util.RandomAccessReader;
import org.apache.cassandra.metrics.RestorableMeter;
import org.apache.cassandra.replication.ImmutableCoordinatorLogOffsets;
import org.apache.cassandra.schema.SchemaConstants;
import org.apache.cassandra.schema.TableMetadataRef;
import org.apache.cassandra.service.ActiveRepairService;
@ -1230,9 +1230,9 @@ public abstract class SSTableReader extends SSTable implements UnfilteredSource,
return sstableMetadata.pendingRepair;
}
public CoordinatorLogBoundaries getCoordinatorLogBoundaries()
public ImmutableCoordinatorLogOffsets getCoordinatorLogOffsets()
{
return sstableMetadata.coordinatorLogBoundaries;
return sstableMetadata.coordinatorLogOffsets;
}
public long getRepairedAt()

View File

@ -25,12 +25,14 @@ import java.util.Collections;
import java.util.HashSet;
import java.util.List;
import java.util.Map;
import java.util.Objects;
import java.util.Set;
import java.util.function.Consumer;
import java.util.function.Supplier;
import javax.annotation.Nullable;
import com.google.common.base.Preconditions;
import com.google.common.collect.ImmutableList;
import com.google.common.collect.ImmutableSet;
import com.google.common.collect.Sets;
@ -38,7 +40,6 @@ import com.google.common.collect.Sets;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.apache.cassandra.db.CoordinatorLogBoundaries;
import org.apache.cassandra.db.DecoratedKey;
import org.apache.cassandra.db.SerializationHeader;
import org.apache.cassandra.db.compression.CompressionDictionaryManager;
@ -60,6 +61,10 @@ import org.apache.cassandra.io.sstable.metadata.MetadataComponent;
import org.apache.cassandra.io.sstable.metadata.MetadataType;
import org.apache.cassandra.io.sstable.metadata.StatsMetadata;
import org.apache.cassandra.io.util.MmappedRegionsCache;
import org.apache.cassandra.replication.ImmutableCoordinatorLogOffsets;
import org.apache.cassandra.replication.MutationTrackingService;
import org.apache.cassandra.service.ActiveRepairService;
import org.apache.cassandra.utils.Clock;
import org.apache.cassandra.utils.Throwables;
import org.apache.cassandra.utils.TimeUUID;
import org.apache.cassandra.utils.concurrent.Transactional;
@ -78,7 +83,7 @@ public abstract class SSTableWriter extends SSTable implements Transactional
protected long repairedAt;
protected TimeUUID pendingRepair;
protected boolean isTransient;
protected CoordinatorLogBoundaries coordinatorLogBoundaries;
protected ImmutableCoordinatorLogOffsets coordinatorLogOffsets;
protected long maxDataAge = -1;
protected final long keyCount;
public final MetadataCollector metadataCollector;
@ -102,13 +107,13 @@ public abstract class SSTableWriter extends SSTable implements Transactional
checkNotNull(builder.getIndexGroups());
checkNotNull(builder.getMetadataCollector());
checkNotNull(builder.getSerializationHeader());
checkNotNull(builder.getCoordinatorLogBoundaries());
checkNotNull(builder.getCoordinatorLogOffsets());
this.keyCount = builder.getKeyCount();
this.repairedAt = builder.getRepairedAt();
this.pendingRepair = builder.getPendingRepair();
this.isTransient = builder.isTransientSSTable();
this.coordinatorLogBoundaries = builder.getCoordinatorLogBoundaries();
this.coordinatorLogOffsets = builder.getCoordinatorLogOffsets();
this.metadataCollector = builder.getMetadataCollector();
this.header = builder.getSerializationHeader();
this.mmappedRegionsCache = builder.getMmappedRegionsCache();
@ -347,12 +352,24 @@ public abstract class SSTableWriter extends SSTable implements Transactional
protected final Map<MetadataType, MetadataComponent> finalizeMetadata()
{
// Migration from incremental repair to mutation tracking will be supported, but support for mixing
// incremental repair and mutation tracking is not planned
if (metadata().replicationType().isTracked() && repairedAt == ActiveRepairService.UNREPAIRED_SSTABLE)
{
Preconditions.checkState(Objects.equals(pendingRepair, ActiveRepairService.NO_PENDING_REPAIR));
if (MutationTrackingService.instance.isDurablyReconciled(getKeyspaceName(), coordinatorLogOffsets))
{
repairedAt = Clock.Global.currentTimeMillis();
logger.debug("Marking SSTable {} as reconciled with repairedAt {}", descriptor, repairedAt);
}
}
return metadataCollector.finalizeMetadata(getPartitioner().getClass().getCanonicalName(),
metadata().params.bloomFilterFpChance,
repairedAt,
pendingRepair,
isTransient,
coordinatorLogBoundaries,
coordinatorLogOffsets,
header,
first.retainable().getKey(),
last.retainable().getKey());
@ -398,13 +415,31 @@ public abstract class SSTableWriter extends SSTable implements Transactional
protected void doPrepare()
{
transactionals.get().forEach(Transactional::prepareToCommit);
new StatsComponent(finalizeMetadata()).save(descriptor);
Map<MetadataType, MetadataComponent> metadata = finalizeMetadata();
new StatsComponent(metadata).save(descriptor);
// save the table of components
TOCComponent.updateTOC(descriptor, components);
if (openResult)
{
finalReader = openFinal(SSTableReader.OpenReason.NORMAL);
/*
When we open above, we re-finalize metadata and may be durably reconciled, but this is after
StatsComponent is saved, so the next reload from disk loses the reconciliation update. We need to save
again to ensure the descriptor file matches what's in memory.
Could move this to the post-flush compaction, but then we couldn't flush directly into the repaired set.
Or see if we can open before the first save, so we never have to write stats to disk twice.
*/
StatsMetadata stale = ((StatsMetadata) metadata.get(MetadataType.STATS));
if (finalReader.getSSTableMetadata().repairedAt != stale.repairedAt)
{
metadata.put(MetadataType.STATS, finalReader.getSSTableMetadata());
new StatsComponent(metadata).save(descriptor);
}
}
}
protected Throwable doCommit(Throwable accumulate)
@ -459,7 +494,7 @@ public abstract class SSTableWriter extends SSTable implements Transactional
private List<Index.Group> indexGroups;
@Nullable
private CompressionDictionaryManager compressionDictionaryManager;
private CoordinatorLogBoundaries coordinatorLogBoundaries;
private ImmutableCoordinatorLogOffsets coordinatorLogOffsets;
public B setMetadataCollector(MetadataCollector metadataCollector)
{
@ -485,9 +520,9 @@ public abstract class SSTableWriter extends SSTable implements Transactional
return (B) this;
}
public B setCoordinatorLogBoundaries(CoordinatorLogBoundaries coordinatorLogBoundaries)
public B setCoordinatorLogOffsets(ImmutableCoordinatorLogOffsets coordinatorLogOffsets)
{
this.coordinatorLogBoundaries = coordinatorLogBoundaries;
this.coordinatorLogOffsets = coordinatorLogOffsets;
return (B) this;
}
@ -581,9 +616,9 @@ public abstract class SSTableWriter extends SSTable implements Transactional
return transientSSTable;
}
public CoordinatorLogBoundaries getCoordinatorLogBoundaries()
public ImmutableCoordinatorLogOffsets getCoordinatorLogOffsets()
{
return coordinatorLogBoundaries;
return coordinatorLogOffsets;
}
public SerializationHeader getSerializationHeader()

View File

@ -190,7 +190,7 @@ public abstract class SortedTableScrubber<R extends SSTableReaderWithFilter> imp
Refs<SSTableReader> refs = Refs.ref(Collections.singleton(sstable)))
{
StatsMetadata metadata = sstable.getSSTableMetadata();
writer.switchWriter(CompactionManager.createWriter(cfs, destination, expectedBloomFilterSize, metadata.repairedAt, metadata.pendingRepair, metadata.isTransient, metadata.coordinatorLogBoundaries, sstable, transaction));
writer.switchWriter(CompactionManager.createWriter(cfs, destination, expectedBloomFilterSize, metadata.repairedAt, metadata.pendingRepair, metadata.isTransient, metadata.coordinatorLogOffsets, sstable, transaction));
scrubInternal(writer);
@ -240,7 +240,7 @@ public abstract class SortedTableScrubber<R extends SSTableReaderWithFilter> imp
// out of order partitions/rows, but no bad partition found - we can keep our repairedAt time
long repairedAt = badPartitions > 0 ? ActiveRepairService.UNREPAIRED_SSTABLE : sstable.getSSTableMetadata().repairedAt;
SSTableReader newInOrderSstable;
try (SSTableWriter inOrderWriter = CompactionManager.createWriter(cfs, destination, expectedBloomFilterSize, repairedAt, metadata.pendingRepair, metadata.isTransient, metadata.coordinatorLogBoundaries, sstable, transaction))
try (SSTableWriter inOrderWriter = CompactionManager.createWriter(cfs, destination, expectedBloomFilterSize, repairedAt, metadata.pendingRepair, metadata.isTransient, metadata.coordinatorLogOffsets, sstable, transaction))
{
for (Partition partition : outOfOrder)
inOrderWriter.append(partition.unfilteredIterator());

View File

@ -32,7 +32,6 @@ import org.apache.cassandra.db.ClusteringBound;
import org.apache.cassandra.db.ClusteringBoundOrBoundary;
import org.apache.cassandra.db.ClusteringComparator;
import org.apache.cassandra.db.ClusteringPrefix;
import org.apache.cassandra.db.CoordinatorLogBoundaries;
import org.apache.cassandra.db.DeletionTime;
import org.apache.cassandra.db.LivenessInfo;
import org.apache.cassandra.db.SerializationHeader;
@ -48,6 +47,7 @@ import org.apache.cassandra.db.rows.Unfiltered;
import org.apache.cassandra.io.sstable.ClusteringDescriptor;
import org.apache.cassandra.io.sstable.SSTable;
import org.apache.cassandra.io.sstable.format.SSTableReader;
import org.apache.cassandra.replication.ImmutableCoordinatorLogOffsets;
import org.apache.cassandra.service.StorageService;
import org.apache.cassandra.utils.EstimatedHistogram;
import org.apache.cassandra.utils.FBUtilities;
@ -413,7 +413,7 @@ public class MetadataCollector implements PartitionStatisticsCollector
return totalRows;
}
public Map<MetadataType, MetadataComponent> finalizeMetadata(String partitioner, double bloomFilterFPChance, long repairedAt, TimeUUID pendingRepair, boolean isTransient, CoordinatorLogBoundaries coordinatorLogBoundaries, SerializationHeader header, ByteBuffer firstKey, ByteBuffer lastKey)
public Map<MetadataType, MetadataComponent> finalizeMetadata(String partitioner, double bloomFilterFPChance, long repairedAt, TimeUUID pendingRepair, boolean isTransient, ImmutableCoordinatorLogOffsets coordinatorLogOffsets, SerializationHeader header, ByteBuffer firstKey, ByteBuffer lastKey)
{
assert minClustering.kind() == ClusteringPrefix.Kind.CLUSTERING || minClustering.kind().isStart();
assert maxClustering.kind() == ClusteringPrefix.Kind.CLUSTERING || maxClustering.kind().isEnd();
@ -450,7 +450,7 @@ public class MetadataCollector implements PartitionStatisticsCollector
pendingRepair,
isTransient,
hasPartitionLevelDeletions,
coordinatorLogBoundaries,
coordinatorLogOffsets,
firstKey,
lastKey));
components.put(MetadataType.COMPACTION, new CompactionMetadata(cardinality));

View File

@ -30,7 +30,6 @@ import org.slf4j.LoggerFactory;
import org.apache.cassandra.db.BufferClusteringBound;
import org.apache.cassandra.db.ClusteringBound;
import org.apache.cassandra.db.CoordinatorLogBoundaries;
import org.apache.cassandra.db.Slice;
import org.apache.cassandra.db.TypeSizes;
import org.apache.cassandra.db.commitlog.CommitLogPosition;
@ -42,6 +41,7 @@ import org.apache.cassandra.io.ISerializer;
import org.apache.cassandra.io.sstable.format.Version;
import org.apache.cassandra.io.util.DataInputPlus;
import org.apache.cassandra.io.util.DataOutputPlus;
import org.apache.cassandra.replication.ImmutableCoordinatorLogOffsets;
import org.apache.cassandra.serializers.AbstractTypeSerializer;
import org.apache.cassandra.utils.ByteBufferUtil;
import org.apache.cassandra.utils.EstimatedHistogram;
@ -80,7 +80,7 @@ public class StatsMetadata extends MetadataComponent
public final UUID originatingHostId;
public final TimeUUID pendingRepair;
public final boolean isTransient;
public final CoordinatorLogBoundaries coordinatorLogBoundaries;
public final ImmutableCoordinatorLogOffsets coordinatorLogOffsets;
// just holds the current encoding stats to avoid allocating - it is not serialized
public final EncodingStats encodingStats;
@ -124,7 +124,7 @@ public class StatsMetadata extends MetadataComponent
TimeUUID pendingRepair,
boolean isTransient,
boolean hasPartitionLevelDeletions,
CoordinatorLogBoundaries coordinatorLogBoundaries,
ImmutableCoordinatorLogOffsets coordinatorLogOffsets,
ByteBuffer firstKey,
ByteBuffer lastKey)
{
@ -150,7 +150,7 @@ public class StatsMetadata extends MetadataComponent
this.originatingHostId = originatingHostId;
this.pendingRepair = pendingRepair;
this.isTransient = isTransient;
this.coordinatorLogBoundaries = coordinatorLogBoundaries;
this.coordinatorLogOffsets = coordinatorLogOffsets;
this.encodingStats = new EncodingStats(minTimestamp, minLocalDeletionTime, minTTL);
this.hasPartitionLevelDeletions = hasPartitionLevelDeletions;
this.firstKey = firstKey;
@ -211,7 +211,7 @@ public class StatsMetadata extends MetadataComponent
pendingRepair,
isTransient,
hasPartitionLevelDeletions,
coordinatorLogBoundaries,
coordinatorLogOffsets,
firstKey,
lastKey);
}
@ -241,7 +241,7 @@ public class StatsMetadata extends MetadataComponent
newPendingRepair,
newIsTransient,
hasPartitionLevelDeletions,
coordinatorLogBoundaries,
coordinatorLogOffsets,
firstKey,
lastKey);
}
@ -275,7 +275,7 @@ public class StatsMetadata extends MetadataComponent
.append(originatingHostId, that.originatingHostId)
.append(pendingRepair, that.pendingRepair)
.append(hasPartitionLevelDeletions, that.hasPartitionLevelDeletions)
.append(coordinatorLogBoundaries, that.coordinatorLogBoundaries)
.append(coordinatorLogOffsets, that.coordinatorLogOffsets)
.append(firstKey, that.firstKey)
.append(lastKey, that.lastKey)
.build();
@ -306,7 +306,7 @@ public class StatsMetadata extends MetadataComponent
.append(originatingHostId)
.append(pendingRepair)
.append(hasPartitionLevelDeletions)
.append(coordinatorLogBoundaries)
.append(coordinatorLogOffsets)
.append(firstKey)
.append(lastKey)
.build();
@ -395,7 +395,7 @@ public class StatsMetadata extends MetadataComponent
}
if (version.hasMutationTrackingMetadata())
size += CoordinatorLogBoundaries.serializer.serializedSize(component.coordinatorLogBoundaries, version.correspondingMessagingVersion());
size += ImmutableCoordinatorLogOffsets.serializer.serializedSize(component.coordinatorLogOffsets, version.correspondingMessagingVersion());
return size;
}
@ -522,7 +522,7 @@ public class StatsMetadata extends MetadataComponent
}
if (version.hasMutationTrackingMetadata())
CoordinatorLogBoundaries.serializer.serialize(component.coordinatorLogBoundaries, out, version.correspondingMessagingVersion());
ImmutableCoordinatorLogOffsets.serializer.serialize(component.coordinatorLogOffsets, out, version.correspondingMessagingVersion());
}
private void serializeImprovedMinMax(Version version, StatsMetadata component, DataOutputPlus out) throws IOException
@ -668,9 +668,9 @@ public class StatsMetadata extends MetadataComponent
tokenSpaceCoverage = in.readDouble();
}
CoordinatorLogBoundaries coordinatorLogBoundaries = CoordinatorLogBoundaries.NONE;
ImmutableCoordinatorLogOffsets coordinatorLogOffsets = ImmutableCoordinatorLogOffsets.NONE;
if (version.hasMutationTrackingMetadata())
coordinatorLogBoundaries = CoordinatorLogBoundaries.serializer.deserialize(in, version.correspondingMessagingVersion());
coordinatorLogOffsets = ImmutableCoordinatorLogOffsets.serializer.deserialize(in, version.correspondingMessagingVersion());
return new StatsMetadata(partitionSizes,
columnCounts,
@ -695,7 +695,7 @@ public class StatsMetadata extends MetadataComponent
pendingRepair,
isTransient,
hasPartitionLevelDeletions,
coordinatorLogBoundaries,
coordinatorLogOffsets,
firstKey,
lastKey);
}

View File

@ -17,6 +17,7 @@
*/
package org.apache.cassandra.replication;
import java.util.Objects;
import java.util.concurrent.atomic.AtomicLong;
import java.util.concurrent.locks.ReadWriteLock;
import java.util.concurrent.locks.ReentrantReadWriteLock;
@ -32,6 +33,7 @@ import org.apache.cassandra.db.PartitionPosition;
import org.apache.cassandra.dht.AbstractBounds;
import org.apache.cassandra.dht.Token;
import org.apache.cassandra.schema.TableId;
import org.apache.cassandra.tcm.ClusterMetadata;
import static org.apache.cassandra.utils.Clock.Global.currentTimeMillis;
@ -74,6 +76,7 @@ public abstract class CoordinatorLog
void receivedWriteResponse(MutationId mutationId, int onHostId)
{
Preconditions.checkArgument(!mutationId.isNone());
Preconditions.checkArgument(!Objects.equals(onHostId, ClusterMetadata.current().myNodeId().id()));
logger.trace("witnessed remote mutation {} from {}", mutationId, onHostId);
lock.writeLock().lock();
try
@ -235,11 +238,37 @@ public abstract class CoordinatorLog
return witnessedOffsets[participants.indexOf(hostId)];
}
boolean isDurablyReconciled(CoordinatorLogOffsets<?> logOffsets)
{
lock.readLock().lock();
try
{
// TODO: reconciledOffsets not necessarily durable, update once durability is implemented
Offsets.RangeIterator durablyReconciled = reconciledOffsets.rangeIterator();
Offsets.RangeIterator difference = Offsets.difference(logOffsets.offsets(logId.asLong()).rangeIterator(), durablyReconciled);
return !difference.tryAdvance();
}
finally
{
lock.readLock().unlock();
}
}
protected Offsets.Mutable getLocal()
{
return witnessedOffsets[participants.indexOf(localHostId)];
}
@Override
public String toString()
{
return "CoordinatorLog{" +
"logId=" + logId +
", localHostId=" + localHostId +
", participants=" + participants +
'}';
}
public static class CoordinatorLogPrimary extends CoordinatorLog
{
private final AtomicLong sequenceId = new AtomicLong(-1);

View File

@ -16,31 +16,24 @@
* limitations under the License.
*/
package org.apache.cassandra.db;
package org.apache.cassandra.replication;
import org.apache.cassandra.replication.MutationId;
import org.apache.cassandra.io.sstable.metadata.StatsMetadata;
public interface MutableCoordinatorLogBoundaries extends CoordinatorLogBoundaries
/**
* Mutation ID offsets present in this SSTable for each coordinator log, to determine whether an SSTable is reconciled
* or not.
* <p>
* Note that peers may have reconciled all mutations included in an SSTable, but {@link StatsMetadata#repairedAt} is
* dependent on compaction timing, so "nodetool repair --validate" may report temporary disagreements on the repaired
* set.
* <p>
* Iterable over {@link CoordinatorLogId}.
*/
public interface CoordinatorLogOffsets<O extends Offsets> extends Iterable<Long>
{
void add(MutationId mutationId);
O offsets(long logId);
int size();
default void addAll(CoordinatorLogBoundaries from)
{
for (long logId : from)
{
MutationId max = from.max(logId);
if (!max.isNone())
add(max);
}
}
static MutableCoordinatorLogBoundaries create()
{
return new CoordinatorLogBoundariesMap();
}
static MutableCoordinatorLogBoundaries create(int size)
{
return new CoordinatorLogBoundariesMap(size);
}
ImmutableCoordinatorLogOffsets NONE = new ImmutableCoordinatorLogOffsets.Builder(0).build();
}

View File

@ -0,0 +1,174 @@
/*
* 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.replication;
import java.io.IOException;
import java.util.Iterator;
import java.util.Map;
import java.util.Objects;
import javax.annotation.concurrent.NotThreadSafe;
import com.google.common.collect.Iterators;
import org.agrona.collections.Long2ObjectHashMap;
import org.apache.cassandra.io.IVersionedSerializer;
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.vint.VIntCoding;
public class ImmutableCoordinatorLogOffsets implements CoordinatorLogOffsets<Offsets.Immutable>
{
private final Long2ObjectHashMap<Offsets.Immutable> ids;
@Override
public Offsets.Immutable offsets(long logId)
{
Offsets.Immutable offsets = ids.get(logId);
if (offsets == null)
return new Offsets.Immutable(new CoordinatorLogId(logId));
return offsets;
}
@Override
public int size()
{
return ids.size();
}
@Override
public Iterator<Long> iterator()
{
return Iterators.unmodifiableIterator(ids.keySet().iterator());
}
@Override
public boolean equals(Object o)
{
if (o == null || getClass() != o.getClass()) return false;
ImmutableCoordinatorLogOffsets longs = (ImmutableCoordinatorLogOffsets) o;
return Objects.equals(ids, longs.ids);
}
@Override
public int hashCode()
{
return Objects.hashCode(ids);
}
public ImmutableCoordinatorLogOffsets(Builder builder)
{
// Important to set shouldAvoidAllocation=false, otherwise iterators are cached and not thread safe, even when
// immutable and read-only
this.ids = new Long2ObjectHashMap<>(builder.ids.size(), 0.9f, false);
for (Map.Entry<Long, Offsets.Immutable.Builder> entry : builder.ids.entrySet())
ids.put(entry.getKey(), entry.getValue().build());
}
@NotThreadSafe
public static class Builder
{
private final Long2ObjectHashMap<Offsets.Immutable.Builder> ids;
public Builder()
{
this(16);
}
public Builder(int size)
{
this.ids = new Long2ObjectHashMap<>(size, 0.9f, false);
}
public Builder add(MutationId mutationId)
{
if (mutationId.isNone())
return this;
ids.computeIfAbsent(mutationId.logId(), logId -> new Offsets.Immutable.Builder(new CoordinatorLogId(logId)))
.add(mutationId.offset());
return this;
}
public Builder addAll(CoordinatorLogOffsets<?> logOffsets)
{
for (long log : logOffsets)
{
Offsets offsets = logOffsets.offsets(log);
ids.computeIfAbsent(log, logId -> new Offsets.Immutable.Builder(new CoordinatorLogId(logId)))
.addAll(offsets);
}
return this;
}
public Builder addAll(Offsets.Immutable offsets)
{
ids.computeIfAbsent(offsets.logId.asLong(), logId -> new Offsets.Immutable.Builder(new CoordinatorLogId(logId)))
.addAll(offsets);
return this;
}
public ImmutableCoordinatorLogOffsets build()
{
return new ImmutableCoordinatorLogOffsets(this);
}
}
public static class Serializer implements IVersionedSerializer<ImmutableCoordinatorLogOffsets>
{
@Override
public void serialize(ImmutableCoordinatorLogOffsets logOffsets, DataOutputPlus out, int version) throws IOException
{
if (version < MessagingService.VERSION_52)
return;
out.writeUnsignedVInt32(logOffsets.size());
for (long logId : logOffsets)
Offsets.serializer.serialize(logOffsets.offsets(logId), out, version);
}
@Override
public ImmutableCoordinatorLogOffsets deserialize(DataInputPlus in, int version) throws IOException
{
if (version < MessagingService.VERSION_52)
return ImmutableCoordinatorLogOffsets.NONE;
int size = in.readUnsignedVInt32();
ImmutableCoordinatorLogOffsets.Builder builder = new ImmutableCoordinatorLogOffsets.Builder(size);
for (int i = 0; i < size; i++)
{
Offsets.Immutable offsets = Offsets.serializer.deserialize(in, version);
builder.addAll(offsets);
}
return builder.build();
}
@Override
public long serializedSize(ImmutableCoordinatorLogOffsets logOffsets, int version)
{
if (version < MessagingService.VERSION_52)
return 0;
long size = 0;
size += VIntCoding.computeUnsignedVIntSize(logOffsets.size());
for (long logId : logOffsets)
size += Offsets.serializer.serializedSize(logOffsets.offsets(logId), version);
return size;
}
}
public static final Serializer serializer = new Serializer();
}

View File

@ -0,0 +1,49 @@
/*
* 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.replication;
import org.apache.cassandra.config.DatabaseDescriptor;
/*
* A replica can only receive writes from another replica it shares ranges with, and tracked writes are executed by
* coordinators, so instances should generally contain up to (2*RF - 1) keys.
*/
public interface MutableCoordinatorLogOffsets extends CoordinatorLogOffsets<Offsets.Mutable>
{
void add(ShortMutationId mutationId);
default void addAll(ImmutableCoordinatorLogOffsets from)
{
for (long logId : from)
{
Offsets offsets = from.offsets(logId);
offsets.forEach(this::add);
}
}
default void addAll(Offsets from)
{
from.forEach(this::add);
}
static MutableCoordinatorLogOffsets create(boolean assumeExclusive)
{
return assumeExclusive ? new NonBlockingCoordinatorLogOffsets.Exclusive() : new NonBlockingCoordinatorLogOffsets.Concurrent(DatabaseDescriptor.getConcurrentWriters());
}
}

View File

@ -180,12 +180,41 @@ public class MutationTrackingService
}
private final AtomicInteger nextHostLogId = new AtomicInteger();
private static class KeyspaceShards
public boolean isDurablyReconciled(String keyspace, ImmutableCoordinatorLogOffsets logOffsets)
{
// Could pass through SSTable bounds to exclude shards for non-overlapping ranges, but this will mostly be
// called on flush for L0 SSTables with wide bounds.
KeyspaceShards keyspaceShards = shards.get(keyspace);
if (keyspaceShards == null)
{
logger.debug("Could not find shards for keyspace {}", keyspace);
return false;
}
for (Long logId : logOffsets)
{
CoordinatorLogId coordinatorLogId = new CoordinatorLogId(logId);
CoordinatorLog log = keyspaceShards.logs.get(coordinatorLogId);
if (log == null)
{
logger.warn("Could not determine lifecycle for unknown logId {}, not marking as durably reconciled", coordinatorLogId);
return false;
}
if (!log.isDurablyReconciled(logOffsets))
return false;
}
return true;
}
private static class KeyspaceShards implements Shard.Subscriber
{
private final String keyspace;
private final Map<Range<Token>, Shard> shards;
private transient final Map<Range<PartitionPosition>, Shard> ppShards;
private transient final Map<CoordinatorLogId, CoordinatorLog> logs;
static KeyspaceShards make(KeyspaceMetadata keyspace, ClusterMetadata cluster, IntSupplier logIdProvider)
{
@ -209,7 +238,11 @@ public class MutationTrackingService
this.shards = shards;
this.ppShards = new HashMap<>();
shards.forEach((range, shard) -> ppShards.put(Range.makeRowRange(range), shard));
this.logs = new HashMap<>();
shards.forEach((range, shard) -> {
ppShards.put(Range.makeRowRange(range), shard);
shard.addSubscriber(this);
});
}
MutationId nextMutationId(Token token)
@ -285,6 +318,20 @@ public class MutationTrackingService
Range<Token> range = ClusterMetadata.current().placements.get(ksm.params.replication).writes.forRange(token).range();
return shards.get(range);
}
@Override
public void onLogCreation(CoordinatorLog log)
{
logger.debug("Indexing created log {}", log);
logs.put(log.logId, log);
}
@Override
public void onSubscribe(CoordinatorLog currentLog)
{
logger.debug("Indexing current log {}", currentLog);
logs.put(currentLog.logId, currentLog);
}
}
// TODO (later): a more intelligent heuristic for offsets included in broadcasts

View File

@ -0,0 +1,206 @@
/*
* 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.replication;
import java.util.Iterator;
import java.util.concurrent.locks.ReentrantLock;
import com.google.common.collect.Iterators;
import org.apache.cassandra.db.memtable.Memtable;
import org.apache.cassandra.db.memtable.SkipListMemtable;
import org.apache.cassandra.db.memtable.TrieMemtable;
import org.jctools.maps.NonBlockingHashMapLong;
import org.jctools.queues.MpscUnboundedArrayQueue;
/**
* This is different from {@link Log2OffsetsMap} because it's focused on supporting fast, frequent updates from multiple
* threads at {@link Memtable#put}, and infrequent reads at {@link Memtable#getFlushSet}.
* <p>
* Concurrent lock-free memtable implementations like {@link SkipListMemtable} should use {@link Concurrent}, which
* performs better for contending updates. Locking memtable implementations like {@link TrieMemtable} should use
* {@link Exclusive} which assumes updates are already protected by a lock at a higher level.
*/
abstract class NonBlockingCoordinatorLogOffsets<E extends NonBlockingCoordinatorLogOffsets.Entry> extends NonBlockingHashMapLong<E> implements MutableCoordinatorLogOffsets
{
interface EntryFactory<E>
{
E create(CoordinatorLogId logId);
}
private final EntryFactory<E> factory;
private NonBlockingCoordinatorLogOffsets(EntryFactory<E> factory)
{
this.factory = factory;
}
abstract static class Entry
{
protected final Offsets.Mutable base;
private Entry(CoordinatorLogId logId)
{
this.base = new Offsets.Mutable(logId);
}
abstract void add(ShortMutationId id);
abstract Offsets.Mutable offsets();
}
public void add(ShortMutationId mutationId)
{
computeIfAbsent(mutationId.logId(), logId -> factory.create(new CoordinatorLogId(logId))).add(mutationId);
}
@Override
public Offsets.Mutable offsets(long logId)
{
E logOffsets = get(logId);
if (logOffsets == null)
return new Offsets.Mutable(new CoordinatorLogId(logId));
return logOffsets.offsets();
}
@Override
public Iterator<Long> iterator()
{
return Iterators.unmodifiableIterator(keys().asIterator());
}
public static class Exclusive extends NonBlockingCoordinatorLogOffsets<Exclusive.Entry>
{
private static final Exclusive.EntryFactory FACTORY = new Exclusive.EntryFactory();
static class Entry extends NonBlockingCoordinatorLogOffsets.Entry
{
private Entry(CoordinatorLogId logId)
{
super(logId);
}
@Override
void add(ShortMutationId id)
{
base.add(id.offset());
}
@Override
Offsets.Mutable offsets()
{
return base;
}
}
static class EntryFactory implements NonBlockingCoordinatorLogOffsets.EntryFactory<Exclusive.Entry>
{
@Override
public Exclusive.Entry create(CoordinatorLogId logId)
{
return new Exclusive.Entry(logId);
}
}
public Exclusive()
{
super(FACTORY);
}
}
public static class Concurrent extends NonBlockingCoordinatorLogOffsets<Concurrent.Entry>
{
static class Entry extends NonBlockingCoordinatorLogOffsets.Entry
{
ReentrantLock lock = new ReentrantLock();
MpscUnboundedArrayQueue<Integer> contended;
private Entry(CoordinatorLogId logId, int contentions)
{
super(logId);
this.contended = new MpscUnboundedArrayQueue<>(Math.max(2, contentions));
}
@Override
void add(ShortMutationId id)
{
int offset = id.offset();
boolean locked = lock.tryLock();
try
{
if (locked)
{
flush();
base.add(offset);
}
else
contended.add(offset);
}
finally
{
if (locked)
lock.unlock();
}
}
@Override
Offsets.Mutable offsets()
{
flush();
return base;
}
private void flush()
{
boolean locked = lock.tryLock();
try
{
if (locked && !contended.isEmpty())
contended.drain(base::add);
}
finally
{
if (locked)
lock.unlock();
}
}
}
static class EntryFactory implements NonBlockingCoordinatorLogOffsets.EntryFactory<Concurrent.Entry>
{
final int contentions;
public EntryFactory(int contentions)
{
this.contentions = contentions;
}
@Override
public Entry create(CoordinatorLogId logId)
{
return new Concurrent.Entry(logId, contentions);
}
}
public Concurrent(int contentions)
{
super(new EntryFactory(contentions));
}
}
}

View File

@ -158,7 +158,25 @@ public class Shard
private CoordinatorLog getOrCreate(long logId)
{
CoordinatorLog log = logs.get(logId);
return log != null
? log : logs.computeIfAbsent(logId, ignore -> CoordinatorLog.create(localHostId, new CoordinatorLogId(logId), participants));
if (log != null)
return log;
CoordinatorLog newLog = logs.computeIfAbsent(logId, ignore -> CoordinatorLog.create(localHostId, new CoordinatorLogId(logId), participants));
for (Subscriber subscriber : subscribers)
subscriber.onLogCreation(newLog);
return newLog;
}
private final List<Subscriber> subscribers = new ArrayList<>();
public interface Subscriber
{
default void onLogCreation(CoordinatorLog log) {}
default void onSubscribe(CoordinatorLog currentLog) {}
}
public void addSubscriber(Subscriber subscriber)
{
subscriber.onSubscribe(currentLocalLog);
subscribers.add(subscriber);
}
}

View File

@ -21,9 +21,9 @@ package org.apache.cassandra.streaming;
import java.io.IOException;
import java.util.List;
import org.apache.cassandra.db.CoordinatorLogBoundaries;
import org.apache.cassandra.dht.Range;
import org.apache.cassandra.dht.Token;
import org.apache.cassandra.replication.ImmutableCoordinatorLogOffsets;
import org.apache.cassandra.schema.TableId;
import org.apache.cassandra.utils.TimeUUID;
@ -49,7 +49,7 @@ public interface OutgoingStream
long getRepairedAt();
TimeUUID getPendingRepair();
CoordinatorLogBoundaries getCoordinatorLogBoundaries();
ImmutableCoordinatorLogOffsets getCoordinatorLogOffsets();
String getName();

View File

@ -75,7 +75,7 @@ public class OutgoingStreamMessage extends StreamMessage
sequenceNumber,
stream.getRepairedAt(),
stream.getPendingRepair(),
stream.getCoordinatorLogBoundaries());
stream.getCoordinatorLogOffsets());
}
public synchronized void serialize(StreamingDataOutputPlus out, int version, StreamSession session) throws IOException

View File

@ -21,11 +21,11 @@ import java.io.IOException;
import com.google.common.base.Objects;
import org.apache.cassandra.db.CoordinatorLogBoundaries;
import org.apache.cassandra.db.TypeSizes;
import org.apache.cassandra.io.util.DataInputPlus;
import org.apache.cassandra.io.util.DataOutputPlus;
import org.apache.cassandra.locator.InetAddressAndPort;
import org.apache.cassandra.replication.ImmutableCoordinatorLogOffsets;
import org.apache.cassandra.schema.TableId;
import org.apache.cassandra.streaming.StreamSession;
import org.apache.cassandra.utils.TimeUUID;
@ -47,7 +47,7 @@ public class StreamMessageHeader
public final int sequenceNumber;
public final long repairedAt;
public final TimeUUID pendingRepair;
public final CoordinatorLogBoundaries coordinatorLogBoundaries;
public final ImmutableCoordinatorLogOffsets coordinatorLogOffsets;
public final InetAddressAndPort sender;
public StreamMessageHeader(TableId tableId,
@ -58,7 +58,7 @@ public class StreamMessageHeader
int sequenceNumber,
long repairedAt,
TimeUUID pendingRepair,
CoordinatorLogBoundaries coordinatorLogBoundaries)
ImmutableCoordinatorLogOffsets coordinatorLogOffsets)
{
this.tableId = tableId;
this.sender = sender;
@ -68,7 +68,7 @@ public class StreamMessageHeader
this.sequenceNumber = sequenceNumber;
this.repairedAt = repairedAt;
this.pendingRepair = pendingRepair;
this.coordinatorLogBoundaries = coordinatorLogBoundaries;
this.coordinatorLogOffsets = coordinatorLogOffsets;
}
@Override
@ -123,7 +123,7 @@ public class StreamMessageHeader
{
header.pendingRepair.serialize(out);
}
CoordinatorLogBoundaries.serializer.serialize(header.coordinatorLogBoundaries, out, version);
ImmutableCoordinatorLogOffsets.serializer.serialize(header.coordinatorLogOffsets, out, version);
}
public StreamMessageHeader deserialize(DataInputPlus in, int version) throws IOException
@ -136,9 +136,9 @@ public class StreamMessageHeader
int sequenceNumber = in.readInt();
long repairedAt = in.readLong();
TimeUUID pendingRepair = in.readBoolean() ? TimeUUID.deserialize(in) : null;
CoordinatorLogBoundaries coordinatorLogBoundaries = CoordinatorLogBoundaries.serializer.deserialize(in, version);
ImmutableCoordinatorLogOffsets coordinatorLogOffsets = ImmutableCoordinatorLogOffsets.serializer.deserialize(in, version);
return new StreamMessageHeader(tableId, sender, planId, sendByFollower, sessionIndex, sequenceNumber, repairedAt, pendingRepair, coordinatorLogBoundaries);
return new StreamMessageHeader(tableId, sender, planId, sendByFollower, sessionIndex, sequenceNumber, repairedAt, pendingRepair, coordinatorLogOffsets);
}
public long serializedSize(StreamMessageHeader header, int version)
@ -152,7 +152,7 @@ public class StreamMessageHeader
size += TypeSizes.sizeof(header.repairedAt);
size += TypeSizes.sizeof(header.pendingRepair != null);
size += header.pendingRepair != null ? TimeUUID.sizeInBytes() : 0;
size += CoordinatorLogBoundaries.serializer.serializedSize(header.coordinatorLogBoundaries, version);
size += ImmutableCoordinatorLogOffsets.serializer.serializedSize(header.coordinatorLogOffsets, version);
return size;
}

View File

@ -204,6 +204,7 @@ public class ClusterMetadataTestHelper
{
KeyspaceAttributes attributes = new KeyspaceAttributes();
attributes.addProperty(KeyspaceParams.Option.REPLICATION.toString(), params.replication.asMap());
attributes.addProperty(KeyspaceParams.Option.REPLICATION_TYPE.toString(), params.replicationType.name());
CreateKeyspaceStatement createKeyspaceStatement = new CreateKeyspaceStatement(name, attributes, false);
try
{

View File

@ -27,6 +27,9 @@ import java.util.Set;
import com.google.common.primitives.Ints;
import org.apache.cassandra.db.Mutation;
import org.apache.cassandra.db.SimpleBuilders;
import org.apache.cassandra.db.partitions.PartitionUpdate;
import org.apache.cassandra.replication.*;
import org.junit.Assert;
import org.junit.Assume;
@ -45,9 +48,12 @@ import org.apache.cassandra.schema.ReplicationType;
import org.apache.cassandra.schema.Schema;
import org.apache.cassandra.schema.TableMetadata;
import org.apache.cassandra.utils.ByteBufferUtil;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
public class MutationTrackingUtils
{
private static final Logger logger = LoggerFactory.getLogger(MutationTrackingUtils.class);
private static final int VERSION = MessagingService.current_version;
public static byte[] encodeId(MutationId id)
@ -298,4 +304,24 @@ public class MutationTrackingUtils
Assert.assertEquals(1, summary.size());
return summary.get(0).logId();
}
public static Mutation createMutation(TableMetadata tableMetadata, int k, int v)
{
DecoratedKey key = tableMetadata.partitioner.decorateKey(ByteBufferUtil.bytes(k));
MutationId mutationId = MutationTrackingService.instance.nextMutationId(tableMetadata.keyspace, key.getToken());
SimpleBuilders.MutationBuilder builder = new SimpleBuilders.MutationBuilder(mutationId, tableMetadata.keyspace, key);
PartitionUpdate.SimpleBuilder partition = builder.update(tableMetadata);
partition.row().add("v", v);
Mutation mutation = builder.build();
Assert.assertFalse(mutation.id().isNone());
return mutation;
}
public static MutationId applyMutation(TableMetadata tableMetadata, int k, int v)
{
Mutation mutation = createMutation(tableMetadata, k, v);
mutation.apply();
logger.debug("Applied mutation {}", mutation.id());
return mutation.id();
}
}

View File

@ -41,7 +41,7 @@ import org.openjdk.jmh.annotations.Warmup;
import org.apache.cassandra.SchemaLoader;
import org.apache.cassandra.db.ColumnFamilyStore;
import org.apache.cassandra.db.CoordinatorLogBoundaries;
import org.apache.cassandra.db.Keyspace;
import org.apache.cassandra.db.RowUpdateBuilder;
import org.apache.cassandra.db.commitlog.CommitLog;
@ -59,6 +59,7 @@ import org.apache.cassandra.io.sstable.format.SSTableReader;
import org.apache.cassandra.locator.InetAddressAndPort;
import org.apache.cassandra.net.AsyncStreamingInputPlus;
import org.apache.cassandra.net.AsyncStreamingOutputPlus;
import org.apache.cassandra.replication.ImmutableCoordinatorLogOffsets;
import org.apache.cassandra.schema.CachingParams;
import org.apache.cassandra.schema.KeyspaceParams;
import org.apache.cassandra.streaming.PreviewKind;
@ -152,7 +153,7 @@ public class ZeroCopyStreamingBench
blockStreamReader = new CassandraEntireSSTableStreamReader(new StreamMessageHeader(sstable.metadata().id,
peer, session.planId(), false,
0, 0, 0,
null, CoordinatorLogBoundaries.NONE), entireSSTableStreamHeader, session);
null, ImmutableCoordinatorLogOffsets.NONE), entireSSTableStreamHeader, session);
List<Range<Token>> requestedRanges = Arrays.asList(new Range<>(sstable.getFirst().minValue().getToken(), sstable.getLast().getToken()));
CassandraStreamHeader partialSSTableStreamHeader =
@ -174,7 +175,7 @@ public class ZeroCopyStreamingBench
partialStreamReader = new CassandraStreamReader(new StreamMessageHeader(sstable.metadata().id,
peer, session.planId(), false,
0, 0, 0,
null, CoordinatorLogBoundaries.NONE),
null, ImmutableCoordinatorLogOffsets.NONE),
partialSSTableStreamHeader, session);
}

View File

@ -1,137 +0,0 @@
/*
* 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;
import org.junit.Test;
import accord.utils.Gen;
import accord.utils.Gens;
import org.apache.cassandra.io.util.DataInputBuffer;
import org.apache.cassandra.io.util.DataOutputBuffer;
import org.apache.cassandra.net.MessagingService;
import org.apache.cassandra.replication.CoordinatorLogId;
import org.apache.cassandra.replication.MutationId;
import org.assertj.core.api.Assertions;
import static accord.utils.Property.qt;
public class CoordinatorLogBoundariesTest
{
private static final Gen<Long> LOG_ID_GEN = rs -> {
int hostId = rs.nextInt(1, 4);
int hostLogId = rs.nextInt(1, 11);
return CoordinatorLogId.asLong(hostId, hostLogId);
};
private static final Gen<Long> SEQUENCE_ID_GEN = rs -> {
int offset = rs.nextBiasedInt(1, 10_000, 1_000_000);
return MutationId.sequenceId(offset, offset);
};
private static final Gen<MutationId> MUTATION_ID_GEN = rs -> new MutationId(LOG_ID_GEN.next(rs), SEQUENCE_ID_GEN.next(rs));
private static final Gen<CoordinatorLogBoundaries> COORDINATOR_LOG_BOUNDARIES_GEN = rs -> {
MutableCoordinatorLogBoundaries boundaries = MutableCoordinatorLogBoundaries.create();
int numIds = rs.nextBiasedInt(0, 10, 1000);
for (int i = 0; i < numIds; i++)
boundaries.add(MUTATION_ID_GEN.next(rs));
return boundaries;
};
@Test
public void roundtripSerde()
{
qt()
.forAll(COORDINATOR_LOG_BOUNDARIES_GEN)
.check(boundaries -> {
try (DataOutputBuffer outputBuffer = DataOutputBuffer.scratchBuffer.get())
{
CoordinatorLogBoundaries.serializer.serialize(boundaries, outputBuffer, MessagingService.current_version);
byte[] bytes = outputBuffer.toByteArray();
try (DataInputBuffer inputBuffer = new DataInputBuffer(bytes))
{
CoordinatorLogBoundaries deserialized = CoordinatorLogBoundaries.serializer.deserialize(inputBuffer, MessagingService.current_version);
Assertions.assertThat(boundaries).isEqualTo(deserialized);
Assertions.assertThat(bytes.length).isEqualTo(CoordinatorLogBoundaries.serializer.serializedSize(boundaries, MessagingService.current_version));
}
}
});
}
@Test
public void monotonicAdd()
{
qt()
.forAll(Gens.lists(MUTATION_ID_GEN).ofSizeBetween(3, 100))
.check(ids -> {
CoordinatorLogBoundariesMap boundaries = new CoordinatorLogBoundariesMap();
for (MutationId id : ids)
{
int originalOffset = boundaries.maxOffset(id.logId());
boundaries.add(id);
int updatedOffset = boundaries.maxOffset(id.logId());
Assertions.assertThat(updatedOffset).isGreaterThanOrEqualTo(originalOffset);
Assertions.assertThat(updatedOffset).isEqualTo(Math.max(originalOffset, id.offset()));
}
});
}
@Test
public void monotonicMerge()
{
qt()
.forAll(COORDINATOR_LOG_BOUNDARIES_GEN, COORDINATOR_LOG_BOUNDARIES_GEN)
.check((left, right) -> {
MutableCoordinatorLogBoundaries boundaries = MutableCoordinatorLogBoundaries.create();
boundaries.addAll(left);
boundaries.addAll(right);
CoordinatorLogBoundaries merged = boundaries;
for (Long logId : merged)
{
int leftOffset = left.maxOffset(logId);
int rightOffset = right.maxOffset(logId);
int mergedOffset = merged.maxOffset(logId);
Assertions.assertThat(mergedOffset).isGreaterThanOrEqualTo(leftOffset);
Assertions.assertThat(mergedOffset).isGreaterThanOrEqualTo(rightOffset);
Assertions.assertThat(mergedOffset).isIn(leftOffset, rightOffset);
}
});
}
@Test
public void builderEquivalentToMutable()
{
qt()
.forAll(Gens.lists(MUTATION_ID_GEN).ofSizeBetween(3, 100))
.check(ids -> {
CoordinatorLogBoundariesMap boundaries = new CoordinatorLogBoundariesMap();
CoordinatorLogBoundariesBuilder builder = new CoordinatorLogBoundariesBuilder();
for (MutationId id : ids)
{
boundaries.add(id);
builder.add(id);
}
CoordinatorLogBoundaries fromBuilder = builder.build();
Assertions.assertThat(fromBuilder).hasSize(boundaries.size());
for (Long logId : boundaries)
Assertions.assertThat(fromBuilder.max(logId)).isEqualTo(boundaries.max(logId));
});
}
}

View File

@ -28,6 +28,7 @@ import org.apache.cassandra.db.memtable.Memtable;
import org.apache.cassandra.db.partitions.PartitionUpdate;
import org.apache.cassandra.exceptions.ConfigurationException;
import org.apache.cassandra.io.sstable.format.SSTableReader;
import org.apache.cassandra.replication.ImmutableCoordinatorLogOffsets;
import org.apache.cassandra.replication.MutationId;
import org.apache.cassandra.replication.MutationJournal;
import org.apache.cassandra.replication.MutationTrackingService;
@ -38,7 +39,6 @@ import org.apache.cassandra.schema.TableMetadata;
import org.apache.cassandra.schema.TableParams;
import org.apache.cassandra.utils.ByteBufferUtil;
import org.assertj.core.api.Assertions;
import org.junit.Assert;
import org.junit.BeforeClass;
import org.junit.Test;
@ -54,7 +54,7 @@ import java.util.List;
import java.util.concurrent.atomic.AtomicInteger;
@RunWith(Parameterized.class)
public class CoordinatorLogBoundariesLifecycleTest
public class CoordinatorLogOffsetsLifecycleTest
{
private static final AtomicInteger keyspaceNumber = new AtomicInteger();
@ -176,9 +176,9 @@ public class CoordinatorLogBoundariesLifecycleTest
assertNumSSTables(view, 0);
Memtable memtable = view.getCurrentMemtable();
CoordinatorLogBoundaries boundaries = memtable.getFlushSet(null, null).coordinatorLogBoundaries();
Assertions.assertThat(boundaries.size()).isEqualTo(1);
Assertions.assertThat(boundaries.max(id2.logId())).isEqualTo(id2);
ImmutableCoordinatorLogOffsets logOffsets = memtable.getFlushSet(null, null).coordinatorLogOffsets();
Assertions.assertThat(logOffsets.size()).isEqualTo(1);
Assertions.assertThat(logOffsets.offsets(id2.logId()).contains(id2.offset())).isTrue();
}
// flush 1
@ -190,9 +190,11 @@ public class CoordinatorLogBoundariesLifecycleTest
assertNumSSTables(view, 1);
SSTableReader sstable = Iterables.getOnlyElement(view.liveSSTables());
CoordinatorLogBoundaries boundaries = sstable.getCoordinatorLogBoundaries();
Assertions.assertThat(boundaries.size()).isEqualTo(1);
Assertions.assertThat(boundaries.max(id2.logId())).isEqualTo(id2);
ImmutableCoordinatorLogOffsets logOffsets = sstable.getCoordinatorLogOffsets();
Assertions.assertThat(logOffsets.size()).isEqualTo(1);
Assertions.assertThat(logOffsets.offsets(id2.logId()).contains(id2.offset())).isTrue();
// Single-participant, so mutations are immediately reconciled once applied
Assertions.assertThat(sstable.isRepaired()).isTrue();
}
MutationId id3;
@ -207,9 +209,9 @@ public class CoordinatorLogBoundariesLifecycleTest
assertNumSSTables(view, 1);
Memtable memtable = view.getCurrentMemtable();
CoordinatorLogBoundaries boundaries = memtable.getFlushSet(null, null).coordinatorLogBoundaries();
Assertions.assertThat(boundaries.size()).isEqualTo(1);
Assertions.assertThat(boundaries.max(id4.logId())).isEqualTo(id4);
ImmutableCoordinatorLogOffsets logOffsets = memtable.getFlushSet(null, null).coordinatorLogOffsets();
Assertions.assertThat(logOffsets.size()).isEqualTo(1);
Assertions.assertThat(logOffsets.offsets(id4.logId()).contains(id4.offset())).isTrue();
}
// flush 2
@ -223,16 +225,18 @@ public class CoordinatorLogBoundariesLifecycleTest
List<SSTableReader> sstables = Lists.newArrayList(view.liveSSTables());
sstables.sort(Comparator.comparing(sst -> sst.descriptor.id.asBytes()));
{
CoordinatorLogBoundaries boundaries = sstables.get(0).getCoordinatorLogBoundaries();
Assertions.assertThat(boundaries.size()).isEqualTo(1);
Assertions.assertThat(boundaries.max(id2.logId())).isEqualTo(id2);
ImmutableCoordinatorLogOffsets logOffsets = sstables.get(0).getCoordinatorLogOffsets();
Assertions.assertThat(logOffsets.size()).isEqualTo(1);
Assertions.assertThat(logOffsets.offsets(id2.logId()).contains(id2.offset())).isTrue();
}
{
CoordinatorLogBoundaries boundaries = sstables.get(1).getCoordinatorLogBoundaries();
Assertions.assertThat(boundaries.size()).isEqualTo(1);
Assertions.assertThat(boundaries.max(id4.logId())).isEqualTo(id4);
ImmutableCoordinatorLogOffsets logOffsets = sstables.get(1).getCoordinatorLogOffsets();
Assertions.assertThat(logOffsets.size()).isEqualTo(1);
Assertions.assertThat(logOffsets.offsets(id4.logId()).contains(id4.offset())).isTrue();
}
for (SSTableReader sstable : sstables)
Assertions.assertThat(sstable.isRepaired()).isTrue();
}
// compaction
@ -244,9 +248,10 @@ public class CoordinatorLogBoundariesLifecycleTest
assertNumSSTables(view, 1);
SSTableReader sstable = Iterables.getOnlyElement(view.liveSSTables());
CoordinatorLogBoundaries boundaries = sstable.getCoordinatorLogBoundaries();
Assertions.assertThat(boundaries.size()).isEqualTo(1);
Assertions.assertThat(boundaries.max(id4.logId())).isEqualTo(id4);
ImmutableCoordinatorLogOffsets logOffsets = sstable.getCoordinatorLogOffsets();
Assertions.assertThat(logOffsets.size()).isEqualTo(1);
Assertions.assertThat(logOffsets.offsets(id4.logId()).contains(id4.offset())).isTrue();
Assertions.assertThat(sstable.isRepaired()).isTrue();
}
}
}

View File

@ -51,6 +51,7 @@ import org.apache.cassandra.io.sstable.format.SSTableWriter;
import org.apache.cassandra.io.sstable.metadata.MetadataCollector;
import org.apache.cassandra.io.util.File;
import org.apache.cassandra.io.util.FileUtils;
import org.apache.cassandra.replication.ImmutableCoordinatorLogOffsets;
import org.apache.cassandra.schema.ColumnMetadata;
import org.apache.cassandra.schema.TableMetadata;
import org.apache.cassandra.schema.TableMetadataRef;
@ -101,7 +102,7 @@ public class SerializationHeaderTest
.setTableMetadataRef(TableMetadataRef.forOfflineTools(schema))
.setKeyCount(1)
.setSerializationHeader(header)
.setCoordinatorLogBoundaries(CoordinatorLogBoundaries.NONE)
.setCoordinatorLogOffsets(ImmutableCoordinatorLogOffsets.NONE)
.setMetadataCollector(new MetadataCollector(schema.comparator))
.addDefaultComponents(Collections.emptySet())
.build(txn, null))

View File

@ -61,6 +61,7 @@ import org.apache.cassandra.locator.InetAddressAndPort;
import org.apache.cassandra.locator.RangesAtEndpoint;
import org.apache.cassandra.locator.Replica;
import org.apache.cassandra.repair.NoSuchRepairSessionException;
import org.apache.cassandra.replication.ImmutableCoordinatorLogOffsets;
import org.apache.cassandra.schema.KeyspaceParams;
import org.apache.cassandra.schema.MockSchema;
import org.apache.cassandra.schema.Schema;
@ -281,7 +282,7 @@ public class AntiCompactionTest
File dir = cfs.getDirectories().getDirectoryForNewSSTables();
Descriptor desc = cfs.newSSTableDescriptor(dir);
try (SSTableTxnWriter writer = SSTableTxnWriter.create(cfs, desc, 0, 0, NO_PENDING_REPAIR, false, CoordinatorLogBoundaries.NONE, new SerializationHeader(true, cfs.metadata(), cfs.metadata().regularAndStaticColumns(), EncodingStats.NO_STATS)))
try (SSTableTxnWriter writer = SSTableTxnWriter.create(cfs, desc, 0, 0, NO_PENDING_REPAIR, false, ImmutableCoordinatorLogOffsets.NONE, new SerializationHeader(true, cfs.metadata(), cfs.metadata().regularAndStaticColumns(), EncodingStats.NO_STATS)))
{
for (int i = 0; i < count; i++)
{

View File

@ -43,7 +43,6 @@ import org.junit.Test;
import org.apache.cassandra.Util;
import org.apache.cassandra.config.DatabaseDescriptor;
import org.apache.cassandra.db.ColumnFamilyStore;
import org.apache.cassandra.db.CoordinatorLogBoundaries;
import org.apache.cassandra.db.DecoratedKey;
import org.apache.cassandra.db.Directories;
import org.apache.cassandra.db.SerializationHeader;
@ -70,6 +69,7 @@ import org.apache.cassandra.io.sstable.metadata.StatsMetadata;
import org.apache.cassandra.io.util.File;
import org.apache.cassandra.io.util.FileHandle;
import org.apache.cassandra.io.util.FileUtils;
import org.apache.cassandra.replication.ImmutableCoordinatorLogOffsets;
import org.apache.cassandra.schema.MockSchema;
import org.apache.cassandra.utils.FilterFactory;
import org.apache.cassandra.utils.Throwables;
@ -1445,7 +1445,7 @@ public class LogTransactionTest extends AbstractTransactionalTest
DecoratedKey key = MockSchema.readerBounds(generation);
SerializationHeader header = SerializationHeader.make(cfs.metadata(), Collections.emptyList());
StatsMetadata metadata = (StatsMetadata) new MetadataCollector(cfs.metadata().comparator)
.finalizeMetadata(cfs.metadata().partitioner.getClass().getCanonicalName(), 0.01f, -1, null, false, CoordinatorLogBoundaries.NONE, header, key.getKey().slice(), key.getKey().slice())
.finalizeMetadata(cfs.metadata().partitioner.getClass().getCanonicalName(), 0.01f, -1, null, false, ImmutableCoordinatorLogOffsets.NONE, header, key.getKey().slice(), key.getKey().slice())
.get(MetadataType.STATS);
SSTableReader reader = new BigTableReader.Builder(descriptor).setComponents(components)
.setTableMetadataRef(cfs.metadata)
@ -1481,7 +1481,7 @@ public class LogTransactionTest extends AbstractTransactionalTest
DecoratedKey key = MockSchema.readerBounds(generation);
SerializationHeader header = SerializationHeader.make(cfs.metadata(), Collections.emptyList());
StatsMetadata metadata = (StatsMetadata) new MetadataCollector(cfs.metadata().comparator)
.finalizeMetadata(cfs.metadata().partitioner.getClass().getCanonicalName(), 0.01f, -1, null, false, CoordinatorLogBoundaries.NONE, header, key.getKey().slice(), key.getKey().slice())
.finalizeMetadata(cfs.metadata().partitioner.getClass().getCanonicalName(), 0.01f, -1, null, false, ImmutableCoordinatorLogOffsets.NONE, header, key.getKey().slice(), key.getKey().slice())
.get(MetadataType.STATS);
SSTableReader reader = new BtiTableReader.Builder(descriptor).setComponents(components)
.setTableMetadataRef(cfs.metadata)

View File

@ -24,7 +24,6 @@ import java.util.List;
import java.util.Set;
import java.util.concurrent.TimeUnit;
import org.apache.cassandra.db.CoordinatorLogBoundaries;
import org.junit.Assert;
import org.junit.BeforeClass;
import org.junit.Test;
@ -44,6 +43,7 @@ import org.apache.cassandra.io.sstable.SSTableRewriter;
import org.apache.cassandra.io.sstable.format.SSTableReader;
import org.apache.cassandra.io.sstable.metadata.MetadataCollector;
import org.apache.cassandra.io.util.File;
import org.apache.cassandra.replication.ImmutableCoordinatorLogOffsets;
import org.apache.cassandra.schema.KeyspaceParams;
import org.apache.cassandra.schema.Schema;
import org.apache.cassandra.schema.TableMetadataRef;
@ -168,7 +168,7 @@ public class RealTransactionsTest extends SchemaLoader
.setSecondaryIndexGroups(cfs.indexManager.listIndexGroups())
.setMetadataCollector(new MetadataCollector(cfs.metadata().comparator))
.addDefaultComponents(cfs.indexManager.listIndexGroups())
.setCoordinatorLogBoundaries(CoordinatorLogBoundaries.NONE)
.setCoordinatorLogOffsets(ImmutableCoordinatorLogOffsets.NONE)
.build(txn, cfs));
while (ci.hasNext())
{

View File

@ -24,7 +24,6 @@ import java.util.Collection;
import java.util.Collections;
import java.util.Queue;
import org.apache.cassandra.db.CoordinatorLogBoundaries;
import org.junit.BeforeClass;
import org.junit.Test;
@ -42,6 +41,7 @@ import org.apache.cassandra.io.util.DataInputBuffer;
import org.apache.cassandra.locator.InetAddressAndPort;
import org.apache.cassandra.net.AsyncStreamingOutputPlus;
import org.apache.cassandra.net.SharedDefaultFileRegion;
import org.apache.cassandra.replication.ImmutableCoordinatorLogOffsets;
import org.apache.cassandra.schema.CachingParams;
import org.apache.cassandra.schema.KeyspaceParams;
import org.apache.cassandra.streaming.PreviewKind;
@ -164,7 +164,7 @@ public class CassandraEntireSSTableStreamWriterTest
.withTableId(sstable.metadata().id)
.build();
CassandraEntireSSTableStreamReader reader = new CassandraEntireSSTableStreamReader(new StreamMessageHeader(sstable.metadata().id, peer, session.planId(), false, 0, 0, 0, null, CoordinatorLogBoundaries.NONE), header, session);
CassandraEntireSSTableStreamReader reader = new CassandraEntireSSTableStreamReader(new StreamMessageHeader(sstable.metadata().id, peer, session.planId(), false, 0, 0, 0, null, ImmutableCoordinatorLogOffsets.NONE), header, session);
SSTableTxnSingleStreamWriter sstableWriter = (SSTableTxnSingleStreamWriter) reader.read(new DataInputBuffer(serializedFile.nioBuffer(), false));
StreamingLifecycleTransaction stt = new StreamingLifecycleTransaction();

View File

@ -44,7 +44,6 @@ import org.junit.runner.RunWith;
import org.apache.cassandra.SchemaLoader;
import org.apache.cassandra.Util;
import org.apache.cassandra.db.ColumnFamilyStore;
import org.apache.cassandra.db.CoordinatorLogBoundaries;
import org.apache.cassandra.db.Keyspace;
import org.apache.cassandra.db.RowUpdateBuilder;
import org.apache.cassandra.db.compaction.CompactionManager;
@ -67,6 +66,7 @@ import org.apache.cassandra.net.AsyncStreamingOutputPlus;
import org.apache.cassandra.net.BufferPoolAllocator;
import org.apache.cassandra.net.MessagingService;
import org.apache.cassandra.net.SharedDefaultFileRegion;
import org.apache.cassandra.replication.ImmutableCoordinatorLogOffsets;
import org.apache.cassandra.schema.KeyspaceParams;
import org.apache.cassandra.schema.SchemaTestUtil;
import org.apache.cassandra.schema.TableMetadata;
@ -235,7 +235,7 @@ public class EntireSSTableStreamConcurrentComponentMutationTest
concurrentMutations.get(3, TimeUnit.MINUTES);
session.prepareReceiving(new StreamSummary(sstable.metadata().id, emptyList(), 1, 5104));
StreamMessageHeader messageHeader = new StreamMessageHeader(sstable.metadata().id, peer, session.planId(), false, 0, 0, 0, null, CoordinatorLogBoundaries.NONE);
StreamMessageHeader messageHeader = new StreamMessageHeader(sstable.metadata().id, peer, session.planId(), false, 0, 0, 0, null, ImmutableCoordinatorLogOffsets.NONE);
try (DataInputBuffer in = new DataInputBuffer(serializedFile.nioBuffer(), false))
{

View File

@ -20,7 +20,6 @@ package org.apache.cassandra.io.sstable;
import java.io.IOException;
import org.apache.cassandra.db.CoordinatorLogBoundaries;
import org.junit.BeforeClass;
import org.junit.Test;
@ -34,6 +33,7 @@ import org.apache.cassandra.db.SerializationHeader;
import org.apache.cassandra.db.compaction.OperationType;
import org.apache.cassandra.db.lifecycle.LifecycleTransaction;
import org.apache.cassandra.dht.Murmur3Partitioner;
import org.apache.cassandra.replication.ImmutableCoordinatorLogOffsets;
import org.apache.cassandra.schema.KeyspaceParams;
import static org.junit.Assert.assertEquals;
@ -76,7 +76,7 @@ public class RangeAwareSSTableWriterTest
0,
null,
false,
CoordinatorLogBoundaries.NONE,
ImmutableCoordinatorLogOffsets.NONE,
DatabaseDescriptor.getSelectedSSTableFormat(),
0,
0,

View File

@ -39,7 +39,6 @@ import org.junit.Test;
import org.apache.cassandra.config.DatabaseDescriptor;
import org.apache.cassandra.db.Clustering;
import org.apache.cassandra.db.CoordinatorLogBoundaries;
import org.apache.cassandra.db.DecoratedKey;
import org.apache.cassandra.db.DeletionTime;
import org.apache.cassandra.db.SerializationHeader;
@ -64,6 +63,7 @@ import org.apache.cassandra.io.sstable.format.SSTableWriter;
import org.apache.cassandra.io.sstable.metadata.MetadataCollector;
import org.apache.cassandra.io.util.File;
import org.apache.cassandra.io.util.FileUtils;
import org.apache.cassandra.replication.ImmutableCoordinatorLogOffsets;
import org.apache.cassandra.schema.ColumnMetadata;
import org.apache.cassandra.schema.TableMetadata;
import org.apache.cassandra.schema.TableMetadataRef;
@ -135,7 +135,7 @@ public class SSTableFlushObserverTest
.setSerializationHeader(new SerializationHeader(true, cfm, cfm.regularAndStaticColumns(), EncodingStats.NO_STATS))
.setSecondaryIndexGroups(Collections.singleton(indexGroup))
.addDefaultComponents(Collections.emptySet())
.setCoordinatorLogBoundaries(CoordinatorLogBoundaries.NONE);
.setCoordinatorLogOffsets(ImmutableCoordinatorLogOffsets.NONE);
assertThat(observer.beginCalled).isFalse();
assertThat(observer.isComplete).isFalse();
@ -220,7 +220,7 @@ public class SSTableFlushObserverTest
.setSerializationHeader(new SerializationHeader(true, cfm, cfm.regularAndStaticColumns(), EncodingStats.NO_STATS))
.setSecondaryIndexGroups(List.of(indexGroup1, indexGroup2))
.addDefaultComponents(Collections.emptySet())
.setCoordinatorLogBoundaries(CoordinatorLogBoundaries.NONE)
.setCoordinatorLogOffsets(ImmutableCoordinatorLogOffsets.NONE)
.build(transaction, null)
).withMessage("Failed to initialize");

View File

@ -37,7 +37,6 @@ import org.apache.cassandra.UpdateBuilder;
import org.apache.cassandra.Util;
import org.apache.cassandra.concurrent.NamedThreadFactory;
import org.apache.cassandra.db.ColumnFamilyStore;
import org.apache.cassandra.db.CoordinatorLogBoundaries;
import org.apache.cassandra.db.DecoratedKey;
import org.apache.cassandra.db.DeletionTime;
import org.apache.cassandra.db.Keyspace;
@ -62,6 +61,7 @@ import org.apache.cassandra.io.sstable.format.SSTableWriter;
import org.apache.cassandra.io.util.File;
import org.apache.cassandra.io.util.FileUtils;
import org.apache.cassandra.metrics.StorageMetrics;
import org.apache.cassandra.replication.ImmutableCoordinatorLogOffsets;
import org.apache.cassandra.utils.ByteBufferUtil;
import org.apache.cassandra.utils.FBUtilities;
@ -940,7 +940,7 @@ public class SSTableRewriterTest extends SSTableWriterTestBase
File dir = cfs.getDirectories().getDirectoryForNewSSTables();
Descriptor desc = cfs.newSSTableDescriptor(dir);
try (SSTableTxnWriter writer = SSTableTxnWriter.create(cfs, desc, 0, 0, null, false, CoordinatorLogBoundaries.NONE, new SerializationHeader(true, cfs.metadata(), cfs.metadata().regularAndStaticColumns(), EncodingStats.NO_STATS)))
try (SSTableTxnWriter writer = SSTableTxnWriter.create(cfs, desc, 0, 0, null, false, ImmutableCoordinatorLogOffsets.NONE, new SerializationHeader(true, cfs.metadata(), cfs.metadata().regularAndStaticColumns(), EncodingStats.NO_STATS)))
{
int end = f == fileCount - 1 ? partitionCount : ((f + 1) * partitionCount) / fileCount;
for ( ; i < end ; i++)

View File

@ -44,6 +44,7 @@ import org.apache.cassandra.io.sstable.format.SSTableFormat.Components;
import org.apache.cassandra.io.sstable.format.SSTableReader;
import org.apache.cassandra.io.util.File;
import org.apache.cassandra.io.util.FileUtils;
import org.apache.cassandra.replication.ImmutableCoordinatorLogOffsets;
import org.apache.cassandra.schema.Schema;
import org.apache.cassandra.schema.TableMetadata;
@ -235,7 +236,7 @@ public class SSTableUtils
if (cfs.metadata().replicationType().isTracked())
throw new IllegalStateException("Can't create writer for table with mutation tracking enabled");
SerializationHeader header = appender.header();
SSTableTxnWriter writer = SSTableTxnWriter.create(cfs, Descriptor.fromFileWithComponent(datafile, false).left, expectedSize, UNREPAIRED_SSTABLE, NO_PENDING_REPAIR, false, CoordinatorLogBoundaries.NONE, header);
SSTableTxnWriter writer = SSTableTxnWriter.create(cfs, Descriptor.fromFileWithComponent(datafile, false).left, expectedSize, UNREPAIRED_SSTABLE, NO_PENDING_REPAIR, false, ImmutableCoordinatorLogOffsets.NONE, header);
while (appender.append(writer)) { /* pass */ }
Collection<SSTableReader> readers = writer.finish(true);

View File

@ -35,7 +35,6 @@ import org.apache.cassandra.ServerTestUtils;
import org.apache.cassandra.config.Config;
import org.apache.cassandra.config.DatabaseDescriptor;
import org.apache.cassandra.db.ColumnFamilyStore;
import org.apache.cassandra.db.CoordinatorLogBoundaries;
import org.apache.cassandra.db.Keyspace;
import org.apache.cassandra.db.SerializationHeader;
import org.apache.cassandra.db.compaction.CompactionManager;
@ -46,6 +45,7 @@ import org.apache.cassandra.io.sstable.format.SSTableReader;
import org.apache.cassandra.io.sstable.format.SSTableWriter;
import org.apache.cassandra.io.sstable.metadata.MetadataCollector;
import org.apache.cassandra.io.util.File;
import org.apache.cassandra.replication.ImmutableCoordinatorLogOffsets;
import org.apache.cassandra.schema.KeyspaceParams;
import org.apache.cassandra.utils.TimeUUID;
@ -169,7 +169,7 @@ public class SSTableWriterTestBase extends SchemaLoader
.setSecondaryIndexGroups(cfs.indexManager.listIndexGroups())
.setMetadataCollector(new MetadataCollector(cfs.metadata().comparator))
.addDefaultComponents(cfs.indexManager.listIndexGroups())
.setCoordinatorLogBoundaries(CoordinatorLogBoundaries.NONE)
.setCoordinatorLogOffsets(ImmutableCoordinatorLogOffsets.NONE)
.build(txn, cfs);
}

View File

@ -21,7 +21,6 @@ package org.apache.cassandra.io.sstable;
import java.io.IOException;
import java.util.Collection;
import org.apache.cassandra.db.CoordinatorLogBoundaries;
import org.junit.Assert;
import org.junit.BeforeClass;
@ -35,6 +34,7 @@ import org.apache.cassandra.db.marshal.Int32Type;
import org.apache.cassandra.db.rows.EncodingStats;
import org.apache.cassandra.io.sstable.format.SSTableFormat.Components;
import org.apache.cassandra.io.util.File;
import org.apache.cassandra.replication.ImmutableCoordinatorLogOffsets;
import org.apache.cassandra.schema.KeyspaceParams;
import org.apache.cassandra.utils.concurrent.AbstractTransactionalTest;
@ -73,7 +73,7 @@ public class SSTableWriterTransactionTest extends AbstractTransactionalTest
private TestableBTW(Descriptor desc)
{
this(desc, SSTableTxnWriter.create(cfs, desc, 0, 0, null, false, CoordinatorLogBoundaries.NONE,
this(desc, SSTableTxnWriter.create(cfs, desc, 0, 0, null, false, ImmutableCoordinatorLogOffsets.NONE,
new SerializationHeader(true, cfs.metadata(),
cfs.metadata().regularAndStaticColumns(),
EncodingStats.NO_STATS)));

View File

@ -64,7 +64,6 @@ import org.apache.cassandra.cql3.QueryProcessor;
import org.apache.cassandra.cql3.UntypedResultSet;
import org.apache.cassandra.db.ColumnFamilyStore;
import org.apache.cassandra.db.ConsistencyLevel;
import org.apache.cassandra.db.CoordinatorLogBoundaries;
import org.apache.cassandra.db.CounterMutation;
import org.apache.cassandra.db.DecoratedKey;
import org.apache.cassandra.db.Keyspace;
@ -98,6 +97,7 @@ import org.apache.cassandra.io.sstable.format.bti.BtiFormat;
import org.apache.cassandra.io.sstable.metadata.MetadataCollector;
import org.apache.cassandra.io.util.File;
import org.apache.cassandra.io.util.FileUtils;
import org.apache.cassandra.replication.ImmutableCoordinatorLogOffsets;
import org.apache.cassandra.schema.KeyspaceParams;
import org.apache.cassandra.utils.ByteBufferUtil;
import org.apache.cassandra.utils.OutputHandler;
@ -843,7 +843,7 @@ public class ScrubTest
.setMetadataCollector(collector)
.setSerializationHeader(header)
.addDefaultComponents(Collections.emptySet())
.setCoordinatorLogBoundaries(CoordinatorLogBoundaries.NONE)
.setCoordinatorLogOffsets(ImmutableCoordinatorLogOffsets.NONE)
.build(txn, cfs);
return new TestMultiWriter(writer, txn);

View File

@ -33,7 +33,6 @@ import org.slf4j.LoggerFactory;
import org.apache.cassandra.config.DatabaseDescriptor;
import org.apache.cassandra.db.Clustering;
import org.apache.cassandra.db.MutableCoordinatorLogBoundaries;
import org.apache.cassandra.db.SerializationHeader;
import org.apache.cassandra.db.commitlog.CommitLogPosition;
import org.apache.cassandra.db.commitlog.IntervalSet;
@ -51,6 +50,7 @@ import org.apache.cassandra.io.util.File;
import org.apache.cassandra.io.util.FileOutputStreamPlus;
import org.apache.cassandra.io.util.FileUtils;
import org.apache.cassandra.io.util.RandomAccessReader;
import org.apache.cassandra.replication.ImmutableCoordinatorLogOffsets;
import org.apache.cassandra.replication.MutationId;
import org.apache.cassandra.schema.TableMetadata;
import org.apache.cassandra.utils.Throwables;
@ -156,10 +156,10 @@ public class MetadataSerializerTest
collector.updateClusteringValues(Clustering.make(UTF8Type.instance.decompose("cba"), withNulls ? null : Int32Type.instance.decompose(234)));
ByteBuffer first = AsciiType.instance.decompose("a");
ByteBuffer last = AsciiType.instance.decompose("b");
MutableCoordinatorLogBoundaries boundaries = MutableCoordinatorLogBoundaries.create();
boundaries.add(new MutationId(1, 12345));
boundaries.add(new MutationId(2, 56789));
return collector.finalizeMetadata(partitioner, bfFpChance, 0, null, false, boundaries, SerializationHeader.make(cfm, Collections.emptyList()), first, last);
ImmutableCoordinatorLogOffsets.Builder logOffsetsBuilder = new ImmutableCoordinatorLogOffsets.Builder();
logOffsetsBuilder.add(new MutationId(1, 12345));
logOffsetsBuilder.add(new MutationId(2, 56789));
return collector.finalizeMetadata(partitioner, bfFpChance, 0, null, false, logOffsetsBuilder.build(), SerializationHeader.make(cfm, Collections.emptyList()), first, last);
}
private void testVersions(List<String> versions) throws Throwable

View File

@ -0,0 +1,334 @@
/*
* 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.replication;
import org.apache.cassandra.config.DatabaseDescriptor;
import org.apache.cassandra.db.Mutation;
import org.apache.cassandra.db.marshal.Int32Type;
import org.apache.cassandra.dht.Murmur3Partitioner;
import org.apache.cassandra.distributed.test.log.ClusterMetadataTestHelper;
import org.apache.cassandra.distributed.test.tracking.MutationTrackingUtils;
import org.apache.cassandra.locator.InetAddressAndPort;
import org.apache.cassandra.schema.DistributedSchema;
import org.apache.cassandra.schema.KeyspaceParams;
import org.apache.cassandra.schema.ReplicationType;
import org.apache.cassandra.schema.SchemaTransformations;
import org.apache.cassandra.schema.TableMetadata;
import org.apache.cassandra.tcm.ClusterMetadata;
import org.apache.cassandra.tcm.ClusterMetadataService;
import org.apache.cassandra.tcm.StubClusterMetadataService;
import org.apache.cassandra.tcm.membership.Directory;
import org.apache.cassandra.tcm.membership.Location;
import org.apache.cassandra.tcm.membership.NodeId;
import org.apache.cassandra.tcm.transformations.AlterSchema;
import org.junit.BeforeClass;
import org.junit.Test;
import accord.utils.Gen;
import accord.utils.Gens;
import org.apache.cassandra.io.util.DataInputBuffer;
import org.apache.cassandra.io.util.DataOutputBuffer;
import org.apache.cassandra.net.MessagingService;
import org.assertj.core.api.Assertions;
import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.Future;
import java.util.function.Supplier;
import com.google.common.collect.Sets;
import static accord.utils.Property.qt;
public class CoordinatorLogOffsetsTest
{
private static final Gen<Long> LOG_ID_GEN = rs -> {
int hostId = rs.nextInt(1, 4);
int hostLogId = rs.nextInt(1, 11);
return CoordinatorLogId.asLong(hostId, hostLogId);
};
private static final Gen<Long> SEQUENCE_ID_GEN = rs -> {
int offset = rs.nextBiasedInt(1, 10_000, 1_000_000);
return MutationId.sequenceId(offset, offset);
};
private static final Gen<MutationId> MUTATION_ID_GEN = rs -> new MutationId(LOG_ID_GEN.next(rs), SEQUENCE_ID_GEN.next(rs));
private static final Gen<ImmutableCoordinatorLogOffsets> COORDINATOR_LOG_OFFSETS_GEN = rs -> {
int numIds = rs.nextBiasedInt(0, 10, 1000);
ImmutableCoordinatorLogOffsets.Builder builder = new ImmutableCoordinatorLogOffsets.Builder(numIds);
for (int i = 0; i < numIds; i++)
builder.add(MUTATION_ID_GEN.next(rs));
return builder.build();
};
@BeforeClass
public static void beforeClass()
{
DatabaseDescriptor.daemonInitialization();
}
@Test
public void roundtripSerde()
{
qt()
.forAll(COORDINATOR_LOG_OFFSETS_GEN)
.check(offsets -> {
try (DataOutputBuffer outputBuffer = DataOutputBuffer.scratchBuffer.get())
{
ImmutableCoordinatorLogOffsets.serializer.serialize(offsets, outputBuffer, MessagingService.current_version);
byte[] bytes = outputBuffer.toByteArray();
try (DataInputBuffer inputBuffer = new DataInputBuffer(bytes))
{
ImmutableCoordinatorLogOffsets deserialized = ImmutableCoordinatorLogOffsets.serializer.deserialize(inputBuffer, MessagingService.current_version);
Assertions.assertThat(Sets.newHashSet(deserialized.iterator())).isEqualTo(Sets.newHashSet(offsets.iterator()));
for (long logId : offsets)
Assertions.assertThat(deserialized.offsets(logId)).isEqualTo(offsets.offsets(logId));
Assertions.assertThat(bytes.length).isEqualTo(ImmutableCoordinatorLogOffsets.serializer.serializedSize(offsets, MessagingService.current_version));
}
}
});
}
@Test
public void monotonicAdd()
{
monotonicAdd(() -> new NonBlockingCoordinatorLogOffsets.Concurrent(16));
monotonicAdd(NonBlockingCoordinatorLogOffsets.Exclusive::new);
}
public void monotonicAdd(Supplier<MutableCoordinatorLogOffsets> ctor)
{
qt()
.forAll(Gens.lists(MUTATION_ID_GEN).ofSizeBetween(3, 100))
.check(ids -> {
MutableCoordinatorLogOffsets logOffsets = ctor.get();
for (MutationId id : ids)
{
Offsets originalOffsets = logOffsets.offsets(id.logId());
boolean existed = originalOffsets.contains(id.offset());
logOffsets.add(id);
Offsets updatedOffsets = logOffsets.offsets(id.logId());
if (existed)
Assertions.assertThat(updatedOffsets).hasSameSizeAs(originalOffsets);
Assertions.assertThat(updatedOffsets.contains(id.offset())).isTrue();
}
});
}
@Test
public void monotonicMerge()
{
monotonicMerge(() -> new NonBlockingCoordinatorLogOffsets.Concurrent(16));
monotonicMerge(NonBlockingCoordinatorLogOffsets.Exclusive::new);
}
public void monotonicMerge(Supplier<MutableCoordinatorLogOffsets> ctor)
{
qt()
.forAll(COORDINATOR_LOG_OFFSETS_GEN, COORDINATOR_LOG_OFFSETS_GEN)
.check((left, right) -> {
MutableCoordinatorLogOffsets merged = ctor.get();
merged.addAll(left);
merged.addAll(right);
for (Long logId : merged)
{
Offsets leftOffsets = left.offsets(logId);
Offsets rightOffsets = right.offsets(logId);
Offsets mergedOffsets = Offsets.Immutable.copy(merged.offsets(logId));
Assertions.assertThat(mergedOffsets).isEqualTo(Offsets.Immutable.union(leftOffsets, rightOffsets));
}
});
}
@Test
public void builderEquivalentToMutable()
{
builderEquivalentToMutable(() -> new NonBlockingCoordinatorLogOffsets.Concurrent(16));
builderEquivalentToMutable(NonBlockingCoordinatorLogOffsets.Exclusive::new);
}
public void builderEquivalentToMutable(Supplier<MutableCoordinatorLogOffsets> ctor)
{
qt()
.forAll(Gens.lists(MUTATION_ID_GEN).ofSizeBetween(3, 100))
.check(ids -> {
MutableCoordinatorLogOffsets logOffsets = ctor.get();
ImmutableCoordinatorLogOffsets.Builder builder = new ImmutableCoordinatorLogOffsets.Builder(ids.size());
for (MutationId id : ids)
{
logOffsets.add(id);
builder.add(id);
}
ImmutableCoordinatorLogOffsets fromBuilder = builder.build();
Assertions.assertThat(fromBuilder).hasSize(logOffsets.size());
for (Long logId : logOffsets)
Assertions.assertThat(fromBuilder.offsets(logId)).isEqualTo(Offsets.Immutable.copy(logOffsets.offsets(logId)));
});
}
@Test
public void mutableImplsEquivalent()
{
class Args
{
public final List<MutationId> ids;
public final int contentions;
public Args(List<MutationId> ids, int contentions)
{
this.ids = ids;
this.contentions = contentions;
}
}
Gen<Args> argsGen = rs -> {
List<MutationId> ids = Gens.lists(MUTATION_ID_GEN).ofSizeBetween(3, 100).next(rs);
int contentions = rs.nextInt(1, 16);
return new Args(ids, contentions);
};
qt()
.forAll(argsGen)
.check(args -> {
NonBlockingCoordinatorLogOffsets.Exclusive exclusive = new NonBlockingCoordinatorLogOffsets.Exclusive();
NonBlockingCoordinatorLogOffsets.Concurrent concurrent = new NonBlockingCoordinatorLogOffsets.Concurrent(args.contentions);
ExecutorService executor = Executors.newFixedThreadPool(args.contentions);
List<Future<?>> concurrentUpdates = new ArrayList<>();
for (MutationId id : args.ids)
{
concurrentUpdates.add(executor.submit(() -> concurrent.add(id)));
exclusive.add(id);
}
for (Future<?> task : concurrentUpdates)
task.get();
Assertions.assertThatIterable(exclusive).hasSameSizeAs(concurrent);
for (Long logId : exclusive)
Assertions.assertThat(exclusive.offsets(logId)).isEqualTo(concurrent.offsets(logId));
});
}
@Test
public void reconciledBounds() throws InterruptedException, ExecutionException {
DatabaseDescriptor.daemonInitialization();
DatabaseDescriptor.setPartitionerUnsafe(Murmur3Partitioner.instance);
MutationJournal.instance.start();
String ks = "ks";
String tbl = "tbl";
TableMetadata tableMetadata = TableMetadata.builder(ks, tbl)
.addPartitionKeyColumn("k", Int32Type.instance)
.addRegularColumn("v", Int32Type.instance)
.build();
InetAddressAndPort addr1 = InetAddressAndPort.getByNameUnchecked("127.0.0.1");
InetAddressAndPort addr2 = InetAddressAndPort.getByNameUnchecked("127.0.0.2");
InetAddressAndPort addr3 = InetAddressAndPort.getByNameUnchecked("127.0.0.3");
Location location = new Location("dc1", "rack1");
ClusterMetadata metadata = new ClusterMetadata(Murmur3Partitioner.instance, Directory.EMPTY, DistributedSchema.empty());
ClusterMetadataService.unsetInstance();
ClusterMetadataService.setInstance(StubClusterMetadataService.forTesting());
// RF=3, all instances are replicas
ClusterMetadataTestHelper.addEndpoint(addr1, new Murmur3Partitioner.LongToken(1), location);
ClusterMetadataTestHelper.addEndpoint(addr2, new Murmur3Partitioner.LongToken(2), location);
ClusterMetadataTestHelper.addEndpoint(addr3, new Murmur3Partitioner.LongToken(3), location);
ClusterMetadataTestHelper.createKeyspace(ks, KeyspaceParams.simple(3, ReplicationType.tracked));
ClusterMetadataTestHelper.commit(new AlterSchema(SchemaTransformations.addTable(tableMetadata, false)));
MutationTrackingService.instance.start(metadata);
// Eventually, will also run perturbations before checking isReconciled (like log truncation, durability, etc.)
// to ensure that we don't prune data required to check what's been reconciled
// Applied at all replicas
{
Mutation mutation = MutationTrackingUtils.createMutation(tableMetadata, 1, 1);
MutationTrackingService.instance.startWriting(mutation);
MutationTrackingService.instance.finishWriting(mutation);
MutationTrackingService.instance.receivedWriteResponse(ks, mutation.key().getToken(), mutation.id(), addr2);
MutationTrackingService.instance.receivedWriteResponse(ks, mutation.key().getToken(), mutation.id(), addr3);
ImmutableCoordinatorLogOffsets logOffsets = new ImmutableCoordinatorLogOffsets.Builder()
.add(mutation.id())
.build();
Assertions.assertThat(MutationTrackingService.instance.isDurablyReconciled(ks, logOffsets)).isTrue();
}
// Applied locally but not on remote replicas
{
Mutation mutation = MutationTrackingUtils.createMutation(tableMetadata, 2, 2);
MutationTrackingService.instance.startWriting(mutation);
MutationTrackingService.instance.finishWriting(mutation);
ImmutableCoordinatorLogOffsets logOffsets = new ImmutableCoordinatorLogOffsets.Builder()
.add(mutation.id())
.build();
Assertions.assertThat(MutationTrackingService.instance.isDurablyReconciled(ks, logOffsets)).isFalse();
}
// Applied on remote replicas but not locally
{
Mutation mutation = MutationTrackingUtils.createMutation(tableMetadata, 3, 3);
MutationTrackingService.instance.startWriting(mutation);
MutationTrackingService.instance.receivedWriteResponse(ks, mutation.key().getToken(), mutation.id(), addr2);
MutationTrackingService.instance.receivedWriteResponse(ks, mutation.key().getToken(), mutation.id(), addr3);
ImmutableCoordinatorLogOffsets logOffsets = new ImmutableCoordinatorLogOffsets.Builder()
.add(mutation.id())
.build();
Assertions.assertThat(MutationTrackingService.instance.isDurablyReconciled(ks, logOffsets)).isFalse();
}
// If no replicas are aware of a log, it should be considered unreconciled out of caution
{
Mutation mutation = MutationTrackingUtils.createMutation(tableMetadata, 4, 4);
MutationTrackingService.instance.startWriting(mutation);
MutationTrackingService.instance.finishWriting(mutation);
MutationTrackingService.instance.receivedWriteResponse(ks, mutation.key().getToken(), mutation.id(), addr2);
MutationTrackingService.instance.receivedWriteResponse(ks, mutation.key().getToken(), mutation.id(), addr3);
MutationId fakeMutationId = new MutationId(CoordinatorLogId.asLong(111, 222), MutationId.sequenceId(333, 444));
Assertions.assertThat(metadata.directory.version(new NodeId(fakeMutationId.hostId()))).isNull();
Offsets.Immutable.Builder offsetsBuilder = new Offsets.Immutable.Builder(new CoordinatorLogId(fakeMutationId.logId()));
offsetsBuilder.add(fakeMutationId.offset());
ImmutableCoordinatorLogOffsets.Builder logOffsetsBuilder = new ImmutableCoordinatorLogOffsets.Builder();
logOffsetsBuilder.add(fakeMutationId);
Assertions.assertThat(MutationTrackingService.instance.isDurablyReconciled(ks, logOffsetsBuilder.build())).isFalse();
}
MutationTrackingService.instance.shutdownBlocking();
}
}

View File

@ -41,7 +41,6 @@ import org.apache.cassandra.Util;
import org.apache.cassandra.config.DatabaseDescriptor;
import org.apache.cassandra.db.BufferDecoratedKey;
import org.apache.cassandra.db.ColumnFamilyStore;
import org.apache.cassandra.db.CoordinatorLogBoundaries;
import org.apache.cassandra.db.DeletionTime;
import org.apache.cassandra.db.Directories;
import org.apache.cassandra.db.Keyspace;
@ -73,6 +72,7 @@ import org.apache.cassandra.io.util.File;
import org.apache.cassandra.io.util.FileHandle;
import org.apache.cassandra.io.util.FileUtils;
import org.apache.cassandra.io.util.Memory;
import org.apache.cassandra.replication.ImmutableCoordinatorLogOffsets;
import org.apache.cassandra.service.CacheService;
import org.apache.cassandra.tcm.ClusterMetadata;
import org.apache.cassandra.utils.ByteBufferUtil;
@ -215,14 +215,14 @@ public class MockSchema
BufferDecoratedKey last = readerBounds(lastToken);
Map<MetadataType, MetadataComponent> metadataComponents = collector.sstableLevel(level)
.finalizeMetadata(cfs.metadata().partitioner.getClass().getCanonicalName(),
0.01f,
UNREPAIRED_SSTABLE,
null,
false,
CoordinatorLogBoundaries.NONE,
header,
first.retainable().getKey().slice(),
last.retainable().getKey().slice());
0.01f,
UNREPAIRED_SSTABLE,
null,
false,
ImmutableCoordinatorLogOffsets.NONE,
header,
first.retainable().getKey().slice(),
last.retainable().getKey().slice());
StatsMetadata statsMetadata = (StatsMetadata) metadataComponents.get(MetadataType.STATS);
try (DataOutputStreamPlus out = descriptor.fileFor(Components.STATS).newOutputStream(File.WriteMode.OVERWRITE))
{
@ -271,7 +271,7 @@ public class MockSchema
BufferDecoratedKey first = readerBounds(firstToken);
BufferDecoratedKey last = readerBounds(lastToken);
StatsMetadata metadata = (StatsMetadata) collector.sstableLevel(level)
.finalizeMetadata(cfs.metadata().partitioner.getClass().getCanonicalName(), 0.01f, UNREPAIRED_SSTABLE, null, false, CoordinatorLogBoundaries.NONE, header, first.retainable().getKey(), last.retainable().getKey())
.finalizeMetadata(cfs.metadata().partitioner.getClass().getCanonicalName(), 0.01f, UNREPAIRED_SSTABLE, null, false, ImmutableCoordinatorLogOffsets.NONE, header, first.retainable().getKey(), last.retainable().getKey())
.get(MetadataType.STATS);
BtiTableReader reader = new BtiTableReader.Builder(descriptor).setComponents(components)
.setTableMetadataRef(cfs.metadata)

View File

@ -38,7 +38,6 @@ import org.apache.cassandra.ServerTestUtils;
import org.apache.cassandra.config.DatabaseDescriptor;
import org.apache.cassandra.db.BufferDecoratedKey;
import org.apache.cassandra.db.ColumnFamilyStore;
import org.apache.cassandra.db.CoordinatorLogBoundaries;
import org.apache.cassandra.db.DecoratedKey;
import org.apache.cassandra.db.DeletionTime;
import org.apache.cassandra.db.Keyspace;
@ -66,6 +65,7 @@ import org.apache.cassandra.locator.InetAddressAndPort;
import org.apache.cassandra.metrics.StorageMetrics;
import org.apache.cassandra.net.AsyncStreamingOutputPlus;
import org.apache.cassandra.net.MessagingService;
import org.apache.cassandra.replication.ImmutableCoordinatorLogOffsets;
import org.apache.cassandra.schema.TableMetadata;
import org.apache.cassandra.streaming.async.StreamCompressionSerializer;
import org.apache.cassandra.streaming.messages.StreamMessageHeader;
@ -466,7 +466,7 @@ public class StreamReaderTest
fakeSeq,
System.currentTimeMillis(),
pendingRepair,
CoordinatorLogBoundaries.NONE);
ImmutableCoordinatorLogOffsets.NONE);
}
private static CassandraStreamHeader streamMessageHeader(int...tokens)
@ -505,9 +505,9 @@ public class StreamReaderTest
super(header, streamHeader, session);
}
protected SSTableTxnSingleStreamWriter createWriter(ColumnFamilyStore cfs, long totalSize, long repairedAt, TimeUUID pendingRepair, CoordinatorLogBoundaries coordinatorLogBoundaries, SSTableFormat<?,?> format) throws IOException
protected SSTableTxnSingleStreamWriter createWriter(ColumnFamilyStore cfs, long totalSize, long repairedAt, TimeUUID pendingRepair, ImmutableCoordinatorLogOffsets coordinatorLogOffsets, SSTableFormat<?,?> format) throws IOException
{
return super.createWriter(cfs, totalSize, repairedAt, pendingRepair, coordinatorLogBoundaries, format);
return super.createWriter(cfs, totalSize, repairedAt, pendingRepair, coordinatorLogOffsets, format);
}
@Override

View File

@ -28,13 +28,13 @@ import org.junit.BeforeClass;
import org.junit.Test;
import org.apache.cassandra.SchemaLoader;
import org.apache.cassandra.db.CoordinatorLogBoundaries;
import org.apache.cassandra.io.util.DataInputBuffer;
import org.apache.cassandra.io.util.DataInputPlus;
import org.apache.cassandra.io.util.DataOutputBuffer;
import org.apache.cassandra.io.util.DataOutputPlus;
import org.apache.cassandra.net.MessagingService;
import org.apache.cassandra.net.TestChannel;
import org.apache.cassandra.replication.ImmutableCoordinatorLogOffsets;
import org.apache.cassandra.schema.TableId;
import org.apache.cassandra.streaming.PreviewKind;
import org.apache.cassandra.streaming.StreamDeserializingTask;
@ -116,7 +116,7 @@ public class StreamingInboundHandlerTest
public void StreamDeserializingTask_deserialize_ISM_NoSession() throws IOException
{
StreamMessageHeader header = new StreamMessageHeader(TableId.generate(), REMOTE_ADDR, nextTimeUUID(), true,
0, 0, 0, nextTimeUUID(), CoordinatorLogBoundaries.NONE);
0, 0, 0, nextTimeUUID(), ImmutableCoordinatorLogOffsets.NONE);
ByteBuffer temp = ByteBuffer.allocate(1024);
DataOutputPlus out = new DataOutputBuffer(temp);
@ -135,7 +135,7 @@ public class StreamingInboundHandlerTest
StreamResultFuture future = StreamResultFuture.createFollower(0, planId, StreamOperation.REPAIR, REMOTE_ADDR, streamingChannel, MessagingService.current_version, nextTimeUUID(), PreviewKind.ALL);
StreamManager.instance.registerFollower(future);
StreamMessageHeader header = new StreamMessageHeader(TableId.generate(), REMOTE_ADDR, planId, false,
0, 0, 0, nextTimeUUID(), CoordinatorLogBoundaries.NONE);
0, 0, 0, nextTimeUUID(), ImmutableCoordinatorLogOffsets.NONE);
// IncomingStreamMessage.serializer.deserialize
StreamSession session = StreamManager.instance.findSession(header.sender, header.planId, header.sessionIndex, header.sendByFollower);