diff --git a/paimon-core/src/main/java/org/apache/paimon/append/dataevolution/DataEvolutionCompactTaskSerializer.java b/paimon-core/src/main/java/org/apache/paimon/append/dataevolution/DataEvolutionCompactTaskSerializer.java index 019cf3458c66..40ae2073dd17 100644 --- a/paimon-core/src/main/java/org/apache/paimon/append/dataevolution/DataEvolutionCompactTaskSerializer.java +++ b/paimon-core/src/main/java/org/apache/paimon/append/dataevolution/DataEvolutionCompactTaskSerializer.java @@ -40,7 +40,7 @@ public class DataEvolutionCompactTaskSerializer implements VersionedSerializer { - private static final int CURRENT_VERSION = 2; + private static final int CURRENT_VERSION = 3; private final DataFileMetaSerializer dataFileSerializer; diff --git a/paimon-core/src/main/java/org/apache/paimon/append/dataevolution/DataEvolutionNormalCompactTask.java b/paimon-core/src/main/java/org/apache/paimon/append/dataevolution/DataEvolutionNormalCompactTask.java index 1db0786d27cc..aebb1c9004e0 100644 --- a/paimon-core/src/main/java/org/apache/paimon/append/dataevolution/DataEvolutionNormalCompactTask.java +++ b/paimon-core/src/main/java/org/apache/paimon/append/dataevolution/DataEvolutionNormalCompactTask.java @@ -28,6 +28,7 @@ import org.apache.paimon.table.FileStoreTable; import org.apache.paimon.table.sink.CommitMessage; import org.apache.paimon.table.source.DataSplit; +import org.apache.paimon.types.DataField; import org.apache.paimon.types.RowType; import org.apache.paimon.utils.FileStorePathFactory; import org.apache.paimon.utils.RecordWriter; @@ -36,13 +37,17 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import java.util.HashMap; import java.util.List; +import java.util.Map; import java.util.Set; import java.util.stream.Collectors; import static org.apache.paimon.types.BlobType.fieldNamesInBlobFile; import static org.apache.paimon.types.VectorType.fieldNamesInVectorFile; import static org.apache.paimon.types.VectorType.isVectorStoreFile; +import static org.apache.paimon.utils.DataEvolutionUtils.fieldMaxSequenceNumber; +import static org.apache.paimon.utils.DataEvolutionUtils.fileFields; import static org.apache.paimon.utils.Preconditions.checkArgument; /** Compacts normal structured files of a data evolution table. */ @@ -122,8 +127,35 @@ public CommitMessage doCompact(FileStoreTable table, String commitUser) throws E dataFileMeta = dataFileMeta.assignSequenceNumber( minSequenceId(compactBefore), maxSequenceId(compactBefore)); + dataFileMeta = + dataFileMeta.withColumnMaxSequenceNumbers( + compactedColumnMaxSequenceNumbers(table, dataFileMeta)); compactAfter.add(dataFileMeta); return commitMessage(compactBefore, compactAfter); } + + private long[] compactedColumnMaxSequenceNumbers( + FileStoreTable table, DataFileMeta outputFile) { + Map fieldMaxSequences = new HashMap<>(); + for (DataFileMeta input : compactBefore) { + List inputFields = fileFields(table.schemaManager()::schema, input); + for (int inputPosition = 0; inputPosition < inputFields.size(); inputPosition++) { + fieldMaxSequences.merge( + inputFields.get(inputPosition).id(), + fieldMaxSequenceNumber(input, inputPosition, inputFields.size()), + (left, right) -> Math.max(left, right)); + } + } + + long fallbackSequence = maxSequenceId(compactBefore); + List outputFields = fileFields(table.schemaManager()::schema, outputFile); + long[] result = new long[outputFields.size()]; + for (int outputPosition = 0; outputPosition < outputFields.size(); outputPosition++) { + result[outputPosition] = + fieldMaxSequences.getOrDefault( + outputFields.get(outputPosition).id(), fallbackSequence); + } + return result; + } } diff --git a/paimon-core/src/main/java/org/apache/paimon/globalindex/DataEvolutionGlobalIndexRefreshPlanner.java b/paimon-core/src/main/java/org/apache/paimon/globalindex/DataEvolutionGlobalIndexRefreshPlanner.java index 44ed39685d24..4f0b14f46c34 100644 --- a/paimon-core/src/main/java/org/apache/paimon/globalindex/DataEvolutionGlobalIndexRefreshPlanner.java +++ b/paimon-core/src/main/java/org/apache/paimon/globalindex/DataEvolutionGlobalIndexRefreshPlanner.java @@ -41,7 +41,8 @@ import java.util.Set; import java.util.TreeMap; -import static org.apache.paimon.utils.DataEvolutionUtils.fileFieldIds; +import static org.apache.paimon.utils.DataEvolutionUtils.fieldMaxSequenceNumber; +import static org.apache.paimon.utils.DataEvolutionUtils.fileFields; /** Plans existing global index files which need refresh after data-evolution updates. */ public final class DataEvolutionGlobalIndexRefreshPlanner { @@ -81,7 +82,7 @@ public static List findIndexesToRefresh( .addIndex(i, indexMeta.rowRange(), scanSnapshotId); } - Map>, Set> fileFieldIdsCache = new HashMap<>(); + Map>, List> fileFieldsCache = new HashMap<>(); for (ManifestEntry dataEntry : dataEntries) { DataFileMeta file = dataEntry.file(); if (dataEntry.kind() != FileKind.ADD || file.firstRowId() == null) { @@ -93,12 +94,21 @@ public static List findIndexesToRefresh( continue; } - Set physicalFieldIds = - fileFieldIdsCache.computeIfAbsent( + List physicalFields = + fileFieldsCache.computeIfAbsent( Pair.of(file.schemaId(), file.writeCols()), - key -> fileFieldIds(schemaManager::schema, file)); - if (!disjoint(indexedFieldIds, physicalFieldIds)) { - group.addDataFile(file); + key -> fileFields(schemaManager::schema, file)); + long indexedMaxSequence = Long.MIN_VALUE; + for (int position = 0; position < physicalFields.size(); position++) { + if (indexedFieldIds.contains(physicalFields.get(position).id())) { + indexedMaxSequence = + Math.max( + indexedMaxSequence, + fieldMaxSequenceNumber(file, position, physicalFields.size())); + } + } + if (indexedMaxSequence != Long.MIN_VALUE) { + group.addDataFile(file, indexedMaxSequence); } } @@ -119,7 +129,7 @@ public static List findIndexesToRefresh( private static final class RefreshGroup { private final List indexes = new ArrayList<>(); - private final List dataFiles = new ArrayList<>(); + private final List dataUpdates = new ArrayList<>(); private final MergedRanges indexedRanges = new MergedRanges(); private long minScanSnapshotId = Long.MAX_VALUE; @@ -134,21 +144,23 @@ private boolean mayContainUpdate(DataFileMeta file) { && indexedRanges.intersects(file.nonNullRowIdRange()); } - private void addDataFile(DataFileMeta file) { - dataFiles.add(file); + private void addDataFile(DataFileMeta file, long maxSequenceNumber) { + if (maxSequenceNumber > minScanSnapshotId) { + dataUpdates.add(new DataUpdate(file.nonNullRowIdRange(), maxSequenceNumber)); + } } private void markIndexesToRefresh(boolean[] result) { // As scan watermarks decrease, eligible data files only grow. - dataFiles.sort(Comparator.comparingLong(DataFileMeta::maxSequenceNumber).reversed()); + dataUpdates.sort(Comparator.comparingLong(DataUpdate::maxSequenceNumber).reversed()); indexes.sort((left, right) -> Long.compare(right.scanSnapshotId, left.scanSnapshotId)); MergedRanges updatedRanges = new MergedRanges(); int nextFile = 0; for (IndexQuery index : indexes) { - while (nextFile < dataFiles.size() - && dataFiles.get(nextFile).maxSequenceNumber() > index.scanSnapshotId) { - updatedRanges.add(dataFiles.get(nextFile).nonNullRowIdRange()); + while (nextFile < dataUpdates.size() + && dataUpdates.get(nextFile).maxSequenceNumber > index.scanSnapshotId) { + updatedRanges.add(dataUpdates.get(nextFile).rowRange); nextFile++; } if (updatedRanges.intersects(index.rowRange)) { @@ -158,6 +170,21 @@ private void markIndexesToRefresh(boolean[] result) { } } + private static final class DataUpdate { + + private final Range rowRange; + private final long maxSequenceNumber; + + private DataUpdate(Range rowRange, long maxSequenceNumber) { + this.rowRange = rowRange; + this.maxSequenceNumber = maxSequenceNumber; + } + + private long maxSequenceNumber() { + return maxSequenceNumber; + } + } + private static final class IndexQuery { private final int ordinal; @@ -218,13 +245,4 @@ private static boolean matchesFields(GlobalIndexMeta meta, List field } return expectedExtraFields != null && Arrays.equals(actualExtraFields, expectedExtraFields); } - - private static boolean disjoint(Set left, Set right) { - for (Integer value : left) { - if (right.contains(value)) { - return false; - } - } - return true; - } } diff --git a/paimon-core/src/main/java/org/apache/paimon/io/BinaryDataFileMeta.java b/paimon-core/src/main/java/org/apache/paimon/io/BinaryDataFileMeta.java index be56deb20727..f48653571e36 100644 --- a/paimon-core/src/main/java/org/apache/paimon/io/BinaryDataFileMeta.java +++ b/paimon-core/src/main/java/org/apache/paimon/io/BinaryDataFileMeta.java @@ -244,6 +244,22 @@ public List writeCols() { return nullableStringArray(Fields.WRITE_COLS); } + @Nullable + @Override + public long[] columnMaxSequenceNumbers() { + int position = requiredPosition(Fields.COLUMN_MAX_SEQUENCE_NUMBERS); + InternalRow row = currentRow(); + if (row.isNullAt(position)) { + return null; + } + InternalArray array = row.getArray(position); + long[] result = new long[array.size()]; + for (int i = 0; i < array.size(); i++) { + result[i] = array.getLong(i); + } + return result; + } + public boolean containsWriteColumn(BinaryString fieldName) { int position = requiredPosition(Fields.WRITE_COLS); InternalRow row = currentRow(); @@ -281,6 +297,11 @@ public DataFileMeta assignSequenceNumber(long minSequenceNumber, long maxSequenc throw unsupportedOperation("assignSequenceNumber(long, long)"); } + @Override + public DataFileMeta withColumnMaxSequenceNumbers(long[] columnMaxSequenceNumbers) { + throw unsupportedOperation("withColumnMaxSequenceNumbers(long[])"); + } + @Override public DataFileMeta assignFirstRowId(long firstRowId) { throw unsupportedOperation("assignFirstRowId(long)"); @@ -378,6 +399,8 @@ private static class Fields { private static final int EXTERNAL_PATH = fieldIndex(DataFileMeta.EXTERNAL_PATH); private static final int FIRST_ROW_ID = fieldIndex(DataFileMeta.FIRST_ROW_ID); private static final int WRITE_COLS = fieldIndex(DataFileMeta.WRITE_COLS); + private static final int COLUMN_MAX_SEQUENCE_NUMBERS = + fieldIndex(DataFileMeta.COLUMN_MAX_SEQUENCE_NUMBERS); } /** Projected data-file schema together with its bound binary field layout. */ diff --git a/paimon-core/src/main/java/org/apache/paimon/io/DataFileMeta.java b/paimon-core/src/main/java/org/apache/paimon/io/DataFileMeta.java index f8b4e6aaf5b7..aa944ae2edbd 100644 --- a/paimon-core/src/main/java/org/apache/paimon/io/DataFileMeta.java +++ b/paimon-core/src/main/java/org/apache/paimon/io/DataFileMeta.java @@ -80,6 +80,7 @@ public interface DataFileMeta { String EXTERNAL_PATH = "_EXTERNAL_PATH"; String FIRST_ROW_ID = "_FIRST_ROW_ID"; String WRITE_COLS = "_WRITE_COLS"; + String COLUMN_MAX_SEQUENCE_NUMBERS = "_COLUMN_MAX_SEQUENCE_NUMBERS"; RowType SCHEMA = new RowType( @@ -109,7 +110,11 @@ public interface DataFileMeta { new DataField(17, EXTERNAL_PATH, newStringType(true)), new DataField(18, FIRST_ROW_ID, new BigIntType(true)), new DataField( - 19, WRITE_COLS, new ArrayType(true, newStringType(false))))); + 19, WRITE_COLS, new ArrayType(true, newStringType(false))), + new DataField( + 20, + COLUMN_MAX_SEQUENCE_NUMBERS, + new ArrayType(true, new BigIntType(false))))); BinaryRow EMPTY_MIN_KEY = EMPTY_ROW; BinaryRow EMPTY_MAX_KEY = EMPTY_ROW; @@ -173,7 +178,7 @@ static DataFileMeta create( @Nullable String externalPath, @Nullable Long firstRowId, @Nullable List writeCols) { - return new PojoDataFileMeta( + return create( fileName, fileSize, rowCount, @@ -193,7 +198,54 @@ static DataFileMeta create( valueStatsCols, externalPath, firstRowId, - writeCols); + writeCols, + null); + } + + static DataFileMeta create( + String fileName, + long fileSize, + long rowCount, + BinaryRow minKey, + BinaryRow maxKey, + SimpleStats keyStats, + SimpleStats valueStats, + long minSequenceNumber, + long maxSequenceNumber, + long schemaId, + int level, + List extraFiles, + Timestamp creationTime, + @Nullable Long deleteRowCount, + @Nullable byte[] embeddedIndex, + @Nullable FileSource fileSource, + @Nullable List valueStatsCols, + @Nullable String externalPath, + @Nullable Long firstRowId, + @Nullable List writeCols, + @Nullable long[] columnMaxSequenceNumbers) { + return new PojoDataFileMeta( + fileName, + fileSize, + rowCount, + minKey, + maxKey, + keyStats, + valueStats, + minSequenceNumber, + maxSequenceNumber, + schemaId, + level, + extraFiles, + creationTime, + deleteRowCount, + embeddedIndex, + fileSource, + valueStatsCols, + externalPath, + firstRowId, + writeCols, + columnMaxSequenceNumbers); } static DataFileMeta create( @@ -354,6 +406,18 @@ default Range nonNullRowIdRange() { @Nullable List writeCols(); + /** + * Maximum sequence number per physical table field after data-evolution compaction. + * + *

Values follow the table-field order selected by {@link #writeCols()} when it is non-null + * (system fields are ignored), or the file schema field order otherwise. A null value means + * that only the file-level sequence range is available. + */ + @Nullable + default long[] columnMaxSequenceNumbers() { + return null; + } + DataFileMeta upgrade(int newLevel); DataFileMeta rename(String newFileName); @@ -362,6 +426,11 @@ default Range nonNullRowIdRange() { DataFileMeta assignSequenceNumber(long minSequenceNumber, long maxSequenceNumber); + default DataFileMeta withColumnMaxSequenceNumbers(long[] columnMaxSequenceNumbers) { + throw new UnsupportedOperationException( + "This DataFileMeta implementation does not support column sequence numbers."); + } + DataFileMeta assignFirstRowId(long firstRowId); DataFileMeta newFirstRowId(@Nullable Long newFirstRowId); diff --git a/paimon-core/src/main/java/org/apache/paimon/io/DataFileMetaFirstRowIdLegacySerializer.java b/paimon-core/src/main/java/org/apache/paimon/io/DataFileMetaFirstRowIdLegacySerializer.java index 59abcc730d38..34ad0dc28d99 100644 --- a/paimon-core/src/main/java/org/apache/paimon/io/DataFileMetaFirstRowIdLegacySerializer.java +++ b/paimon-core/src/main/java/org/apache/paimon/io/DataFileMetaFirstRowIdLegacySerializer.java @@ -36,7 +36,7 @@ public class DataFileMetaFirstRowIdLegacySerializer extends ObjectSerializer { + + private static final long serialVersionUID = 1L; + + static final RowType SCHEMA = + DataFileMeta.SCHEMA.project( + DataFileMeta.FILE_NAME, + DataFileMeta.FILE_SIZE, + DataFileMeta.ROW_COUNT, + DataFileMeta.MIN_KEY, + DataFileMeta.MAX_KEY, + DataFileMeta.KEY_STATS, + DataFileMeta.VALUE_STATS, + DataFileMeta.MIN_SEQUENCE_NUMBER, + DataFileMeta.MAX_SEQUENCE_NUMBER, + DataFileMeta.SCHEMA_ID, + DataFileMeta.LEVEL, + DataFileMeta.EXTRA_FILES, + DataFileMeta.CREATION_TIME, + DataFileMeta.DELETE_ROW_COUNT, + DataFileMeta.EMBEDDED_FILE_INDEX, + DataFileMeta.FILE_SOURCE, + DataFileMeta.VALUE_STATS_COLS, + DataFileMeta.EXTERNAL_PATH, + DataFileMeta.FIRST_ROW_ID, + DataFileMeta.WRITE_COLS); + + public DataFileMetaWriteColsLegacySerializer() { + super(SCHEMA); + } + + @Override + public InternalRow toRow(DataFileMeta meta) { + return GenericRow.of( + BinaryString.fromString(meta.fileName()), + meta.fileSize(), + meta.rowCount(), + serializeBinaryRow(meta.minKey()), + serializeBinaryRow(meta.maxKey()), + meta.keyStats().toRow(), + meta.valueStats().toRow(), + meta.minSequenceNumber(), + meta.maxSequenceNumber(), + meta.schemaId(), + meta.level(), + toStringArrayData(meta.extraFiles()), + meta.creationTime(), + meta.deleteRowCount().orElse(null), + meta.embeddedIndex(), + meta.fileSource().map(FileSource::toByteValue).orElse(null), + toStringArrayData(meta.valueStatsCols()), + meta.externalPath().map(BinaryString::fromString).orElse(null), + meta.firstRowId(), + meta.writeCols() == null ? null : toStringArrayData(meta.writeCols())); + } + + @Override + public DataFileMeta fromRow(InternalRow row) { + return DataFileMeta.create( + row.getString(0).toString(), + row.getLong(1), + row.getLong(2), + deserializeBinaryRow(row.getBinary(3)), + deserializeBinaryRow(row.getBinary(4)), + SimpleStats.fromRow(row.getRow(5, 3)), + SimpleStats.fromRow(row.getRow(6, 3)), + row.getLong(7), + row.getLong(8), + row.getLong(9), + row.getInt(10), + fromStringArrayData(row.getArray(11)), + row.getTimestamp(12, 3), + row.isNullAt(13) ? null : row.getLong(13), + row.isNullAt(14) ? null : row.getBinary(14), + row.isNullAt(15) ? null : FileSource.fromByteValue(row.getByte(15)), + row.isNullAt(16) ? null : fromStringArrayData(row.getArray(16)), + row.isNullAt(17) ? null : row.getString(17).toString(), + row.isNullAt(18) ? null : row.getLong(18), + row.isNullAt(19) ? null : fromStringArrayData(row.getArray(19))); + } +} diff --git a/paimon-core/src/main/java/org/apache/paimon/io/PojoDataFileMeta.java b/paimon-core/src/main/java/org/apache/paimon/io/PojoDataFileMeta.java index 9b288d1d5f7f..dcbd650c8b1c 100644 --- a/paimon-core/src/main/java/org/apache/paimon/io/PojoDataFileMeta.java +++ b/paimon-core/src/main/java/org/apache/paimon/io/PojoDataFileMeta.java @@ -78,6 +78,8 @@ public class PojoDataFileMeta implements DataFileMeta { private final @Nullable List writeCols; + private final @Nullable long[] columnMaxSequenceNumbers; + public PojoDataFileMeta( String fileName, long fileSize, @@ -99,6 +101,52 @@ public PojoDataFileMeta( @Nullable String externalPath, @Nullable Long firstRowId, @Nullable List writeCols) { + this( + fileName, + fileSize, + rowCount, + minKey, + maxKey, + keyStats, + valueStats, + minSequenceNumber, + maxSequenceNumber, + schemaId, + level, + extraFiles, + creationTime, + deleteRowCount, + embeddedIndex, + fileSource, + valueStatsCols, + externalPath, + firstRowId, + writeCols, + null); + } + + public PojoDataFileMeta( + String fileName, + long fileSize, + long rowCount, + BinaryRow minKey, + BinaryRow maxKey, + SimpleStats keyStats, + SimpleStats valueStats, + long minSequenceNumber, + long maxSequenceNumber, + long schemaId, + int level, + List extraFiles, + Timestamp creationTime, + @Nullable Long deleteRowCount, + @Nullable byte[] embeddedIndex, + @Nullable FileSource fileSource, + @Nullable List valueStatsCols, + @Nullable String externalPath, + @Nullable Long firstRowId, + @Nullable List writeCols, + @Nullable long[] columnMaxSequenceNumbers) { this.fileName = fileName; this.fileSize = fileSize; @@ -123,6 +171,8 @@ public PojoDataFileMeta( this.externalPath = externalPath; this.firstRowId = firstRowId; this.writeCols = writeCols; + this.columnMaxSequenceNumbers = + columnMaxSequenceNumbers == null ? null : columnMaxSequenceNumbers.clone(); } @Override @@ -239,6 +289,12 @@ public List writeCols() { return writeCols; } + @Nullable + @Override + public long[] columnMaxSequenceNumbers() { + return columnMaxSequenceNumbers == null ? null : columnMaxSequenceNumbers.clone(); + } + @Override public PojoDataFileMeta upgrade(int newLevel) { checkArgument(newLevel > this.level); @@ -262,7 +318,8 @@ public PojoDataFileMeta upgrade(int newLevel) { valueStatsCols, externalPath, firstRowId, - writeCols); + writeCols, + columnMaxSequenceNumbers); } @Override @@ -288,7 +345,8 @@ public PojoDataFileMeta rename(String newFileName) { valueStatsCols, newExternalPath, firstRowId, - writeCols); + writeCols, + columnMaxSequenceNumbers); } @Override @@ -313,7 +371,8 @@ public PojoDataFileMeta copyWithoutStats() { Collections.emptyList(), externalPath, firstRowId, - writeCols); + writeCols, + columnMaxSequenceNumbers); } @Override @@ -338,7 +397,34 @@ public PojoDataFileMeta assignSequenceNumber(long minSequenceNumber, long maxSeq valueStatsCols, externalPath, firstRowId, - writeCols); + writeCols, + columnMaxSequenceNumbers); + } + + @Override + public PojoDataFileMeta withColumnMaxSequenceNumbers(long[] columnMaxSequenceNumbers) { + return new PojoDataFileMeta( + fileName, + fileSize, + rowCount, + minKey, + maxKey, + keyStats, + valueStats, + minSequenceNumber, + maxSequenceNumber, + schemaId, + level, + extraFiles, + creationTime, + deleteRowCount, + embeddedIndex, + fileSource, + valueStatsCols, + externalPath, + firstRowId, + writeCols, + columnMaxSequenceNumbers); } @Override @@ -363,7 +449,8 @@ public PojoDataFileMeta assignFirstRowId(long firstRowId) { valueStatsCols, externalPath, firstRowId, - writeCols); + writeCols, + columnMaxSequenceNumbers); } @Override @@ -388,7 +475,8 @@ public PojoDataFileMeta newFirstRowId(@Nullable Long newFirstRowId) { valueStatsCols, externalPath, newFirstRowId, - writeCols); + writeCols, + columnMaxSequenceNumbers); } @Override @@ -413,7 +501,8 @@ public PojoDataFileMeta copy(List newExtraFiles) { valueStatsCols, externalPath, firstRowId, - writeCols); + writeCols, + columnMaxSequenceNumbers); } @Override @@ -438,7 +527,8 @@ public PojoDataFileMeta newExternalPath(String newExternalPath) { valueStatsCols, newExternalPath, firstRowId, - writeCols); + writeCols, + columnMaxSequenceNumbers); } @Override @@ -463,7 +553,8 @@ public PojoDataFileMeta copy(byte[] newEmbeddedIndex) { valueStatsCols, externalPath, firstRowId, - writeCols); + writeCols, + columnMaxSequenceNumbers); } @Override @@ -494,7 +585,8 @@ public boolean equals(Object o) { && Objects.equals(valueStatsCols, that.valueStatsCols()) && Objects.equals(externalPath, that.externalPath().orElse(null)) && Objects.equals(firstRowId, that.firstRowId()) - && Objects.equals(writeCols, that.writeCols()); + && Objects.equals(writeCols, that.writeCols()) + && Arrays.equals(columnMaxSequenceNumbers, that.columnMaxSequenceNumbers()); } @Override @@ -519,7 +611,8 @@ public int hashCode() { valueStatsCols, externalPath, firstRowId, - writeCols); + writeCols, + Arrays.hashCode(columnMaxSequenceNumbers)); } @Override @@ -529,7 +622,8 @@ public String toString() { + "minKey: %s, maxKey: %s, keyStats: %s, valueStats: %s, " + "minSequenceNumber: %d, maxSequenceNumber: %d, " + "schemaId: %d, level: %d, extraFiles: %s, creationTime: %s, " - + "deleteRowCount: %d, fileSource: %s, valueStatsCols: %s, externalPath: %s, firstRowId: %s, writeCols: %s}", + + "deleteRowCount: %d, fileSource: %s, valueStatsCols: %s, externalPath: %s, " + + "firstRowId: %s, writeCols: %s, columnMaxSequenceNumbers: %s}", fileName, fileSize, rowCount, @@ -549,6 +643,7 @@ public String toString() { valueStatsCols, externalPath, firstRowId, - writeCols); + writeCols, + Arrays.toString(columnMaxSequenceNumbers)); } } diff --git a/paimon-core/src/main/java/org/apache/paimon/table/sink/AppendCompactTaskSerializer.java b/paimon-core/src/main/java/org/apache/paimon/table/sink/AppendCompactTaskSerializer.java index ed78ec2f7f53..ec6df446a869 100644 --- a/paimon-core/src/main/java/org/apache/paimon/table/sink/AppendCompactTaskSerializer.java +++ b/paimon-core/src/main/java/org/apache/paimon/table/sink/AppendCompactTaskSerializer.java @@ -37,7 +37,7 @@ /** Serializer for {@link AppendCompactTask}. */ public class AppendCompactTaskSerializer implements VersionedSerializer { - private static final int CURRENT_VERSION = 2; + private static final int CURRENT_VERSION = 3; private final DataFileMetaSerializer dataFileSerializer; diff --git a/paimon-core/src/main/java/org/apache/paimon/table/sink/CommitMessageSerializer.java b/paimon-core/src/main/java/org/apache/paimon/table/sink/CommitMessageSerializer.java index 2593a2a68ac2..8222b07c8c8c 100644 --- a/paimon-core/src/main/java/org/apache/paimon/table/sink/CommitMessageSerializer.java +++ b/paimon-core/src/main/java/org/apache/paimon/table/sink/CommitMessageSerializer.java @@ -34,6 +34,7 @@ import org.apache.paimon.io.DataFileMeta12LegacySerializer; import org.apache.paimon.io.DataFileMetaFirstRowIdLegacySerializer; import org.apache.paimon.io.DataFileMetaSerializer; +import org.apache.paimon.io.DataFileMetaWriteColsLegacySerializer; import org.apache.paimon.io.DataIncrement; import org.apache.paimon.io.DataInputDeserializer; import org.apache.paimon.io.DataInputView; @@ -52,12 +53,13 @@ /** {@link VersionedSerializer} for {@link CommitMessage}. */ public class CommitMessageSerializer implements VersionedSerializer { - public static final int CURRENT_VERSION = 12; + public static final int CURRENT_VERSION = 13; private final DataFileMetaSerializer dataFileSerializer; private final IndexFileMetaSerializer indexEntrySerializer; private DataFileMetaFirstRowIdLegacySerializer dataFileMetaFirstRowIdLegacySerializer; + private DataFileMetaWriteColsLegacySerializer dataFileMetaWriteColsLegacySerializer; private DataFileMeta12LegacySerializer dataFileMeta12LegacySerializer; private DataFileMeta10LegacySerializer dataFileMeta10LegacySerializer; private DataFileMeta09Serializer dataFile09Serializer; @@ -186,8 +188,13 @@ private CommitMessage deserialize(int version, DataInputView view) throws IOExce private IOExceptionSupplier> fileDeserializer( int version, DataInputView view) { - if (version >= 9) { + if (version >= 13) { return () -> dataFileSerializer.deserializeList(view); + } else if (version >= 9) { + if (dataFileMetaWriteColsLegacySerializer == null) { + dataFileMetaWriteColsLegacySerializer = new DataFileMetaWriteColsLegacySerializer(); + } + return () -> dataFileMetaWriteColsLegacySerializer.deserializeList(view); } else if (version == 8) { if (dataFileMetaFirstRowIdLegacySerializer == null) { dataFileMetaFirstRowIdLegacySerializer = diff --git a/paimon-core/src/main/java/org/apache/paimon/table/sink/MultiTableCompactionTaskSerializer.java b/paimon-core/src/main/java/org/apache/paimon/table/sink/MultiTableCompactionTaskSerializer.java index 8478d04ea3b2..d1212db99f60 100644 --- a/paimon-core/src/main/java/org/apache/paimon/table/sink/MultiTableCompactionTaskSerializer.java +++ b/paimon-core/src/main/java/org/apache/paimon/table/sink/MultiTableCompactionTaskSerializer.java @@ -39,7 +39,7 @@ public class MultiTableCompactionTaskSerializer implements VersionedSerializer { - private static final int CURRENT_VERSION = 1; + private static final int CURRENT_VERSION = 2; private final DataFileMetaSerializer dataFileSerializer; diff --git a/paimon-core/src/main/java/org/apache/paimon/table/source/ChainSplit.java b/paimon-core/src/main/java/org/apache/paimon/table/source/ChainSplit.java index bad2dc1c7bf8..39edeb6c233b 100644 --- a/paimon-core/src/main/java/org/apache/paimon/table/source/ChainSplit.java +++ b/paimon-core/src/main/java/org/apache/paimon/table/source/ChainSplit.java @@ -21,10 +21,12 @@ import org.apache.paimon.data.BinaryRow; import org.apache.paimon.io.DataFileMeta; import org.apache.paimon.io.DataFileMetaSerializer; +import org.apache.paimon.io.DataFileMetaWriteColsLegacySerializer; import org.apache.paimon.io.DataInputView; import org.apache.paimon.io.DataInputViewStreamWrapper; import org.apache.paimon.io.DataOutputView; import org.apache.paimon.io.DataOutputViewStreamWrapper; +import org.apache.paimon.utils.ObjectSerializer; import org.apache.paimon.utils.SerializationUtils; import javax.annotation.Nullable; @@ -49,7 +51,7 @@ public class ChainSplit implements Split { private static final long serialVersionUID = 1L; private static final int VERSION_1 = 1; - private static final int VERSION = 2; + private static final int VERSION = 3; private BinaryRow logicalPartition; private List dataFiles; @@ -217,7 +219,10 @@ public static ChainSplit deserialize(DataInputView in) throws IOException { int n = in.readInt(); List dataFiles = new ArrayList<>(n); - DataFileMetaSerializer dataFileSer = new DataFileMetaSerializer(); + ObjectSerializer dataFileSer = + version <= 2 + ? new DataFileMetaWriteColsLegacySerializer() + : new DataFileMetaSerializer(); for (int i = 0; i < n; i++) { dataFiles.add(dataFileSer.deserialize(in)); } diff --git a/paimon-core/src/main/java/org/apache/paimon/table/source/DataSplit.java b/paimon-core/src/main/java/org/apache/paimon/table/source/DataSplit.java index df31763a3001..88bf60f019c8 100644 --- a/paimon-core/src/main/java/org/apache/paimon/table/source/DataSplit.java +++ b/paimon-core/src/main/java/org/apache/paimon/table/source/DataSplit.java @@ -28,6 +28,7 @@ import org.apache.paimon.io.DataFileMeta12LegacySerializer; import org.apache.paimon.io.DataFileMetaFirstRowIdLegacySerializer; import org.apache.paimon.io.DataFileMetaSerializer; +import org.apache.paimon.io.DataFileMetaWriteColsLegacySerializer; import org.apache.paimon.io.DataInputView; import org.apache.paimon.io.DataInputViewStreamWrapper; import org.apache.paimon.io.DataOutputView; @@ -63,7 +64,7 @@ public class DataSplit implements Split { private static final long serialVersionUID = 7L; private static final long MAGIC = -2394839472490812314L; - private static final int VERSION = 8; + private static final int VERSION = 9; private long snapshotId = 0; private BinaryRow partition; @@ -509,6 +510,10 @@ private static FunctionWithIOException getFileMetaS new DataFileMetaFirstRowIdLegacySerializer(); return serializer::deserialize; } else if (version == 8) { + DataFileMetaWriteColsLegacySerializer serializer = + new DataFileMetaWriteColsLegacySerializer(); + return serializer::deserialize; + } else if (version == 9) { DataFileMetaSerializer serializer = new DataFileMetaSerializer(); return serializer::deserialize; } else { diff --git a/paimon-core/src/main/java/org/apache/paimon/table/source/IncrementalSplit.java b/paimon-core/src/main/java/org/apache/paimon/table/source/IncrementalSplit.java index 5b672a4fefd5..2672bc6d0fbe 100644 --- a/paimon-core/src/main/java/org/apache/paimon/table/source/IncrementalSplit.java +++ b/paimon-core/src/main/java/org/apache/paimon/table/source/IncrementalSplit.java @@ -21,10 +21,12 @@ import org.apache.paimon.data.BinaryRow; import org.apache.paimon.io.DataFileMeta; import org.apache.paimon.io.DataFileMetaSerializer; +import org.apache.paimon.io.DataFileMetaWriteColsLegacySerializer; import org.apache.paimon.io.DataInputView; import org.apache.paimon.io.DataInputViewStreamWrapper; import org.apache.paimon.io.DataOutputViewStreamWrapper; import org.apache.paimon.utils.FunctionWithIOException; +import org.apache.paimon.utils.ObjectSerializer; import javax.annotation.Nullable; @@ -44,7 +46,7 @@ public class IncrementalSplit implements Split { private static final long serialVersionUID = 1L; - private static final int VERSION = 1; + private static final int VERSION = 2; private long snapshotId; private BinaryRow partition; @@ -220,7 +222,7 @@ private void readObject(ObjectInputStream objectInputStream) throws IOException, ClassNotFoundException { DataInputViewStreamWrapper in = new DataInputViewStreamWrapper(objectInputStream); int version = in.readInt(); - if (version != VERSION) { + if (version < 1 || version > VERSION) { throw new UnsupportedOperationException("Unsupported version: " + version); } @@ -229,7 +231,10 @@ private void readObject(ObjectInputStream objectInputStream) bucket = in.readInt(); totalBuckets = in.readInt(); - DataFileMetaSerializer dataFileMetaSerializer = new DataFileMetaSerializer(); + ObjectSerializer dataFileMetaSerializer = + version == 1 + ? new DataFileMetaWriteColsLegacySerializer() + : new DataFileMetaSerializer(); FunctionWithIOException deletionFileSerializer = DeletionFile::deserialize; diff --git a/paimon-core/src/main/java/org/apache/paimon/table/source/SplitSerializer.java b/paimon-core/src/main/java/org/apache/paimon/table/source/SplitSerializer.java index 070956b811d3..fabfb5604e1d 100644 --- a/paimon-core/src/main/java/org/apache/paimon/table/source/SplitSerializer.java +++ b/paimon-core/src/main/java/org/apache/paimon/table/source/SplitSerializer.java @@ -23,12 +23,14 @@ import org.apache.paimon.globalindex.IndexedSplit; import org.apache.paimon.io.DataFileMeta; import org.apache.paimon.io.DataFileMetaSerializer; +import org.apache.paimon.io.DataFileMetaWriteColsLegacySerializer; import org.apache.paimon.io.DataInputDeserializer; import org.apache.paimon.io.DataInputView; import org.apache.paimon.io.DataOutputView; import org.apache.paimon.io.DataOutputViewStreamWrapper; import org.apache.paimon.table.FallbackReadFileStoreTable; import org.apache.paimon.utils.FunctionWithIOException; +import org.apache.paimon.utils.ObjectSerializer; import javax.annotation.Nullable; @@ -54,7 +56,7 @@ public class SplitSerializer { private static final long MAGIC = 0x53504C49545F5631L; // "SPLIT_V1" - private static final int VERSION = 1; + private static final int VERSION = 2; private static final int DATA_SPLIT = 1; private static final int INCREMENTAL_SPLIT = 2; @@ -113,7 +115,7 @@ public static Split deserialize(DataInputView in) throws IOException { } int version = in.readInt(); - if (version != VERSION) { + if (version < 1 || version > VERSION) { throw new IOException("Unsupported split serializer version: " + version); } @@ -122,7 +124,7 @@ public static Split deserialize(DataInputView in) throws IOException { case DATA_SPLIT: return DataSplit.deserialize(in); case INCREMENTAL_SPLIT: - return readIncrementalSplit(in); + return readIncrementalSplit(in, version); case INDEXED_SPLIT: return IndexedSplit.deserialize(in); case CHAIN_SPLIT: @@ -151,17 +153,18 @@ private static void writeIncrementalSplit(IncrementalSplit split, DataOutputView out.writeBoolean(split.isStreaming()); } - private static IncrementalSplit readIncrementalSplit(DataInputView in) throws IOException { + private static IncrementalSplit readIncrementalSplit(DataInputView in, int version) + throws IOException { long snapshotId = in.readLong(); BinaryRow partition = deserializeBinaryRow(in); int bucket = in.readInt(); int totalBuckets = in.readInt(); - List beforeFiles = readDataFiles(in); + List beforeFiles = readDataFiles(in, version); FunctionWithIOException deletionFileSerializer = DeletionFile::deserialize; List beforeDeletionFiles = DeletionFile.deserializeList(in, deletionFileSerializer); - List afterFiles = readDataFiles(in); + List afterFiles = readDataFiles(in, version); List afterDeletionFiles = DeletionFile.deserializeList(in, deletionFileSerializer); boolean isStreaming = in.readBoolean(); @@ -232,10 +235,14 @@ private static void writeDataFiles(List files, DataOutputView out) } } - private static List readDataFiles(DataInputView in) throws IOException { + private static List readDataFiles(DataInputView in, int version) + throws IOException { int size = in.readInt(); List files = new ArrayList<>(size); - DataFileMetaSerializer serializer = new DataFileMetaSerializer(); + ObjectSerializer serializer = + version == 1 + ? new DataFileMetaWriteColsLegacySerializer() + : new DataFileMetaSerializer(); for (int i = 0; i < size; i++) { files.add(serializer.deserialize(in)); } diff --git a/paimon-core/src/main/java/org/apache/paimon/utils/DataEvolutionUtils.java b/paimon-core/src/main/java/org/apache/paimon/utils/DataEvolutionUtils.java index c1294f39462f..96c25850e411 100644 --- a/paimon-core/src/main/java/org/apache/paimon/utils/DataEvolutionUtils.java +++ b/paimon-core/src/main/java/org/apache/paimon/utils/DataEvolutionUtils.java @@ -22,10 +22,12 @@ import org.apache.paimon.schema.TableSchema; import org.apache.paimon.types.DataField; +import java.util.ArrayList; import java.util.Collection; import java.util.Comparator; -import java.util.HashSet; +import java.util.HashMap; import java.util.List; +import java.util.Map; import java.util.Set; import java.util.function.Function; import java.util.stream.Collectors; @@ -43,17 +45,43 @@ public class DataEvolutionUtils { */ public static Set fileFieldIds( Function scanTableSchema, DataFileMeta file) { + return fileFields(scanTableSchema, file).stream() + .map(DataField::id) + .collect(Collectors.toSet()); + } + + /** Table fields physically present in a file, in their physical write order. */ + public static List fileFields( + Function scanTableSchema, DataFileMeta file) { TableSchema schema = scanTableSchema.apply(file.schemaId()); List writeCols = file.writeCols(); - Set writeColNames = writeCols == null ? null : new HashSet<>(writeCols); - Set ids = new HashSet<>(); + if (writeCols == null) { + return schema.fields(); + } + + Map fieldsByName = new HashMap<>(); for (DataField field : schema.fields()) { + fieldsByName.put(field.name(), field); + } + List fields = new ArrayList<>(); + for (String writeCol : writeCols) { // writeCols may also contain physical row-tracking fields outside the table schema. - if (writeColNames == null || writeColNames.contains(field.name())) { - ids.add(field.id()); + DataField field = fieldsByName.get(writeCol); + if (field != null) { + fields.add(field); } } - return ids; + return fields; + } + + /** Returns the latest sequence known for a physical field position in the file. */ + public static long fieldMaxSequenceNumber( + DataFileMeta file, int fieldPosition, int physicalFieldCount) { + long[] columnSequences = file.columnMaxSequenceNumbers(); + if (columnSequences == null || columnSequences.length != physicalFieldCount) { + return file.maxSequenceNumber(); + } + return columnSequences[fieldPosition]; } /** diff --git a/paimon-core/src/test/java/org/apache/paimon/globalindex/DataEvolutionGlobalIndexRefreshPlannerTest.java b/paimon-core/src/test/java/org/apache/paimon/globalindex/DataEvolutionGlobalIndexRefreshPlannerTest.java index 68d1a26218de..3fa8491ff640 100644 --- a/paimon-core/src/test/java/org/apache/paimon/globalindex/DataEvolutionGlobalIndexRefreshPlannerTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/globalindex/DataEvolutionGlobalIndexRefreshPlannerTest.java @@ -152,6 +152,76 @@ void testUsesStableFieldIdAcrossRenameAndFullWrites() { .containsExactly(index); } + @Test + void testUsesColumnSequenceNumbersForCompactedFullFile() { + IndexManifestEntry index = index("index", 0, 99, 5L, BinaryRow.EMPTY_ROW, 0); + + assertThat( + plan( + Collections.singletonList( + dataWithColumnSequences( + "unrelated-compact", + 0, + 100, + 10, + new long[] {5L, 10L, 10L})), + index)) + .isEmpty(); + assertThat( + plan( + Collections.singletonList( + dataWithColumnSequences( + "index-compact", + 0, + 100, + 10, + new long[] {6L, 10L, 10L})), + index)) + .containsExactly(index); + + // Legacy compacted files have no column metadata and remain conservative. + assertThat(plan(Collections.singletonList(data("legacy", 0, 100, 10, 1)), index)) + .containsExactly(index); + } + + @Test + void testColumnSequenceNumbersFollowWriteColsOrder() { + IndexManifestEntry index = index("index", 0, 99, 5L, BinaryRow.EMPTY_ROW, 0); + + assertThat( + plan( + Collections.singletonList( + dataWithColumnSequences( + "reordered-compact", + 0, + 100, + 10, + new long[] {10L, 5L}, + "unrelated", + "vector")), + index)) + .isEmpty(); + } + + @Test + void testMalformedColumnSequenceNumbersFallBackToFileSequence() { + IndexManifestEntry index = index("index", 0, 99, 5L, BinaryRow.EMPTY_ROW, 0); + + assertThat( + plan( + Collections.singletonList( + dataWithColumnSequences( + "malformed-compact", + 0, + 100, + 10, + new long[] {5L}, + "vector", + "other")), + index)) + .containsExactly(index); + } + @Test void testRefreshesFromUpdateLayerOverBaseSchemaWithoutIndexColumn() { IndexManifestEntry index = index("index", 0, 99, 5L, BinaryRow.EMPTY_ROW, 0); @@ -491,4 +561,18 @@ private ManifestEntry dataWithWriteCols( writeCols); return ManifestEntry.create(FileKind.ADD, BinaryRow.EMPTY_ROW, 0, 1, file); } + + private ManifestEntry dataWithColumnSequences( + String fileName, + long firstRowId, + long rowCount, + long maxSequenceNumber, + long[] columnSequences, + String... writeCols) { + DataFileMeta file = + data(fileName, firstRowId, rowCount, maxSequenceNumber, 1, writeCols) + .file() + .withColumnMaxSequenceNumbers(columnSequences); + return ManifestEntry.create(FileKind.ADD, BinaryRow.EMPTY_ROW, 0, 1, file); + } } diff --git a/paimon-core/src/test/java/org/apache/paimon/globalindex/sorted/SortedGlobalIndexScannerTest.java b/paimon-core/src/test/java/org/apache/paimon/globalindex/sorted/SortedGlobalIndexScannerTest.java index f918b1728981..91373f4b0b50 100644 --- a/paimon-core/src/test/java/org/apache/paimon/globalindex/sorted/SortedGlobalIndexScannerTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/globalindex/sorted/SortedGlobalIndexScannerTest.java @@ -20,6 +20,9 @@ import org.apache.paimon.CoreOptions; import org.apache.paimon.Snapshot; +import org.apache.paimon.append.dataevolution.DataEvolutionCompactCoordinator; +import org.apache.paimon.append.dataevolution.DataEvolutionCompactTask; +import org.apache.paimon.append.dataevolution.DataEvolutionCompactionCommitPreparation; import org.apache.paimon.data.BinaryRow; import org.apache.paimon.data.BinaryString; import org.apache.paimon.data.BlobData; @@ -35,6 +38,8 @@ import org.apache.paimon.manifest.IndexManifestEntry; import org.apache.paimon.manifest.ManifestEntry; import org.apache.paimon.memory.MemorySlice; +import org.apache.paimon.options.ExpireConfig; +import org.apache.paimon.options.Options; import org.apache.paimon.partition.PartitionPredicate; import org.apache.paimon.predicate.Predicate; import org.apache.paimon.schema.Schema; @@ -44,7 +49,9 @@ import org.apache.paimon.table.sink.BatchTableWrite; import org.apache.paimon.table.sink.BatchWriteBuilder; import org.apache.paimon.table.sink.CommitMessage; +import org.apache.paimon.table.sink.CommitMessageImpl; import org.apache.paimon.table.source.DataSplit; +import org.apache.paimon.types.DataField; import org.apache.paimon.types.DataTypes; import org.apache.paimon.types.RowType; import org.apache.paimon.utils.Pair; @@ -52,6 +59,7 @@ import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.Test; +import java.time.Duration; import java.util.ArrayList; import java.util.Collections; import java.util.Comparator; @@ -59,7 +67,9 @@ import java.util.List; import java.util.Map; import java.util.Optional; +import java.util.stream.Collectors; +import static org.apache.paimon.utils.DataEvolutionUtils.fileFields; import static org.assertj.core.api.Assertions.assertThat; /** Test class for {@link SortedGlobalIndexScanner}. */ @@ -277,6 +287,134 @@ public void testIncrementalScanWithNewData() throws Exception { 500, totalRowCount, "incrementalScan should only return the newly written rows"); } + @Test + public void testIncrementalScanIgnoresNonIndexColumnCompactionAfterSnapshotExpiration() + throws Exception { + write(); + createIndex(null); + + DataFileMeta firstCompact = updateColumnAndCompact("f1", 1); + int f0Id = getTableDefault().rowType().getField("f0").id(); + long f0Sequence = columnSequence(firstCompact, f0Id); + assertThat(f0Sequence).isLessThan(firstCompact.maxSequenceNumber()); + + DataFileMeta secondCompact = updateColumnAndCompact("f1", 2); + assertThat(columnSequence(secondCompact, f0Id)).isEqualTo(f0Sequence); + + FileStoreTable table = getTableDefault(); + table.newExpireSnapshots() + .config( + ExpireConfig.builder() + .snapshotRetainMax(1) + .snapshotRetainMin(1) + .snapshotTimeRetain(Duration.ZERO) + .build()) + .expire(); + assertThat(table.snapshotManager().earliestSnapshotId()) + .isEqualTo(table.snapshotManager().latestSnapshotId()); + + assertThat(dataEvolutionScanner(table).withIndexField("f0").incrementalScan()).isEmpty(); + } + + @Test + public void testIncrementalScanRefreshesIndexColumnCompaction() throws Exception { + write(); + createIndex(null); + DataFileMeta compacted = updateColumnAndCompact("f0", 1); + + int f0Id = getTableDefault().rowType().getField("f0").id(); + assertThat(columnSequence(compacted, f0Id)).isEqualTo(compacted.maxSequenceNumber()); + + Optional> scanResult = + dataEvolutionScanner(getTableDefault()).withIndexField("f0").incrementalScan(); + assertThat(scanResult).isPresent(); + assertThat(scanResult.get().deletedIndexEntries()).isNotEmpty(); + } + + private long columnSequence(DataFileMeta file, int fieldId) throws Exception { + List fields = fileFields(getTableDefault().schemaManager()::schema, file); + long[] sequences = file.columnMaxSequenceNumbers(); + assertThat(sequences).hasSize(fields.size()); + for (int i = 0; i < fields.size(); i++) { + if (fields.get(i).id() == fieldId) { + return sequences[i]; + } + } + throw new IllegalArgumentException("Field not found in data file: " + fieldId); + } + + private SortedGlobalIndexScanner dataEvolutionScanner(FileStoreTable table) { + Options options = new Options(); + options.set( + CoreOptions.GLOBAL_INDEX_COLUMN_UPDATE_ACTION, + CoreOptions.GlobalIndexColumnUpdateAction.IGNORE); + return new SortedGlobalIndexScanner(table, "btree", options); + } + + private DataFileMeta updateColumnAndCompact(String column, int updateRound) throws Exception { + Map writeOptions = new HashMap<>(); + writeOptions.put( + CoreOptions.GLOBAL_INDEX_COLUMN_UPDATE_ACTION.key(), + CoreOptions.GlobalIndexColumnUpdateAction.IGNORE.toString()); + FileStoreTable table = getTableDefault().copy(writeOptions); + RowType writeType = table.rowType().project("dt", column); + BatchWriteBuilder writeBuilder = table.newBatchWriteBuilder(); + try (BatchTableWrite batchWrite = writeBuilder.newWrite().withWriteType(writeType)) { + for (int i = 0; i < PART_ROW_NUM; i++) { + Object value = + "f0".equals(column) + ? i + updateRound * (int) PART_ROW_NUM + : BinaryString.fromString("updated_" + updateRound + "_" + i); + batchWrite.write(GenericRow.of(BinaryString.fromString("p0"), value)); + } + List messages = batchWrite.prepareCommit(); + assignFirstRowId(messages, 0L); + try (BatchTableCommit commit = writeBuilder.newCommit()) { + commit.commit(messages); + } + } + + writeOptions.put(CoreOptions.COMPACTION_MIN_FILE_NUM.key(), "2"); + table = getTableDefault().copy(writeOptions); + Snapshot compactSnapshot = table.snapshotManager().latestSnapshot(); + DataEvolutionCompactCoordinator coordinator = + new DataEvolutionCompactCoordinator(table, false, false, compactSnapshot); + List compactMessages = new ArrayList<>(); + for (DataEvolutionCompactTask task : coordinator.plan()) { + compactMessages.add(task.doCompact(table, "test-compact")); + } + assertThat(compactMessages).isNotEmpty(); + compactMessages.addAll( + new DataEvolutionCompactionCommitPreparation(table, compactSnapshot) + .prepare(compactMessages)); + try (BatchTableCommit commit = table.newBatchWriteBuilder().newCommit()) { + commit.commit(compactMessages); + } + + List rowRangeFiles = + getTableDefault().store().newScan().plan().files().stream() + .map(ManifestEntry::file) + .filter(file -> file.firstRowId() != null && file.firstRowId() == 0L) + .collect(Collectors.toList()); + assertThat(rowRangeFiles).hasSize(1); + assertThat(rowRangeFiles.get(0).columnMaxSequenceNumbers()).isNotNull(); + return rowRangeFiles.get(0); + } + + private void assignFirstRowId(List messages, long firstRowId) { + for (CommitMessage message : messages) { + CommitMessageImpl impl = (CommitMessageImpl) message; + List files = new ArrayList<>(impl.newFilesIncrement().newFiles()); + impl.newFilesIncrement().newFiles().clear(); + impl.newFilesIncrement() + .newFiles() + .addAll( + files.stream() + .map(file -> file.assignFirstRowId(firstRowId)) + .collect(Collectors.toList())); + } + } + @Test public void testIncrementalScanWithPartitionPredicate() throws Exception { write(); diff --git a/paimon-core/src/test/java/org/apache/paimon/io/BinaryDataFileMetaTest.java b/paimon-core/src/test/java/org/apache/paimon/io/BinaryDataFileMetaTest.java index df7c7a8f17ef..2518fdfe1d79 100644 --- a/paimon-core/src/test/java/org/apache/paimon/io/BinaryDataFileMetaTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/io/BinaryDataFileMetaTest.java @@ -42,20 +42,21 @@ public class BinaryDataFileMetaTest { void testImplementsProjectedDataFileMeta() { DataFileMeta expected = DataFileMeta.forAppend( - "data.parquet", - 123L, - 5L, - SimpleStats.EMPTY_STATS, - 2L, - 3L, - 4L, - Arrays.asList("extra-1", "extra-2"), - new byte[] {1, 2}, - FileSource.COMPACT, - Collections.singletonList("value_col"), - "external/dir/data.parquet", - 10L, - Collections.singletonList("write_col")); + "data.parquet", + 123L, + 5L, + SimpleStats.EMPTY_STATS, + 2L, + 3L, + 4L, + Arrays.asList("extra-1", "extra-2"), + new byte[] {1, 2}, + FileSource.COMPACT, + Collections.singletonList("value_col"), + "external/dir/data.parquet", + 10L, + Collections.singletonList("write_col")) + .withColumnMaxSequenceNumbers(new long[] {11L}); BinaryDataFileMeta actual = BinaryDataFileMeta.Projection.create(DataFileMeta.SCHEMA) .createDataFile() @@ -91,6 +92,7 @@ void testImplementsProjectedDataFileMeta() { assertThat(actual.firstRowId()).isEqualTo(10L); assertThat(actual.nonNullFirstRowId()).isEqualTo(10L); assertThat(actual.writeCols()).containsExactly("write_col"); + assertThat(actual.columnMaxSequenceNumbers()).containsExactly(11L); assertThat(actual.containsWriteColumn(BinaryString.fromString("write_col"))).isTrue(); assertThat(actual.containsWriteColumn(BinaryString.fromString("other"))).isFalse(); assertThat(actual.toFileSelection(Collections.singletonList(new Range(11L, 12L)))) @@ -101,6 +103,9 @@ void testImplementsProjectedDataFileMeta() { assertUnsupported(actual::copyWithoutStats, "copyWithoutStats()"); assertUnsupported( () -> actual.assignSequenceNumber(4L, 5L), "assignSequenceNumber(long, long)"); + assertUnsupported( + () -> actual.withColumnMaxSequenceNumbers(new long[] {2L}), + "withColumnMaxSequenceNumbers(long[])"); assertUnsupported(() -> actual.assignFirstRowId(20L), "assignFirstRowId(long)"); assertUnsupported(() -> actual.newFirstRowId(20L), "newFirstRowId(Long)"); assertUnsupported(() -> actual.copy(Collections.emptyList()), "copy(List)"); diff --git a/paimon-core/src/test/java/org/apache/paimon/io/DataFileMetaSerializerTest.java b/paimon-core/src/test/java/org/apache/paimon/io/DataFileMetaSerializerTest.java index 5074ccf1cf7d..be87c063aa8b 100644 --- a/paimon-core/src/test/java/org/apache/paimon/io/DataFileMetaSerializerTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/io/DataFileMetaSerializerTest.java @@ -20,7 +20,12 @@ import org.apache.paimon.utils.ObjectSerializerTestBase; +import org.junit.jupiter.api.Test; + import java.util.Arrays; +import java.util.Collections; + +import static org.assertj.core.api.Assertions.assertThat; /** Tests for {@link DataFileMetaSerializer}. */ public class DataFileMetaSerializerTest extends ObjectSerializerTestBase { @@ -34,6 +39,47 @@ protected DataFileMetaSerializer serializer() { @Override protected DataFileMeta object() { - return gen.next().meta.copy(Arrays.asList("extra1", "extra2")); + return gen.next() + .meta + .copy(Arrays.asList("extra1", "extra2")) + .withColumnMaxSequenceNumbers(new long[] {3L, 42L}); + } + + @Test + void testCopyOperationsPreserveColumnSequences() { + DataFileMeta file = object(); + assertColumnSequences(file.upgrade(file.level() + 1)); + assertColumnSequences(file.rename("renamed.parquet")); + assertColumnSequences(file.copyWithoutStats()); + assertColumnSequences(file.assignSequenceNumber(1L, 2L)); + assertColumnSequences(file.assignFirstRowId(1L)); + assertColumnSequences(file.newFirstRowId(null)); + assertColumnSequences(file.copy(Collections.emptyList())); + assertColumnSequences(file.newExternalPath("external/renamed.parquet")); + assertColumnSequences(file.copy(new byte[] {1})); + } + + @Test + void testLegacySerializerDropsColumnSequences() { + DataFileMetaWriteColsLegacySerializer legacy = new DataFileMetaWriteColsLegacySerializer(); + DataFileMeta file = legacy.fromRow(legacy.toRow(object())); + assertThat(file.columnMaxSequenceNumbers()).isNull(); + } + + @Test + void testColumnSequencesAreDefensivelyCopied() { + long[] sequences = {3L, 42L}; + DataFileMeta file = gen.next().meta.withColumnMaxSequenceNumbers(sequences); + + sequences[0] = 100L; + assertColumnSequences(file); + + long[] returned = file.columnMaxSequenceNumbers(); + returned[1] = 100L; + assertColumnSequences(file); + } + + private void assertColumnSequences(DataFileMeta file) { + assertThat(file.columnMaxSequenceNumbers()).containsExactly(3L, 42L); } } diff --git a/paimon-core/src/test/java/org/apache/paimon/manifest/ManifestCommittableSerializerCompatibilityTest.java b/paimon-core/src/test/java/org/apache/paimon/manifest/ManifestCommittableSerializerCompatibilityTest.java index 20eac11da7e8..b10b6f7dcb94 100644 --- a/paimon-core/src/test/java/org/apache/paimon/manifest/ManifestCommittableSerializerCompatibilityTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/manifest/ManifestCommittableSerializerCompatibilityTest.java @@ -27,6 +27,7 @@ import org.apache.paimon.io.DataIncrement; import org.apache.paimon.stats.SimpleStats; import org.apache.paimon.table.sink.CommitMessageImpl; +import org.apache.paimon.utils.CompatibilityUtils; import org.apache.paimon.utils.IOUtils; import org.junit.jupiter.api.Test; @@ -45,6 +46,62 @@ /** Compatibility Test for {@link ManifestCommittableSerializer}. */ public class ManifestCommittableSerializerCompatibilityTest { + private static final String GENERATE_GOLDEN_FILES_PROPERTY = + "generateManifestCommittableGoldenFiles"; + + @Test + public void testCompatibilityToV5CommitV13() throws IOException { + DataFileMeta dataFile = + DataFileMeta.create( + "column-sequence-file", + 1024L, + 10L, + singleColumn("min_key"), + singleColumn("max_key"), + SimpleStats.EMPTY_STATS, + SimpleStats.EMPTY_STATS, + 1L, + 5L, + 1L, + 0, + Collections.emptyList(), + Timestamp.fromLocalDateTime( + LocalDateTime.parse("2026-08-07T00:00:00")), + 0L, + null, + FileSource.COMPACT, + null, + null, + 1L, + Arrays.asList("a", "b")) + .withColumnMaxSequenceNumbers(new long[] {3L, 5L}); + IndexFileMeta indexFile = + new IndexFileMeta( + "index-type", "index-file", 100L, 10L, (GlobalIndexMeta) null, null); + ManifestCommittable committable = + createManifestCommittable( + Collections.singletonList(dataFile), indexFile, indexFile); + + ManifestCommittableSerializer serializer = new ManifestCommittableSerializer(); + byte[] current = serializer.serialize(committable); + byte[] serialized; + if (Boolean.parseBoolean( + System.getProperties().getProperty(GENERATE_GOLDEN_FILES_PROPERTY))) { + CompatibilityUtils.writeCompatibilityFile("manifest-committable-v13-v5", current); + serialized = current; + } else { + serialized = + IOUtils.readFully( + ManifestCommittableSerializerCompatibilityTest.class + .getClassLoader() + .getResourceAsStream( + "compatibility/manifest-committable-v13-v5"), + true); + } + + assertThat(serializer.deserialize(5, serialized)).isEqualTo(committable); + } + @Test public void testCompatibilityToV5CommitV11() throws IOException { String fileName = "manifest-committable-v11-v5"; diff --git a/paimon-core/src/test/java/org/apache/paimon/table/sink/CommitMessageSerializerTest.java b/paimon-core/src/test/java/org/apache/paimon/table/sink/CommitMessageSerializerTest.java index c4f519e84e52..bc36deedca9a 100644 --- a/paimon-core/src/test/java/org/apache/paimon/table/sink/CommitMessageSerializerTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/table/sink/CommitMessageSerializerTest.java @@ -40,6 +40,14 @@ public void test() throws IOException { CommitMessageSerializer serializer = new CommitMessageSerializer(); DataIncrement dataIncrement = randomNewFilesIncrement(); + dataIncrement + .newFiles() + .set( + 0, + dataIncrement + .newFiles() + .get(0) + .withColumnMaxSequenceNumbers(new long[] {3L, 42L})); dataIncrement.newIndexFiles().addAll(Arrays.asList(randomIndexFile(), randomIndexFile())); dataIncrement .deletedIndexFiles() diff --git a/paimon-core/src/test/java/org/apache/paimon/table/source/DataSplitCompatibleTest.java b/paimon-core/src/test/java/org/apache/paimon/table/source/DataSplitCompatibleTest.java index a0c76537ab10..b37f25946a5b 100644 --- a/paimon-core/src/test/java/org/apache/paimon/table/source/DataSplitCompatibleTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/table/source/DataSplitCompatibleTest.java @@ -36,6 +36,7 @@ import org.apache.paimon.types.IntType; import org.apache.paimon.types.SmallIntType; import org.apache.paimon.types.TimestampType; +import org.apache.paimon.utils.CompatibilityUtils; import org.apache.paimon.utils.IOUtils; import org.apache.paimon.utils.InstantiationUtil; @@ -61,6 +62,8 @@ /** Test for {@link DataSplit}. */ public class DataSplitCompatibleTest { + private static final String GENERATE_GOLDEN_FILES_PROPERTY = "generateDataSplitGoldenFiles"; + @Test public void testSplitMergedRowCount() { // not rawConvertible @@ -218,6 +221,7 @@ public void testSerializer() throws IOException { DataFileTestDataGenerator gen = DataFileTestDataGenerator.builder().build(); DataFileTestDataGenerator.Data data = gen.next(); List files = new ArrayList<>(); + files.add(gen.next().meta.withColumnMaxSequenceNumbers(new long[] {3L, 42L})); for (int i = 0; i < ThreadLocalRandom.current().nextInt(10); i++) { files.add(gen.next().meta); } @@ -784,6 +788,84 @@ public void testSerializerCompatibleV8() throws Exception { assertThat(actual).isEqualTo(split); } + @Test + public void testSerializerCompatibleV9() throws Exception { + SimpleStats keyStats = + new SimpleStats( + singleColumn("min_key"), + singleColumn("max_key"), + fromLongArray(new Long[] {0L})); + SimpleStats valueStats = + new SimpleStats( + singleColumn("min_value"), + singleColumn("max_value"), + fromLongArray(new Long[] {0L})); + + DataFileMeta dataFile = + DataFileMeta.create( + "my_file", + 1024 * 1024, + 1024, + singleColumn("min_key"), + singleColumn("max_key"), + keyStats, + valueStats, + 15, + 200, + 5, + 3, + Arrays.asList("extra1", "extra2"), + Timestamp.fromLocalDateTime( + LocalDateTime.parse("2022-03-02T20:20:12")), + 11L, + new byte[] {1, 2, 4}, + FileSource.COMPACT, + Arrays.asList("field1", "field2", "field3"), + "hdfs:///path/to/warehouse", + 12L, + Arrays.asList("a", "b", "c", "f")) + .withColumnMaxSequenceNumbers(new long[] {15L, 100L, 150L, 200L}); + List dataFiles = Collections.singletonList(dataFile); + + DeletionFile deletionFile = new DeletionFile("deletion_file", 100, 22, 33L); + List deletionFiles = Collections.singletonList(deletionFile); + + BinaryRow partition = new BinaryRow(1); + BinaryRowWriter binaryRowWriter = new BinaryRowWriter(partition); + binaryRowWriter.writeString(0, BinaryString.fromString("aaaaa")); + binaryRowWriter.complete(); + + DataSplit split = + DataSplit.builder() + .withSnapshot(18) + .withPartition(partition) + .withBucket(20) + .withTotalBuckets(32) + .withDataFiles(dataFiles) + .withDataDeletionFiles(deletionFiles) + .withBucketPath("my path") + .build(); + + byte[] current = InstantiationUtil.serializeObject(split); + byte[] serialized; + if (Boolean.parseBoolean( + System.getProperties().getProperty(GENERATE_GOLDEN_FILES_PROPERTY))) { + CompatibilityUtils.writeCompatibilityFile("datasplit-v9", current); + serialized = current; + } else { + serialized = + IOUtils.readFully( + DataSplitCompatibleTest.class + .getClassLoader() + .getResourceAsStream("compatibility/datasplit-v9"), + true); + } + + DataSplit actual = + InstantiationUtil.deserializeObject(serialized, DataSplit.class.getClassLoader()); + assertThat(actual).isEqualTo(split); + } + private DataFileMeta newDataFile(long rowCount) { return newDataFile(rowCount, null, null); } diff --git a/paimon-core/src/test/java/org/apache/paimon/table/source/SplitSerializerTest.java b/paimon-core/src/test/java/org/apache/paimon/table/source/SplitSerializerTest.java index 0fd6f982b527..eea7a817f2cf 100644 --- a/paimon-core/src/test/java/org/apache/paimon/table/source/SplitSerializerTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/table/source/SplitSerializerTest.java @@ -56,31 +56,42 @@ public class SplitSerializerTest { @Test public void testRoundTrip() throws IOException { - for (GoldenCase goldenCase : goldenCases()) { + for (GoldenCase goldenCase : goldenCases(2)) { Split actual = SplitSerializer.deserialize(SplitSerializer.serialize(goldenCase.split)); assertSplitEquals(goldenCase.split, actual); } } @Test - public void testGoldenFiles() throws IOException { + public void testVersion1GoldenFiles() throws IOException { + for (GoldenCase goldenCase : goldenCases(1)) { + assertSplitEquals( + goldenCase.split, + SplitSerializer.deserialize(readGoldenFile(goldenCase.fileName))); + } + } + + @Test + public void testVersion2GoldenFiles() throws IOException { boolean generateGoldenFiles = Boolean.parseBoolean( System.getProperties().getProperty(GENERATE_GOLDEN_FILES_PROPERTY)); - for (GoldenCase goldenCase : goldenCases()) { - byte[] actual = SplitSerializer.serialize(goldenCase.split); + for (GoldenCase goldenCase : goldenCases(2)) { + byte[] current = SplitSerializer.serialize(goldenCase.split); + byte[] serialized; if (generateGoldenFiles) { - CompatibilityUtils.writeCompatibilityFile(goldenCase.fileName, actual); + CompatibilityUtils.writeCompatibilityFile(goldenCase.fileName, current); + serialized = current; } else { - assertThat(actual).isEqualTo(readGoldenFile(goldenCase.fileName)); + serialized = readGoldenFile(goldenCase.fileName); } - assertSplitEquals(goldenCase.split, SplitSerializer.deserialize(actual)); + assertSplitEquals(goldenCase.split, SplitSerializer.deserialize(serialized)); } } @Test public void testFallbackSplitImplSerializeAndDeserialize() throws IOException { - FallbackReadFileStoreTable.FallbackSplitImpl split = fallbackSplit(); + FallbackReadFileStoreTable.FallbackSplitImpl split = fallbackSplit(true); ByteArrayOutputStream out = new ByteArrayOutputStream(); split.serialize(new DataOutputViewStreamWrapper(out)); @@ -93,7 +104,7 @@ public void testFallbackSplitImplSerializeAndDeserialize() throws IOException { @Test public void testFallbackSplitImplJavaSerializeAndDeserialize() throws Exception { - FallbackReadFileStoreTable.FallbackSplitImpl split = fallbackSplit(); + FallbackReadFileStoreTable.FallbackSplitImpl split = fallbackSplit(true); byte[] bytes = InstantiationUtil.serializeObject(split); FallbackReadFileStoreTable.FallbackSplitImpl deserialized = @@ -111,32 +122,36 @@ private static byte[] readGoldenFile(String fileName) throws IOException { } } - private static List goldenCases() { + private static List goldenCases(int version) { + boolean withColumnSequences = version >= 2; List cases = new ArrayList<>(); - DataSplit dataSplit = dataSplit(); - IncrementalSplit incrementalSplit = incrementalSplit(); + DataSplit dataSplit = dataSplit(withColumnSequences); + IncrementalSplit incrementalSplit = incrementalSplit(withColumnSequences); IndexedSplit indexedSplit = indexedSplit(dataSplit); - ChainSplit chainSplit = chainSplit(); + ChainSplit chainSplit = chainSplit(withColumnSequences); QueryAuthSplit queryAuthSplit = queryAuthSplit(dataSplit); FallbackReadFileStoreTable.FallbackDataSplit fallbackDataSplit = (FallbackReadFileStoreTable.FallbackDataSplit) FallbackReadFileStoreTable.toFallbackSplit(dataSplit, true); - FallbackReadFileStoreTable.FallbackSplitImpl fallbackSplit = fallbackSplit(); - - cases.add(new GoldenCase("split-v1-data", dataSplit)); - cases.add(new GoldenCase("split-v1-incremental", incrementalSplit)); - cases.add(new GoldenCase("split-v1-indexed", indexedSplit)); - cases.add(new GoldenCase("split-v1-chain", chainSplit)); - cases.add(new GoldenCase("split-v1-query-auth", queryAuthSplit)); - cases.add(new GoldenCase("split-v1-fallback-data", fallbackDataSplit)); - cases.add(new GoldenCase("split-v1-fallback", fallbackSplit)); + FallbackReadFileStoreTable.FallbackSplitImpl fallbackSplit = + fallbackSplit(withColumnSequences); + + String prefix = "split-v" + version + "-"; + cases.add(new GoldenCase(prefix + "data", dataSplit)); + cases.add(new GoldenCase(prefix + "incremental", incrementalSplit)); + cases.add(new GoldenCase(prefix + "indexed", indexedSplit)); + cases.add(new GoldenCase(prefix + "chain", chainSplit)); + cases.add(new GoldenCase(prefix + "query-auth", queryAuthSplit)); + cases.add(new GoldenCase(prefix + "fallback-data", fallbackDataSplit)); + cases.add(new GoldenCase(prefix + "fallback", fallbackSplit)); return cases; } - private static DataSplit dataSplit() { + private static DataSplit dataSplit(boolean withColumnSequences) { List files = Arrays.asList( - dataFile("file-a", 0, 1, 10, 100L), dataFile("file-b", 1, 11, 20, 200L)); + dataFile("file-a", 0, 1, 10, 100L, withColumnSequences), + dataFile("file-b", 1, 11, 20, 200L, withColumnSequences)); List deletionFiles = Arrays.asList(null, new DeletionFile("dv/file-b", 2L, 10L, 3L)); return DataSplit.builder() @@ -151,10 +166,13 @@ private static DataSplit dataSplit() { .build(); } - private static IncrementalSplit incrementalSplit() { + private static IncrementalSplit incrementalSplit(boolean withColumnSequences) { List before = - Collections.singletonList(dataFile("before-file", 0, 1, 5, 10L)); - List after = Collections.singletonList(dataFile("after-file", 0, 6, 12, 20L)); + Collections.singletonList( + dataFile("before-file", 0, 1, 5, 10L, withColumnSequences)); + List after = + Collections.singletonList( + dataFile("after-file", 0, 6, 12, 20L, withColumnSequences)); return new IncrementalSplit( 43L, DataFileTestUtils.row(2026, 8), @@ -174,8 +192,8 @@ private static IndexedSplit indexedSplit(DataSplit dataSplit) { new float[] {0.5f, 0.25f, 0.125f}); } - private static ChainSplit chainSplit() { - DataSplit left = dataSplit(); + private static ChainSplit chainSplit(boolean withColumnSequences) { + DataSplit left = dataSplit(withColumnSequences); DataSplit right = DataSplit.builder() .withSnapshot(44L) @@ -184,7 +202,14 @@ private static ChainSplit chainSplit() { .withTotalBuckets(8) .withBucketPath("dt=20260707/bucket-5") .withDataFiles( - Collections.singletonList(dataFile("chain-file", 0, 21, 30, 300L))) + Collections.singletonList( + dataFile( + "chain-file", + 0, + 21, + 30, + 300L, + withColumnSequences))) .withDataDeletionFiles( Collections.singletonList( new DeletionFile("deletion_file", 100, 22, null))) @@ -224,33 +249,44 @@ private static QueryAuthSplit queryAuthSplit(Split split) { Arrays.asList("filter-json-1", "filter-json-2"), columnMasking)); } - private static FallbackReadFileStoreTable.FallbackSplitImpl fallbackSplit() { - return new FallbackReadFileStoreTable.FallbackSplitImpl(chainSplit(), true); + private static FallbackReadFileStoreTable.FallbackSplitImpl fallbackSplit( + boolean withColumnSequences) { + return new FallbackReadFileStoreTable.FallbackSplitImpl( + chainSplit(withColumnSequences), true); } private static DataFileMeta dataFile( - String name, int level, int minKey, int maxKey, long maxSequence) { - return DataFileMeta.create( - name, - maxKey - minKey + 1, - maxKey - minKey + 1, - DataFileTestUtils.row(minKey), - DataFileTestUtils.row(maxKey), - SimpleStats.EMPTY_STATS, - SimpleStats.EMPTY_STATS, - 0L, - maxSequence, - 0L, - level, - Collections.emptyList(), - Timestamp.fromEpochMillis(100), - 0L, - null, - FileSource.APPEND, - null, - null, - null, - null); + String name, + int level, + int minKey, + int maxKey, + long maxSequence, + boolean withColumnSequences) { + DataFileMeta file = + DataFileMeta.create( + name, + maxKey - minKey + 1, + maxKey - minKey + 1, + DataFileTestUtils.row(minKey), + DataFileTestUtils.row(maxKey), + SimpleStats.EMPTY_STATS, + SimpleStats.EMPTY_STATS, + 0L, + maxSequence, + 0L, + level, + Collections.emptyList(), + Timestamp.fromEpochMillis(100), + 0L, + null, + FileSource.APPEND, + null, + null, + null, + null); + return withColumnSequences + ? file.withColumnMaxSequenceNumbers(new long[] {maxSequence}) + : file; } private static void assertSplitEquals(Split expected, Split actual) { diff --git a/paimon-core/src/test/java/org/apache/paimon/utils/DataEvolutionUtilsTest.java b/paimon-core/src/test/java/org/apache/paimon/utils/DataEvolutionUtilsTest.java index 33feb9d850e1..b05907306214 100644 --- a/paimon-core/src/test/java/org/apache/paimon/utils/DataEvolutionUtilsTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/utils/DataEvolutionUtilsTest.java @@ -109,6 +109,48 @@ public void testFileFieldIdsHandlesFullEmptyAndUnrelatedWrites() { .containsExactly(2); } + @Test + public void testFileFieldsFollowWriteColsOrderAndIgnoreSystemFields() { + TableSchema schema = + new TableSchema( + 1L, + Arrays.asList( + new DataField(1, "indexed", new IntType()), + new DataField(2, "other", new IntType())), + 2, + Collections.emptyList(), + Collections.emptyList(), + new HashMap<>(), + ""); + + assertThat( + DataEvolutionUtils.fileFields( + ignored -> schema, + dataFile( + "reordered.parquet", + 1, + Arrays.asList( + "other", SpecialFields.ROW_ID.name(), "indexed")))) + .extracting(DataField::id) + .containsExactly(2, 1); + } + + @Test + public void testFieldMaxSequenceNumberFallsBackForMissingOrMalformedArray() { + DataFileMeta legacy = dataFile("legacy.parquet", 10, null); + DataFileMeta malformed = + dataFile("malformed.parquet", 10, null) + .withColumnMaxSequenceNumbers(new long[] {5L}); + DataFileMeta valid = + dataFile("valid.parquet", 10, null) + .withColumnMaxSequenceNumbers(new long[] {5L, 8L}); + + assertThat(DataEvolutionUtils.fieldMaxSequenceNumber(legacy, 0, 2)).isEqualTo(10L); + assertThat(DataEvolutionUtils.fieldMaxSequenceNumber(malformed, 0, 2)).isEqualTo(10L); + assertThat(DataEvolutionUtils.fieldMaxSequenceNumber(valid, 0, 2)).isEqualTo(5L); + assertThat(DataEvolutionUtils.fieldMaxSequenceNumber(valid, 1, 2)).isEqualTo(8L); + } + @Test public void testRetrieveAnchorFileSkipsSpecialFiles() { DataFileMeta blobFile = dataFile("blob-file.blob", 1); diff --git a/paimon-core/src/test/resources/compatibility/datasplit-v9 b/paimon-core/src/test/resources/compatibility/datasplit-v9 new file mode 100644 index 000000000000..6277c55665ef Binary files /dev/null and b/paimon-core/src/test/resources/compatibility/datasplit-v9 differ diff --git a/paimon-core/src/test/resources/compatibility/manifest-committable-v13-v5 b/paimon-core/src/test/resources/compatibility/manifest-committable-v13-v5 new file mode 100644 index 000000000000..d919649c29e3 Binary files /dev/null and b/paimon-core/src/test/resources/compatibility/manifest-committable-v13-v5 differ diff --git a/paimon-core/src/test/resources/compatibility/split-v2-chain b/paimon-core/src/test/resources/compatibility/split-v2-chain new file mode 100644 index 000000000000..0a4cb5a29fda Binary files /dev/null and b/paimon-core/src/test/resources/compatibility/split-v2-chain differ diff --git a/paimon-core/src/test/resources/compatibility/split-v2-data b/paimon-core/src/test/resources/compatibility/split-v2-data new file mode 100644 index 000000000000..81204d780982 Binary files /dev/null and b/paimon-core/src/test/resources/compatibility/split-v2-data differ diff --git a/paimon-core/src/test/resources/compatibility/split-v2-fallback b/paimon-core/src/test/resources/compatibility/split-v2-fallback new file mode 100644 index 000000000000..c92ca5a3eba6 Binary files /dev/null and b/paimon-core/src/test/resources/compatibility/split-v2-fallback differ diff --git a/paimon-core/src/test/resources/compatibility/split-v2-fallback-data b/paimon-core/src/test/resources/compatibility/split-v2-fallback-data new file mode 100644 index 000000000000..eab303006fcb Binary files /dev/null and b/paimon-core/src/test/resources/compatibility/split-v2-fallback-data differ diff --git a/paimon-core/src/test/resources/compatibility/split-v2-incremental b/paimon-core/src/test/resources/compatibility/split-v2-incremental new file mode 100644 index 000000000000..a3ca67a57d4b Binary files /dev/null and b/paimon-core/src/test/resources/compatibility/split-v2-incremental differ diff --git a/paimon-core/src/test/resources/compatibility/split-v2-indexed b/paimon-core/src/test/resources/compatibility/split-v2-indexed new file mode 100644 index 000000000000..76652e33d4ba Binary files /dev/null and b/paimon-core/src/test/resources/compatibility/split-v2-indexed differ diff --git a/paimon-core/src/test/resources/compatibility/split-v2-query-auth b/paimon-core/src/test/resources/compatibility/split-v2-query-auth new file mode 100644 index 000000000000..f8a06462a45f Binary files /dev/null and b/paimon-core/src/test/resources/compatibility/split-v2-query-auth differ diff --git a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/compact/changelog/ChangelogCompactTaskSerializer.java b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/compact/changelog/ChangelogCompactTaskSerializer.java index c7f56dc5de15..e796b92054ba 100644 --- a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/compact/changelog/ChangelogCompactTaskSerializer.java +++ b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/compact/changelog/ChangelogCompactTaskSerializer.java @@ -40,7 +40,7 @@ public class ChangelogCompactTaskSerializer implements SimpleVersionedSerializer { - private static final int CURRENT_VERSION = 2; + private static final int CURRENT_VERSION = 3; private final DataFileMetaSerializer dataFileSerializer; diff --git a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/compact/changelog/ChangelogCompactTaskSerializerTest.java b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/compact/changelog/ChangelogCompactTaskSerializerTest.java index df62115f2750..2dba94abd828 100644 --- a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/compact/changelog/ChangelogCompactTaskSerializerTest.java +++ b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/compact/changelog/ChangelogCompactTaskSerializerTest.java @@ -87,22 +87,23 @@ private List newFiles(int num) { private DataFileMeta newFile() { return DataFileMeta.create( - UUID.randomUUID().toString(), - 0, - 1, - row(0), - row(0), - newSimpleStats(0, 1), - newSimpleStats(0, 1), - 0, - 1, - 0, - 0, - 0L, - null, - FileSource.APPEND, - null, - null, - null); + UUID.randomUUID().toString(), + 0, + 1, + row(0), + row(0), + newSimpleStats(0, 1), + newSimpleStats(0, 1), + 0, + 1, + 0, + 0, + 0L, + null, + FileSource.APPEND, + null, + null, + null) + .withColumnMaxSequenceNumbers(new long[] {1L}); } } diff --git a/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/copy/CopyFilesUtil.java b/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/copy/CopyFilesUtil.java index 1a1e12efbe33..88dd1fb40311 100644 --- a/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/copy/CopyFilesUtil.java +++ b/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/copy/CopyFilesUtil.java @@ -73,7 +73,10 @@ public static DataFileMeta toNewDataFileMeta( oldFileMeta.valueStatsCols(), newExternalPath, oldFileMeta.firstRowId(), - oldFileMeta.writeCols()); + oldFileMeta.writeCols(), + // Column sequence numbers are positional and cannot be safely reused after + // changing the schema id. A null value makes readers fall back conservatively. + null); } public static IndexFileMeta toNewIndexFileMeta(IndexFileMeta oldFileMeta, String newFileName) { diff --git a/paimon-spark/paimon-spark-common/src/test/java/org/apache/paimon/spark/copy/CopyFilesUtilTest.java b/paimon-spark/paimon-spark-common/src/test/java/org/apache/paimon/spark/copy/CopyFilesUtilTest.java new file mode 100644 index 000000000000..2f973ffce3ad --- /dev/null +++ b/paimon-spark/paimon-spark-common/src/test/java/org/apache/paimon/spark/copy/CopyFilesUtilTest.java @@ -0,0 +1,61 @@ +/* + * 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.paimon.spark.copy; + +import org.apache.paimon.io.DataFileMeta; +import org.apache.paimon.stats.SimpleStats; + +import org.junit.jupiter.api.Test; + +import java.util.Arrays; +import java.util.Collections; + +import static org.assertj.core.api.Assertions.assertThat; + +/** Tests for {@link CopyFilesUtil}. */ +public class CopyFilesUtilTest { + + @Test + void testClearColumnSequencesWhenChangingSchemaId() { + DataFileMeta source = + DataFileMeta.forAppend( + "source.parquet", + 10L, + 2L, + SimpleStats.EMPTY_STATS, + 1L, + 3L, + 5L, + Collections.emptyList(), + null, + null, + null, + null, + null, + Arrays.asList("a", "b")) + .withColumnMaxSequenceNumbers(new long[] {2L, 3L}); + + DataFileMeta copied = CopyFilesUtil.toNewDataFileMeta(source, "copied.parquet", 6L); + + assertThat(copied.fileName()).isEqualTo("copied.parquet"); + assertThat(copied.schemaId()).isEqualTo(6L); + assertThat(copied.writeCols()).containsExactly("a", "b"); + assertThat(copied.columnMaxSequenceNumbers()).isNull(); + } +}