Snapshot the last-written clustering instead of a stale reusable reference

CursorCompactor.mergePartition tracked the partition's last-written unfiltered
via a direct reference to the source cursor's descriptor (sstableCursors[0].unfiltered()).
That descriptor is reusable and gets overwritten in place as soon as the cursor reads
its next unfiltered -- including ones that are subsequently merged away or purged and
never actually written. By the time the partition closed and updateClusteringMetadata
was called, the reference could describe a clustering that was never written, corrupting
the sstable's min/max clustering metadata.

Fixed by copying the clustering into a dedicated, owned ClusteringDescriptor at the moment
something is genuinely written, rather than holding a reference into a cursor's mutable
scratch state. SSTableCursorWriter.updateClusteringMetadata now takes a ClusteringDescriptor
instead of an UnfilteredDescriptor to make that ownership explicit.

Found by an extended randomized differential soak run (hundreds of examples rather than
the default handful): a single-byte Statistics.db divergence between the iterator and
cursor paths that the logical JSON dump and coarse stats summary both missed, since
neither checks min/max clustering metadata directly.
This commit is contained in:
Jon Haddad 2026-07-28 17:45:04 -07:00
parent 8f88969e61
commit fe617688df
2 changed files with 15 additions and 6 deletions

View File

@ -45,6 +45,7 @@ import org.apache.cassandra.db.ReusableLivenessInfo;
import org.apache.cassandra.db.SerializationHeader; import org.apache.cassandra.db.SerializationHeader;
import org.apache.cassandra.db.SystemKeyspace; import org.apache.cassandra.db.SystemKeyspace;
import org.apache.cassandra.db.compaction.writers.CompactionAwareWriter; import org.apache.cassandra.db.compaction.writers.CompactionAwareWriter;
import org.apache.cassandra.db.marshal.AbstractType;
import org.apache.cassandra.db.partitions.UnfilteredPartitionIterators; import org.apache.cassandra.db.partitions.UnfilteredPartitionIterators;
import org.apache.cassandra.db.rows.BTreeRow; import org.apache.cassandra.db.rows.BTreeRow;
import org.apache.cassandra.db.rows.Cell; import org.apache.cassandra.db.rows.Cell;
@ -54,6 +55,7 @@ import org.apache.cassandra.db.rows.Row;
import org.apache.cassandra.db.rows.UnfilteredRowIterators; import org.apache.cassandra.db.rows.UnfilteredRowIterators;
import org.apache.cassandra.db.rows.UnfilteredSerializer; import org.apache.cassandra.db.rows.UnfilteredSerializer;
import org.apache.cassandra.dht.ReusableDecoratedKey; import org.apache.cassandra.dht.ReusableDecoratedKey;
import org.apache.cassandra.io.sstable.ClusteringDescriptor;
import org.apache.cassandra.io.sstable.ISSTableScanner; import org.apache.cassandra.io.sstable.ISSTableScanner;
import org.apache.cassandra.io.sstable.PartitionDescriptor; import org.apache.cassandra.io.sstable.PartitionDescriptor;
import org.apache.cassandra.io.sstable.SSTableCursorReader; import org.apache.cassandra.io.sstable.SSTableCursorReader;
@ -245,6 +247,13 @@ public class CursorCompactor extends CompactionInfo.Holder
private StatefulCursor lastSource = null; private StatefulCursor lastSource = null;
private final ReusableDecoratedKey detachedLastWrittenKey; private final ReusableDecoratedKey detachedLastWrittenKey;
// Snapshot of the clustering of the last unfiltered WRITTEN to the current output partition.
// sstableCursors[*].unfiltered() descriptors are reusable: a cursor overwrites its descriptor
// in place as soon as it reads its next unfiltered including unfiltereds that then merge
// away or purge and never reach the output so a reference captured during the merge loop
// can, by the time the partition closes, describe a clustering that was never written. A copy
// per written unfiltered keeps the true last written clustering for the min/max metadata.
private final ClusteringDescriptor lastWrittenClustering;
// Partition state. Writes can be delayed if the deletion is purged, or live and partition is empty -> LIVE deletion. // Partition state. Writes can be delayed if the deletion is purged, or live and partition is empty -> LIVE deletion.
PartitionDescriptor partitionDescriptor; PartitionDescriptor partitionDescriptor;
@ -327,6 +336,7 @@ public class CursorCompactor extends CompactionInfo.Holder
purger = new Purger(type, controller); purger = new Purger(type, controller);
detachedLastWrittenKey = metadata.partitioner.createReusableKey(128); detachedLastWrittenKey = metadata.partitioner.createReusableKey(128);
lastWrittenClustering = new ClusteringDescriptor(metadata.comparator.subtypes().toArray(AbstractType[]::new));
} }
/** /**
@ -432,7 +442,6 @@ public class CursorCompactor extends CompactionInfo.Holder
int unfilteredMergeLimit = partitionMergeLimit; int unfilteredMergeLimit = partitionMergeLimit;
boolean isFirstUnfiltered = true; boolean isFirstUnfiltered = true;
int unfilteredCount = 0; int unfilteredCount = 0;
UnfilteredDescriptor lastClustering = null;
while (true) while (true)
{ {
unfilteredMergeLimit = prepareAndSortUnfilteredForMerge(partitionMergeLimit, unfilteredMergeLimit); unfilteredMergeLimit = prepareAndSortUnfilteredForMerge(partitionMergeLimit, unfilteredMergeLimit);
@ -445,7 +454,7 @@ public class CursorCompactor extends CompactionInfo.Holder
{ {
isFirstUnfiltered = false; isFirstUnfiltered = false;
unfilteredCount++; unfilteredCount++;
lastClustering = sstableCursors[0].unfiltered(); lastWrittenClustering.copy(sstableCursors[0].unfiltered());
} }
} }
else if (UnfilteredSerializer.isTombstoneMarker(flags)) { else if (UnfilteredSerializer.isTombstoneMarker(flags)) {
@ -454,7 +463,7 @@ public class CursorCompactor extends CompactionInfo.Holder
{ {
isFirstUnfiltered = false; isFirstUnfiltered = false;
unfilteredCount++; unfilteredCount++;
lastClustering = sstableCursors[0].unfiltered(); lastWrittenClustering.copy(sstableCursors[0].unfiltered());
} }
if (activeOpenRangeDeletion == DeletionTime.LIVE) { if (activeOpenRangeDeletion == DeletionTime.LIVE) {
activeDeletion = mergedDeletion; activeDeletion = mergedDeletion;
@ -476,7 +485,7 @@ public class CursorCompactor extends CompactionInfo.Holder
ssTableCursorWriter.writePartitionEnd(partitionDescriptor.keyBytes(), partitionDescriptor.keyLength(), toWritePartitionDeletion, partitionHeaderLength); ssTableCursorWriter.writePartitionEnd(partitionDescriptor.keyBytes(), partitionDescriptor.keyLength(), toWritePartitionDeletion, partitionHeaderLength);
// update metadata tracking of min/max clustering on last unfiltered // update metadata tracking of min/max clustering on last unfiltered
if (unfilteredCount > 1) { if (unfilteredCount > 1) {
ssTableCursorWriter.updateClusteringMetadata(lastClustering); ssTableCursorWriter.updateClusteringMetadata(lastWrittenClustering);
} }
} }
// move along // move along

View File

@ -503,9 +503,9 @@ public class SSTableCursorWriter implements AutoCloseable
cursorIndexWriter.rowWritten(unfilteredDescriptor, unfilteredStartPosition, unfilteredEndPosition, openMarker); cursorIndexWriter.rowWritten(unfilteredDescriptor, unfilteredStartPosition, unfilteredEndPosition, openMarker);
} }
public void updateClusteringMetadata(UnfilteredDescriptor unfilteredDescriptor) public void updateClusteringMetadata(ClusteringDescriptor clusteringDescriptor)
{ {
metadataCollector.updateClusteringValues(unfilteredDescriptor); metadataCollector.updateClusteringValues(clusteringDescriptor);
} }
/** /**