From fe617688df964df7cd1276fed6ff608d23d317df Mon Sep 17 00:00:00 2001 From: Jon Haddad Date: Tue, 28 Jul 2026 17:45:04 -0700 Subject: [PATCH] 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. --- .../db/compaction/CursorCompactor.java | 17 +++++++++++++---- .../io/sstable/SSTableCursorWriter.java | 4 ++-- 2 files changed, 15 insertions(+), 6 deletions(-) diff --git a/src/java/org/apache/cassandra/db/compaction/CursorCompactor.java b/src/java/org/apache/cassandra/db/compaction/CursorCompactor.java index 2a5c47f0db..81e85e5bb8 100644 --- a/src/java/org/apache/cassandra/db/compaction/CursorCompactor.java +++ b/src/java/org/apache/cassandra/db/compaction/CursorCompactor.java @@ -45,6 +45,7 @@ import org.apache.cassandra.db.ReusableLivenessInfo; import org.apache.cassandra.db.SerializationHeader; import org.apache.cassandra.db.SystemKeyspace; 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.rows.BTreeRow; 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.UnfilteredSerializer; 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.PartitionDescriptor; import org.apache.cassandra.io.sstable.SSTableCursorReader; @@ -245,6 +247,13 @@ public class CursorCompactor extends CompactionInfo.Holder private StatefulCursor lastSource = null; 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. PartitionDescriptor partitionDescriptor; @@ -327,6 +336,7 @@ public class CursorCompactor extends CompactionInfo.Holder purger = new Purger(type, controller); 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; boolean isFirstUnfiltered = true; int unfilteredCount = 0; - UnfilteredDescriptor lastClustering = null; while (true) { unfilteredMergeLimit = prepareAndSortUnfilteredForMerge(partitionMergeLimit, unfilteredMergeLimit); @@ -445,7 +454,7 @@ public class CursorCompactor extends CompactionInfo.Holder { isFirstUnfiltered = false; unfilteredCount++; - lastClustering = sstableCursors[0].unfiltered(); + lastWrittenClustering.copy(sstableCursors[0].unfiltered()); } } else if (UnfilteredSerializer.isTombstoneMarker(flags)) { @@ -454,7 +463,7 @@ public class CursorCompactor extends CompactionInfo.Holder { isFirstUnfiltered = false; unfilteredCount++; - lastClustering = sstableCursors[0].unfiltered(); + lastWrittenClustering.copy(sstableCursors[0].unfiltered()); } if (activeOpenRangeDeletion == DeletionTime.LIVE) { activeDeletion = mergedDeletion; @@ -476,7 +485,7 @@ public class CursorCompactor extends CompactionInfo.Holder ssTableCursorWriter.writePartitionEnd(partitionDescriptor.keyBytes(), partitionDescriptor.keyLength(), toWritePartitionDeletion, partitionHeaderLength); // update metadata tracking of min/max clustering on last unfiltered if (unfilteredCount > 1) { - ssTableCursorWriter.updateClusteringMetadata(lastClustering); + ssTableCursorWriter.updateClusteringMetadata(lastWrittenClustering); } } // move along diff --git a/src/java/org/apache/cassandra/io/sstable/SSTableCursorWriter.java b/src/java/org/apache/cassandra/io/sstable/SSTableCursorWriter.java index 467532f5a3..41cffc8429 100644 --- a/src/java/org/apache/cassandra/io/sstable/SSTableCursorWriter.java +++ b/src/java/org/apache/cassandra/io/sstable/SSTableCursorWriter.java @@ -503,9 +503,9 @@ public class SSTableCursorWriter implements AutoCloseable cursorIndexWriter.rowWritten(unfilteredDescriptor, unfilteredStartPosition, unfilteredEndPosition, openMarker); } - public void updateClusteringMetadata(UnfilteredDescriptor unfilteredDescriptor) + public void updateClusteringMetadata(ClusteringDescriptor clusteringDescriptor) { - metadataCollector.updateClusteringValues(unfilteredDescriptor); + metadataCollector.updateClusteringValues(clusteringDescriptor); } /**