mirror of https://github.com/apache/cassandra
Merge branch 'cassandra-5.0' into cassandra-6.0
* cassandra-5.0: SAI Component Checksum Validation Should be Segment-Aware
This commit is contained in:
commit
27cc24bc79
|
|
@ -63,6 +63,7 @@
|
||||||
* Fix a removed TTLed row re-appearance in a materialized view after a cursor compaction (CASSANDRA-21152)
|
* Fix a removed TTLed row re-appearance in a materialized view after a cursor compaction (CASSANDRA-21152)
|
||||||
* Rework ZSTD dictionary compression logic to create a trainer per training (CASSANDRA-21209)
|
* Rework ZSTD dictionary compression logic to create a trainer per training (CASSANDRA-21209)
|
||||||
Merged from 5.0:
|
Merged from 5.0:
|
||||||
|
* SAI Component Checksum Validation Should be Segment-Aware (CASSANDRA-21516)
|
||||||
* Support Python 3.12 and 3.13 in cqlsh (CASSANDRA-20997)
|
* Support Python 3.12 and 3.13 in cqlsh (CASSANDRA-20997)
|
||||||
* Make synchronization on VectorMemoryIndex inserts more granular (CASSANDRA-21160)
|
* Make synchronization on VectorMemoryIndex inserts more granular (CASSANDRA-21160)
|
||||||
* putShortVolatile is not volatile in InMemoryTrie (CASSANDRA-21353)
|
* putShortVolatile is not volatile in InMemoryTrie (CASSANDRA-21353)
|
||||||
|
|
|
||||||
|
|
@ -20,6 +20,8 @@ package org.apache.cassandra.index.sai.disk.io;
|
||||||
|
|
||||||
import java.io.IOException;
|
import java.io.IOException;
|
||||||
|
|
||||||
|
import javax.annotation.concurrent.NotThreadSafe;
|
||||||
|
|
||||||
import org.apache.lucene.store.DataInput;
|
import org.apache.lucene.store.DataInput;
|
||||||
import org.apache.lucene.store.IndexInput;
|
import org.apache.lucene.store.IndexInput;
|
||||||
|
|
||||||
|
|
@ -30,10 +32,13 @@ import org.apache.cassandra.io.util.RandomAccessReader;
|
||||||
* This is a wrapper over a Cassandra {@link RandomAccessReader} that provides an {@link IndexInput}
|
* This is a wrapper over a Cassandra {@link RandomAccessReader} that provides an {@link IndexInput}
|
||||||
* interface for Lucene classes that need {@link IndexInput}. This is an optimisation because the
|
* interface for Lucene classes that need {@link IndexInput}. This is an optimisation because the
|
||||||
* Lucene {@link DataInput} reads bytes one at a time whereas the {@link RandomAccessReader} is
|
* Lucene {@link DataInput} reads bytes one at a time whereas the {@link RandomAccessReader} is
|
||||||
* optimised to read multibyte objects faster.
|
* optimized to read multibyte objects faster.
|
||||||
*/
|
*/
|
||||||
|
@NotThreadSafe
|
||||||
public class IndexInputReader extends IndexInput
|
public class IndexInputReader extends IndexInput
|
||||||
{
|
{
|
||||||
|
public static final Runnable NO_OP_ON_CLOSE = () -> {};
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* the byte order of `input`'s native readX operations doesn't matter,
|
* the byte order of `input`'s native readX operations doesn't matter,
|
||||||
* because we only use `readFully` and `readByte` methods. IndexInput calls these
|
* because we only use `readFully` and `readByte` methods. IndexInput calls these
|
||||||
|
|
@ -42,27 +47,47 @@ public class IndexInputReader extends IndexInput
|
||||||
private final RandomAccessReader input;
|
private final RandomAccessReader input;
|
||||||
private final Runnable doOnClose;
|
private final Runnable doOnClose;
|
||||||
|
|
||||||
private IndexInputReader(RandomAccessReader input, Runnable doOnClose)
|
/** Absolute offset in the underlying file that this input's position 0 refers to. */
|
||||||
|
private final long offset;
|
||||||
|
|
||||||
|
/** Bounded length of this input, in bytes. */
|
||||||
|
private final long length;
|
||||||
|
|
||||||
|
private IndexInputReader(RandomAccessReader input, Runnable doOnClose, long offset, long length)
|
||||||
{
|
{
|
||||||
super(input.getPath());
|
super(input.getPath());
|
||||||
this.input = input;
|
this.input = input;
|
||||||
this.doOnClose = doOnClose;
|
this.doOnClose = doOnClose;
|
||||||
|
this.offset = offset;
|
||||||
|
this.length = length;
|
||||||
}
|
}
|
||||||
|
|
||||||
public static IndexInputReader create(RandomAccessReader input)
|
public static IndexInputReader create(RandomAccessReader input)
|
||||||
{
|
{
|
||||||
return new IndexInputReader(input, () -> {});
|
// Top-level inputs own the underlying reader; folding its close into doOnClose lets us
|
||||||
|
// avoid a separate ownership flag on the class.
|
||||||
|
return new IndexInputReader(input, input::close, 0L, input.length());
|
||||||
}
|
}
|
||||||
|
|
||||||
public static IndexInputReader create(RandomAccessReader input, Runnable doOnClose)
|
public static IndexInputReader create(RandomAccessReader input, Runnable doOnClose)
|
||||||
{
|
{
|
||||||
return new IndexInputReader(input, doOnClose);
|
Runnable close = () -> {
|
||||||
|
try
|
||||||
|
{
|
||||||
|
input.close();
|
||||||
|
}
|
||||||
|
finally
|
||||||
|
{
|
||||||
|
doOnClose.run();
|
||||||
|
}
|
||||||
|
};
|
||||||
|
return new IndexInputReader(input, close, 0L, input.length());
|
||||||
}
|
}
|
||||||
|
|
||||||
public static IndexInputReader create(FileHandle handle)
|
public static IndexInputReader create(FileHandle handle)
|
||||||
{
|
{
|
||||||
RandomAccessReader reader = handle.createReader();
|
RandomAccessReader reader = handle.createReader();
|
||||||
return new IndexInputReader(reader, () -> {});
|
return new IndexInputReader(reader, reader::close, 0L, reader.length());
|
||||||
}
|
}
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
|
|
@ -79,38 +104,43 @@ public class IndexInputReader extends IndexInput
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
public void close()
|
public void close()
|
||||||
{
|
|
||||||
try
|
|
||||||
{
|
|
||||||
input.close();
|
|
||||||
}
|
|
||||||
finally
|
|
||||||
{
|
{
|
||||||
doOnClose.run();
|
doOnClose.run();
|
||||||
}
|
}
|
||||||
}
|
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
public long getFilePointer()
|
public long getFilePointer()
|
||||||
{
|
{
|
||||||
return input.getFilePointer();
|
return input.getFilePointer() - offset;
|
||||||
}
|
}
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
public void seek(long position)
|
public void seek(long position)
|
||||||
{
|
{
|
||||||
input.seek(position);
|
if (position > length)
|
||||||
|
throw new IllegalArgumentException("Cannot seek to position " + position + " past length of " + length);
|
||||||
|
|
||||||
|
input.seek(offset + position);
|
||||||
}
|
}
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
public long length()
|
public long length()
|
||||||
{
|
{
|
||||||
return input.length();
|
return length;
|
||||||
}
|
}
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
public IndexInput slice(String sliceDescription, long offset, long length)
|
public IndexInput slice(String sliceDescription, long offset, long length)
|
||||||
{
|
{
|
||||||
throw new UnsupportedOperationException("Slice operations are not supported");
|
if (offset < 0 || length < 0 || offset + length > this.length)
|
||||||
|
throw new IllegalArgumentException("Invalid slice: offset=" + offset + ", length=" + length + ", parent length=" + this.length + " for " + sliceDescription);
|
||||||
|
|
||||||
|
// Slices share the underlying reader with their parent; the no-op close keeps the parent's lifecycle intact.
|
||||||
|
IndexInputReader slice = new IndexInputReader(input, NO_OP_ON_CLOSE, this.offset + offset, length);
|
||||||
|
|
||||||
|
// Seek to the beginning of the slice...
|
||||||
|
slice.seek(0);
|
||||||
|
|
||||||
|
return slice;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -21,11 +21,14 @@ package org.apache.cassandra.index.sai.disk.v1;
|
||||||
import java.io.IOException;
|
import java.io.IOException;
|
||||||
import java.io.UncheckedIOException;
|
import java.io.UncheckedIOException;
|
||||||
import java.util.EnumSet;
|
import java.util.EnumSet;
|
||||||
|
import java.util.List;
|
||||||
import java.util.Set;
|
import java.util.Set;
|
||||||
|
|
||||||
import com.codahale.metrics.Gauge;
|
import com.codahale.metrics.Gauge;
|
||||||
import com.google.common.annotations.VisibleForTesting;
|
import com.google.common.annotations.VisibleForTesting;
|
||||||
|
|
||||||
|
import org.apache.lucene.codecs.CodecUtil;
|
||||||
|
import org.apache.lucene.index.CorruptIndexException;
|
||||||
import org.apache.lucene.store.IndexInput;
|
import org.apache.lucene.store.IndexInput;
|
||||||
import org.slf4j.Logger;
|
import org.slf4j.Logger;
|
||||||
import org.slf4j.LoggerFactory;
|
import org.slf4j.LoggerFactory;
|
||||||
|
|
@ -44,6 +47,7 @@ import org.apache.cassandra.index.sai.disk.format.IndexComponent;
|
||||||
import org.apache.cassandra.index.sai.disk.format.IndexDescriptor;
|
import org.apache.cassandra.index.sai.disk.format.IndexDescriptor;
|
||||||
import org.apache.cassandra.index.sai.disk.format.OnDiskFormat;
|
import org.apache.cassandra.index.sai.disk.format.OnDiskFormat;
|
||||||
import org.apache.cassandra.index.sai.disk.v1.segment.SegmentBuilder;
|
import org.apache.cassandra.index.sai.disk.v1.segment.SegmentBuilder;
|
||||||
|
import org.apache.cassandra.index.sai.disk.v1.segment.SegmentMetadata;
|
||||||
import org.apache.cassandra.index.sai.metrics.AbstractMetrics;
|
import org.apache.cassandra.index.sai.metrics.AbstractMetrics;
|
||||||
import org.apache.cassandra.index.sai.utils.IndexIdentifier;
|
import org.apache.cassandra.index.sai.utils.IndexIdentifier;
|
||||||
import org.apache.cassandra.index.sai.utils.IndexTermType;
|
import org.apache.cassandra.index.sai.utils.IndexTermType;
|
||||||
|
|
@ -96,6 +100,18 @@ public class V1OnDiskFormat implements OnDiskFormat
|
||||||
IndexComponent.TERMS_DATA,
|
IndexComponent.TERMS_DATA,
|
||||||
IndexComponent.POSTING_LISTS);
|
IndexComponent.POSTING_LISTS);
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Per-column components whose files are written in append mode with one SAI codec footer
|
||||||
|
* per segment (see {@link org.apache.cassandra.index.sai.disk.v1.bbtree.NumericIndexWriter},
|
||||||
|
* {@link org.apache.cassandra.index.sai.disk.v1.trie.TrieTermsDictionaryWriter},
|
||||||
|
* {@link org.apache.cassandra.index.sai.disk.v1.postings.PostingsWriter}, and
|
||||||
|
* {@link org.apache.cassandra.index.sai.disk.v1.vector.OnHeapGraph}).
|
||||||
|
*/
|
||||||
|
private static final Set<IndexComponent> SEGMENTED_COMPONENTS = EnumSet.of(IndexComponent.BALANCED_TREE,
|
||||||
|
IndexComponent.POSTING_LISTS,
|
||||||
|
IndexComponent.TERMS_DATA,
|
||||||
|
IndexComponent.COMPRESSED_VECTORS);
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Global limit on heap consumed by all index segment building that occurs outside the context of Memtable flush.
|
* Global limit on heap consumed by all index segment building that occurs outside the context of Memtable flush.
|
||||||
* <p>
|
* <p>
|
||||||
|
|
@ -219,15 +235,90 @@ public class V1OnDiskFormat implements OnDiskFormat
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
if (isEmptyIndex)
|
||||||
|
return;
|
||||||
|
|
||||||
|
// Safely read the segment metadata so we can validate per-segment checksums below...
|
||||||
|
List<SegmentMetadata> segments = null;
|
||||||
|
if (checksum)
|
||||||
|
{
|
||||||
|
validateIndexComponent(indexDescriptor, indexIdentifier, IndexComponent.META, true);
|
||||||
|
try
|
||||||
|
{
|
||||||
|
segments = SegmentMetadata.load(MetadataSource.loadColumnMetadata(indexDescriptor, indexIdentifier), indexDescriptor.primaryKeyFactory);
|
||||||
|
}
|
||||||
|
catch (IOException e)
|
||||||
|
{
|
||||||
|
rethrowIOException(e);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
for (IndexComponent indexComponent : perColumnIndexComponents(indexTermType))
|
for (IndexComponent indexComponent : perColumnIndexComponents(indexTermType))
|
||||||
{
|
{
|
||||||
if (!isEmptyIndex && isNotBuildCompletionMarker(indexComponent))
|
if (isNotBuildCompletionMarker(indexComponent))
|
||||||
{
|
{
|
||||||
|
// META was validated up-front in CHECKSUM mode; don't validate it twice.
|
||||||
|
if (checksum && indexComponent == IndexComponent.META)
|
||||||
|
continue;
|
||||||
|
|
||||||
|
if (checksum && SEGMENTED_COMPONENTS.contains(indexComponent))
|
||||||
|
{
|
||||||
|
assert segments != null : "No segment metadata available!";
|
||||||
|
validateSegmentedIndexComponent(indexDescriptor, indexIdentifier, indexComponent, segments, indexTermType.isVector());
|
||||||
|
}
|
||||||
|
else
|
||||||
validateIndexComponent(indexDescriptor, indexIdentifier, indexComponent, checksum);
|
validateIndexComponent(indexDescriptor, indexIdentifier, indexComponent, checksum);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
private static void validateSegmentedIndexComponent(IndexDescriptor indexDescriptor,
|
||||||
|
IndexIdentifier indexIdentifier,
|
||||||
|
IndexComponent indexComponent,
|
||||||
|
List<SegmentMetadata> segments,
|
||||||
|
boolean payloadOnlyMetadata)
|
||||||
|
{
|
||||||
|
try (IndexInput input = indexDescriptor.openPerIndexInput(indexComponent, indexIdentifier))
|
||||||
|
{
|
||||||
|
long fileLength = input.length();
|
||||||
|
long frameStart = 0;
|
||||||
|
|
||||||
|
for (SegmentMetadata segment : segments)
|
||||||
|
{
|
||||||
|
SegmentMetadata.ComponentMetadata cm = segment.componentMetadatas.get(indexComponent);
|
||||||
|
|
||||||
|
// Non-vector writers record offsets as the codec-framed segment starts (before
|
||||||
|
// the header) and length as the full framed length (through the footer). The vector
|
||||||
|
// writer (OnHeapGraph#writeData) instead records the offset as the payload start (after
|
||||||
|
// the header) and length as just the payload length, because vector readers seek
|
||||||
|
// directly at the payload. Segments are written contiguously in append mode, so we can
|
||||||
|
// recover the vector-path frame extent by walking segment ends and adding the trailing
|
||||||
|
// 16-byte codec footer.
|
||||||
|
long frameEnd = payloadOnlyMetadata ? cm.offset + cm.length + CodecUtil.footerLength() : cm.offset + cm.length;
|
||||||
|
|
||||||
|
if (frameEnd > fileLength || frameEnd < frameStart)
|
||||||
|
throw new CorruptIndexException(String.format("Segment frame [%d, %d) is inconsistent with component file length %d",
|
||||||
|
frameStart, frameEnd, fileLength),
|
||||||
|
indexComponent.name + '@' + frameStart);
|
||||||
|
|
||||||
|
IndexInput slice = input.slice(indexComponent.name + '@' + frameStart, frameStart, frameEnd - frameStart);
|
||||||
|
SAICodecUtils.validateChecksum(slice);
|
||||||
|
frameStart = frameEnd;
|
||||||
|
}
|
||||||
|
|
||||||
|
if (frameStart != fileLength)
|
||||||
|
throw new CorruptIndexException(String.format("Component file length %d does not match combined frame length of all segments %d",
|
||||||
|
fileLength, frameStart),
|
||||||
|
indexComponent.name);
|
||||||
|
}
|
||||||
|
catch (Exception e)
|
||||||
|
{
|
||||||
|
logger.warn(indexDescriptor.logMessage("Segmented checksum validation failed for index component {} on SSTable {}"),
|
||||||
|
indexComponent, indexDescriptor.sstableDescriptor);
|
||||||
|
rethrowIOException(e);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
private static void validateIndexComponent(IndexDescriptor indexDescriptor,
|
private static void validateIndexComponent(IndexDescriptor indexDescriptor,
|
||||||
IndexIdentifier indexContext,
|
IndexIdentifier indexContext,
|
||||||
IndexComponent indexComponent,
|
IndexComponent indexComponent,
|
||||||
|
|
@ -245,9 +336,7 @@ public class V1OnDiskFormat implements OnDiskFormat
|
||||||
catch (Exception e)
|
catch (Exception e)
|
||||||
{
|
{
|
||||||
logger.warn(indexDescriptor.logMessage("{} failed for index component {} on SSTable {}"),
|
logger.warn(indexDescriptor.logMessage("{} failed for index component {} on SSTable {}"),
|
||||||
checksum ? "Checksum validation" : "Validation",
|
checksum ? "Checksum validation" : "Validation", indexComponent, indexDescriptor.sstableDescriptor);
|
||||||
indexComponent,
|
|
||||||
indexDescriptor.sstableDescriptor);
|
|
||||||
rethrowIOException(e);
|
rethrowIOException(e);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -55,8 +55,10 @@ import org.apache.cassandra.db.SystemKeyspace;
|
||||||
import org.apache.cassandra.db.compaction.CompactionManager;
|
import org.apache.cassandra.db.compaction.CompactionManager;
|
||||||
import org.apache.cassandra.db.compaction.OperationType;
|
import org.apache.cassandra.db.compaction.OperationType;
|
||||||
import org.apache.cassandra.db.marshal.AbstractType;
|
import org.apache.cassandra.db.marshal.AbstractType;
|
||||||
|
import org.apache.cassandra.db.marshal.FloatType;
|
||||||
import org.apache.cassandra.db.marshal.Int32Type;
|
import org.apache.cassandra.db.marshal.Int32Type;
|
||||||
import org.apache.cassandra.db.marshal.UTF8Type;
|
import org.apache.cassandra.db.marshal.UTF8Type;
|
||||||
|
import org.apache.cassandra.db.marshal.VectorType;
|
||||||
import org.apache.cassandra.exceptions.InvalidRequestException;
|
import org.apache.cassandra.exceptions.InvalidRequestException;
|
||||||
import org.apache.cassandra.index.Index;
|
import org.apache.cassandra.index.Index;
|
||||||
import org.apache.cassandra.index.SecondaryIndexManager;
|
import org.apache.cassandra.index.SecondaryIndexManager;
|
||||||
|
|
@ -1429,6 +1431,27 @@ public class StorageAttachedIndexDDLTest extends SAITester
|
||||||
assertEquals(Arrays.asList(2L, 1L), toSize.apply(iterator.next()));
|
assertEquals(Arrays.asList(2L, 1L), toSize.apply(iterator.next()));
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
public void multiSegmentVectorIndexPassesChecksumValidation()
|
||||||
|
{
|
||||||
|
createTable("CREATE TABLE %s (pk int, val vector<float, 3>, PRIMARY KEY(pk))");
|
||||||
|
|
||||||
|
int vectorCount = 100;
|
||||||
|
for (int pk = 0; pk < vectorCount; pk++)
|
||||||
|
execute("INSERT INTO %s (pk, val) VALUES (" + pk + ", [" + pk + ".0, " + (pk + 1) + ".0, " + (pk + 2) + ".0])");
|
||||||
|
|
||||||
|
flush();
|
||||||
|
|
||||||
|
SegmentBuilder.updateLastValidSegmentRowId(17); // 17 rows per segment -> multi-segment build
|
||||||
|
IndexIdentifier vectorIndexIdentifier = createIndexIdentifier(createIndex("CREATE CUSTOM INDEX ON %s(val) USING 'StorageAttachedIndex'"));
|
||||||
|
IndexTermType vectorIndexTermType = createIndexTermType(VectorType.getInstance(FloatType.instance, 3));
|
||||||
|
|
||||||
|
// A vector index writes CompressedVectors.db, TermsData.db, and PostingLists.db in append
|
||||||
|
// mode with one SAI codec footer per segment (see OnHeapGraph.writeData). Multi-segment
|
||||||
|
// builds therefore need segment-aware checksum validation.
|
||||||
|
assertTrue(verifyChecksum(vectorIndexTermType, vectorIndexIdentifier));
|
||||||
|
}
|
||||||
|
|
||||||
private void assertZeroSegmentBuilderUsage()
|
private void assertZeroSegmentBuilderUsage()
|
||||||
{
|
{
|
||||||
assertEquals("Segment memory limiter should revert to zero.", 0L, getSegmentBufferUsedBytes());
|
assertEquals("Segment memory limiter should revert to zero.", 0L, getSegmentBufferUsedBytes());
|
||||||
|
|
|
||||||
|
|
@ -18,15 +18,19 @@
|
||||||
package org.apache.cassandra.index.sai.disk.v1;
|
package org.apache.cassandra.index.sai.disk.v1;
|
||||||
|
|
||||||
import java.io.IOException;
|
import java.io.IOException;
|
||||||
|
import java.io.RandomAccessFile;
|
||||||
|
import java.io.UncheckedIOException;
|
||||||
import java.nio.ByteBuffer;
|
import java.nio.ByteBuffer;
|
||||||
import java.nio.file.Files;
|
import java.nio.file.Files;
|
||||||
import java.nio.file.Path;
|
import java.nio.file.Path;
|
||||||
import java.util.Arrays;
|
import java.util.Arrays;
|
||||||
import java.util.Collections;
|
import java.util.Collections;
|
||||||
|
import java.util.Date;
|
||||||
import java.util.List;
|
import java.util.List;
|
||||||
|
|
||||||
import com.google.common.base.Stopwatch;
|
import com.google.common.base.Stopwatch;
|
||||||
|
|
||||||
|
import org.apache.lucene.index.CorruptIndexException;
|
||||||
import org.junit.After;
|
import org.junit.After;
|
||||||
import org.junit.BeforeClass;
|
import org.junit.BeforeClass;
|
||||||
import org.junit.Test;
|
import org.junit.Test;
|
||||||
|
|
@ -34,11 +38,13 @@ import org.junit.Test;
|
||||||
import org.apache.cassandra.config.DatabaseDescriptor;
|
import org.apache.cassandra.config.DatabaseDescriptor;
|
||||||
import org.apache.cassandra.db.Clustering;
|
import org.apache.cassandra.db.Clustering;
|
||||||
import org.apache.cassandra.db.DecoratedKey;
|
import org.apache.cassandra.db.DecoratedKey;
|
||||||
|
import org.apache.cassandra.db.marshal.TimestampType;
|
||||||
import org.apache.cassandra.db.marshal.UTF8Type;
|
import org.apache.cassandra.db.marshal.UTF8Type;
|
||||||
import org.apache.cassandra.db.rows.BTreeRow;
|
import org.apache.cassandra.db.rows.BTreeRow;
|
||||||
import org.apache.cassandra.db.rows.BufferCell;
|
import org.apache.cassandra.db.rows.BufferCell;
|
||||||
import org.apache.cassandra.db.rows.Row;
|
import org.apache.cassandra.db.rows.Row;
|
||||||
import org.apache.cassandra.dht.Murmur3Partitioner;
|
import org.apache.cassandra.dht.Murmur3Partitioner;
|
||||||
|
import org.apache.cassandra.index.sai.IndexValidation;
|
||||||
import org.apache.cassandra.index.sai.SAITester;
|
import org.apache.cassandra.index.sai.SAITester;
|
||||||
import org.apache.cassandra.index.sai.StorageAttachedIndex;
|
import org.apache.cassandra.index.sai.StorageAttachedIndex;
|
||||||
import org.apache.cassandra.index.sai.disk.format.IndexComponent;
|
import org.apache.cassandra.index.sai.disk.format.IndexComponent;
|
||||||
|
|
@ -61,6 +67,8 @@ import org.apache.cassandra.utils.bytecomparable.ByteSource;
|
||||||
import static org.apache.cassandra.Util.dk;
|
import static org.apache.cassandra.Util.dk;
|
||||||
import static org.junit.Assert.assertEquals;
|
import static org.junit.Assert.assertEquals;
|
||||||
import static org.junit.Assert.assertFalse;
|
import static org.junit.Assert.assertFalse;
|
||||||
|
import static org.junit.Assert.assertTrue;
|
||||||
|
import static org.junit.Assert.fail;
|
||||||
|
|
||||||
public class SegmentFlushTest
|
public class SegmentFlushTest
|
||||||
{
|
{
|
||||||
|
|
@ -87,6 +95,147 @@ public class SegmentFlushTest
|
||||||
SegmentBuilder.updateLastValidSegmentRowId(-1); // reset
|
SegmentBuilder.updateLastValidSegmentRowId(-1); // reset
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
public void multiSegmentBalancedTreePassesChecksumValidation() throws IOException
|
||||||
|
{
|
||||||
|
Path tmpDir = Files.createTempDirectory("SegmentFlushTest");
|
||||||
|
IndexDescriptor indexDescriptor = IndexDescriptor.create(new Descriptor(new File(tmpDir.toFile()), "ks", "cf", new SequenceBasedSSTableId(1)), Murmur3Partitioner.instance, SAITester.EMPTY_COMPARATOR);
|
||||||
|
ColumnMetadata column = ColumnMetadata.regularColumn("sai", "internal", "ts", TimestampType.instance, 1);
|
||||||
|
StorageAttachedIndex index = SAITester.createMockIndex(column);
|
||||||
|
|
||||||
|
SSTableIndexWriter writer = new SSTableIndexWriter(indexDescriptor, index, V1OnDiskFormat.SEGMENT_BUILD_MEMORY_LIMITER, () -> true);
|
||||||
|
|
||||||
|
List<DecoratedKey> keys = Arrays.asList(dk("1"), dk("2"));
|
||||||
|
Collections.sort(keys);
|
||||||
|
|
||||||
|
writer.addRow(SAITester.TEST_FACTORY.create(keys.get(0)), createRow(column, TimestampType.instance.decompose(new Date(1_000L))), 0L);
|
||||||
|
writer.addRow(SAITester.TEST_FACTORY.create(keys.get(1)), createRow(column, TimestampType.instance.decompose(new Date(2_000L))), SegmentBuilder.LAST_VALID_SEGMENT_ROW_ID + 1);
|
||||||
|
writer.complete(Stopwatch.createStarted());
|
||||||
|
|
||||||
|
// Will throw if checksum validation fails:
|
||||||
|
indexDescriptor.validatePerIndexComponents(index.termType(), index.identifier(), IndexValidation.CHECKSUM, true, true);
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
public void multiSegmentTermsDataPassesChecksumValidation() throws IOException
|
||||||
|
{
|
||||||
|
Path tmpDir = Files.createTempDirectory("SegmentFlushTest");
|
||||||
|
IndexDescriptor indexDescriptor = IndexDescriptor.create(new Descriptor(new File(tmpDir.toFile()), "ks", "cf", new SequenceBasedSSTableId(1)), Murmur3Partitioner.instance, SAITester.EMPTY_COMPARATOR);
|
||||||
|
ColumnMetadata column = ColumnMetadata.regularColumn("sai", "internal", "name", UTF8Type.instance, 1);
|
||||||
|
StorageAttachedIndex index = SAITester.createMockIndex(column);
|
||||||
|
|
||||||
|
SSTableIndexWriter writer = new SSTableIndexWriter(indexDescriptor, index, V1OnDiskFormat.SEGMENT_BUILD_MEMORY_LIMITER, () -> true);
|
||||||
|
|
||||||
|
List<DecoratedKey> keys = Arrays.asList(dk("1"), dk("2"));
|
||||||
|
Collections.sort(keys);
|
||||||
|
|
||||||
|
writer.addRow(SAITester.TEST_FACTORY.create(keys.get(0)), createRow(column, UTF8Type.instance.decompose("a")), 0L);
|
||||||
|
writer.addRow(SAITester.TEST_FACTORY.create(keys.get(1)), createRow(column, UTF8Type.instance.decompose("b")), SegmentBuilder.LAST_VALID_SEGMENT_ROW_ID + 1);
|
||||||
|
writer.complete(Stopwatch.createStarted());
|
||||||
|
|
||||||
|
// Will throw if checksum validation fails:
|
||||||
|
indexDescriptor.validatePerIndexComponents(index.termType(), index.identifier(), IndexValidation.CHECKSUM, true, true);
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
public void multiSegmentBalancedTreePassesHeaderFooterValidation() throws IOException
|
||||||
|
{
|
||||||
|
Path tmpDir = Files.createTempDirectory("SegmentFlushTest");
|
||||||
|
IndexDescriptor indexDescriptor = IndexDescriptor.create(new Descriptor(new File(tmpDir.toFile()), "ks", "cf", new SequenceBasedSSTableId(1)), Murmur3Partitioner.instance, SAITester.EMPTY_COMPARATOR);
|
||||||
|
ColumnMetadata column = ColumnMetadata.regularColumn("sai", "internal", "ts", TimestampType.instance, 1);
|
||||||
|
StorageAttachedIndex index = SAITester.createMockIndex(column);
|
||||||
|
|
||||||
|
SSTableIndexWriter writer = new SSTableIndexWriter(indexDescriptor, index, V1OnDiskFormat.SEGMENT_BUILD_MEMORY_LIMITER, () -> true);
|
||||||
|
|
||||||
|
List<DecoratedKey> keys = Arrays.asList(dk("1"), dk("2"));
|
||||||
|
Collections.sort(keys);
|
||||||
|
|
||||||
|
writer.addRow(SAITester.TEST_FACTORY.create(keys.get(0)), createRow(column, TimestampType.instance.decompose(new Date(1_000L))), 0L);
|
||||||
|
writer.addRow(SAITester.TEST_FACTORY.create(keys.get(1)), createRow(column, TimestampType.instance.decompose(new Date(2_000L))), SegmentBuilder.LAST_VALID_SEGMENT_ROW_ID + 1);
|
||||||
|
writer.complete(Stopwatch.createStarted());
|
||||||
|
|
||||||
|
indexDescriptor.validatePerIndexComponents(index.termType(), index.identifier(), IndexValidation.HEADER_FOOTER, false, true);
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
public void multiSegmentBalancedTreeFailsChecksumOnFirstSegmentByteFlip() throws IOException
|
||||||
|
{
|
||||||
|
Path tmpDir = Files.createTempDirectory("SegmentFlushTest");
|
||||||
|
IndexDescriptor indexDescriptor = IndexDescriptor.create(new Descriptor(new File(tmpDir.toFile()), "ks", "cf", new SequenceBasedSSTableId(1)), Murmur3Partitioner.instance, SAITester.EMPTY_COMPARATOR);
|
||||||
|
ColumnMetadata column = ColumnMetadata.regularColumn("sai", "internal", "ts", TimestampType.instance, 1);
|
||||||
|
StorageAttachedIndex index = SAITester.createMockIndex(column);
|
||||||
|
|
||||||
|
SSTableIndexWriter writer = new SSTableIndexWriter(indexDescriptor, index, V1OnDiskFormat.SEGMENT_BUILD_MEMORY_LIMITER, () -> true);
|
||||||
|
|
||||||
|
List<DecoratedKey> keys = Arrays.asList(dk("1"), dk("2"));
|
||||||
|
Collections.sort(keys);
|
||||||
|
|
||||||
|
writer.addRow(SAITester.TEST_FACTORY.create(keys.get(0)), createRow(column, TimestampType.instance.decompose(new Date(1_000L))), 0L);
|
||||||
|
writer.addRow(SAITester.TEST_FACTORY.create(keys.get(1)), createRow(column, TimestampType.instance.decompose(new Date(2_000L))), SegmentBuilder.LAST_VALID_SEGMENT_ROW_ID + 1);
|
||||||
|
writer.complete(Stopwatch.createStarted());
|
||||||
|
|
||||||
|
// Locate segment 0's payload extent so the flip lands inside it. Corrupting the FIRST
|
||||||
|
// (not last) segment specifically proves the validator inspects every segment -- a
|
||||||
|
// validator that only checked the trailing footer would miss this and silently pass.
|
||||||
|
MetadataSource source = MetadataSource.loadColumnMetadata(indexDescriptor, index.identifier());
|
||||||
|
List<SegmentMetadata> segments = SegmentMetadata.load(source, indexDescriptor.primaryKeyFactory);
|
||||||
|
assertEquals(2, segments.size());
|
||||||
|
SegmentMetadata.ComponentMetadata cm = segments.get(0).componentMetadatas.get(IndexComponent.BALANCED_TREE);
|
||||||
|
long flipPosition = cm.offset + cm.length / 2;
|
||||||
|
|
||||||
|
File balancedTree = indexDescriptor.fileFor(IndexComponent.BALANCED_TREE, index.identifier());
|
||||||
|
try (RandomAccessFile raf = new RandomAccessFile(balancedTree.toJavaIOFile(), "rw"))
|
||||||
|
{
|
||||||
|
raf.seek(flipPosition);
|
||||||
|
int original = raf.readByte();
|
||||||
|
raf.seek(flipPosition);
|
||||||
|
raf.writeByte(original ^ 0xFF);
|
||||||
|
}
|
||||||
|
|
||||||
|
try
|
||||||
|
{
|
||||||
|
indexDescriptor.validatePerIndexComponents(index.termType(), index.identifier(), IndexValidation.CHECKSUM, true, true);
|
||||||
|
fail("Expected corrupted first segment to fail checksum validation");
|
||||||
|
}
|
||||||
|
catch (UncheckedIOException expected)
|
||||||
|
{
|
||||||
|
assertTrue("Expected CorruptIndexException cause; got " + expected.getCause(), expected.getCause() instanceof CorruptIndexException);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
public void multiSegmentBalancedTreeFailsChecksumOnAppendedGarbage() throws IOException
|
||||||
|
{
|
||||||
|
Path tmpDir = Files.createTempDirectory("SegmentFlushTest");
|
||||||
|
IndexDescriptor indexDescriptor = IndexDescriptor.create(new Descriptor(new File(tmpDir.toFile()), "ks", "cf", new SequenceBasedSSTableId(1)), Murmur3Partitioner.instance, SAITester.EMPTY_COMPARATOR);
|
||||||
|
ColumnMetadata column = ColumnMetadata.regularColumn("sai", "internal", "ts", TimestampType.instance, 1);
|
||||||
|
StorageAttachedIndex index = SAITester.createMockIndex(column);
|
||||||
|
|
||||||
|
SSTableIndexWriter writer = new SSTableIndexWriter(indexDescriptor, index, V1OnDiskFormat.SEGMENT_BUILD_MEMORY_LIMITER, () -> true);
|
||||||
|
|
||||||
|
List<DecoratedKey> keys = Arrays.asList(dk("1"), dk("2"));
|
||||||
|
Collections.sort(keys);
|
||||||
|
|
||||||
|
writer.addRow(SAITester.TEST_FACTORY.create(keys.get(0)), createRow(column, TimestampType.instance.decompose(new Date(1_000L))), 0L);
|
||||||
|
writer.addRow(SAITester.TEST_FACTORY.create(keys.get(1)), createRow(column, TimestampType.instance.decompose(new Date(2_000L))), SegmentBuilder.LAST_VALID_SEGMENT_ROW_ID + 1);
|
||||||
|
writer.complete(Stopwatch.createStarted());
|
||||||
|
|
||||||
|
// Append 100 random bytes past the last segment's footer. The per-segment slice loop
|
||||||
|
// walks each segment's declared frame, then asserts that frameStart == input.length()
|
||||||
|
// once done. This corruption trips that post-loop invariant, not a per-segment CRC.
|
||||||
|
SAITester.CorruptionType.APPENDED_DATA.corrupt(indexDescriptor.fileFor(IndexComponent.BALANCED_TREE, index.identifier()));
|
||||||
|
|
||||||
|
try
|
||||||
|
{
|
||||||
|
indexDescriptor.validatePerIndexComponents(index.termType(), index.identifier(), IndexValidation.CHECKSUM, true, true);
|
||||||
|
fail("Expected trailing garbage to fail checksum validation");
|
||||||
|
}
|
||||||
|
catch (UncheckedIOException expected)
|
||||||
|
{
|
||||||
|
assertTrue("Expected CorruptIndexException cause; got " + expected.getCause(), expected.getCause() instanceof CorruptIndexException);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
@Test
|
@Test
|
||||||
public void testFlushBetweenRowIds() throws Exception
|
public void testFlushBetweenRowIds() throws Exception
|
||||||
{
|
{
|
||||||
|
|
|
||||||
Loading…
Reference in New Issue