From 6e0d8c983359d8632ce7504e41776f30f69b4dca Mon Sep 17 00:00:00 2001 From: 1fanwang <1fannnw@gmail.com> Date: Wed, 26 Aug 2026 07:44:41 -0400 Subject: [PATCH 1/7] CASSANDRA-21607: Fix compaction of retained complex cells Iterator compaction can encounter complex cells newer than a column drop even when the current schema has no complex columns. Signed-off-by: 1fanwang <1fannnw@gmail.com> --- .../org/apache/cassandra/db/rows/Row.java | 9 ++++++-- ...oppedColumnDifferentialCompactionTest.java | 21 +++++++++++++++++++ 2 files changed, 28 insertions(+), 2 deletions(-) diff --git a/src/java/org/apache/cassandra/db/rows/Row.java b/src/java/org/apache/cassandra/db/rows/Row.java index 4137c42ca242..a727246cd555 100644 --- a/src/java/org/apache/cassandra/db/rows/Row.java +++ b/src/java/org/apache/cassandra/db/rows/Row.java @@ -859,8 +859,8 @@ private static class ColumnDataReducer extends RowMergeIterator.Reducer>> complexCells; + private ComplexColumnData.Builder complexBuilder; + private List>> complexCells; private final CellReducer cellReducer; public ColumnDataReducer(int size, boolean hasComplex) @@ -916,6 +916,11 @@ protected ColumnData getReduced() } else { + if (complexBuilder == null) + { + complexBuilder = ComplexColumnData.builder(); + complexCells = new ArrayList<>(versions.length); + } complexBuilder.newColumn(column); complexCells.clear(); DeletionTime complexDeletion = DeletionTime.LIVE; diff --git a/test/unit/org/apache/cassandra/db/compaction/differential/DroppedColumnDifferentialCompactionTest.java b/test/unit/org/apache/cassandra/db/compaction/differential/DroppedColumnDifferentialCompactionTest.java index f1dc9ea0d21e..6d2b0242bb1c 100644 --- a/test/unit/org/apache/cassandra/db/compaction/differential/DroppedColumnDifferentialCompactionTest.java +++ b/test/unit/org/apache/cassandra/db/compaction/differential/DroppedColumnDifferentialCompactionTest.java @@ -24,6 +24,7 @@ import org.apache.cassandra.db.LivenessInfo; import org.apache.cassandra.db.Mutation; import org.apache.cassandra.db.partitions.PartitionUpdate; +import org.apache.cassandra.utils.FBUtilities; import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertTrue; @@ -95,6 +96,26 @@ public void cellsNewerThanDropRetained() throws Exception assertCursorMatchesIterator(cfs); } + @Test + public void complexCellsNewerThanDropRetained() throws Exception + { + createTable("CREATE TABLE %s (pk bigint PRIMARY KEY, m map)"); + ColumnFamilyStore cfs = getCurrentColumnFamilyStore(); + cfs.disableAutoCompaction(); + + execute("UPDATE %s USING TIMESTAMP " + FUTURE_TS + " SET m['a'] = 1 WHERE pk = 0"); + flush(); + execute("UPDATE %s USING TIMESTAMP " + FUTURE_TS + " SET m['b'] = 2 WHERE pk = 0"); + flush(); + + alterTable("ALTER TABLE %s DROP m"); + + commitCompaction(cfs, cfs.getLiveSSTables(), false, cfs.getDefaultGcBefore(FBUtilities.nowInSeconds())); + + alterTable("ALTER TABLE %s ADD m map"); + assertRows(execute("SELECT m FROM %s WHERE pk = 0"), row(map("a", 1L, "b", 2L))); + } + /** DROP then ADD: pre-drop cells are filtered, post-re-add cells survive — the resurrection shape. */ @Test public void droppedColumnReAdded() throws Exception From 78cfeff9f3f32fe6f81fba851ee5f7171ccfe382 Mon Sep 17 00:00:00 2001 From: 1fanwang <1fannnw@gmail.com> Date: Wed, 26 Aug 2026 12:45:29 -0400 Subject: [PATCH 2/7] CASSANDRA-21607: Verify compacted table remains readable Signed-off-by: 1fanwang <1fannnw@gmail.com> --- .../DroppedColumnDifferentialCompactionTest.java | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/test/unit/org/apache/cassandra/db/compaction/differential/DroppedColumnDifferentialCompactionTest.java b/test/unit/org/apache/cassandra/db/compaction/differential/DroppedColumnDifferentialCompactionTest.java index 6d2b0242bb1c..b967d8fbb9cc 100644 --- a/test/unit/org/apache/cassandra/db/compaction/differential/DroppedColumnDifferentialCompactionTest.java +++ b/test/unit/org/apache/cassandra/db/compaction/differential/DroppedColumnDifferentialCompactionTest.java @@ -99,11 +99,11 @@ public void cellsNewerThanDropRetained() throws Exception @Test public void complexCellsNewerThanDropRetained() throws Exception { - createTable("CREATE TABLE %s (pk bigint PRIMARY KEY, m map)"); + createTable("CREATE TABLE %s (pk bigint PRIMARY KEY, v bigint, m map)"); ColumnFamilyStore cfs = getCurrentColumnFamilyStore(); cfs.disableAutoCompaction(); - execute("UPDATE %s USING TIMESTAMP " + FUTURE_TS + " SET m['a'] = 1 WHERE pk = 0"); + execute("UPDATE %s USING TIMESTAMP " + FUTURE_TS + " SET v = 7, m['a'] = 1 WHERE pk = 0"); flush(); execute("UPDATE %s USING TIMESTAMP " + FUTURE_TS + " SET m['b'] = 2 WHERE pk = 0"); flush(); @@ -111,6 +111,7 @@ public void complexCellsNewerThanDropRetained() throws Exception alterTable("ALTER TABLE %s DROP m"); commitCompaction(cfs, cfs.getLiveSSTables(), false, cfs.getDefaultGcBefore(FBUtilities.nowInSeconds())); + assertRows(execute("SELECT * FROM %s"), row(0L, 7L)); alterTable("ALTER TABLE %s ADD m map"); assertRows(execute("SELECT m FROM %s WHERE pk = 0"), row(map("a", 1L, "b", 2L))); From 3e493c09afab68c5348d9622c6d76cab9543943f Mon Sep 17 00:00:00 2001 From: 1fanwang <1fannnw@gmail.com> Date: Wed, 26 Aug 2026 14:12:26 -0400 Subject: [PATCH 3/7] CASSANDRA-21607: Cover reads before compaction Signed-off-by: 1fanwang <1fannnw@gmail.com> --- src/java/org/apache/cassandra/db/rows/Row.java | 15 ++++++++------- .../DroppedColumnDifferentialCompactionTest.java | 1 + 2 files changed, 9 insertions(+), 7 deletions(-) diff --git a/src/java/org/apache/cassandra/db/rows/Row.java b/src/java/org/apache/cassandra/db/rows/Row.java index a727246cd555..6a0f9abcda50 100644 --- a/src/java/org/apache/cassandra/db/rows/Row.java +++ b/src/java/org/apache/cassandra/db/rows/Row.java @@ -859,8 +859,8 @@ private static class ColumnDataReducer extends RowMergeIterator.Reducer>> complexCells; + private final ComplexColumnData.Builder complexBuilder; + private final List>> complexCells; private final CellReducer cellReducer; public ColumnDataReducer(int size, boolean hasComplex) @@ -916,11 +916,12 @@ protected ColumnData getReduced() } else { - if (complexBuilder == null) - { - complexBuilder = ComplexColumnData.builder(); - complexCells = new ArrayList<>(versions.length); - } + ComplexColumnData.Builder complexBuilder = this.complexBuilder != null + ? this.complexBuilder + : ComplexColumnData.builder(); + List>> complexCells = this.complexCells != null + ? this.complexCells + : new ArrayList<>(versions.length); complexBuilder.newColumn(column); complexCells.clear(); DeletionTime complexDeletion = DeletionTime.LIVE; diff --git a/test/unit/org/apache/cassandra/db/compaction/differential/DroppedColumnDifferentialCompactionTest.java b/test/unit/org/apache/cassandra/db/compaction/differential/DroppedColumnDifferentialCompactionTest.java index b967d8fbb9cc..deff9e19ec98 100644 --- a/test/unit/org/apache/cassandra/db/compaction/differential/DroppedColumnDifferentialCompactionTest.java +++ b/test/unit/org/apache/cassandra/db/compaction/differential/DroppedColumnDifferentialCompactionTest.java @@ -110,6 +110,7 @@ public void complexCellsNewerThanDropRetained() throws Exception alterTable("ALTER TABLE %s DROP m"); + assertRows(execute("SELECT * FROM %s"), row(0L, 7L)); commitCompaction(cfs, cfs.getLiveSSTables(), false, cfs.getDefaultGcBefore(FBUtilities.nowInSeconds())); assertRows(execute("SELECT * FROM %s"), row(0L, 7L)); From edc35c2368fae7fa7175e9907fe0fe27f8e84d35 Mon Sep 17 00:00:00 2001 From: 1fanwang <1fannnw@gmail.com> Date: Wed, 26 Aug 2026 18:02:32 -0400 Subject: [PATCH 4/7] CASSANDRA-21607: Filter dropped columns from read responses Signed-off-by: 1fanwang <1fannnw@gmail.com> --- src/java/org/apache/cassandra/db/ReadCommand.java | 13 +++++++++++++ .../DroppedColumnDifferentialCompactionTest.java | 8 ++++---- 2 files changed, 17 insertions(+), 4 deletions(-) diff --git a/src/java/org/apache/cassandra/db/ReadCommand.java b/src/java/org/apache/cassandra/db/ReadCommand.java index d848e23b6b61..aff6a0dc251c 100644 --- a/src/java/org/apache/cassandra/db/ReadCommand.java +++ b/src/java/org/apache/cassandra/db/ReadCommand.java @@ -448,6 +448,19 @@ public ReadResponse createResponse(UnfilteredPartitionIterator iterator, Repaire // ends equal, and there are no dangling RT bound in any partition. iterator = RTBoundValidator.validate(iterator, Stage.PROCESSED, true); + if (!metadata().droppedColumns.isEmpty()) + { + ColumnFilter selection = columnFilter(); + iterator = Transformation.apply(iterator, new Transformation() + { + @Override + protected UnfilteredRowIterator applyToPartition(UnfilteredRowIterator partition) + { + return ReadCommand.this.clusteringIndexFilter(partition.partitionKey()).filterNotIndexed(selection, partition); + } + }); + } + return isDigestQuery() ? ReadResponse.createDigestResponse(iterator, this) : ReadResponse.createDataResponse(iterator, this, rdi); diff --git a/test/unit/org/apache/cassandra/db/compaction/differential/DroppedColumnDifferentialCompactionTest.java b/test/unit/org/apache/cassandra/db/compaction/differential/DroppedColumnDifferentialCompactionTest.java index deff9e19ec98..cb9b503a49b1 100644 --- a/test/unit/org/apache/cassandra/db/compaction/differential/DroppedColumnDifferentialCompactionTest.java +++ b/test/unit/org/apache/cassandra/db/compaction/differential/DroppedColumnDifferentialCompactionTest.java @@ -99,20 +99,20 @@ public void cellsNewerThanDropRetained() throws Exception @Test public void complexCellsNewerThanDropRetained() throws Exception { - createTable("CREATE TABLE %s (pk bigint PRIMARY KEY, v bigint, m map)"); + createTable("CREATE TABLE %s (pk bigint PRIMARY KEY, m map)"); ColumnFamilyStore cfs = getCurrentColumnFamilyStore(); cfs.disableAutoCompaction(); - execute("UPDATE %s USING TIMESTAMP " + FUTURE_TS + " SET v = 7, m['a'] = 1 WHERE pk = 0"); + execute("UPDATE %s USING TIMESTAMP " + FUTURE_TS + " SET m['a'] = 1 WHERE pk = 0"); flush(); execute("UPDATE %s USING TIMESTAMP " + FUTURE_TS + " SET m['b'] = 2 WHERE pk = 0"); flush(); alterTable("ALTER TABLE %s DROP m"); - assertRows(execute("SELECT * FROM %s"), row(0L, 7L)); + assertTrue(executeNet("SELECT * FROM %s").all().isEmpty()); commitCompaction(cfs, cfs.getLiveSSTables(), false, cfs.getDefaultGcBefore(FBUtilities.nowInSeconds())); - assertRows(execute("SELECT * FROM %s"), row(0L, 7L)); + assertTrue(executeNet("SELECT * FROM %s").all().isEmpty()); alterTable("ALTER TABLE %s ADD m map"); assertRows(execute("SELECT m FROM %s WHERE pk = 0"), row(map("a", 1L, "b", 2L))); From 4920e8007893f89f8f265181a5b2396a9cd8e650 Mon Sep 17 00:00:00 2001 From: 1fanwang <1fannnw@gmail.com> Date: Thu, 27 Aug 2026 00:08:45 -0400 Subject: [PATCH 5/7] CASSANDRA-21607: Filter dropped cells before read limits Signed-off-by: 1fanwang <1fannnw@gmail.com> --- .../org/apache/cassandra/db/ReadCommand.java | 26 +++++++++---------- ...oppedColumnDifferentialCompactionTest.java | 17 ++++++++++++ 2 files changed, 30 insertions(+), 13 deletions(-) diff --git a/src/java/org/apache/cassandra/db/ReadCommand.java b/src/java/org/apache/cassandra/db/ReadCommand.java index aff6a0dc251c..a0e7f86c6958 100644 --- a/src/java/org/apache/cassandra/db/ReadCommand.java +++ b/src/java/org/apache/cassandra/db/ReadCommand.java @@ -448,19 +448,6 @@ public ReadResponse createResponse(UnfilteredPartitionIterator iterator, Repaire // ends equal, and there are no dangling RT bound in any partition. iterator = RTBoundValidator.validate(iterator, Stage.PROCESSED, true); - if (!metadata().droppedColumns.isEmpty()) - { - ColumnFilter selection = columnFilter(); - iterator = Transformation.apply(iterator, new Transformation() - { - @Override - protected UnfilteredRowIterator applyToPartition(UnfilteredRowIterator partition) - { - return ReadCommand.this.clusteringIndexFilter(partition.partitionKey()).filterNotIndexed(selection, partition); - } - }); - } - return isDigestQuery() ? ReadResponse.createDigestResponse(iterator, this) : ReadResponse.createDataResponse(iterator, this, rdi); @@ -575,6 +562,19 @@ public UnfilteredPartitionIterator executeLocally(ReadExecutionController execut */ iterator = filter.filter(iterator, nowInSec()); + if (!metadata().droppedColumns.isEmpty()) + { + ColumnFilter selection = columnFilter(); + iterator = Transformation.apply(iterator, new Transformation() + { + @Override + protected UnfilteredRowIterator applyToPartition(UnfilteredRowIterator partition) + { + return ReadCommand.this.clusteringIndexFilter(partition.partitionKey()).filterNotIndexed(selection, partition); + } + }); + } + // apply the limits/row counter; this transformation is stopping and would close the iterator as soon // as the count is observed; if that happens in the middle of an open RT, its end bound will not be included. // If tracking repaired data, the counter is needed for overreading repaired data, otherwise we can diff --git a/test/unit/org/apache/cassandra/db/compaction/differential/DroppedColumnDifferentialCompactionTest.java b/test/unit/org/apache/cassandra/db/compaction/differential/DroppedColumnDifferentialCompactionTest.java index cb9b503a49b1..c237463e1751 100644 --- a/test/unit/org/apache/cassandra/db/compaction/differential/DroppedColumnDifferentialCompactionTest.java +++ b/test/unit/org/apache/cassandra/db/compaction/differential/DroppedColumnDifferentialCompactionTest.java @@ -118,6 +118,23 @@ public void complexCellsNewerThanDropRetained() throws Exception assertRows(execute("SELECT m FROM %s WHERE pk = 0"), row(map("a", 1L, "b", 2L))); } + @Test + public void droppedComplexCellsDoNotConsumeReadLimit() throws Throwable + { + createTable("CREATE TABLE %s (pk bigint, ck bigint, v bigint, m map, PRIMARY KEY (pk, ck))"); + + execute("UPDATE %s USING TIMESTAMP " + FUTURE_TS + " SET m['a'] = 1 WHERE pk = 0 AND ck = 0"); + execute("UPDATE %s USING TIMESTAMP " + FUTURE_TS + " SET m['a'] = 1 WHERE pk = 0 AND ck = 2"); + execute("UPDATE %s SET v = 7 WHERE pk = 0 AND ck = 1"); + flush(); + + alterTable("ALTER TABLE %s DROP m"); + + assertRowsNet(executeNet("SELECT * FROM %s WHERE pk = 0 LIMIT 1"), row(0L, 1L, 7L)); + assertRowsNet(executeNet("SELECT * FROM %s WHERE pk = 0 ORDER BY ck DESC LIMIT 1"), row(0L, 1L, 7L)); + assertRowsNet(executeNetWithPaging("SELECT * FROM %s WHERE pk = 0", 1), row(0L, 1L, 7L)); + } + /** DROP then ADD: pre-drop cells are filtered, post-re-add cells survive — the resurrection shape. */ @Test public void droppedColumnReAdded() throws Exception From 40e5a0c3bb28b1dbb1e47dca72fe570efd0c23e8 Mon Sep 17 00:00:00 2001 From: 1fanwang <1fannnw@gmail.com> Date: Thu, 27 Aug 2026 01:15:20 -0400 Subject: [PATCH 6/7] CASSANDRA-21607: Use current schema for dropped columns Signed-off-by: 1fanwang <1fannnw@gmail.com> --- .../org/apache/cassandra/db/ReadCommand.java | 2 +- ...oppedColumnDifferentialCompactionTest.java | 19 +++++++++++++++++++ 2 files changed, 20 insertions(+), 1 deletion(-) diff --git a/src/java/org/apache/cassandra/db/ReadCommand.java b/src/java/org/apache/cassandra/db/ReadCommand.java index a0e7f86c6958..bcb14e0aab8e 100644 --- a/src/java/org/apache/cassandra/db/ReadCommand.java +++ b/src/java/org/apache/cassandra/db/ReadCommand.java @@ -562,7 +562,7 @@ public UnfilteredPartitionIterator executeLocally(ReadExecutionController execut */ iterator = filter.filter(iterator, nowInSec()); - if (!metadata().droppedColumns.isEmpty()) + if (!cfs.metadata().droppedColumns.isEmpty()) { ColumnFilter selection = columnFilter(); iterator = Transformation.apply(iterator, new Transformation() diff --git a/test/unit/org/apache/cassandra/db/compaction/differential/DroppedColumnDifferentialCompactionTest.java b/test/unit/org/apache/cassandra/db/compaction/differential/DroppedColumnDifferentialCompactionTest.java index c237463e1751..74faa4493702 100644 --- a/test/unit/org/apache/cassandra/db/compaction/differential/DroppedColumnDifferentialCompactionTest.java +++ b/test/unit/org/apache/cassandra/db/compaction/differential/DroppedColumnDifferentialCompactionTest.java @@ -23,7 +23,16 @@ import org.apache.cassandra.db.ColumnFamilyStore; import org.apache.cassandra.db.LivenessInfo; import org.apache.cassandra.db.Mutation; +import org.apache.cassandra.db.ReadCommandVerbHandler; +import org.apache.cassandra.db.SinglePartitionReadCommand; +import org.apache.cassandra.db.Slices; +import org.apache.cassandra.db.filter.ClusteringIndexSliceFilter; +import org.apache.cassandra.db.filter.ColumnFilter; +import org.apache.cassandra.db.filter.DataLimits; +import org.apache.cassandra.db.filter.RowFilter; import org.apache.cassandra.db.partitions.PartitionUpdate; +import org.apache.cassandra.schema.TableMetadata; +import org.apache.cassandra.utils.ByteBufferUtil; import org.apache.cassandra.utils.FBUtilities; import static org.junit.Assert.assertEquals; @@ -102,6 +111,7 @@ public void complexCellsNewerThanDropRetained() throws Exception createTable("CREATE TABLE %s (pk bigint PRIMARY KEY, m map)"); ColumnFamilyStore cfs = getCurrentColumnFamilyStore(); cfs.disableAutoCompaction(); + TableMetadata metadataBeforeDrop = cfs.metadata(); execute("UPDATE %s USING TIMESTAMP " + FUTURE_TS + " SET m['a'] = 1 WHERE pk = 0"); flush(); @@ -110,6 +120,15 @@ public void complexCellsNewerThanDropRetained() throws Exception alterTable("ALTER TABLE %s DROP m"); + SinglePartitionReadCommand commandDeserializedBeforeDrop = + SinglePartitionReadCommand.create(metadataBeforeDrop, + FBUtilities.nowInSeconds(), + ColumnFilter.all(cfs.metadata()), + RowFilter.none(), + DataLimits.NONE, + metadataBeforeDrop.partitioner.decorateKey(ByteBufferUtil.bytes(0L)), + new ClusteringIndexSliceFilter(Slices.ALL, false)); + ReadCommandVerbHandler.instance.doRead(commandDeserializedBeforeDrop, false); assertTrue(executeNet("SELECT * FROM %s").all().isEmpty()); commitCompaction(cfs, cfs.getLiveSSTables(), false, cfs.getDefaultGcBefore(FBUtilities.nowInSeconds())); assertTrue(executeNet("SELECT * FROM %s").all().isEmpty()); From 9b96953a063e1e303be44a76c6fde946adf20c14 Mon Sep 17 00:00:00 2001 From: 1fanwang <1fannnw@gmail.com> Date: Sat, 29 Aug 2026 10:18:32 -0400 Subject: [PATCH 7/7] Make wildcard ColumnFilter honour its fixed column set fetches() answered true for every column, including ones dropped after the filter was built, while fetchedColumns() returned the fixed set it was built with. A read crossing a schema change then deserialized cells it could not serialize: IllegalStateException: [m] is not a subset of [] Skipping them here lets ReadCommand stop filtering. Signed-off-by: 1fanwang <1fannnw@gmail.com> --- src/java/org/apache/cassandra/db/ReadCommand.java | 13 ------------- .../apache/cassandra/db/filter/ColumnFilter.java | 5 ++++- 2 files changed, 4 insertions(+), 14 deletions(-) diff --git a/src/java/org/apache/cassandra/db/ReadCommand.java b/src/java/org/apache/cassandra/db/ReadCommand.java index bcb14e0aab8e..d848e23b6b61 100644 --- a/src/java/org/apache/cassandra/db/ReadCommand.java +++ b/src/java/org/apache/cassandra/db/ReadCommand.java @@ -562,19 +562,6 @@ public UnfilteredPartitionIterator executeLocally(ReadExecutionController execut */ iterator = filter.filter(iterator, nowInSec()); - if (!cfs.metadata().droppedColumns.isEmpty()) - { - ColumnFilter selection = columnFilter(); - iterator = Transformation.apply(iterator, new Transformation() - { - @Override - protected UnfilteredRowIterator applyToPartition(UnfilteredRowIterator partition) - { - return ReadCommand.this.clusteringIndexFilter(partition.partitionKey()).filterNotIndexed(selection, partition); - } - }); - } - // apply the limits/row counter; this transformation is stopping and would close the iterator as soon // as the count is observed; if that happens in the middle of an open RT, its end bound will not be included. // If tracking repaired data, the counter is needed for overreading repaired data, otherwise we can diff --git a/src/java/org/apache/cassandra/db/filter/ColumnFilter.java b/src/java/org/apache/cassandra/db/filter/ColumnFilter.java index e28d1e0161bb..64f7a8df4e83 100644 --- a/src/java/org/apache/cassandra/db/filter/ColumnFilter.java +++ b/src/java/org/apache/cassandra/db/filter/ColumnFilter.java @@ -598,7 +598,10 @@ public boolean allFetchedColumnsAreQueried() @Override public boolean fetches(ColumnMetadata column) { - return true; + // A column dropped after this filter was built is absent from fetchedAndQueried, and + // answering true for it would fetch cells that fetchedColumns() cannot then encode as a + // superset. Use allEver() to read dropped columns deliberately. + return fetchedAndQueried.contains(column); } @Override