From 1e4af17df04df5620735ccab8ce11d62414112c6 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E4=BB=9F=E5=BC=8B?= Date: Thu, 20 Aug 2026 14:28:38 +0800 Subject: [PATCH 1/4] [core] Persist column sequence numbers in data file metadata --- .../java/org/apache/paimon/CoreOptions.java | 4 + .../DataEvolutionCompactTaskSerializer.java | 2 +- .../DataEvolutionNormalCompactTask.java | 67 +++++ .../org/apache/paimon/io/DataFileMeta.java | 36 ++- .../paimon/io/DataFileMeta08Serializer.java | 1 + .../paimon/io/DataFileMeta09Serializer.java | 1 + .../io/DataFileMeta10LegacySerializer.java | 1 + .../io/DataFileMeta12LegacySerializer.java | 1 + ...ataFileMetaFirstRowIdLegacySerializer.java | 3 +- .../paimon/io/DataFileMetaSerializer.java | 19 +- ...DataFileMetaWriteColsLegacySerializer.java | 116 +++++++++ .../apache/paimon/io/PojoDataFileMeta.java | 78 ++++-- .../paimon/io/ProjectedDataFileMeta.java | 23 ++ ...anifestEntryWriteColsLegacySerializer.java | 81 ++++++ .../sink/AppendCompactTaskSerializer.java | 2 +- .../sink/CommitMessageLegacyV2Serializer.java | 1 + .../table/sink/CommitMessageSerializer.java | 11 +- .../MultiTableCompactionTaskSerializer.java | 2 +- .../paimon/table/source/ChainSplit.java | 9 +- .../apache/paimon/table/source/DataSplit.java | 7 +- .../paimon/table/source/IncrementalSplit.java | 11 +- .../paimon/utils/DataEvolutionUtils.java | 38 +++ .../org/apache/paimon/CoreOptionsTest.java | 5 +- .../apache/paimon/append/BlobUpdateTest.java | 3 +- .../DataEvolutionCompactCoordinatorTest.java | 3 +- .../DataEvolutionCompactRangePlannerTest.java | 3 +- .../DataEvolutionNormalCompactTaskTest.java | 231 ++++++++++++++++++ .../crosspartition/IndexBootstrapTest.java | 1 + .../GlobalIndexBuilderUtilsTest.java | 2 + .../paimon/io/DataFileMetaSerializerTest.java | 48 +++- .../apache/paimon/io/DataFileTestUtils.java | 1 + .../paimon/io/ProjectedDataFileMetaTest.java | 33 +-- ...ommittableSerializerCompatibilityTest.java | 80 +++++- .../manifest/ManifestEntrySerializerTest.java | 22 ++ .../paimon/manifest/ManifestFileMetaTest.java | 7 +- .../manifest/ManifestFileMetaTestBase.java | 1 + .../paimon/manifest/ManifestFileTest.java | 45 ++++ .../NoPartitionManifestFileMetaTest.java | 3 +- .../mergetree/LevelsOrderingFixTest.java | 1 + .../compact/IntervalPartitionTest.java | 1 + .../BlobFallbackRecordReaderTest.java | 3 +- .../operation/DataEvolutionReadTest.java | 10 +- .../paimon/operation/ExpireSnapshotsTest.java | 2 + .../operation/ManifestEntryRunMergeTest.java | 3 +- .../operation/ManifestRewriteCleanupTest.java | 3 +- .../sink/CommitMessageSerializerTest.java | 8 + .../table/source/DataSplitCompatibleTest.java | 94 ++++++- .../table/source/SplitSerializerTest.java | 158 ++++++++---- .../paimon/utils/ChainTableUtilsTest.java | 1 + .../paimon/utils/DataEvolutionUtilsTest.java | 51 ++++ .../test/resources/compatibility/datasplit-v9 | Bin 0 -> 1034 bytes .../compatibility/manifest-committable-v13-v5 | Bin 0 -> 3146 bytes .../resources/compatibility/split-v2-chain | Bin 0 -> 1491 bytes .../resources/compatibility/split-v2-data | Bin 0 -> 944 bytes .../resources/compatibility/split-v2-fallback | Bin 0 -> 1508 bytes .../compatibility/split-v2-fallback-data | Bin 0 -> 945 bytes .../compatibility/split-v2-incremental | Bin 0 -> 978 bytes .../resources/compatibility/split-v2-indexed | Bin 0 -> 1009 bytes .../compatibility/split-v2-query-auth | Bin 0 -> 1059 bytes .../ChangelogCompactTaskSerializer.java | 2 +- .../ChangelogCompactSortOperatorTest.java | 1 + .../ChangelogCompactTaskSerializerTest.java | 35 +-- .../GenericIndexTopoBuilderTest.java | 1 + ...PendingSplitsCheckpointSerializerTest.java | 154 ++++++++++++ .../pending-splits-incremental-v1-chain-v2 | Bin 0 -> 2659 bytes .../paimon/spark/copy/CopyFilesUtil.java | 5 +- .../paimon/spark/copy/CopyFilesUtilTest.java | 61 +++++ .../CreateGlobalIndexProcedureTest.java | 1 + 68 files changed, 1456 insertions(+), 141 deletions(-) create mode 100644 paimon-core/src/main/java/org/apache/paimon/io/DataFileMetaWriteColsLegacySerializer.java create mode 100644 paimon-core/src/main/java/org/apache/paimon/manifest/ManifestEntryWriteColsLegacySerializer.java create mode 100644 paimon-core/src/test/java/org/apache/paimon/append/dataevolution/DataEvolutionNormalCompactTaskTest.java create mode 100644 paimon-core/src/test/resources/compatibility/datasplit-v9 create mode 100644 paimon-core/src/test/resources/compatibility/manifest-committable-v13-v5 create mode 100644 paimon-core/src/test/resources/compatibility/split-v2-chain create mode 100644 paimon-core/src/test/resources/compatibility/split-v2-data create mode 100644 paimon-core/src/test/resources/compatibility/split-v2-fallback create mode 100644 paimon-core/src/test/resources/compatibility/split-v2-fallback-data create mode 100644 paimon-core/src/test/resources/compatibility/split-v2-incremental create mode 100644 paimon-core/src/test/resources/compatibility/split-v2-indexed create mode 100644 paimon-core/src/test/resources/compatibility/split-v2-query-auth create mode 100644 paimon-flink/paimon-flink-common/src/test/resources/compatibility/pending-splits-incremental-v1-chain-v2 create mode 100644 paimon-spark/paimon-spark-common/src/test/java/org/apache/paimon/spark/copy/CopyFilesUtilTest.java diff --git a/paimon-api/src/main/java/org/apache/paimon/CoreOptions.java b/paimon-api/src/main/java/org/apache/paimon/CoreOptions.java index 9a6fd514dfbc..a19d1818c247 100644 --- a/paimon-api/src/main/java/org/apache/paimon/CoreOptions.java +++ b/paimon-api/src/main/java/org/apache/paimon/CoreOptions.java @@ -3734,6 +3734,10 @@ public GlobalIndexColumnUpdateAction globalIndexColumnUpdateAction() { return options.get(GLOBAL_INDEX_COLUMN_UPDATE_ACTION); } + public boolean ignoreIndexColumnUpdate() { + return globalIndexColumnUpdateAction() == GlobalIndexColumnUpdateAction.IGNORE; + } + public LookupStrategy lookupStrategy() { return LookupStrategy.from( mergeEngine().equals(MergeEngine.FIRST_ROW), 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 9f4e41c34dcf..10a368a49a49 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 @@ -26,25 +26,36 @@ import org.apache.paimon.manifest.FileSource; import org.apache.paimon.operation.AppendFileStoreWrite; import org.apache.paimon.reader.RecordReader; +import org.apache.paimon.schema.SchemaManager; +import org.apache.paimon.schema.TableSchema; 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.Pair; import org.apache.paimon.utils.RecordWriter; import org.apache.paimon.utils.SetUtils; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import javax.annotation.Nullable; + +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; 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.checkContiguousRowRange; +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. */ @@ -126,8 +137,64 @@ public CommitMessage doCompact(FileStoreTable table, String commitUser) throws E dataFileMeta = dataFileMeta.assignSequenceNumber( minSequenceId(compactBefore), maxSequenceId(compactBefore)); + if (options.ignoreIndexColumnUpdate()) { + long[] columnMaxSequenceNumbers = + compactedColumnMaxSequenceNumbers(table, dataFileMeta); + if (columnMaxSequenceNumbers != null) { + dataFileMeta = dataFileMeta.withColumnMaxSequenceNumbers(columnMaxSequenceNumbers); + } + } compactAfter.add(dataFileMeta); return commitMessage(compactBefore, compactAfter); } + + @Nullable + private long[] compactedColumnMaxSequenceNumbers( + FileStoreTable table, DataFileMeta outputFile) { + SchemaManager schemaManager = table.schemaManager(); + Map schemaCache = new HashMap<>(); + Function schemaLoader = + schemaId -> schemaCache.computeIfAbsent(schemaId, schemaManager::schema); + Map>, List> fileFieldsCache = new HashMap<>(); + + Map fieldMaxSequences = new HashMap<>(); + for (DataFileMeta input : compactBefore) { + List inputFields = + fileFieldsCache.computeIfAbsent( + Pair.of(input.schemaId(), input.writeCols()), + key -> fileFields(schemaLoader, input)); + long[] inputColumnSequences = input.columnMaxSequenceNumbers(); + for (int inputPosition = 0; inputPosition < inputFields.size(); inputPosition++) { + fieldMaxSequences.merge( + inputFields.get(inputPosition).id(), + fieldMaxSequenceNumber( + input, inputColumnSequences, inputPosition, inputFields.size()), + Math::max); + } + } + + long fallbackSequence = outputFile.maxSequenceNumber(); + List outputFields = + fileFieldsCache.computeIfAbsent( + Pair.of(outputFile.schemaId(), outputFile.writeCols()), + key -> fileFields(schemaLoader, outputFile)); + boolean allEqualToFileMax = + outputFields.stream() + .allMatch( + field -> + fieldMaxSequences.getOrDefault(field.id(), fallbackSequence) + == fallbackSequence); + if (allEqualToFileMax) { + return null; + } + + 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/io/DataFileMeta.java b/paimon-core/src/main/java/org/apache/paimon/io/DataFileMeta.java index f8b4e6aaf5b7..912c2e9f8b1e 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; @@ -150,7 +155,8 @@ static DataFileMeta forAppend( valueStatsCols, externalPath, firstRowId, - writeCols); + writeCols, + null); } static DataFileMeta create( @@ -173,7 +179,7 @@ static DataFileMeta create( @Nullable String externalPath, @Nullable Long firstRowId, @Nullable List writeCols) { - return new PojoDataFileMeta( + return create( fileName, fileSize, rowCount, @@ -193,7 +199,8 @@ static DataFileMeta create( valueStatsCols, externalPath, firstRowId, - writeCols); + writeCols, + null); } static DataFileMeta create( @@ -234,7 +241,8 @@ static DataFileMeta create( valueStatsCols, null, firstRowId, - writeCols); + writeCols, + null); } static DataFileMeta create( @@ -257,7 +265,8 @@ static DataFileMeta create( @Nullable List valueStatsCols, @Nullable String externalPath, @Nullable Long firstRowId, - @Nullable List writeCols) { + @Nullable List writeCols, + @Nullable long[] columnMaxSequenceNumbers) { return new PojoDataFileMeta( fileName, fileSize, @@ -278,7 +287,8 @@ static DataFileMeta create( valueStatsCols, externalPath, firstRowId, - writeCols); + writeCols, + columnMaxSequenceNumbers); } String fileName(); @@ -354,6 +364,16 @@ 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 + long[] columnMaxSequenceNumbers(); + DataFileMeta upgrade(int newLevel); DataFileMeta rename(String newFileName); @@ -362,6 +382,8 @@ default Range nonNullRowIdRange() { DataFileMeta assignSequenceNumber(long minSequenceNumber, long maxSequenceNumber); + DataFileMeta withColumnMaxSequenceNumbers(long[] columnMaxSequenceNumbers); + DataFileMeta assignFirstRowId(long firstRowId); DataFileMeta newFirstRowId(@Nullable Long newFirstRowId); diff --git a/paimon-core/src/main/java/org/apache/paimon/io/DataFileMeta08Serializer.java b/paimon-core/src/main/java/org/apache/paimon/io/DataFileMeta08Serializer.java index b646ef08ca39..c69a319f2d83 100644 --- a/paimon-core/src/main/java/org/apache/paimon/io/DataFileMeta08Serializer.java +++ b/paimon-core/src/main/java/org/apache/paimon/io/DataFileMeta08Serializer.java @@ -136,6 +136,7 @@ public DataFileMeta deserialize(DataInputView in) throws IOException { null, null, null, + null, null); } } diff --git a/paimon-core/src/main/java/org/apache/paimon/io/DataFileMeta09Serializer.java b/paimon-core/src/main/java/org/apache/paimon/io/DataFileMeta09Serializer.java index 662f1276c8d6..41e93bbcba5e 100644 --- a/paimon-core/src/main/java/org/apache/paimon/io/DataFileMeta09Serializer.java +++ b/paimon-core/src/main/java/org/apache/paimon/io/DataFileMeta09Serializer.java @@ -142,6 +142,7 @@ public DataFileMeta deserialize(DataInputView in) throws IOException { null, null, null, + null, null); } } diff --git a/paimon-core/src/main/java/org/apache/paimon/io/DataFileMeta10LegacySerializer.java b/paimon-core/src/main/java/org/apache/paimon/io/DataFileMeta10LegacySerializer.java index dca1aa528f33..2171085da850 100644 --- a/paimon-core/src/main/java/org/apache/paimon/io/DataFileMeta10LegacySerializer.java +++ b/paimon-core/src/main/java/org/apache/paimon/io/DataFileMeta10LegacySerializer.java @@ -147,6 +147,7 @@ public DataFileMeta deserialize(DataInputView in) throws IOException { row.isNullAt(16) ? null : fromStringArrayData(row.getArray(16)), null, null, + null, null); } } diff --git a/paimon-core/src/main/java/org/apache/paimon/io/DataFileMeta12LegacySerializer.java b/paimon-core/src/main/java/org/apache/paimon/io/DataFileMeta12LegacySerializer.java index e888c1ca7462..b56a3396e330 100644 --- a/paimon-core/src/main/java/org/apache/paimon/io/DataFileMeta12LegacySerializer.java +++ b/paimon-core/src/main/java/org/apache/paimon/io/DataFileMeta12LegacySerializer.java @@ -149,6 +149,7 @@ public DataFileMeta deserialize(DataInputView in) throws IOException { row.isNullAt(16) ? null : fromStringArrayData(row.getArray(16)), row.isNullAt(17) ? null : row.getString(17).toString(), null, + null, null); } } 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..63981a1e299b 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; + + public 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)), + null); + } +} 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..6cf1eaa0d59c 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, @@ -98,7 +100,8 @@ public PojoDataFileMeta( @Nullable List valueStatsCols, @Nullable String externalPath, @Nullable Long firstRowId, - @Nullable List writeCols) { + @Nullable List writeCols, + @Nullable long[] columnMaxSequenceNumbers) { this.fileName = fileName; this.fileSize = fileSize; @@ -123,6 +126,8 @@ public PojoDataFileMeta( this.externalPath = externalPath; this.firstRowId = firstRowId; this.writeCols = writeCols; + this.columnMaxSequenceNumbers = + columnMaxSequenceNumbers == null ? null : columnMaxSequenceNumbers.clone(); } @Override @@ -239,6 +244,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 +273,8 @@ public PojoDataFileMeta upgrade(int newLevel) { valueStatsCols, externalPath, firstRowId, - writeCols); + writeCols, + columnMaxSequenceNumbers); } @Override @@ -288,7 +300,8 @@ public PojoDataFileMeta rename(String newFileName) { valueStatsCols, newExternalPath, firstRowId, - writeCols); + writeCols, + columnMaxSequenceNumbers); } @Override @@ -313,7 +326,8 @@ public PojoDataFileMeta copyWithoutStats() { Collections.emptyList(), externalPath, firstRowId, - writeCols); + writeCols, + columnMaxSequenceNumbers); } @Override @@ -338,7 +352,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 +404,8 @@ public PojoDataFileMeta assignFirstRowId(long firstRowId) { valueStatsCols, externalPath, firstRowId, - writeCols); + writeCols, + columnMaxSequenceNumbers); } @Override @@ -388,7 +430,8 @@ public PojoDataFileMeta newFirstRowId(@Nullable Long newFirstRowId) { valueStatsCols, externalPath, newFirstRowId, - writeCols); + writeCols, + columnMaxSequenceNumbers); } @Override @@ -413,7 +456,8 @@ public PojoDataFileMeta copy(List newExtraFiles) { valueStatsCols, externalPath, firstRowId, - writeCols); + writeCols, + columnMaxSequenceNumbers); } @Override @@ -438,7 +482,8 @@ public PojoDataFileMeta newExternalPath(String newExternalPath) { valueStatsCols, newExternalPath, firstRowId, - writeCols); + writeCols, + columnMaxSequenceNumbers); } @Override @@ -463,7 +508,8 @@ public PojoDataFileMeta copy(byte[] newEmbeddedIndex) { valueStatsCols, externalPath, firstRowId, - writeCols); + writeCols, + columnMaxSequenceNumbers); } @Override @@ -494,7 +540,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 +566,8 @@ public int hashCode() { valueStatsCols, externalPath, firstRowId, - writeCols); + writeCols, + Arrays.hashCode(columnMaxSequenceNumbers)); } @Override @@ -529,7 +577,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 +598,7 @@ public String toString() { valueStatsCols, externalPath, firstRowId, - writeCols); + writeCols, + Arrays.toString(columnMaxSequenceNumbers)); } } diff --git a/paimon-core/src/main/java/org/apache/paimon/io/ProjectedDataFileMeta.java b/paimon-core/src/main/java/org/apache/paimon/io/ProjectedDataFileMeta.java index e525c4e2fc53..bc5edc413ede 100644 --- a/paimon-core/src/main/java/org/apache/paimon/io/ProjectedDataFileMeta.java +++ b/paimon-core/src/main/java/org/apache/paimon/io/ProjectedDataFileMeta.java @@ -254,6 +254,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(); @@ -291,6 +307,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)"); @@ -388,6 +409,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/manifest/ManifestEntryWriteColsLegacySerializer.java b/paimon-core/src/main/java/org/apache/paimon/manifest/ManifestEntryWriteColsLegacySerializer.java new file mode 100644 index 000000000000..3640fe65eee3 --- /dev/null +++ b/paimon-core/src/main/java/org/apache/paimon/manifest/ManifestEntryWriteColsLegacySerializer.java @@ -0,0 +1,81 @@ +/* + * 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.manifest; + +import org.apache.paimon.data.GenericRow; +import org.apache.paimon.data.InternalRow; +import org.apache.paimon.io.DataFileMetaWriteColsLegacySerializer; +import org.apache.paimon.types.RowType; +import org.apache.paimon.utils.ObjectSerializer; +import org.apache.paimon.utils.OffsetRow; + +import java.util.Arrays; + +import static org.apache.paimon.utils.SerializationUtils.deserializeBinaryRow; +import static org.apache.paimon.utils.SerializationUtils.serializeBinaryRow; + +/** Legacy serializer for {@link ManifestEntry} before column sequence numbers were introduced. */ +public class ManifestEntryWriteColsLegacySerializer extends ObjectSerializer { + + private static final long serialVersionUID = 1L; + private static final int FORMAT_IDENTIFIER = 2; + + private static final RowType SCHEMA = + new RowType( + false, + Arrays.asList( + ManifestEntry.SCHEMA.getField(ManifestEntry.KIND), + ManifestEntry.SCHEMA.getField(ManifestEntry.PARTITION), + ManifestEntry.SCHEMA.getField(ManifestEntry.BUCKET), + ManifestEntry.SCHEMA.getField(ManifestEntry.TOTAL_BUCKETS), + ManifestEntry.SCHEMA + .getField(ManifestEntry.FILE) + .newType(DataFileMetaWriteColsLegacySerializer.SCHEMA))); + + private final DataFileMetaWriteColsLegacySerializer dataFileMetaSerializer; + + public ManifestEntryWriteColsLegacySerializer() { + super(ManifestSchemaUtils.withFormatIdentifier(SCHEMA)); + this.dataFileMetaSerializer = new DataFileMetaWriteColsLegacySerializer(); + } + + @Override + public InternalRow toRow(ManifestEntry entry) { + return GenericRow.of( + FORMAT_IDENTIFIER, + entry.kind().toByteValue(), + serializeBinaryRow(entry.partition()), + entry.bucket(), + entry.totalBuckets(), + dataFileMetaSerializer.toRow(entry.file())); + } + + @Override + public ManifestEntry fromRow(InternalRow row) { + ManifestEntrySerializer.checkFormatIdentifier(row.getInt(0)); + InternalRow dataRow = new OffsetRow(row.getFieldCount() - 1, 1).replace(row); + return ManifestEntry.create( + FileKind.fromByteValue(dataRow.getByte(0)), + deserializeBinaryRow(dataRow.getBinary(1)), + dataRow.getInt(2), + dataRow.getInt(3), + dataFileMetaSerializer.fromRow( + dataRow.getRow(4, dataFileMetaSerializer.numFields()))); + } +} 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/CommitMessageLegacyV2Serializer.java b/paimon-core/src/main/java/org/apache/paimon/table/sink/CommitMessageLegacyV2Serializer.java index f60415bc8017..3f8622033830 100644 --- a/paimon-core/src/main/java/org/apache/paimon/table/sink/CommitMessageLegacyV2Serializer.java +++ b/paimon-core/src/main/java/org/apache/paimon/table/sink/CommitMessageLegacyV2Serializer.java @@ -164,6 +164,7 @@ public DataFileMeta fromRow(InternalRow row) { null, null, null, + null, null); } } 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 b4bd8be490ce..7bb2436d929b 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,11 +21,13 @@ 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.FunctionWithIOException; +import org.apache.paimon.utils.ObjectSerializer; import javax.annotation.Nullable; @@ -45,7 +47,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; @@ -239,7 +241,7 @@ public void serialize(DataOutputView out) throws IOException { public static IncrementalSplit deserialize(DataInputView in) throws IOException { int version = in.readInt(); - if (version != VERSION) { + if (version < 1 || version > VERSION) { throw new UnsupportedOperationException("Unsupported version: " + version); } @@ -248,7 +250,10 @@ public static IncrementalSplit deserialize(DataInputView in) throws IOException int bucket = in.readInt(); int 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/utils/DataEvolutionUtils.java b/paimon-core/src/main/java/org/apache/paimon/utils/DataEvolutionUtils.java index 262da197aea9..b8751ff81b8d 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 @@ -24,6 +24,8 @@ import org.apache.paimon.table.source.DataSplit; import org.apache.paimon.types.DataField; +import javax.annotation.Nullable; + import java.util.ArrayList; import java.util.Collection; import java.util.Collections; @@ -119,6 +121,42 @@ private static Set resolveFileFieldIds( return ids; } + /** 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(); + 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. + DataField field = fieldsByName.get(writeCol); + if (field != null) { + fields.add(field); + } + } + return fields; + } + + /** Returns the latest sequence known for a physical field position in the file. */ + public static long fieldMaxSequenceNumber( + DataFileMeta file, + @Nullable long[] columnSequences, + int fieldPosition, + int physicalFieldCount) { + if (columnSequences == null || columnSequences.length != physicalFieldCount) { + return file.maxSequenceNumber(); + } + return columnSequences[fieldPosition]; + } + /** * Retrieve the anchor file of a row range group. Always the oldest normal file. Files are * compared by (max_seq, fileName) pairs. diff --git a/paimon-core/src/test/java/org/apache/paimon/CoreOptionsTest.java b/paimon-core/src/test/java/org/apache/paimon/CoreOptionsTest.java index 266ff305f406..ac71ff2552aa 100644 --- a/paimon-core/src/test/java/org/apache/paimon/CoreOptionsTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/CoreOptionsTest.java @@ -121,8 +121,11 @@ public void testIgnoreGlobalIndexColumnUpdateAction() { Options conf = new Options(); conf.setString(CoreOptions.GLOBAL_INDEX_COLUMN_UPDATE_ACTION.key(), "IGNORE"); - assertThat(new CoreOptions(conf).globalIndexColumnUpdateAction()) + CoreOptions options = new CoreOptions(conf); + assertThat(options.globalIndexColumnUpdateAction()) .isEqualTo(CoreOptions.GlobalIndexColumnUpdateAction.IGNORE); + assertThat(options.ignoreIndexColumnUpdate()).isTrue(); + assertThat(new CoreOptions(new Options()).ignoreIndexColumnUpdate()).isFalse(); } @Test diff --git a/paimon-core/src/test/java/org/apache/paimon/append/BlobUpdateTest.java b/paimon-core/src/test/java/org/apache/paimon/append/BlobUpdateTest.java index b86b020cd199..9d1f05aa6537 100644 --- a/paimon-core/src/test/java/org/apache/paimon/append/BlobUpdateTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/append/BlobUpdateTest.java @@ -502,7 +502,8 @@ private static DataFileMeta writeBlobFile( null, null, firstRowId, - writeCols); + writeCols, + null); } @Override diff --git a/paimon-core/src/test/java/org/apache/paimon/append/dataevolution/DataEvolutionCompactCoordinatorTest.java b/paimon-core/src/test/java/org/apache/paimon/append/dataevolution/DataEvolutionCompactCoordinatorTest.java index 8e8c95193a71..cd784064bbae 100644 --- a/paimon-core/src/test/java/org/apache/paimon/append/dataevolution/DataEvolutionCompactCoordinatorTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/append/dataevolution/DataEvolutionCompactCoordinatorTest.java @@ -857,7 +857,8 @@ private DataFileMeta createDataFileMeta( null, null, firstRowId, - writeCols); + writeCols, + null); } private DataEvolutionCompactCoordinator.CompactPlanner blobPlanner( diff --git a/paimon-core/src/test/java/org/apache/paimon/append/dataevolution/DataEvolutionCompactRangePlannerTest.java b/paimon-core/src/test/java/org/apache/paimon/append/dataevolution/DataEvolutionCompactRangePlannerTest.java index 31003beb7501..88278ec40f54 100644 --- a/paimon-core/src/test/java/org/apache/paimon/append/dataevolution/DataEvolutionCompactRangePlannerTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/append/dataevolution/DataEvolutionCompactRangePlannerTest.java @@ -498,6 +498,7 @@ private DataFileMeta dataFile( null, null, firstRowId, - writeColumns); + writeColumns, + null); } } diff --git a/paimon-core/src/test/java/org/apache/paimon/append/dataevolution/DataEvolutionNormalCompactTaskTest.java b/paimon-core/src/test/java/org/apache/paimon/append/dataevolution/DataEvolutionNormalCompactTaskTest.java new file mode 100644 index 000000000000..abf036b97497 --- /dev/null +++ b/paimon-core/src/test/java/org/apache/paimon/append/dataevolution/DataEvolutionNormalCompactTaskTest.java @@ -0,0 +1,231 @@ +/* + * 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.append.dataevolution; + +import org.apache.paimon.CoreOptions; +import org.apache.paimon.Snapshot; +import org.apache.paimon.data.BinaryString; +import org.apache.paimon.data.GenericRow; +import org.apache.paimon.io.DataFileMeta; +import org.apache.paimon.manifest.ManifestEntry; +import org.apache.paimon.schema.Schema; +import org.apache.paimon.table.FileStoreTable; +import org.apache.paimon.table.TableTestBase; +import org.apache.paimon.table.sink.BatchTableCommit; +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.types.DataField; +import org.apache.paimon.types.DataTypes; +import org.apache.paimon.types.RowType; + +import org.junit.jupiter.api.Test; + +import java.util.ArrayList; +import java.util.Arrays; +import java.util.Collections; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.stream.Collectors; + +import static org.apache.paimon.utils.DataEvolutionUtils.fileFields; +import static org.assertj.core.api.Assertions.assertThat; + +/** Tests for column sequence propagation in {@link DataEvolutionNormalCompactTask}. */ +public class DataEvolutionNormalCompactTaskTest extends TableTestBase { + + private static final int ROW_COUNT = 100; + + @Override + public Schema schemaDefault() { + return Schema.newBuilder() + .column("dt", DataTypes.STRING()) + .column("f0", DataTypes.INT()) + .column("f1", DataTypes.STRING()) + .partitionKeys(Collections.singletonList("dt")) + .option(CoreOptions.ROW_TRACKING_ENABLED.key(), "true") + .option(CoreOptions.DATA_EVOLUTION_ENABLED.key(), "true") + .build(); + } + + @Test + public void testPropagateColumnSequencesAcrossCompactions() throws Exception { + write(); + + int f0Id = getTableDefault().rowType().getField("f0").id(); + int f1Id = getTableDefault().rowType().getField("f1").id(); + DataFileMeta firstCompact = + updateColumnsAndCompact( + Collections.singletonList("f1"), + 1, + CoreOptions.GlobalIndexColumnUpdateAction.IGNORE); + + long f0Sequence = columnSequence(firstCompact, f0Id); + assertThat(f0Sequence).isLessThan(firstCompact.maxSequenceNumber()); + assertThat(columnSequence(firstCompact, f1Id)).isEqualTo(firstCompact.maxSequenceNumber()); + + DataFileMeta secondCompact = + updateColumnsAndCompact( + Collections.singletonList("f1"), + 2, + CoreOptions.GlobalIndexColumnUpdateAction.IGNORE); + assertThat(columnSequence(secondCompact, f0Id)).isEqualTo(f0Sequence); + assertThat(columnSequence(secondCompact, f1Id)) + .isEqualTo(secondCompact.maxSequenceNumber()); + } + + @Test + public void testOmitAndReconstructRedundantColumnSequences() throws Exception { + write(); + + DataFileMeta fullUpdate = + updateColumnsAndCompact( + Arrays.asList("f0", "f1"), + 1, + CoreOptions.GlobalIndexColumnUpdateAction.IGNORE); + assertThat(fullUpdate.columnMaxSequenceNumbers()).isNull(); + + int f0Id = getTableDefault().rowType().getField("f0").id(); + DataFileMeta partialUpdate = + updateColumnsAndCompact( + Collections.singletonList("f1"), + 2, + CoreOptions.GlobalIndexColumnUpdateAction.IGNORE); + assertThat(columnSequence(partialUpdate, f0Id)) + .isEqualTo(fullUpdate.maxSequenceNumber()) + .isLessThan(partialUpdate.maxSequenceNumber()); + } + + @Test + public void testOmitColumnSequencesUnlessUpdatesAreIgnored() throws Exception { + write(); + + DataFileMeta compacted = + updateColumnsAndCompact( + Collections.singletonList("f1"), + 1, + CoreOptions.GlobalIndexColumnUpdateAction.THROW_ERROR); + assertThat(compacted.columnMaxSequenceNumbers()).isNull(); + } + + private void write() throws Exception { + createTableDefault(); + + BatchWriteBuilder builder = getTableDefault().newBatchWriteBuilder(); + try (BatchTableWrite write = builder.newWrite()) { + for (int i = 0; i < ROW_COUNT; i++) { + write.write( + GenericRow.of( + BinaryString.fromString("p0"), + i, + BinaryString.fromString("f1_" + i))); + } + try (BatchTableCommit commit = builder.newCommit()) { + commit.commit(write.prepareCommit()); + } + } + } + + private DataFileMeta updateColumnsAndCompact( + List columns, + int updateRound, + CoreOptions.GlobalIndexColumnUpdateAction updateAction) + throws Exception { + Map writeOptions = new HashMap<>(); + writeOptions.put( + CoreOptions.GLOBAL_INDEX_COLUMN_UPDATE_ACTION.key(), updateAction.toString()); + FileStoreTable table = getTableDefault().copy(writeOptions); + List writeColumns = new ArrayList<>(); + writeColumns.add("dt"); + writeColumns.addAll(columns); + RowType writeType = table.rowType().project(writeColumns); + BatchWriteBuilder writeBuilder = table.newBatchWriteBuilder(); + try (BatchTableWrite batchWrite = writeBuilder.newWrite().withWriteType(writeType)) { + for (int i = 0; i < ROW_COUNT; i++) { + List values = new ArrayList<>(); + values.add(BinaryString.fromString("p0")); + for (String column : columns) { + values.add( + "f0".equals(column) + ? i + updateRound * ROW_COUNT + : BinaryString.fromString("updated_" + updateRound + "_" + i)); + } + batchWrite.write(GenericRow.of(values.toArray())); + } + 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); + return rowRangeFiles.get(0); + } + + 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 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())); + } + } +} diff --git a/paimon-core/src/test/java/org/apache/paimon/crosspartition/IndexBootstrapTest.java b/paimon-core/src/test/java/org/apache/paimon/crosspartition/IndexBootstrapTest.java index 0e352b63ec7d..58dad8079f51 100644 --- a/paimon-core/src/test/java/org/apache/paimon/crosspartition/IndexBootstrapTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/crosspartition/IndexBootstrapTest.java @@ -339,6 +339,7 @@ private static DataFileMeta newFile(long timeMillis) { null, null, null, + null, null); } diff --git a/paimon-core/src/test/java/org/apache/paimon/globalindex/GlobalIndexBuilderUtilsTest.java b/paimon-core/src/test/java/org/apache/paimon/globalindex/GlobalIndexBuilderUtilsTest.java index 549814fafee6..063114a99611 100644 --- a/paimon-core/src/test/java/org/apache/paimon/globalindex/GlobalIndexBuilderUtilsTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/globalindex/GlobalIndexBuilderUtilsTest.java @@ -331,6 +331,7 @@ private ManifestEntry createEntry(Long firstRowId, long rowCount) { null, null, firstRowId, + null, null); return ManifestEntry.create(FileKind.ADD, BinaryRow.EMPTY_ROW, 0, 1, file); } @@ -362,6 +363,7 @@ private static DataFileMeta createDataFileMeta(long firstRowId, long rowCount) { null, null, firstRowId, + null, null); } } 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/io/DataFileTestUtils.java b/paimon-core/src/test/java/org/apache/paimon/io/DataFileTestUtils.java index 974fe7aaf8d8..556e4de335b4 100644 --- a/paimon-core/src/test/java/org/apache/paimon/io/DataFileTestUtils.java +++ b/paimon-core/src/test/java/org/apache/paimon/io/DataFileTestUtils.java @@ -60,6 +60,7 @@ public static DataFileMeta newFile(long minSeq, long maxSeq) { null, null, null, + null, null); } diff --git a/paimon-core/src/test/java/org/apache/paimon/io/ProjectedDataFileMetaTest.java b/paimon-core/src/test/java/org/apache/paimon/io/ProjectedDataFileMetaTest.java index d684b447392c..744aa14ed0e4 100644 --- a/paimon-core/src/test/java/org/apache/paimon/io/ProjectedDataFileMetaTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/io/ProjectedDataFileMetaTest.java @@ -42,20 +42,21 @@ public class ProjectedDataFileMetaTest { 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}); ProjectedDataFileMeta actual = ProjectedDataFileMeta.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/manifest/ManifestCommittableSerializerCompatibilityTest.java b/paimon-core/src/test/java/org/apache/paimon/manifest/ManifestCommittableSerializerCompatibilityTest.java index 20eac11da7e8..dbe1ebfab09e 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,64 @@ /** 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"), + null) + .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(current).isEqualTo(serialized); + assertThat(serializer.deserialize(5, serialized)).isEqualTo(committable); + } + @Test public void testCompatibilityToV5CommitV11() throws IOException { String fileName = "manifest-committable-v11-v5"; @@ -80,7 +139,8 @@ public void testCompatibilityToV5CommitV11() throws IOException { Arrays.asList("field1", "field2", "field3"), "hdfs://localhost:9000/path/to/file", 1L, - Arrays.asList("asdf", "qwer", "zxcv")); + Arrays.asList("asdf", "qwer", "zxcv"), + null); List dataFiles = Collections.singletonList(dataFile); GlobalIndexMeta globalIndexMeta = new GlobalIndexMeta( @@ -212,7 +272,8 @@ public void testCompatibilityToV4CommitV11() throws IOException { Arrays.asList("field1", "field2", "field3"), "hdfs://localhost:9000/path/to/file", 1L, - Arrays.asList("asdf", "qwer", "zxcv")); + Arrays.asList("asdf", "qwer", "zxcv"), + null); List dataFiles = Collections.singletonList(dataFile); GlobalIndexMeta globalIndexMeta = new GlobalIndexMeta(1L, 2L, 3, new int[] {5, 6, 7}, new byte[] {0x23, 0x45}); @@ -311,7 +372,8 @@ public void testCompatibilityToV4CommitV10() throws IOException { Arrays.asList("field1", "field2", "field3"), "hdfs://localhost:9000/path/to/file", 1L, - Arrays.asList("asdf", "qwer", "zxcv")); + Arrays.asList("asdf", "qwer", "zxcv"), + null); List dataFiles = Collections.singletonList(dataFile); IndexFileMeta hashIndexFile = new IndexFileMeta( @@ -400,7 +462,8 @@ public void testCompatibilityToV4CommitV9() throws IOException { Arrays.asList("field1", "field2", "field3"), "hdfs://localhost:9000/path/to/file", 1L, - Arrays.asList("asdf", "qwer", "zxcv")); + Arrays.asList("asdf", "qwer", "zxcv"), + null); List dataFiles = Collections.singletonList(dataFile); LinkedHashMap dvRanges = new LinkedHashMap<>(); @@ -484,6 +547,7 @@ public void testCompatibilityToV4CommitV8() throws IOException { Arrays.asList("field1", "field2", "field3"), "hdfs://localhost:9000/path/to/file", 1L, + null, null); List dataFiles = Collections.singletonList(dataFile); @@ -563,6 +627,7 @@ public void testCompatibilityToV4CommitV7() throws IOException { Arrays.asList("field1", "field2", "field3"), "hdfs://localhost:9000/path/to/file", null, + null, null); List dataFiles = Collections.singletonList(dataFile); @@ -641,6 +706,7 @@ public void testCompatibilityToV3CommitV7() throws IOException { Arrays.asList("field1", "field2", "field3"), "hdfs://localhost:9000/path/to/file", null, + null, null); List dataFiles = Collections.singletonList(dataFile); @@ -716,6 +782,7 @@ public void testCompatibilityToV3CommitV6() throws IOException { Arrays.asList("field1", "field2", "field3"), "hdfs://localhost:9000/path/to/file", null, + null, null); List dataFiles = Collections.singletonList(dataFile); @@ -791,6 +858,7 @@ public void testCompatibilityToV3CommitV5() throws IOException { Arrays.asList("field1", "field2", "field3"), null, null, + null, null); List dataFiles = Collections.singletonList(dataFile); @@ -865,6 +933,7 @@ public void testCompatibilityToV3CommitV4() throws IOException { Arrays.asList("field1", "field2", "field3"), null, null, + null, null); List dataFiles = Collections.singletonList(dataFile); @@ -940,6 +1009,7 @@ public void testCompatibilityToV3CommitV3() throws IOException { null, null, null, + null, null); List dataFiles = Collections.singletonList(dataFile); @@ -1015,6 +1085,7 @@ public void testCompatibilityToV2CommitV2() throws IOException { null, null, null, + null, null); List dataFiles = Collections.singletonList(dataFile); @@ -1090,6 +1161,7 @@ public void testCompatibilityToVersion2PaimonV07() throws IOException { null, null, null, + null, null); List dataFiles = Collections.singletonList(dataFile); diff --git a/paimon-core/src/test/java/org/apache/paimon/manifest/ManifestEntrySerializerTest.java b/paimon-core/src/test/java/org/apache/paimon/manifest/ManifestEntrySerializerTest.java index 0547c9acce25..5005f04cc91a 100644 --- a/paimon-core/src/test/java/org/apache/paimon/manifest/ManifestEntrySerializerTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/manifest/ManifestEntrySerializerTest.java @@ -23,6 +23,8 @@ import org.junit.jupiter.api.Test; +import java.io.IOException; + import static org.assertj.core.api.Assertions.assertThat; /** Tests for {@link ManifestEntrySerializer}. */ @@ -35,6 +37,26 @@ void testFormatIdentifier() { assertThat(new ManifestEntrySerializer().toRow(gen.next()).getInt(0)).isEqualTo(2); } + @Test + void testWriteColsLegacySerializer() throws IOException { + ManifestEntry expected = gen.next(); + ManifestEntry withColumnSequences = + ManifestEntry.create( + expected.kind(), + expected.partition(), + expected.bucket(), + expected.totalBuckets(), + expected.file().withColumnMaxSequenceNumbers(new long[] {3L, 42L})); + ManifestEntryWriteColsLegacySerializer serializer = + new ManifestEntryWriteColsLegacySerializer(); + + ManifestEntry actual = + serializer.deserializeFromBytes(serializer.serializeToBytes(withColumnSequences)); + + assertThat(actual).isEqualTo(expected); + assertThat(actual.file().columnMaxSequenceNumbers()).isNull(); + } + @Override protected ObjectSerializer serializer() { return new ManifestEntrySerializer(); diff --git a/paimon-core/src/test/java/org/apache/paimon/manifest/ManifestFileMetaTest.java b/paimon-core/src/test/java/org/apache/paimon/manifest/ManifestFileMetaTest.java index 6da1b8d2f7a0..0120471b4d06 100644 --- a/paimon-core/src/test/java/org/apache/paimon/manifest/ManifestFileMetaTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/manifest/ManifestFileMetaTest.java @@ -2737,6 +2737,7 @@ private ManifestEntry makeMultiPartEntry( null, null, null, + null, null)); } @@ -2804,7 +2805,8 @@ private ManifestEntry makeMultiPartRowIdEntry( null, null, firstRowId, - Collections.singletonList("f0"))); + Collections.singletonList("f0"), + null)); } private List readFileNames( @@ -2887,7 +2889,8 @@ private ManifestEntry makeRowIdEntry( null, externalPath, firstRowId, - Collections.singletonList("f0"))); + Collections.singletonList("f0"), + null)); } private List readEntries(List manifestMetas) { diff --git a/paimon-core/src/test/java/org/apache/paimon/manifest/ManifestFileMetaTestBase.java b/paimon-core/src/test/java/org/apache/paimon/manifest/ManifestFileMetaTestBase.java index 649ebb73ec99..b7eff0526baf 100644 --- a/paimon-core/src/test/java/org/apache/paimon/manifest/ManifestFileMetaTestBase.java +++ b/paimon-core/src/test/java/org/apache/paimon/manifest/ManifestFileMetaTestBase.java @@ -98,6 +98,7 @@ protected ManifestEntry makeEntry( null, null, null, + null, null)); } diff --git a/paimon-core/src/test/java/org/apache/paimon/manifest/ManifestFileTest.java b/paimon-core/src/test/java/org/apache/paimon/manifest/ManifestFileTest.java index d3ba2527de25..b551edc959fc 100644 --- a/paimon-core/src/test/java/org/apache/paimon/manifest/ManifestFileTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/manifest/ManifestFileTest.java @@ -30,6 +30,7 @@ import org.apache.paimon.fs.PositionOutputStream; import org.apache.paimon.fs.local.LocalFileIO; import org.apache.paimon.io.DataFileMeta; +import org.apache.paimon.io.DataFileMetaWriteColsLegacySerializer; import org.apache.paimon.options.Options; import org.apache.paimon.partition.PartitionPredicate; import org.apache.paimon.schema.SchemaManager; @@ -530,6 +531,50 @@ void testAvroReaderReadsLegacyDataFileMetaWithFewerFields() throws Exception { } } + @Test + void testLegacyAvroReaderSkipsColumnSequenceNumbers() throws Exception { + ManifestEntry expected = gen.next(); + ManifestEntry source = + ManifestEntry.create( + expected.kind(), + expected.partition(), + expected.bucket(), + expected.totalBuckets(), + expected.file().withColumnMaxSequenceNumbers(new long[] {3L, 42L})); + List legacyManifestFields = + ManifestEntry.MANIFEST_ROW_TYPE.getFields().stream() + .map( + field -> + ManifestEntry.FILE.equals(field.name()) + ? field.newType( + DataFileMetaWriteColsLegacySerializer + .SCHEMA) + : field) + .collect(Collectors.toList()); + RowType legacyManifestType = new RowType(false, legacyManifestFields); + Path path = new Path(new Path(tempDir.toUri()), "new-manifest.avro"); + LocalFileIO fileIO = LocalFileIO.create(); + ManifestEntrySerializer serializer = new ManifestEntrySerializer(); + + try (PositionOutputStream out = fileIO.newOutputStream(path, false); + FormatWriter writer = + avro.createWriterFactory(ManifestEntry.MANIFEST_ROW_TYPE) + .create(out, "zstd")) { + writer.addElement(serializer.toRow(source)); + } + + ManifestEntry actual; + try (ManifestAvroReader reader = new ManifestAvroReader(fileIO.newInputStream(path)); + CloseableIterator rows = reader.read(legacyManifestType, null, null)) { + assertThat(rows.hasNext()).isTrue(); + actual = new ManifestEntryWriteColsLegacySerializer().fromRow(rows.next()); + assertThat(rows.hasNext()).isFalse(); + } + + assertThat(actual).isEqualTo(expected); + assertThat(actual.file().columnMaxSequenceNumbers()).isNull(); + } + @Test void testAvroReaderRejectsReorderedTopLevelFields() throws Exception { ManifestEntry entry = gen.next(); diff --git a/paimon-core/src/test/java/org/apache/paimon/manifest/NoPartitionManifestFileMetaTest.java b/paimon-core/src/test/java/org/apache/paimon/manifest/NoPartitionManifestFileMetaTest.java index 52ac56608b86..3170c08d02c9 100644 --- a/paimon-core/src/test/java/org/apache/paimon/manifest/NoPartitionManifestFileMetaTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/manifest/NoPartitionManifestFileMetaTest.java @@ -191,6 +191,7 @@ private ManifestEntry makeRowIdEntry( null, null, firstRowId, - Collections.singletonList("f0"))); + Collections.singletonList("f0"), + null)); } } diff --git a/paimon-core/src/test/java/org/apache/paimon/mergetree/LevelsOrderingFixTest.java b/paimon-core/src/test/java/org/apache/paimon/mergetree/LevelsOrderingFixTest.java index 02d52df9494f..d34cc1cfd9d3 100644 --- a/paimon-core/src/test/java/org/apache/paimon/mergetree/LevelsOrderingFixTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/mergetree/LevelsOrderingFixTest.java @@ -119,6 +119,7 @@ private DataFileMeta createTestFile( null, null, null, + null, null); } } diff --git a/paimon-core/src/test/java/org/apache/paimon/mergetree/compact/IntervalPartitionTest.java b/paimon-core/src/test/java/org/apache/paimon/mergetree/compact/IntervalPartitionTest.java index 464b26b944f0..ede8b678e946 100644 --- a/paimon-core/src/test/java/org/apache/paimon/mergetree/compact/IntervalPartitionTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/mergetree/compact/IntervalPartitionTest.java @@ -187,6 +187,7 @@ private DataFileMeta makeInterval(int left, int right) { null, null, null, + null, null); } diff --git a/paimon-core/src/test/java/org/apache/paimon/operation/BlobFallbackRecordReaderTest.java b/paimon-core/src/test/java/org/apache/paimon/operation/BlobFallbackRecordReaderTest.java index 1c41741b701f..859e8bc5b885 100644 --- a/paimon-core/src/test/java/org/apache/paimon/operation/BlobFallbackRecordReaderTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/operation/BlobFallbackRecordReaderTest.java @@ -560,7 +560,8 @@ private static DataFileMeta blobFile( null, null, firstRowId, - Arrays.asList(BLOB_FIELD)); + Arrays.asList(BLOB_FIELD), + null); } private static List ranges(long... bounds) { diff --git a/paimon-core/src/test/java/org/apache/paimon/operation/DataEvolutionReadTest.java b/paimon-core/src/test/java/org/apache/paimon/operation/DataEvolutionReadTest.java index 51a0a93aa233..24d299495753 100644 --- a/paimon-core/src/test/java/org/apache/paimon/operation/DataEvolutionReadTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/operation/DataEvolutionReadTest.java @@ -354,7 +354,8 @@ private DataFileMeta createBlobFileWithSchema( null, null, firstRowId, - Arrays.asList("blob_col")); + Arrays.asList("blob_col"), + null); } /** Creates a blob file with custom write columns. */ @@ -384,7 +385,8 @@ private DataFileMeta createBlobFileWithCols( null, null, firstRowId, - writeCols); + writeCols, + null); } private DataFileMeta createVectorFile( @@ -445,7 +447,8 @@ private DataFileMeta createFile( null, null, firstRowId, - writeCols); + writeCols, + null); } @Test @@ -523,6 +526,7 @@ private DataFileMeta createNormalFile( null, null, firstRowId, + null, null); } diff --git a/paimon-core/src/test/java/org/apache/paimon/operation/ExpireSnapshotsTest.java b/paimon-core/src/test/java/org/apache/paimon/operation/ExpireSnapshotsTest.java index 143f8445b8bc..46c7626f1108 100644 --- a/paimon-core/src/test/java/org/apache/paimon/operation/ExpireSnapshotsTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/operation/ExpireSnapshotsTest.java @@ -299,6 +299,7 @@ public void testExpireExtraFiles() throws IOException { null, null, null, + null, null); ManifestEntry add = ManifestEntry.create(FileKind.ADD, partition, 0, 1, dataFile); ManifestEntry delete = ManifestEntry.create(FileKind.DELETE, partition, 0, 1, dataFile); @@ -360,6 +361,7 @@ public void testExpireExtraFilesWithExternalPath() throws IOException { null, myDataFile.toString(), null, + null, null); ManifestEntry add = ManifestEntry.create(FileKind.ADD, partition, 0, 1, dataFile); ManifestEntry delete = ManifestEntry.create(FileKind.DELETE, partition, 0, 1, dataFile); diff --git a/paimon-core/src/test/java/org/apache/paimon/operation/ManifestEntryRunMergeTest.java b/paimon-core/src/test/java/org/apache/paimon/operation/ManifestEntryRunMergeTest.java index f93d1e713453..2d08a158d5ca 100644 --- a/paimon-core/src/test/java/org/apache/paimon/operation/ManifestEntryRunMergeTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/operation/ManifestEntryRunMergeTest.java @@ -137,7 +137,8 @@ private ManifestEntry rowIdEntry(String fileName, long firstRowId) { null, null, firstRowId, - Collections.singletonList("f0"))); + Collections.singletonList("f0"), + null)); } @Override diff --git a/paimon-core/src/test/java/org/apache/paimon/operation/ManifestRewriteCleanupTest.java b/paimon-core/src/test/java/org/apache/paimon/operation/ManifestRewriteCleanupTest.java index d54e29e91f6c..658507da8601 100644 --- a/paimon-core/src/test/java/org/apache/paimon/operation/ManifestRewriteCleanupTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/operation/ManifestRewriteCleanupTest.java @@ -454,7 +454,8 @@ private ManifestEntry rowIdEntry(FileKind kind, String fileName, long firstRowId null, null, firstRowId, - Collections.singletonList("f0"))); + Collections.singletonList("f0"), + null)); } private ManifestFile createManifestFile(long suggestedFileSize) { 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..e3f5403c68be 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); } @@ -271,6 +275,7 @@ public void testSerializerCompatibleV1() throws Exception { null, null, null, + null, null); List dataFiles = Collections.singletonList(dataFile); @@ -336,6 +341,7 @@ public void testSerializerCompatibleV2() throws Exception { null, null, null, + null, null); List dataFiles = Collections.singletonList(dataFile); @@ -401,6 +407,7 @@ public void testSerializerCompatibleV3() throws Exception { Arrays.asList("field1", "field2", "field3"), null, null, + null, null); List dataFiles = Collections.singletonList(dataFile); @@ -470,6 +477,7 @@ public void testSerializerCompatibleV4() throws Exception { Arrays.asList("field1", "field2", "field3"), null, null, + null, null); List dataFiles = Collections.singletonList(dataFile); @@ -539,6 +547,7 @@ public void testSerializerCompatibleV5() throws Exception { Arrays.asList("field1", "field2", "field3"), "hdfs:///path/to/warehouse", null, + null, null); List dataFiles = Collections.singletonList(dataFile); @@ -608,6 +617,7 @@ public void testSerializerCompatibleV6() throws Exception { Arrays.asList("field1", "field2", "field3"), "hdfs:///path/to/warehouse", null, + null, null); List dataFiles = Collections.singletonList(dataFile); @@ -678,6 +688,7 @@ public void testSerializerCompatibleV7() throws Exception { Arrays.asList("field1", "field2", "field3"), "hdfs:///path/to/warehouse", 12L, + null, null); List dataFiles = Collections.singletonList(dataFile); @@ -748,7 +759,8 @@ public void testSerializerCompatibleV8() throws Exception { Arrays.asList("field1", "field2", "field3"), "hdfs:///path/to/warehouse", 12L, - Arrays.asList("a", "b", "c", "f")); + Arrays.asList("a", "b", "c", "f"), + null); List dataFiles = Collections.singletonList(dataFile); DeletionFile deletionFile = new DeletionFile("deletion_file", 100, 22, 33L); @@ -784,6 +796,86 @@ 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"), + null) + .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); + } + + assertThat(current).isEqualTo(serialized); + 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..e21983248efb 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 @@ -38,6 +38,7 @@ import java.io.ByteArrayOutputStream; import java.io.IOException; import java.io.InputStream; +import java.nio.ByteBuffer; import java.util.ArrayList; import java.util.Arrays; import java.util.Collections; @@ -47,6 +48,7 @@ import java.util.Optional; import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; /** Test for {@link SplitSerializer}. */ public class SplitSerializerTest { @@ -56,31 +58,53 @@ 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); + assertThat(current).isEqualTo(serialized); } - assertSplitEquals(goldenCase.split, SplitSerializer.deserialize(actual)); + assertSplitEquals(goldenCase.split, SplitSerializer.deserialize(serialized)); } } + @Test + public void testUnsupportedVersion() throws IOException { + byte[] serialized = SplitSerializer.serialize(dataSplit(true)); + ByteBuffer.wrap(serialized, Long.BYTES, Integer.BYTES).putInt(3); + + assertThatThrownBy(() -> SplitSerializer.deserialize(serialized)) + .isInstanceOf(IOException.class) + .hasMessage("Unsupported split serializer version: 3"); + } + @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 +117,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 +135,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 +179,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 +205,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 +215,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 +262,45 @@ 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, + 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/ChainTableUtilsTest.java b/paimon-core/src/test/java/org/apache/paimon/utils/ChainTableUtilsTest.java index 9ddcf7dfda2c..9797af2a91b6 100644 --- a/paimon-core/src/test/java/org/apache/paimon/utils/ChainTableUtilsTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/utils/ChainTableUtilsTest.java @@ -867,6 +867,7 @@ private static DataFileMeta makeFile(String name, String min, String max, long s null, null, null, + null, null); } 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 c499c5f27651..99c4508bd661 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 @@ -242,6 +242,57 @@ public void testCollectWrittenColumnIdsExpandsLegacyFileSchema() { .hasValue(Arrays.asList(1, 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, legacy.columnMaxSequenceNumbers(), 0, 2)) + .isEqualTo(10L); + assertThat( + DataEvolutionUtils.fieldMaxSequenceNumber( + malformed, malformed.columnMaxSequenceNumbers(), 0, 2)) + .isEqualTo(10L); + long[] validSequences = valid.columnMaxSequenceNumbers(); + assertThat(DataEvolutionUtils.fieldMaxSequenceNumber(valid, validSequences, 0, 2)) + .isEqualTo(5L); + assertThat(DataEvolutionUtils.fieldMaxSequenceNumber(valid, validSequences, 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 0000000000000000000000000000000000000000..6277c55665ef8369e01e4458cf2ca4f1739174d9 GIT binary patch literal 1034 zcma)5yKWOf6uou`k&^&PKq%=9qDX-}KpI4vC`1THTSTNwlU#ep-XZJVWq16d03w9> z1HOWSPe7ufrlLZK8hT33Jl2k&;Yw$Z=gfW1jNkvF`68#yH19Sz<8~w)8LM8JG&Hwj z*(lO}-jAwybz7Low4*=LK?mzNP zon7qodkeqnfZ@l0$)5qRX^6Qiqno}59QUQ!o!P&BnB#B1sgsX$M?xe=I_JAIv3!pv z<@tA%j6>*_pPYBR3^|uk+AqqU?CqciyEJIewRxX1=DJ_sAR7Gv5m>cf literal 0 HcmV?d00001 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 0000000000000000000000000000000000000000..d919649c29e31701f71f2510d46b07ac53c7b638 GIT binary patch literal 3146 zcmZQz00UMC#Vo1Wg9Xul>u?K5s)?litz$5BM^h+K{Ocf07(HL&aI3uNGvMJ zEXmBzgUB#&gA{?}fHdO)1_lRakc0>jZvZh^Km-E=7lZX$O!A zKsE?KoeAedL>U-3fQ$~XqSl-9VTxgrKpKk>SQaKO0pv{pazOM7Af23_Q<|HnTbx>0 znwpoKs+*RXlM2>=WK?cuUVL_HWjls=Vg<4|SQ->i3P22^4S*QrUI!p16;R6@CaA@U z5V|o5O2fp#kp&ckg*T8143*J%q*nM-%N$sIjm9Ie5E^m$ks*Ph0FrS*<){M?gXjPt zP61*NfR(Qh9wwQYmy%kcTT)p7E!`k|q|y#j-qFdP(ei~_M$7U|<4b79eH@Vi6!_1Y$6F#SUYEM$7U|<4b#(Re^8R$s9NdvMuffxj|fLI2ILGobmiXB2jDG-O5K_sQb z*2uuf%)s2hOh2hKIXksP*O(C~#Q{+a*3a0%z~I0S5=hI;N!3kcXyby)!e|j7e*+MM z=o3Kf0mKiWv~Dj@d3ZG|%+y5L-x+mQ|90df?8 z7(^QYG00sGKujv2mN_uTBEkhNE|SQP3s5{E1DIX3j0{zq?S1_ f$HKx5#AgC}fitB{ADl0epgw@*I9LK>W?%#W8yP0? literal 0 HcmV?d00001 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 0000000000000000000000000000000000000000..c92ca5a3eba6f550771799d10f3009938fa3609a GIT binary patch literal 1508 zcmWFz@bL_Z4>M$7U|<4bcE(^-0T!SjGZ2daF(VLz!7Fwc3na(b!NB0a4-!es%t_Tv zWN71pO2cT7<_$m$qE7&^2M|Ai(i(7685mN4V%Pu&P_O`~4wpPOJ=nwqfPxhe`{1_1 z84wL{F3j!7=I{VH3P22^4S*QrE(ahc6;R6@m}3#)f)*D^Q0+g0O=^vJ+U^K{I2Y?tv-vD9{ApQWQEno%$`8r5! zAQuimanS&jL&QH^iX;XoUvvP?NzO>j%+m$sVz_p=&7yE2Fas_whbj){Q7cTTWe&`- zNa4uLzy>X?L{ds@jSP&;49pG8^pi@Hvr|iSjiIth4A^D4kU|Pcg*i;WDKHoiHgFW@ zB^DHCN6aWC{U}ReW literal 0 HcmV?d00001 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 0000000000000000000000000000000000000000..eab303006fcbae62473c737f4affb04f15915895 GIT binary patch literal 945 zcmWFz@bL_Z4>M$7U|<4bwtI&!8R$s9NdvMuffxj|fLI2ILGobmiXB2jDG-O5K_sQb z*2uuf%)s2hOh2hKIXksP*O(C~#Q{+a*3a0%z~I0S5=hI;N!3kcXyby)!e|j7e*+MM z=o3Kf0mKiWv~Dj@d3ZG|%+y5L-x+mQ|90df?8 z7(^QYG00sGKujv2mN_uTBEkhNE|SQP3s5{E1DIX3j0M$7U|@n`AjO~!#4<>HhF9z$VFm^c2n8Zppj<`<2F4i-3=aGtL2e-4 z0K}{y4iLa-5g_{j5QFF&KM$7U|<4b=1)H#ggB|_f%uH~4qr0Rk$jT|WOD*B2xtMZ3=o6Vg25|x z2o0q`9A*ZQloDGb10yp7a|1K|q|)T<)Dm4|MxYc2L@`)DV+R9+13yS0Ei)%oH<6)@ z3n~kvMS%PbKn$W!0I>%UKY-F2AmgAQ1;#+5LADhD)!~vyQ;w4>0F2&VcBG zb75{rHjD?zQ2=5PZ2-g|cR2ttseoGMz#NMR7qqxYB0nxb@q`RucF{5}xREUcrdt>n zly1-gwZf2E=D-{a3pWs-3FrmRlrnvAzDR=l0G8um35*%+Pnamo7#65%Sdj37B$(}i Kgk1v=18D$pvn)0M literal 0 HcmV?d00001 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 0000000000000000000000000000000000000000..f8a06462a45f641c4314b1f6bfe9ff949291963f GIT binary patch literal 1059 zcmWFz@bL_Z4>M$7U|<4b)?idV%ifx&?vB#@Swld7A@(8dLoh0!8F z{stfh(I=#+pwS@P3V`Zx$)hR9$rS*~R6y*5+X`nubiuhWw<8w-979Y7J*f{gq;kQ6^y3d{q`0RX0%LOuWh literal 0 HcmV?d00001 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/ChangelogCompactSortOperatorTest.java b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/compact/changelog/ChangelogCompactSortOperatorTest.java index 70530de9dc49..bd706147a598 100644 --- a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/compact/changelog/ChangelogCompactSortOperatorTest.java +++ b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/compact/changelog/ChangelogCompactSortOperatorTest.java @@ -193,6 +193,7 @@ private DataFileMeta createDataFileMeta(int mb, long creationMillis) { null, null, null, + null, null); } 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-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/globalindex/GenericIndexTopoBuilderTest.java b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/globalindex/GenericIndexTopoBuilderTest.java index 57d58b34c6b0..343c74419440 100644 --- a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/globalindex/GenericIndexTopoBuilderTest.java +++ b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/globalindex/GenericIndexTopoBuilderTest.java @@ -257,6 +257,7 @@ private static ManifestEntry createEntry(BinaryRow partition, Long firstRowId, l null, null, firstRowId, + null, null); return ManifestEntry.create(FileKind.ADD, partition, 0, 1, file); } diff --git a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/source/PendingSplitsCheckpointSerializerTest.java b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/source/PendingSplitsCheckpointSerializerTest.java index 2e4b66983e83..e89ef16bb2c5 100644 --- a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/source/PendingSplitsCheckpointSerializerTest.java +++ b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/source/PendingSplitsCheckpointSerializerTest.java @@ -18,12 +18,23 @@ package org.apache.paimon.flink.source; +import org.apache.paimon.io.DataFileMeta; +import org.apache.paimon.table.source.ChainSplit; +import org.apache.paimon.table.source.DeletionFile; +import org.apache.paimon.table.source.IncrementalSplit; +import org.apache.paimon.utils.IOUtils; + import org.apache.flink.core.io.SimpleVersionedSerialization; import org.junit.jupiter.api.Test; import java.io.IOException; +import java.io.InputStream; +import java.util.ArrayList; import java.util.Arrays; import java.util.Collections; +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Map; import static org.apache.paimon.flink.source.FileStoreSourceSplitSerializerTest.newFile; import static org.apache.paimon.flink.source.FileStoreSourceSplitSerializerTest.newSourceSplit; @@ -33,6 +44,9 @@ /** Unit tests for the {@link PendingSplitsCheckpointSerializer}. */ public class PendingSplitsCheckpointSerializerTest { + private static final String LEGACY_CHECKPOINT_RESOURCE = + "compatibility/pending-splits-incremental-v1-chain-v2"; + @Test public void serializeEmptyCheckpoint() throws Exception { final PendingSplitsCheckpoint checkpoint = @@ -77,6 +91,107 @@ public void repeatedSerialization() throws Exception { assertCheckpointsEqual(checkpoint, deSerialized); } + @Test + public void restoreIncrementalAndChainSplits() throws Exception { + DataFileMeta before = file("before.parquet", 1L); + DataFileMeta after = file("after.parquet", 2L); + IncrementalSplit incremental = + new IncrementalSplit( + 10L, + row(1), + 2, + 8, + Collections.singletonList(before), + Collections.singletonList(new DeletionFile("before.dv", 0L, 1L, 1L)), + Collections.singletonList(after), + Collections.singletonList(new DeletionFile("after.dv", 1L, 1L, 1L)), + true); + + Map branchMapping = new LinkedHashMap<>(); + branchMapping.put(before.fileName(), "snapshot"); + branchMapping.put(after.fileName(), "delta"); + Map bucketPathMapping = new LinkedHashMap<>(); + bucketPathMapping.put(before.fileName(), "dt=1/bucket-2"); + bucketPathMapping.put(after.fileName(), "dt=1/bucket-2"); + ChainSplit chain = + new ChainSplit( + row(1), + Arrays.asList(before, after), + branchMapping, + bucketPathMapping, + Arrays.asList(null, new DeletionFile("chain.dv", 2L, 1L, 1L))); + + PendingSplitsCheckpoint restored = + serializeAndDeserialize( + new PendingSplitsCheckpoint( + Arrays.asList( + new FileStoreSourceSplit("incremental", incremental, 11L), + new FileStoreSourceSplit("chain", chain, 12L)), + 20L)); + + assertThat(restored.currentSnapshotId()).isEqualTo(20L); + assertThat(restored.splits()).hasSize(2); + List restoredSplits = new ArrayList<>(restored.splits()); + assertIncrementalSplit(restoredSplits.get(0), 11L); + assertChainSplit(restoredSplits.get(1), 12L, branchMapping, bucketPathMapping); + } + + @Test + public void restoreLegacyIncrementalAndChainSplits() throws Exception { + byte[] bytes; + try (InputStream in = + PendingSplitsCheckpointSerializerTest.class + .getClassLoader() + .getResourceAsStream(LEGACY_CHECKPOINT_RESOURCE)) { + bytes = IOUtils.readFully(in, false); + } + PendingSplitsCheckpointSerializer serializer = + new PendingSplitsCheckpointSerializer(new FileStoreSourceSplitSerializer()); + PendingSplitsCheckpoint restored = + SimpleVersionedSerialization.readVersionAndDeSerialize(serializer, bytes); + + assertThat(restored.currentSnapshotId()).isEqualTo(20L); + assertThat(restored.splits()).hasSize(2); + List restoredSplits = new ArrayList<>(restored.splits()); + + FileStoreSourceSplit incrementalSource = restoredSplits.get(0); + assertThat(incrementalSource.recordsToSkip()).isEqualTo(11L); + assertThat(incrementalSource.split()).isInstanceOf(IncrementalSplit.class); + IncrementalSplit incremental = (IncrementalSplit) incrementalSource.split(); + assertThat(incremental.beforeFiles()) + .extracting(DataFileMeta::fileName) + .containsExactly("before.parquet"); + assertThat(incremental.afterFiles()) + .extracting(DataFileMeta::fileName) + .containsExactly("after.parquet"); + assertThat(incremental.beforeFiles().get(0).columnMaxSequenceNumbers()).isNull(); + assertThat(incremental.afterFiles().get(0).columnMaxSequenceNumbers()).isNull(); + assertThat(incremental.beforeDeletionFiles()) + .containsExactly(new DeletionFile("before.dv", 0L, 1L, 1L)); + assertThat(incremental.afterDeletionFiles()) + .containsExactly(new DeletionFile("after.dv", 1L, 1L, 1L)); + + FileStoreSourceSplit chainSource = restoredSplits.get(1); + assertThat(chainSource.recordsToSkip()).isEqualTo(12L); + assertThat(chainSource.split()).isInstanceOf(ChainSplit.class); + ChainSplit chain = (ChainSplit) chainSource.split(); + assertThat(chain.dataFiles()) + .extracting(DataFileMeta::fileName) + .containsExactly("before.parquet", "after.parquet"); + assertThat(chain.dataFiles()) + .allSatisfy(file -> assertThat(file.columnMaxSequenceNumbers()).isNull()); + assertThat(chain.fileBranchMapping()) + .containsOnly( + org.assertj.core.data.MapEntry.entry("before.parquet", "snapshot"), + org.assertj.core.data.MapEntry.entry("after.parquet", "delta")); + assertThat(chain.fileBucketPathMapping()) + .containsOnly( + org.assertj.core.data.MapEntry.entry("before.parquet", "dt=1/bucket-2"), + org.assertj.core.data.MapEntry.entry("after.parquet", "dt=1/bucket-2")); + assertThat(chain.deletionFiles()) + .hasValue(Arrays.asList(null, new DeletionFile("chain.dv", 2L, 1L, 1L))); + } + // ------------------------------------------------------------------------ // test utils // ------------------------------------------------------------------------ @@ -93,6 +208,45 @@ private static FileStoreSourceSplit testSplit3() { return newSourceSplit("id3", row(3), 4, Arrays.asList(newFile(5), newFile(6))); } + private static DataFileMeta file(String fileName, long sequence) { + return newFile(0).rename(fileName).withColumnMaxSequenceNumbers(new long[] {sequence}); + } + + private static void assertIncrementalSplit( + FileStoreSourceSplit sourceSplit, long recordsToSkip) { + assertThat(sourceSplit.recordsToSkip()).isEqualTo(recordsToSkip); + assertThat(sourceSplit.split()).isInstanceOf(IncrementalSplit.class); + IncrementalSplit split = (IncrementalSplit) sourceSplit.split(); + assertColumnSequences(split.beforeFiles(), 1L); + assertColumnSequences(split.afterFiles(), 2L); + assertThat(split.beforeDeletionFiles()) + .containsExactly(new DeletionFile("before.dv", 0L, 1L, 1L)); + assertThat(split.afterDeletionFiles()) + .containsExactly(new DeletionFile("after.dv", 1L, 1L, 1L)); + } + + private static void assertChainSplit( + FileStoreSourceSplit sourceSplit, + long recordsToSkip, + Map branchMapping, + Map bucketPathMapping) { + assertThat(sourceSplit.recordsToSkip()).isEqualTo(recordsToSkip); + assertThat(sourceSplit.split()).isInstanceOf(ChainSplit.class); + ChainSplit split = (ChainSplit) sourceSplit.split(); + assertColumnSequences(split.dataFiles(), 1L, 2L); + assertThat(split.fileBranchMapping()).isEqualTo(branchMapping); + assertThat(split.fileBucketPathMapping()).isEqualTo(bucketPathMapping); + assertThat(split.deletionFiles()) + .hasValue(Arrays.asList(null, new DeletionFile("chain.dv", 2L, 1L, 1L))); + } + + private static void assertColumnSequences(List files, long... sequences) { + assertThat(files).hasSize(sequences.length); + for (int i = 0; i < sequences.length; i++) { + assertThat(files.get(i).columnMaxSequenceNumbers()).containsExactly(sequences[i]); + } + } + private static PendingSplitsCheckpoint serializeAndDeserialize( final PendingSplitsCheckpoint split) throws IOException { diff --git a/paimon-flink/paimon-flink-common/src/test/resources/compatibility/pending-splits-incremental-v1-chain-v2 b/paimon-flink/paimon-flink-common/src/test/resources/compatibility/pending-splits-incremental-v1-chain-v2 new file mode 100644 index 0000000000000000000000000000000000000000..dffb60857f20650e9df60f588216fc5220841e76 GIT binary patch literal 2659 zcmeHJ&2G~`5FV!(<)>+rQVvKQDuj?YV5^*vDkKyTBFIP)y}^Zzy-5}v+g)$k!hvVu z*aHWSJO+=!1zvy)Gwa_(wh%%-Aa$kju6Mqj^?bAbYydC-x~~DC1z_HQo(NpYmpo9+ z|NH~YSb`nOxy2&pF1Qsju?z!Cv8m6kI9y4WTjOHIapVhyv8Wka&>6$k>B@b_)hi4f zA1le(QUvqo(2WBY#fwmly)kU75O*7CVC=vin<*}zaGxs?22X|0V+8}}EjwlQN(tX~ zM68Y+=xUtypTI{j9Jn^+vrzK2rKzizxXS2G#H13mhk{!UCTS;0+DVeO#}uzJo*5z^^-Ew`5|rC@0ad#2hJgz&`x~D7l22*IEEdx3mT*?3q(P_qEY`1?G4)Z zw11X?@V+#}yQ0xX7w8=Rx_$Fg9jJiUjuMWr$ns>xlRDN#tDD!cDUSm*>K~vD&?acu zjG81W=j%3UIzH`d7|==i@J{uk7nkh0G!CJ8f%}MQvcP-)SYZ9etkt)s9Q9{Sv(R(_ znvS0qxSjBk*SB0(%F{5;b-DmO6h{k8cfmxE%HrKW$l_zdZV003z>s8RO`$5qtwXS~ zS`B9?gd@oUckzNn5vr(y-I6HcBx>SyjnYawVJbtC2DTg+_~Bf*%%A+d?p+}sNYkh? z8BN7N`-<1#QsAv6F0 literal 0 HcmV?d00001 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(); + } +} diff --git a/paimon-spark/paimon-spark-common/src/test/java/org/apache/paimon/spark/procedure/CreateGlobalIndexProcedureTest.java b/paimon-spark/paimon-spark-common/src/test/java/org/apache/paimon/spark/procedure/CreateGlobalIndexProcedureTest.java index f40c751fe245..f25bde2ba517 100644 --- a/paimon-spark/paimon-spark-common/src/test/java/org/apache/paimon/spark/procedure/CreateGlobalIndexProcedureTest.java +++ b/paimon-spark/paimon-spark-common/src/test/java/org/apache/paimon/spark/procedure/CreateGlobalIndexProcedureTest.java @@ -481,6 +481,7 @@ private DataFileMeta createDataFileMeta(Long firstRowId, Long rowCount) { null, null, firstRowId, + null, null); } From 68f069c755202d5b59312b76fc08c027f9ca7951 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E4=BB=9F=E5=BC=8B?= Date: Thu, 20 Aug 2026 15:27:11 +0800 Subject: [PATCH 2/4] [core] Keep split serializer at version 1 --- .../table/source/SplitSerializerTest.java | 159 ++++++------------ .../resources/compatibility/split-v1-chain | Bin 1419 -> 1443 bytes .../resources/compatibility/split-v1-data | Bin 896 -> 912 bytes .../resources/compatibility/split-v1-fallback | Bin 1436 -> 1460 bytes .../compatibility/split-v1-fallback-data | Bin 897 -> 913 bytes .../compatibility/split-v1-incremental | Bin 934 -> 946 bytes .../resources/compatibility/split-v1-indexed | Bin 961 -> 977 bytes .../compatibility/split-v1-query-auth | Bin 1011 -> 1027 bytes .../resources/compatibility/split-v2-chain | Bin 1491 -> 0 bytes .../resources/compatibility/split-v2-data | Bin 944 -> 0 bytes .../resources/compatibility/split-v2-fallback | Bin 1508 -> 0 bytes .../compatibility/split-v2-fallback-data | Bin 945 -> 0 bytes .../compatibility/split-v2-incremental | Bin 978 -> 0 bytes .../resources/compatibility/split-v2-indexed | Bin 1009 -> 0 bytes .../compatibility/split-v2-query-auth | Bin 1059 -> 0 bytes 15 files changed, 55 insertions(+), 104 deletions(-) delete mode 100644 paimon-core/src/test/resources/compatibility/split-v2-chain delete mode 100644 paimon-core/src/test/resources/compatibility/split-v2-data delete mode 100644 paimon-core/src/test/resources/compatibility/split-v2-fallback delete mode 100644 paimon-core/src/test/resources/compatibility/split-v2-fallback-data delete mode 100644 paimon-core/src/test/resources/compatibility/split-v2-incremental delete mode 100644 paimon-core/src/test/resources/compatibility/split-v2-indexed delete mode 100644 paimon-core/src/test/resources/compatibility/split-v2-query-auth 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 e21983248efb..3b6afc3acfe3 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 @@ -38,7 +38,6 @@ import java.io.ByteArrayOutputStream; import java.io.IOException; import java.io.InputStream; -import java.nio.ByteBuffer; import java.util.ArrayList; import java.util.Arrays; import java.util.Collections; @@ -48,7 +47,6 @@ import java.util.Optional; import static org.assertj.core.api.Assertions.assertThat; -import static org.assertj.core.api.Assertions.assertThatThrownBy; /** Test for {@link SplitSerializer}. */ public class SplitSerializerTest { @@ -58,53 +56,31 @@ public class SplitSerializerTest { @Test public void testRoundTrip() throws IOException { - for (GoldenCase goldenCase : goldenCases(2)) { + for (GoldenCase goldenCase : goldenCases()) { Split actual = SplitSerializer.deserialize(SplitSerializer.serialize(goldenCase.split)); assertSplitEquals(goldenCase.split, actual); } } @Test - public void testVersion1GoldenFiles() throws IOException { - for (GoldenCase goldenCase : goldenCases(1)) { - assertSplitEquals( - goldenCase.split, - SplitSerializer.deserialize(readGoldenFile(goldenCase.fileName))); - } - } - - @Test - public void testVersion2GoldenFiles() throws IOException { + public void testGoldenFiles() throws IOException { boolean generateGoldenFiles = Boolean.parseBoolean( System.getProperties().getProperty(GENERATE_GOLDEN_FILES_PROPERTY)); - for (GoldenCase goldenCase : goldenCases(2)) { - byte[] current = SplitSerializer.serialize(goldenCase.split); - byte[] serialized; + for (GoldenCase goldenCase : goldenCases()) { + byte[] actual = SplitSerializer.serialize(goldenCase.split); if (generateGoldenFiles) { - CompatibilityUtils.writeCompatibilityFile(goldenCase.fileName, current); - serialized = current; + CompatibilityUtils.writeCompatibilityFile(goldenCase.fileName, actual); } else { - serialized = readGoldenFile(goldenCase.fileName); - assertThat(current).isEqualTo(serialized); + assertThat(actual).isEqualTo(readGoldenFile(goldenCase.fileName)); } - assertSplitEquals(goldenCase.split, SplitSerializer.deserialize(serialized)); + assertSplitEquals(goldenCase.split, SplitSerializer.deserialize(actual)); } } - @Test - public void testUnsupportedVersion() throws IOException { - byte[] serialized = SplitSerializer.serialize(dataSplit(true)); - ByteBuffer.wrap(serialized, Long.BYTES, Integer.BYTES).putInt(3); - - assertThatThrownBy(() -> SplitSerializer.deserialize(serialized)) - .isInstanceOf(IOException.class) - .hasMessage("Unsupported split serializer version: 3"); - } - @Test public void testFallbackSplitImplSerializeAndDeserialize() throws IOException { - FallbackReadFileStoreTable.FallbackSplitImpl split = fallbackSplit(true); + FallbackReadFileStoreTable.FallbackSplitImpl split = fallbackSplit(); ByteArrayOutputStream out = new ByteArrayOutputStream(); split.serialize(new DataOutputViewStreamWrapper(out)); @@ -117,7 +93,7 @@ public void testFallbackSplitImplSerializeAndDeserialize() throws IOException { @Test public void testFallbackSplitImplJavaSerializeAndDeserialize() throws Exception { - FallbackReadFileStoreTable.FallbackSplitImpl split = fallbackSplit(true); + FallbackReadFileStoreTable.FallbackSplitImpl split = fallbackSplit(); byte[] bytes = InstantiationUtil.serializeObject(split); FallbackReadFileStoreTable.FallbackSplitImpl deserialized = @@ -135,36 +111,32 @@ private static byte[] readGoldenFile(String fileName) throws IOException { } } - private static List goldenCases(int version) { - boolean withColumnSequences = version >= 2; + private static List goldenCases() { List cases = new ArrayList<>(); - DataSplit dataSplit = dataSplit(withColumnSequences); - IncrementalSplit incrementalSplit = incrementalSplit(withColumnSequences); + DataSplit dataSplit = dataSplit(); + IncrementalSplit incrementalSplit = incrementalSplit(); IndexedSplit indexedSplit = indexedSplit(dataSplit); - ChainSplit chainSplit = chainSplit(withColumnSequences); + ChainSplit chainSplit = chainSplit(); QueryAuthSplit queryAuthSplit = queryAuthSplit(dataSplit); FallbackReadFileStoreTable.FallbackDataSplit fallbackDataSplit = (FallbackReadFileStoreTable.FallbackDataSplit) FallbackReadFileStoreTable.toFallbackSplit(dataSplit, true); - 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)); + 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)); return cases; } - private static DataSplit dataSplit(boolean withColumnSequences) { + private static DataSplit dataSplit() { List files = Arrays.asList( - dataFile("file-a", 0, 1, 10, 100L, withColumnSequences), - dataFile("file-b", 1, 11, 20, 200L, withColumnSequences)); + dataFile("file-a", 0, 1, 10, 100L), dataFile("file-b", 1, 11, 20, 200L)); List deletionFiles = Arrays.asList(null, new DeletionFile("dv/file-b", 2L, 10L, 3L)); return DataSplit.builder() @@ -179,13 +151,10 @@ private static DataSplit dataSplit(boolean withColumnSequences) { .build(); } - private static IncrementalSplit incrementalSplit(boolean withColumnSequences) { + private static IncrementalSplit incrementalSplit() { List before = - Collections.singletonList( - dataFile("before-file", 0, 1, 5, 10L, withColumnSequences)); - List after = - Collections.singletonList( - dataFile("after-file", 0, 6, 12, 20L, withColumnSequences)); + Collections.singletonList(dataFile("before-file", 0, 1, 5, 10L)); + List after = Collections.singletonList(dataFile("after-file", 0, 6, 12, 20L)); return new IncrementalSplit( 43L, DataFileTestUtils.row(2026, 8), @@ -205,8 +174,8 @@ private static IndexedSplit indexedSplit(DataSplit dataSplit) { new float[] {0.5f, 0.25f, 0.125f}); } - private static ChainSplit chainSplit(boolean withColumnSequences) { - DataSplit left = dataSplit(withColumnSequences); + private static ChainSplit chainSplit() { + DataSplit left = dataSplit(); DataSplit right = DataSplit.builder() .withSnapshot(44L) @@ -215,14 +184,7 @@ private static ChainSplit chainSplit(boolean withColumnSequences) { .withTotalBuckets(8) .withBucketPath("dt=20260707/bucket-5") .withDataFiles( - Collections.singletonList( - dataFile( - "chain-file", - 0, - 21, - 30, - 300L, - withColumnSequences))) + Collections.singletonList(dataFile("chain-file", 0, 21, 30, 300L))) .withDataDeletionFiles( Collections.singletonList( new DeletionFile("deletion_file", 100, 22, null))) @@ -262,45 +224,34 @@ private static QueryAuthSplit queryAuthSplit(Split split) { Arrays.asList("filter-json-1", "filter-json-2"), columnMasking)); } - private static FallbackReadFileStoreTable.FallbackSplitImpl fallbackSplit( - boolean withColumnSequences) { - return new FallbackReadFileStoreTable.FallbackSplitImpl( - chainSplit(withColumnSequences), true); + private static FallbackReadFileStoreTable.FallbackSplitImpl fallbackSplit() { + return new FallbackReadFileStoreTable.FallbackSplitImpl(chainSplit(), true); } private static DataFileMeta dataFile( - 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, - null); - return withColumnSequences - ? file.withColumnMaxSequenceNumbers(new long[] {maxSequence}) - : file; + 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, + null); } private static void assertSplitEquals(Split expected, Split actual) { diff --git a/paimon-core/src/test/resources/compatibility/split-v1-chain b/paimon-core/src/test/resources/compatibility/split-v1-chain index 08f61c52faec6af460b3bee156a25b47c9e03d7e..d7fa63bf04ea360a0e3a306ef7c188e3f5509d9b 100644 GIT binary patch delta 245 zcmeC?Ud(L~9N^;_5+7#Bz`(!=#4JF}48$T(K9FKyc*PE;Km<^Zv4VlYL4Klzg5(CE z07&5pAoc*_2T)pLVxc-?!Q@6}=ZPD*CoyqMe13oxtZDKZMm)MEKV?MGHaVKTSpfrs1OG$|1<4ga z0g%E2Klahlx2#5tLfaX-YE$wwI#BypNE`5&V?D~iTiru~8qAp1eaa{)0(@8pY&B9c3R f43J|k05QlhFQBx+dnvr9oOfYK!1A_zq#1aL`6$}gvB0ziqh&_P#29#EqxK^DpV{#yq+vNF7 XoRf7J_e1nd{=kSw(_~F1gsx8jAodrw diff --git a/paimon-core/src/test/resources/compatibility/split-v1-incremental b/paimon-core/src/test/resources/compatibility/split-v1-incremental index cb057e42d72f52a32f41d4361c578d5908fa0c1c..50fd72f825be878faf0ae4a3a475e0130d121b1f 100644 GIT binary patch delta 182 zcmZ3+zKLBRIKam$QFgVBqISkxDya9+=K^!1}(IPN;pbC&_>_7|xFdDCEldYKG#%*58#K;H$ DJoOu{ delta 162 zcmdnQzKmTYIKamCjdjRnZC~YvYRh_Y7;#%j8KR6iq8em3o0i{<==3^3<#Hn$z9TP<7=8a5@i~th~ B7}fv) diff --git a/paimon-core/src/test/resources/compatibility/split-v1-indexed b/paimon-core/src/test/resources/compatibility/split-v1-indexed index 0d20df101234947fef1436014942786e1cb21a11..eab6fbe8f82a59013a8ea83137aa4d4541fb5909 100644 GIT binary patch delta 126 zcmX@eevy5GIwR*qjbPRa1_lTDi46*p8yFZEM1c4N5PJae11POAai=2p3lTO*?@6BMCar$jCiz8)@4G`yLmp75+eW}+8lxa delta 116 zcmcb}evo~FIwQwKjbPRS1_lTIi46*pD;O9UM1c4J5PJae4JfTJai=9tilLMGI bCYv)JfM}WghY^pS$)-#QO`A6}$uR-|-EkS9 diff --git a/paimon-core/src/test/resources/compatibility/split-v1-query-auth b/paimon-core/src/test/resources/compatibility/split-v1-query-auth index fce1aa5fb9aebf2d06269c5ed880bf765c9d1a84..5e5acea55d31c3ead85e44fdc463618b8a4098a9 100644 GIT binary patch delta 128 zcmey&-pnyUn~`&(PB3c)1A~M7#1;j~4GatnB0zish&_P#0hHF5xL2LEfRTYAWpW~u j-6Tey$@7^wCz~+thiIMrgAtG3$%ag*nm5m93Sk5Q+|wGD delta 113 zcmZqX_{=^*n~`IpPB3c$1A_zq#1;j~6$}gvB0ziqh&_P#29#EqxL2JqV{#&s+vN33 aoRe)B_e1ndV#1?mvLzE#)8_R|A&dY|I2TR; diff --git a/paimon-core/src/test/resources/compatibility/split-v2-chain b/paimon-core/src/test/resources/compatibility/split-v2-chain deleted file mode 100644 index 0a4cb5a29fdaa6bac43719911d595b8388c60978..0000000000000000000000000000000000000000 GIT binary patch literal 0 HcmV?d00001 literal 1491 zcmWFz@bL_Z4>M$7U|<4b79eH@Vi6!_1Y$6F#SUYEM$7U|<4b#(Re^8R$s9NdvMuffxj|fLI2ILGobmiXB2jDG-O5K_sQb z*2uuf%)s2hOh2hKIXksP*O(C~#Q{+a*3a0%z~I0S5=hI;N!3kcXyby)!e|j7e*+MM z=o3Kf0mKiWv~Dj@d3ZG|%+y5L-x+mQ|90df?8 z7(^QYG00sGKujv2mN_uTBEkhNE|SQP3s5{E1DIX3j0{zq?S1_ f$HKx5#AgC}fitB{ADl0epgw@*I9LK>W?%#W8yP0? diff --git a/paimon-core/src/test/resources/compatibility/split-v2-fallback b/paimon-core/src/test/resources/compatibility/split-v2-fallback deleted file mode 100644 index c92ca5a3eba6f550771799d10f3009938fa3609a..0000000000000000000000000000000000000000 GIT binary patch literal 0 HcmV?d00001 literal 1508 zcmWFz@bL_Z4>M$7U|<4bcE(^-0T!SjGZ2daF(VLz!7Fwc3na(b!NB0a4-!es%t_Tv zWN71pO2cT7<_$m$qE7&^2M|Ai(i(7685mN4V%Pu&P_O`~4wpPOJ=nwqfPxhe`{1_1 z84wL{F3j!7=I{VH3P22^4S*QrE(ahc6;R6@m}3#)f)*D^Q0+g0O=^vJ+U^K{I2Y?tv-vD9{ApQWQEno%$`8r5! zAQuimanS&jL&QH^iX;XoUvvP?NzO>j%+m$sVz_p=&7yE2Fas_whbj){Q7cTTWe&`- zNa4uLzy>X?L{ds@jSP&;49pG8^pi@Hvr|iSjiIth4A^D4kU|Pcg*i;WDKHoiHgFW@ zB^DHCN6aWC{U}ReW diff --git a/paimon-core/src/test/resources/compatibility/split-v2-fallback-data b/paimon-core/src/test/resources/compatibility/split-v2-fallback-data deleted file mode 100644 index eab303006fcbae62473c737f4affb04f15915895..0000000000000000000000000000000000000000 GIT binary patch literal 0 HcmV?d00001 literal 945 zcmWFz@bL_Z4>M$7U|<4bwtI&!8R$s9NdvMuffxj|fLI2ILGobmiXB2jDG-O5K_sQb z*2uuf%)s2hOh2hKIXksP*O(C~#Q{+a*3a0%z~I0S5=hI;N!3kcXyby)!e|j7e*+MM z=o3Kf0mKiWv~Dj@d3ZG|%+y5L-x+mQ|90df?8 z7(^QYG00sGKujv2mN_uTBEkhNE|SQP3s5{E1DIX3j0M$7U|@n`AjO~!#4<>HhF9z$VFm^c2n8Zppj<`<2F4i-3=aGtL2e-4 z0K}{y4iLa-5g_{j5QFF&KM$7U|<4b=1)H#ggB|_f%uH~4qr0Rk$jT|WOD*B2xtMZ3=o6Vg25|x z2o0q`9A*ZQloDGb10yp7a|1K|q|)T<)Dm4|MxYc2L@`)DV+R9+13yS0Ei)%oH<6)@ z3n~kvMS%PbKn$W!0I>%UKY-F2AmgAQ1;#+5LADhD)!~vyQ;w4>0F2&VcBG zb75{rHjD?zQ2=5PZ2-g|cR2ttseoGMz#NMR7qqxYB0nxb@q`RucF{5}xREUcrdt>n zly1-gwZf2E=D-{a3pWs-3FrmRlrnvAzDR=l0G8um35*%+Pnamo7#65%Sdj37B$(}i Kgk1v=18D$pvn)0M diff --git a/paimon-core/src/test/resources/compatibility/split-v2-query-auth b/paimon-core/src/test/resources/compatibility/split-v2-query-auth deleted file mode 100644 index f8a06462a45f641c4314b1f6bfe9ff949291963f..0000000000000000000000000000000000000000 GIT binary patch literal 0 HcmV?d00001 literal 1059 zcmWFz@bL_Z4>M$7U|<4b)?idV%ifx&?vB#@Swld7A@(8dLoh0!8F z{stfh(I=#+pwS@P3V`Zx$)hR9$rS*~R6y*5+X`nubiuhWw<8w-979Y7J*f{gq;kQ6^y3d{q`0RX0%LOuWh From 91c674bac02cca0cd1aa24e1b85fb2d26cb86bec Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E4=BB=9F=E5=BC=8B?= Date: Thu, 20 Aug 2026 15:37:13 +0800 Subject: [PATCH 3/4] [core] Encapsulate incremental split serialization --- .../compatibility/split-v1-incremental | Bin 946 -> 950 bytes 1 file changed, 0 insertions(+), 0 deletions(-) diff --git a/paimon-core/src/test/resources/compatibility/split-v1-incremental b/paimon-core/src/test/resources/compatibility/split-v1-incremental index 50fd72f825be878faf0ae4a3a475e0130d121b1f..e5368a7ea610726179a1e275250c02f5d108ffc1 100644 GIT binary patch delta 25 fcmdnQzKvZVIKamEMo=$QBVb9 delta 21 ccmdnSzKNYDIKamqrfs|06{JVU;qFB From 133e672eb679d6e469e70cb10b3d241f9c266f51 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E4=BB=9F=E5=BC=8B?= Date: Thu, 20 Aug 2026 17:10:42 +0800 Subject: [PATCH 4/4] [core] Rename write column sequence metadata --- .../src/main/java/org/apache/paimon/io/DataFileMeta.java | 4 ++-- .../java/org/apache/paimon/io/ProjectedDataFileMeta.java | 6 +++--- 2 files changed, 5 insertions(+), 5 deletions(-) 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 912c2e9f8b1e..e9874917a17c 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,7 +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"; + String WRITE_COLS_SEQUENCES = "_WRITE_COLS_SEQUENCES"; RowType SCHEMA = new RowType( @@ -113,7 +113,7 @@ public interface DataFileMeta { 19, WRITE_COLS, new ArrayType(true, newStringType(false))), new DataField( 20, - COLUMN_MAX_SEQUENCE_NUMBERS, + WRITE_COLS_SEQUENCES, new ArrayType(true, new BigIntType(false))))); BinaryRow EMPTY_MIN_KEY = EMPTY_ROW; diff --git a/paimon-core/src/main/java/org/apache/paimon/io/ProjectedDataFileMeta.java b/paimon-core/src/main/java/org/apache/paimon/io/ProjectedDataFileMeta.java index bc5edc413ede..70da63dba7f3 100644 --- a/paimon-core/src/main/java/org/apache/paimon/io/ProjectedDataFileMeta.java +++ b/paimon-core/src/main/java/org/apache/paimon/io/ProjectedDataFileMeta.java @@ -257,7 +257,7 @@ public List writeCols() { @Nullable @Override public long[] columnMaxSequenceNumbers() { - int position = requiredPosition(Fields.COLUMN_MAX_SEQUENCE_NUMBERS); + int position = requiredPosition(Fields.WRITE_COLS_SEQUENCES); InternalRow row = currentRow(); if (row.isNullAt(position)) { return null; @@ -409,8 +409,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); + private static final int WRITE_COLS_SEQUENCES = + fieldIndex(DataFileMeta.WRITE_COLS_SEQUENCES); } /** Projected data-file schema together with its bound binary field layout. */