Filter dropped-column cells in cursor compaction; read static-row presence from sstable headers

After ALTER TABLE ... DROP of a column, older sstables can still carry cells for it with
timestamps at or before the drop. The iterator path discards those cells at
deserialization (UnfilteredSerializer.readSimpleColumn, via DeserializationHelper.isDropped);
the cursor path did not, so dropped data could survive cursor compaction and reappear if
the same column name were re-added later. CellCursor now tracks each column's drop
horizon (droppedTimesArray, built once per superset) and discards a cell's value when its
timestamp is at or before that horizon, mirroring the iterator exactly.
DroppedColumnDifferentialCompactionTest pins dropped regular columns, cells newer than
the drop surviving, a drop-then-re-add resurrection shape, and dropped static columns.

Also fixes hasStaticColumns: it was read from current table metadata, but after ALTER
TABLE ... DROP of the last static column, current metadata has no static columns while
older input sstables legitimately still carry static rows. It's now derived from the
union of the input sstables' own headers instead.

The dropped-column filter initially used a plain boolean return from
CellCursor.readCellHeader, conflating two distinct outcomes: a genuine valueless cell
(tombstone) versus no cell at all (every remaining column in the row was
dropped-filtered). When a dropped column was also a row's last column, the reader
returned false without clearing cellFlags, leaving CursorCompactor.mergeCells reading a
stale hasValue=true flag against a state that said no value was available ("Flags and
state contradict"). Fixed with a tri-state return (-1/0/1): 1 = has value, 0 = valueless
cell (still a real cell to merge), -1 = no cell remains at all. The outer state machine
handles -1 by skipping straight past the row via continueReading(), while 0 keeps the
existing last-column end-of-cells handling. Verified by rerunning the full differential
suite, including DroppedColumnDifferentialCompactionTest.

CellCursor.cellDroppedTime was a field but only ever set and read within
readCellHeader() itself; made it a local variable.
This commit is contained in:
Jon Haddad 2026-07-27 11:23:09 -07:00
parent 627b83cba8
commit 272123524a
3 changed files with 250 additions and 49 deletions

View File

@ -297,7 +297,14 @@ public class CursorCompactor extends CompactionInfo.Holder
this.activeCompactions.beginCompaction(this); // note that CompactionTask also calls this, but CT only creates CompactionIterator with a NOOP ActiveCompactions
TableMetadata metadata = metadata();
this.hasStaticColumns = metadata.hasStaticColumns();
// the INPUT headers decide whether static rows can occur in this merge (and the output
// header, SerializationHeader.make, is their union): after ALTER TABLE ... DROP of the
// last static column, current metadata has no static columns but older sstables
// legitimately still carry static rows
boolean anyStaticColumns = false;
for (SSTableReader sstable : this.sstables)
anyStaticColumns |= sstable.header.hasStatic();
this.hasStaticColumns = anyStaticColumns;
/**
* Pipeline should end up similar to the one in {@link CompactionIterator}:
* [MERGED -> ?TopPartitionTracker -> GarbageSkipper -> Purger -> org.apache.cassandra.db.transform.DuplicateRowChecker -> Abortable] -> next()

View File

@ -19,6 +19,8 @@
package org.apache.cassandra.io.sstable;
import java.io.IOException;
import java.nio.ByteBuffer;
import java.util.Map;
import com.google.common.collect.ImmutableList;
@ -44,6 +46,7 @@ import org.apache.cassandra.io.util.DataOutputPlus;
import org.apache.cassandra.io.util.RandomAccessReader;
import org.apache.cassandra.io.util.ResizableByteBuffer;
import org.apache.cassandra.schema.ColumnMetadata;
import org.apache.cassandra.schema.DroppedColumn;
import org.apache.cassandra.schema.TableMetadata;
import org.apache.cassandra.schema.TableMetadataRef;
import org.apache.cassandra.tools.Util;
@ -105,6 +108,10 @@ public class SSTableCursorReader implements AutoCloseable
public ColumnMetadata cellColumn;
private ColumnMetadata[] columnsArray;
private AbstractType<?>[] cellTypeArray;
// Per-column drop horizon (microseconds), Long.MIN_VALUE when the column was never
// dropped; cells with timestamp <= the horizon are discarded, mirroring
// DeserializationHelper.isDropped. Built once per superset, like cellTypeArray.
private long[] droppedTimesArray;
// Remaining PRESENT columns of this row as a bitmask over columnsArray indices.
// Garbage-free sparse-row iteration: rows that do not contain every header column
@ -126,10 +133,13 @@ public class SSTableCursorReader implements AutoCloseable
this.columns = columns;
columnsArray = columns.toArray(COLUMN_METADATA_TYPE);
cellTypeArray = new AbstractType<?>[columnsArray.length];
droppedTimesArray = new long[columnsArray.length];
for (int i = 0; i < columnsArray.length; i++)
{
ColumnMetadata cellColumn = columnsArray[i];
cellTypeArray[i] = serializationHeader.getType(cellColumn);
DroppedColumn dropped = droppedColumns.get(cellColumn.name.bytes);
droppedTimesArray[i] = dropped == null ? Long.MIN_VALUE : dropped.droppedTime;
}
columnsSize = columns.size();
}
@ -191,58 +201,84 @@ public class SSTableCursorReader implements AutoCloseable
/**
* For Cell deserialization see {@link Cell.Serializer#deserialize}
*
* @return true if the cell has a value, false otherwise
* Dropped-column filtering happens here, mirroring the iterator's deserialization:
* cells of a dropped column written at or before the drop are consumed and never
* surfaced; the loop advances to the next column in that case. A dropped column
* that turns out to be the row's last remaining column leaves NO cell at all for
* this position (distinct from a genuine valueless cell/tombstone), which is why
* this returns a tri-state rather than a plain hasValue boolean: the caller must
* skip straight past the row/unfiltered end rather than stopping at a cell that
* doesn't exist.
*
* @return 1 if the next cell has a value, 0 if it has none (tombstone), -1 if no
* cell remains in this row (all trailing columns were dropped-filtered)
*/
boolean readCellHeader() throws IOException
int readCellHeader() throws IOException
{
if (!hasNext()) throw new IllegalStateException();
// HOTSPOT: suprisingly expensive
int currIndex;
if (columnsSize >= 64)
for (;;)
{
// columnsRemain() (via hasNext() above) parked presentWordIndex on a
// non-empty word; same low-to-high bit walk as the single-mask path below
long word = presentWords[presentWordIndex];
currIndex = (presentWordIndex << 6) + Long.numberOfTrailingZeros(word);
presentWords[presentWordIndex] = word & (word - 1);
// HOTSPOT: suprisingly expensive
int currIndex;
if (columnsSize >= 64)
{
// columnsRemain() (via hasNext() above, or the loop-continue path below)
// parked presentWordIndex on a non-empty word; same low-to-high bit walk
// as the single-mask path below
long word = presentWords[presentWordIndex];
currIndex = (presentWordIndex << 6) + Long.numberOfTrailingZeros(word);
presentWords[presentWordIndex] = word & (word - 1);
}
else
{
// Bit i of presentMask corresponds to the i-th column of the superset in
// its iteration order the SAME order the serializer assigned bits and
// the same order cells appear on disk. Walking bits low-to-high therefore
// visits cells in exactly their on-disk order:
// numberOfTrailingZeros = index of the lowest set bit (next present column)
// x & (x - 1) = clears that lowest set bit (subtracting 1 borrows
// through the trailing zeros; the AND kills both)
currIndex = Long.numberOfTrailingZeros(presentMask);
presentMask &= presentMask - 1;
}
cellColumn = columnsArray[currIndex];
cellType = cellTypeArray[currIndex];
long cellDroppedTime = droppedTimesArray[currIndex];
cellFlags = dataReader.readUnsignedByte();
// TODO: specialize common case where flags == HAS_VALUE | USE_ROW_TS?
boolean hasValue = Cell.Serializer.hasValue(cellFlags);
boolean isDeleted = Cell.Serializer.isDeleted(cellFlags);
boolean isExpiring = Cell.Serializer.isExpiring(cellFlags);
boolean useRowTimestamp = Cell.Serializer.useRowTimestamp(cellFlags);
boolean useRowTTL = Cell.Serializer.useRowTTL(cellFlags);
long timestamp = useRowTimestamp ? rowLiveness.timestamp() : serializationHeader.readTimestamp(dataReader);
long localDeletionTime = useRowTTL
? rowLiveness.localExpirationTime()
: (isDeleted || isExpiring ? serializationHeader.readLocalDeletionTime(dataReader) : Cell.NO_DELETION_TIME);
int ttl = useRowTTL ? rowLiveness.ttl() : (isExpiring ? serializationHeader.readTTL(dataReader) : Cell.NO_TTL);
localDeletionTime = Cell.decodeLocalDeletionTime(localDeletionTime, ttl, deserializationHelper);
cellLiveness.reset(timestamp, ttl, localDeletionTime);
cellPath = cellColumn.isComplex()
? cellColumn.cellPathSerializer().deserialize(dataReader)
: null;
if (hasDroppedColumns && timestamp <= cellDroppedTime)
{
// mirror UnfilteredSerializer.readSimpleColumn: cells of a dropped column
// written at or before the drop are discarded on read
if (hasValue)
cellType.skipValue(dataReader);
if (!hasNext())
return -1; // no cell remains; caller must skip past the row/unfiltered end
continue;
}
return hasValue ? 1 : 0;
}
else
{
// Bit i of presentMask corresponds to the i-th column of the superset in
// its iteration order the SAME order the serializer assigned bits and
// the same order cells appear on disk. Walking bits low-to-high therefore
// visits cells in exactly their on-disk order:
// numberOfTrailingZeros = index of the lowest set bit (next present column)
// x & (x - 1) = clears that lowest set bit (subtracting 1 borrows
// through the trailing zeros; the AND kills both)
currIndex = Long.numberOfTrailingZeros(presentMask);
presentMask &= presentMask - 1;
}
cellColumn = columnsArray[currIndex];
cellType = cellTypeArray[currIndex];
cellFlags = dataReader.readUnsignedByte();
// TODO: specialize common case where flags == HAS_VALUE | USE_ROW_TS?
boolean hasValue = Cell.Serializer.hasValue(cellFlags);
boolean isDeleted = Cell.Serializer.isDeleted(cellFlags);
boolean isExpiring = Cell.Serializer.isExpiring(cellFlags);
boolean useRowTimestamp = Cell.Serializer.useRowTimestamp(cellFlags);
boolean useRowTTL = Cell.Serializer.useRowTTL(cellFlags);
long timestamp = useRowTimestamp ? rowLiveness.timestamp() : serializationHeader.readTimestamp(dataReader);
long localDeletionTime = useRowTTL
? rowLiveness.localExpirationTime()
: (isDeleted || isExpiring ? serializationHeader.readLocalDeletionTime(dataReader) : Cell.NO_DELETION_TIME);
int ttl = useRowTTL ? rowLiveness.ttl() : (isExpiring ? serializationHeader.readTTL(dataReader) : Cell.NO_TTL);
localDeletionTime = Cell.decodeLocalDeletionTime(localDeletionTime, ttl, deserializationHelper);
cellLiveness.reset(timestamp, ttl, localDeletionTime);
cellPath = cellColumn.isComplex()
? cellColumn.cellPathSerializer().deserialize(dataReader)
: null;
return hasValue;
}
}
@ -250,6 +286,13 @@ public class SSTableCursorReader implements AutoCloseable
private final AbstractType<?>[] clusteringColumnTypes;
private final DeserializationHelper deserializationHelper;
private final SerializationHeader serializationHeader;
// Dropped-column filtering, mirroring DeserializationHelper.isDropped /
// isDroppedComplexDeletion: cells (and complex deletions) of a dropped column written at or
// before the drop are discarded at deserialization on the iterator path
// (UnfilteredSerializer.readSimpleColumn/readComplexColumn), so the cursor must never
// surface them either.
private final Map<ByteBuffer, DroppedColumn> droppedColumns;
private final boolean hasDroppedColumns;
// need to be closed
private final SSTableReader ssTableReader;
@ -300,9 +343,14 @@ public class SSTableCursorReader implements AutoCloseable
}
deserializationHelper = new DeserializationHelper(metadata, version.correspondingMessagingVersion(), DeserializationHelper.Flag.LOCAL, null);
serializationHeader = reader.header;
droppedColumns = metadata.droppedColumns;
hasDroppedColumns = !droppedColumns.isEmpty();
dataReader = reader.openDataReaderForScan(diskAccessMode);
hasStaticColumns = metadata.hasStaticColumns();
// the HEADER decides whether this sstable can contain static rows: after
// ALTER TABLE ... DROP of the last static column, current metadata has no static
// columns but older sstables legitimately still carry static rows
hasStaticColumns = serializationHeader.hasStatic();
}
@Override
@ -546,7 +594,16 @@ public class SSTableCursorReader implements AutoCloseable
if (state != State.CELL_HEADER_START) throw new IllegalStateException();
try
{
if (cellCursor.readCellHeader())
int cell = cellCursor.readCellHeader();
if (cell < 0)
{
// no cell surfaced at all (every remaining column was dropped-column
// filtered): nothing is current, so advance straight past the would-be
// CELL_END stop instead of surfacing a cell-less position
checkNextFlagsAfterCellValuesEnd();
return continueReading();
}
if (cell > 0)
{
return state = State.CELL_VALUE_START;
}

View File

@ -0,0 +1,137 @@
/*
* 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.compaction.differential;
import org.junit.Test;
import org.apache.cassandra.db.ColumnFamilyStore;
/**
* Dropped-column scenarios: the iterator path discards cells (and complex deletions) of a
* dropped column at deserialization when their timestamp is at or before the drop
* (UnfilteredSerializer.readSimpleColumn/readComplexColumn via DeserializationHelper.isDropped /
* isDroppedComplexDeletion). The cursor path must filter identically, or dropped data survives
* cursor compaction and is resurrected by ALTER TABLE ... ADD of the same column.
*
* Timestamp determinism: ClientState timestamps are strictly increasing per node, so a write
* executed before ALTER TABLE DROP always carries a timestamp strictly below droppedTime;
* writes that must survive the drop use an explicit far-future USING TIMESTAMP.
*/
public class DroppedColumnDifferentialCompactionTest extends DifferentialCompactionTester
{
/** Far in the future relative to any wall-clock droppedTime (year ~2100, microseconds). */
private static final long FUTURE_TS = 4102444800_000000L;
@Test
public void droppedRegularColumn() throws Exception
{
createTable("CREATE TABLE %s (pk bigint, ck bigint, v1 bigint, v2 text, PRIMARY KEY (pk, ck))");
ColumnFamilyStore cfs = getCurrentColumnFamilyStore();
cfs.disableAutoCompaction();
for (long pk = 0; pk < 5; pk++)
for (long ck = 0; ck < 10; ck++)
execute("INSERT INTO %s (pk, ck, v1, v2) VALUES (?, ?, ?, ?)", pk, ck, ck, "dropped-" + ck);
flush();
alterTable("ALTER TABLE %s DROP v2");
// second generation written after the drop, overlapping the first
for (long pk = 0; pk < 5; pk++)
for (long ck = 5; ck < 15; ck++)
execute("INSERT INTO %s (pk, ck, v1) VALUES (?, ?, ?)", pk, ck, ck + 100);
flush();
assertCursorMatchesIterator(cfs);
}
/**
* Cells written with a timestamp ABOVE droppedTime survive compaction on both paths
* (isDropped is timestamp-gated, not column-gated) pins the comparison direction.
*/
@Test
public void cellsNewerThanDropRetained() throws Exception
{
createTable("CREATE TABLE %s (pk bigint, ck bigint, v1 bigint, v2 text, PRIMARY KEY (pk, ck))");
ColumnFamilyStore cfs = getCurrentColumnFamilyStore();
cfs.disableAutoCompaction();
for (long ck = 0; ck < 10; ck++)
execute("INSERT INTO %s (pk, ck, v1, v2) VALUES (0, ?, ?, ?) USING TIMESTAMP " + FUTURE_TS,
ck, ck, "survives-" + ck);
// and some cells that do not survive
for (long ck = 0; ck < 10; ck++)
execute("INSERT INTO %s (pk, ck, v2) VALUES (1, ?, ?)", ck, "dropped-" + ck);
flush();
alterTable("ALTER TABLE %s DROP v2");
for (long ck = 0; ck < 10; ck++)
execute("INSERT INTO %s (pk, ck, v1) VALUES (2, ?, ?)", ck, ck);
flush();
assertCursorMatchesIterator(cfs);
}
/** DROP then ADD: pre-drop cells are filtered, post-re-add cells survive — the resurrection shape. */
@Test
public void droppedColumnReAdded() throws Exception
{
createTable("CREATE TABLE %s (pk bigint, ck bigint, v1 bigint, v2 text, PRIMARY KEY (pk, ck))");
ColumnFamilyStore cfs = getCurrentColumnFamilyStore();
cfs.disableAutoCompaction();
for (long ck = 0; ck < 10; ck++)
execute("INSERT INTO %s (pk, ck, v1, v2) VALUES (0, ?, ?, ?)", ck, ck, "old-" + ck);
flush();
alterTable("ALTER TABLE %s DROP v2");
alterTable("ALTER TABLE %s ADD v2 text");
for (long ck = 5; ck < 15; ck++)
execute("INSERT INTO %s (pk, ck, v2) VALUES (0, ?, ?)", ck, "new-" + ck);
flush();
assertCursorMatchesIterator(cfs);
}
@Test
public void droppedStaticColumn() throws Exception
{
createTable("CREATE TABLE %s (pk bigint, s text static, ck bigint, v1 bigint, PRIMARY KEY (pk, ck))");
ColumnFamilyStore cfs = getCurrentColumnFamilyStore();
cfs.disableAutoCompaction();
for (long pk = 0; pk < 5; pk++)
{
execute("INSERT INTO %s (pk, s) VALUES (?, ?)", pk, "static-" + pk);
for (long ck = 0; ck < 5; ck++)
execute("INSERT INTO %s (pk, ck, v1) VALUES (?, ?, ?)", pk, ck, ck);
}
flush();
alterTable("ALTER TABLE %s DROP s");
for (long pk = 0; pk < 5; pk++)
execute("INSERT INTO %s (pk, ck, v1) VALUES (?, ?, ?)", pk, 99L, 99L);
flush();
assertCursorMatchesIterator(cfs);
}
}