This is an automated email from the ASF dual-hosted git repository.
JingsongLi pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/paimon.git
The following commit(s) were added to refs/heads/master by this push:
new 7a4643b907 [core] Add an optional positional long[] of per-column
maximum sequence numbers to DataFileMeta (#9313)
7a4643b907 is described below
commit 7a4643b907a84987605df8451581529e2f06b3ce
Author: YeJunHao <[email protected]>
AuthorDate: Thu Aug 20 18:00:01 2026 +0800
[core] Add an optional positional long[] of per-column maximum sequence
numbers to DataFileMeta (#9313)
---
.../main/java/org/apache/paimon/CoreOptions.java | 4 +
.../DataEvolutionCompactTaskSerializer.java | 2 +-
.../DataEvolutionNormalCompactTask.java | 67 ++++++
.../java/org/apache/paimon/io/DataFileMeta.java | 36 +++-
.../apache/paimon/io/DataFileMeta08Serializer.java | 1 +
.../apache/paimon/io/DataFileMeta09Serializer.java | 1 +
.../paimon/io/DataFileMeta10LegacySerializer.java | 1 +
.../paimon/io/DataFileMeta12LegacySerializer.java | 1 +
.../io/DataFileMetaFirstRowIdLegacySerializer.java | 3 +-
.../apache/paimon/io/DataFileMetaSerializer.java | 19 +-
... => DataFileMetaWriteColsLegacySerializer.java} | 35 +++-
.../org/apache/paimon/io/PojoDataFileMeta.java | 78 +++++--
.../apache/paimon/io/ProjectedDataFileMeta.java | 23 ++
.../ManifestEntryWriteColsLegacySerializer.java | 81 ++++++++
.../table/sink/AppendCompactTaskSerializer.java | 2 +-
.../sink/CommitMessageLegacyV2Serializer.java | 1 +
.../paimon/table/sink/CommitMessageSerializer.java | 11 +-
.../sink/MultiTableCompactionTaskSerializer.java | 2 +-
.../org/apache/paimon/table/source/ChainSplit.java | 9 +-
.../org/apache/paimon/table/source/DataSplit.java | 7 +-
.../paimon/table/source/IncrementalSplit.java | 11 +-
.../apache/paimon/utils/DataEvolutionUtils.java | 38 ++++
.../java/org/apache/paimon/CoreOptionsTest.java | 5 +-
.../org/apache/paimon/append/BlobUpdateTest.java | 3 +-
.../DataEvolutionCompactCoordinatorTest.java | 3 +-
.../DataEvolutionCompactRangePlannerTest.java | 3 +-
.../DataEvolutionNormalCompactTaskTest.java | 231 +++++++++++++++++++++
.../paimon/crosspartition/IndexBootstrapTest.java | 1 +
.../globalindex/GlobalIndexBuilderUtilsTest.java | 2 +
.../paimon/io/DataFileMetaSerializerTest.java | 48 ++++-
.../org/apache/paimon/io/DataFileTestUtils.java | 1 +
.../paimon/io/ProjectedDataFileMetaTest.java | 33 +--
...festCommittableSerializerCompatibilityTest.java | 80 ++++++-
.../manifest/ManifestEntrySerializerTest.java | 22 ++
.../paimon/manifest/ManifestFileMetaTest.java | 7 +-
.../paimon/manifest/ManifestFileMetaTestBase.java | 1 +
.../apache/paimon/manifest/ManifestFileTest.java | 45 ++++
.../manifest/NoPartitionManifestFileMetaTest.java | 3 +-
.../paimon/mergetree/LevelsOrderingFixTest.java | 1 +
.../mergetree/compact/IntervalPartitionTest.java | 1 +
.../operation/BlobFallbackRecordReaderTest.java | 3 +-
.../paimon/operation/DataEvolutionReadTest.java | 10 +-
.../paimon/operation/ExpireSnapshotsTest.java | 2 +
.../operation/ManifestEntryRunMergeTest.java | 3 +-
.../operation/ManifestRewriteCleanupTest.java | 3 +-
.../table/sink/CommitMessageSerializerTest.java | 8 +
.../table/source/DataSplitCompatibleTest.java | 94 ++++++++-
.../paimon/table/source/SplitSerializerTest.java | 1 +
.../apache/paimon/utils/ChainTableUtilsTest.java | 1 +
.../paimon/utils/DataEvolutionUtilsTest.java | 51 +++++
.../src/test/resources/compatibility/datasplit-v9 | Bin 0 -> 1034 bytes
.../compatibility/manifest-committable-v13-v5 | Bin 0 -> 3146 bytes
.../test/resources/compatibility/split-v1-chain | Bin 1419 -> 1443 bytes
.../src/test/resources/compatibility/split-v1-data | Bin 896 -> 912 bytes
.../test/resources/compatibility/split-v1-fallback | Bin 1436 -> 1460 bytes
.../resources/compatibility/split-v1-fallback-data | Bin 897 -> 913 bytes
.../resources/compatibility/split-v1-incremental | Bin 934 -> 950 bytes
.../test/resources/compatibility/split-v1-indexed | Bin 961 -> 977 bytes
.../resources/compatibility/split-v1-query-auth | Bin 1011 -> 1027 bytes
.../changelog/ChangelogCompactTaskSerializer.java | 2 +-
.../ChangelogCompactSortOperatorTest.java | 1 +
.../ChangelogCompactTaskSerializerTest.java | 35 ++--
.../globalindex/GenericIndexTopoBuilderTest.java | 1 +
.../PendingSplitsCheckpointSerializerTest.java | 154 ++++++++++++++
.../pending-splits-incremental-v1-chain-v2 | Bin 0 -> 2659 bytes
.../apache/paimon/spark/copy/CopyFilesUtil.java | 5 +-
.../paimon/spark/copy/CopyFilesUtilTest.java | 61 ++++++
.../procedure/CreateGlobalIndexProcedureTest.java | 1 +
68 files changed, 1267 insertions(+), 92 deletions(-)
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 9a6fd514df..a19d1818c2 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 class CoreOptions implements Serializable {
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 019cf3458c..40ae2073dd 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 @@ import static
org.apache.paimon.utils.SerializationUtils.serializeBinaryRow;
public class DataEvolutionCompactTaskSerializer
implements VersionedSerializer<DataEvolutionCompactTask> {
- 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 9f4e41c34d..10a368a49a 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.io.DataFileMeta;
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 class DataEvolutionNormalCompactTask extends
DataEvolutionCompactTask {
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<Long, TableSchema> schemaCache = new HashMap<>();
+ Function<Long, TableSchema> schemaLoader =
+ schemaId -> schemaCache.computeIfAbsent(schemaId,
schemaManager::schema);
+ Map<Pair<Long, List<String>>, List<DataField>> fileFieldsCache = new
HashMap<>();
+
+ Map<Integer, Long> fieldMaxSequences = new HashMap<>();
+ for (DataFileMeta input : compactBefore) {
+ List<DataField> 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<DataField> 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 f8b4e6aaf5..e9874917a1 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 WRITE_COLS_SEQUENCES = "_WRITE_COLS_SEQUENCES";
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,
+ WRITE_COLS_SEQUENCES,
+ new ArrayType(true, new
BigIntType(false)))));
BinaryRow EMPTY_MIN_KEY = EMPTY_ROW;
BinaryRow EMPTY_MAX_KEY = EMPTY_ROW;
@@ -150,7 +155,8 @@ public interface DataFileMeta {
valueStatsCols,
externalPath,
firstRowId,
- writeCols);
+ writeCols,
+ null);
}
static DataFileMeta create(
@@ -173,7 +179,7 @@ public interface DataFileMeta {
@Nullable String externalPath,
@Nullable Long firstRowId,
@Nullable List<String> writeCols) {
- return new PojoDataFileMeta(
+ return create(
fileName,
fileSize,
rowCount,
@@ -193,7 +199,8 @@ public interface DataFileMeta {
valueStatsCols,
externalPath,
firstRowId,
- writeCols);
+ writeCols,
+ null);
}
static DataFileMeta create(
@@ -234,7 +241,8 @@ public interface DataFileMeta {
valueStatsCols,
null,
firstRowId,
- writeCols);
+ writeCols,
+ null);
}
static DataFileMeta create(
@@ -257,7 +265,8 @@ public interface DataFileMeta {
@Nullable List<String> valueStatsCols,
@Nullable String externalPath,
@Nullable Long firstRowId,
- @Nullable List<String> writeCols) {
+ @Nullable List<String> writeCols,
+ @Nullable long[] columnMaxSequenceNumbers) {
return new PojoDataFileMeta(
fileName,
fileSize,
@@ -278,7 +287,8 @@ public interface DataFileMeta {
valueStatsCols,
externalPath,
firstRowId,
- writeCols);
+ writeCols,
+ columnMaxSequenceNumbers);
}
String fileName();
@@ -354,6 +364,16 @@ public interface DataFileMeta {
@Nullable
List<String> writeCols();
+ /**
+ * Maximum sequence number per physical table field after data-evolution
compaction.
+ *
+ * <p>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 @@ public interface DataFileMeta {
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 b646ef08ca..c69a319f2d 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 class DataFileMeta08Serializer implements
Serializable {
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 662f1276c8..41e93bbcba 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 class DataFileMeta09Serializer implements
Serializable {
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 dca1aa528f..2171085da8 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 class DataFileMeta10LegacySerializer implements
Serializable {
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 e888c1ca74..b56a3396e3 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 class DataFileMeta12LegacySerializer implements
Serializable {
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 59abcc730d..63981a1e29 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<Dat
private static final long serialVersionUID = 1L;
public DataFileMetaFirstRowIdLegacySerializer() {
- super(DataFileMeta.SCHEMA);
+ super(DataFileMetaWriteColsLegacySerializer.SCHEMA);
}
@Override
@@ -86,6 +86,7 @@ public class DataFileMetaFirstRowIdLegacySerializer extends
ObjectSerializer<Dat
row.isNullAt(16) ? null :
fromStringArrayData(row.getArray(16)),
row.isNullAt(17) ? null : row.getString(17).toString(),
row.isNullAt(18) ? null : row.getLong(18),
+ null,
null);
}
}
diff --git
a/paimon-core/src/main/java/org/apache/paimon/io/DataFileMetaSerializer.java
b/paimon-core/src/main/java/org/apache/paimon/io/DataFileMetaSerializer.java
index afed7265d4..f148bb397e 100644
--- a/paimon-core/src/main/java/org/apache/paimon/io/DataFileMetaSerializer.java
+++ b/paimon-core/src/main/java/org/apache/paimon/io/DataFileMetaSerializer.java
@@ -19,7 +19,9 @@
package org.apache.paimon.io;
import org.apache.paimon.data.BinaryString;
+import org.apache.paimon.data.GenericArray;
import org.apache.paimon.data.GenericRow;
+import org.apache.paimon.data.InternalArray;
import org.apache.paimon.data.InternalRow;
import org.apache.paimon.manifest.FileSource;
import org.apache.paimon.stats.SimpleStats;
@@ -41,6 +43,7 @@ public class DataFileMetaSerializer extends
ObjectSerializer<DataFileMeta> {
@Override
public InternalRow toRow(DataFileMeta meta) {
+ long[] columnMaxSequenceNumbers = meta.columnMaxSequenceNumbers();
return GenericRow.of(
BinaryString.fromString(meta.fileName()),
meta.fileSize(),
@@ -61,7 +64,10 @@ public class DataFileMetaSerializer extends
ObjectSerializer<DataFileMeta> {
toStringArrayData(meta.valueStatsCols()),
meta.externalPath().map(BinaryString::fromString).orElse(null),
meta.firstRowId(),
- meta.writeCols() == null ? null :
toStringArrayData(meta.writeCols()));
+ meta.writeCols() == null ? null :
toStringArrayData(meta.writeCols()),
+ columnMaxSequenceNumbers == null
+ ? null
+ : new GenericArray(columnMaxSequenceNumbers));
}
@Override
@@ -86,6 +92,15 @@ public class DataFileMetaSerializer extends
ObjectSerializer<DataFileMeta> {
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)));
+ row.isNullAt(19) ? null :
fromStringArrayData(row.getArray(19)),
+ row.isNullAt(20) ? null : fromLongArray(row.getArray(20)));
+ }
+
+ private static long[] fromLongArray(InternalArray array) {
+ long[] result = new long[array.size()];
+ for (int i = 0; i < array.size(); i++) {
+ result[i] = array.getLong(i);
+ }
+ return result;
}
}
diff --git
a/paimon-core/src/main/java/org/apache/paimon/io/DataFileMetaSerializer.java
b/paimon-core/src/main/java/org/apache/paimon/io/DataFileMetaWriteColsLegacySerializer.java
similarity index 73%
copy from
paimon-core/src/main/java/org/apache/paimon/io/DataFileMetaSerializer.java
copy to
paimon-core/src/main/java/org/apache/paimon/io/DataFileMetaWriteColsLegacySerializer.java
index afed7265d4..062bcf4fd1 100644
--- a/paimon-core/src/main/java/org/apache/paimon/io/DataFileMetaSerializer.java
+++
b/paimon-core/src/main/java/org/apache/paimon/io/DataFileMetaWriteColsLegacySerializer.java
@@ -23,6 +23,7 @@ import org.apache.paimon.data.GenericRow;
import org.apache.paimon.data.InternalRow;
import org.apache.paimon.manifest.FileSource;
import org.apache.paimon.stats.SimpleStats;
+import org.apache.paimon.types.RowType;
import org.apache.paimon.utils.ObjectSerializer;
import static org.apache.paimon.utils.InternalRowUtils.fromStringArrayData;
@@ -30,13 +31,36 @@ import static
org.apache.paimon.utils.InternalRowUtils.toStringArrayData;
import static org.apache.paimon.utils.SerializationUtils.deserializeBinaryRow;
import static org.apache.paimon.utils.SerializationUtils.serializeBinaryRow;
-/** Serializer for {@link DataFileMeta}. */
-public class DataFileMetaSerializer extends ObjectSerializer<DataFileMeta> {
+/** Legacy serializer for {@link DataFileMeta} before column sequence numbers
were introduced. */
+public class DataFileMetaWriteColsLegacySerializer extends
ObjectSerializer<DataFileMeta> {
private static final long serialVersionUID = 1L;
- public DataFileMetaSerializer() {
- super(DataFileMeta.SCHEMA);
+ 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
@@ -86,6 +110,7 @@ public class DataFileMetaSerializer extends
ObjectSerializer<DataFileMeta> {
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)));
+ 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 9b288d1d5f..6cf1eaa0d5 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<String> writeCols;
+ private final @Nullable long[] columnMaxSequenceNumbers;
+
public PojoDataFileMeta(
String fileName,
long fileSize,
@@ -98,7 +100,8 @@ public class PojoDataFileMeta implements DataFileMeta {
@Nullable List<String> valueStatsCols,
@Nullable String externalPath,
@Nullable Long firstRowId,
- @Nullable List<String> writeCols) {
+ @Nullable List<String> writeCols,
+ @Nullable long[] columnMaxSequenceNumbers) {
this.fileName = fileName;
this.fileSize = fileSize;
@@ -123,6 +126,8 @@ public class PojoDataFileMeta implements DataFileMeta {
this.externalPath = externalPath;
this.firstRowId = firstRowId;
this.writeCols = writeCols;
+ this.columnMaxSequenceNumbers =
+ columnMaxSequenceNumbers == null ? null :
columnMaxSequenceNumbers.clone();
}
@Override
@@ -239,6 +244,12 @@ public class PojoDataFileMeta implements DataFileMeta {
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 class PojoDataFileMeta implements DataFileMeta {
valueStatsCols,
externalPath,
firstRowId,
- writeCols);
+ writeCols,
+ columnMaxSequenceNumbers);
}
@Override
@@ -288,7 +300,8 @@ public class PojoDataFileMeta implements DataFileMeta {
valueStatsCols,
newExternalPath,
firstRowId,
- writeCols);
+ writeCols,
+ columnMaxSequenceNumbers);
}
@Override
@@ -313,7 +326,8 @@ public class PojoDataFileMeta implements DataFileMeta {
Collections.emptyList(),
externalPath,
firstRowId,
- writeCols);
+ writeCols,
+ columnMaxSequenceNumbers);
}
@Override
@@ -338,7 +352,34 @@ public class PojoDataFileMeta implements DataFileMeta {
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 class PojoDataFileMeta implements DataFileMeta {
valueStatsCols,
externalPath,
firstRowId,
- writeCols);
+ writeCols,
+ columnMaxSequenceNumbers);
}
@Override
@@ -388,7 +430,8 @@ public class PojoDataFileMeta implements DataFileMeta {
valueStatsCols,
externalPath,
newFirstRowId,
- writeCols);
+ writeCols,
+ columnMaxSequenceNumbers);
}
@Override
@@ -413,7 +456,8 @@ public class PojoDataFileMeta implements DataFileMeta {
valueStatsCols,
externalPath,
firstRowId,
- writeCols);
+ writeCols,
+ columnMaxSequenceNumbers);
}
@Override
@@ -438,7 +482,8 @@ public class PojoDataFileMeta implements DataFileMeta {
valueStatsCols,
newExternalPath,
firstRowId,
- writeCols);
+ writeCols,
+ columnMaxSequenceNumbers);
}
@Override
@@ -463,7 +508,8 @@ public class PojoDataFileMeta implements DataFileMeta {
valueStatsCols,
externalPath,
firstRowId,
- writeCols);
+ writeCols,
+ columnMaxSequenceNumbers);
}
@Override
@@ -494,7 +540,8 @@ public class PojoDataFileMeta implements DataFileMeta {
&& 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 class PojoDataFileMeta implements DataFileMeta {
valueStatsCols,
externalPath,
firstRowId,
- writeCols);
+ writeCols,
+ Arrays.hashCode(columnMaxSequenceNumbers));
}
@Override
@@ -529,7 +577,8 @@ public class PojoDataFileMeta implements DataFileMeta {
+ "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 class PojoDataFileMeta implements DataFileMeta {
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 e525c4e2fc..70da63dba7 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 final class ProjectedDataFileMeta implements
DataFileMeta {
return nullableStringArray(Fields.WRITE_COLS);
}
+ @Nullable
+ @Override
+ public long[] columnMaxSequenceNumbers() {
+ int position = requiredPosition(Fields.WRITE_COLS_SEQUENCES);
+ 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 final class ProjectedDataFileMeta implements
DataFileMeta {
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 @@ public final class ProjectedDataFileMeta implements
DataFileMeta {
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 WRITE_COLS_SEQUENCES =
+ fieldIndex(DataFileMeta.WRITE_COLS_SEQUENCES);
}
/** 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 0000000000..3640fe65ee
--- /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<ManifestEntry> {
+
+ 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 ed78ec2f7f..ec6df446a8 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 @@ import static
org.apache.paimon.utils.SerializationUtils.serializeBinaryRow;
/** Serializer for {@link AppendCompactTask}. */
public class AppendCompactTaskSerializer implements
VersionedSerializer<AppendCompactTask> {
- 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 f60415bc80..3f86220338 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 class CommitMessageLegacyV2Serializer {
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 2593a2a68a..8222b07c8c 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.DataFileMeta10LegacySerializer;
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 @@ import static
org.apache.paimon.utils.SerializationUtils.serializeBinaryRow;
/** {@link VersionedSerializer} for {@link CommitMessage}. */
public class CommitMessageSerializer implements
VersionedSerializer<CommitMessage> {
- 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 @@ public class CommitMessageSerializer implements
VersionedSerializer<CommitMessag
private IOExceptionSupplier<List<DataFileMeta>> 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 8478d04ea3..d1212db99f 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 @@ import static
org.apache.paimon.utils.SerializationUtils.serializeBinaryRow;
public class MultiTableCompactionTaskSerializer
implements VersionedSerializer<MultiTableAppendCompactTask> {
- 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 bad2dc1c7b..39edeb6c23 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 @@ package org.apache.paimon.table.source;
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<DataFileMeta> dataFiles;
@@ -217,7 +219,10 @@ public class ChainSplit implements Split {
int n = in.readInt();
List<DataFileMeta> dataFiles = new ArrayList<>(n);
- DataFileMetaSerializer dataFileSer = new DataFileMetaSerializer();
+ ObjectSerializer<DataFileMeta> 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 df31763a30..88bf60f019 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.DataFileMeta10LegacySerializer;
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 @@ public class DataSplit implements Split {
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 b4bd8be490..7bb2436d92 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 @@ package org.apache.paimon.table.source;
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 class IncrementalSplit implements Split {
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 class IncrementalSplit implements Split {
int bucket = in.readInt();
int totalBuckets = in.readInt();
- DataFileMetaSerializer dataFileMetaSerializer = new
DataFileMetaSerializer();
+ ObjectSerializer<DataFileMeta> dataFileMetaSerializer =
+ version == 1
+ ? new DataFileMetaWriteColsLegacySerializer()
+ : new DataFileMetaSerializer();
FunctionWithIOException<DataInputView, DeletionFile>
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 262da197ae..b8751ff81b 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.SpecialFields;
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 @@ public class DataEvolutionUtils {
return ids;
}
+ /** Table fields physically present in a file, in their physical write
order. */
+ public static List<DataField> fileFields(
+ Function<Long, TableSchema> scanTableSchema, DataFileMeta file) {
+ TableSchema schema = scanTableSchema.apply(file.schemaId());
+ List<String> writeCols = file.writeCols();
+ if (writeCols == null) {
+ return schema.fields();
+ }
+
+ Map<String, DataField> fieldsByName = new HashMap<>();
+ for (DataField field : schema.fields()) {
+ fieldsByName.put(field.name(), field);
+ }
+ List<DataField> 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 266ff305f4..ac71ff2552 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 class CoreOptionsTest {
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 b86b020cd1..9d1f05aa65 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 @@ public class BlobUpdateTest extends TableTestBase {
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 8e8c95193a..cd784064bb 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 @@ public class DataEvolutionCompactCoordinatorTest {
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 31003beb75..88278ec40f 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 @@ class DataEvolutionCompactRangePlannerTest extends
TableTestBase {
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 0000000000..abf036b974
--- /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<String> columns,
+ int updateRound,
+ CoreOptions.GlobalIndexColumnUpdateAction updateAction)
+ throws Exception {
+ Map<String, String> writeOptions = new HashMap<>();
+ writeOptions.put(
+ CoreOptions.GLOBAL_INDEX_COLUMN_UPDATE_ACTION.key(),
updateAction.toString());
+ FileStoreTable table = getTableDefault().copy(writeOptions);
+ List<String> 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<Object> 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<CommitMessage> 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<CommitMessage> 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<DataFileMeta> 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<DataField> 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<CommitMessage> messages, long
firstRowId) {
+ for (CommitMessage message : messages) {
+ CommitMessageImpl impl = (CommitMessageImpl) message;
+ List<DataFileMeta> 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 0e352b63ec..58dad8079f 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 @@ public class IndexBootstrapTest extends TableTestBase {
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 549814fafe..063114a996 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 @@ class GlobalIndexBuilderUtilsTest {
null,
null,
firstRowId,
+ null,
null);
return ManifestEntry.create(FileKind.ADD, BinaryRow.EMPTY_ROW, 0, 1,
file);
}
@@ -362,6 +363,7 @@ class GlobalIndexBuilderUtilsTest {
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 5074ccf1cf..be87c063aa 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 @@ package org.apache.paimon.io;
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<DataFileMeta> {
@@ -34,6 +39,47 @@ public class DataFileMetaSerializerTest extends
ObjectSerializerTestBase<DataFil
@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 974fe7aaf8..556e4de335 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 class DataFileTestUtils {
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 d684b44739..744aa14ed0 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 @@ public class ProjectedDataFileMetaTest {
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 @@ public class ProjectedDataFileMetaTest {
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 20eac11da7..dbe1ebfab0 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.DataFileMeta;
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 @@ import static org.assertj.core.api.Assertions.assertThat;
/** 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 class ManifestCommittableSerializerCompatibilityTest {
Arrays.asList("field1", "field2", "field3"),
"hdfs://localhost:9000/path/to/file",
1L,
- Arrays.asList("asdf", "qwer", "zxcv"));
+ Arrays.asList("asdf", "qwer", "zxcv"),
+ null);
List<DataFileMeta> dataFiles = Collections.singletonList(dataFile);
GlobalIndexMeta globalIndexMeta =
new GlobalIndexMeta(
@@ -212,7 +272,8 @@ public class ManifestCommittableSerializerCompatibilityTest
{
Arrays.asList("field1", "field2", "field3"),
"hdfs://localhost:9000/path/to/file",
1L,
- Arrays.asList("asdf", "qwer", "zxcv"));
+ Arrays.asList("asdf", "qwer", "zxcv"),
+ null);
List<DataFileMeta> 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 class ManifestCommittableSerializerCompatibilityTest
{
Arrays.asList("field1", "field2", "field3"),
"hdfs://localhost:9000/path/to/file",
1L,
- Arrays.asList("asdf", "qwer", "zxcv"));
+ Arrays.asList("asdf", "qwer", "zxcv"),
+ null);
List<DataFileMeta> dataFiles = Collections.singletonList(dataFile);
IndexFileMeta hashIndexFile =
new IndexFileMeta(
@@ -400,7 +462,8 @@ public class ManifestCommittableSerializerCompatibilityTest
{
Arrays.asList("field1", "field2", "field3"),
"hdfs://localhost:9000/path/to/file",
1L,
- Arrays.asList("asdf", "qwer", "zxcv"));
+ Arrays.asList("asdf", "qwer", "zxcv"),
+ null);
List<DataFileMeta> dataFiles = Collections.singletonList(dataFile);
LinkedHashMap<String, DeletionVectorMeta> dvRanges = new
LinkedHashMap<>();
@@ -484,6 +547,7 @@ public class ManifestCommittableSerializerCompatibilityTest
{
Arrays.asList("field1", "field2", "field3"),
"hdfs://localhost:9000/path/to/file",
1L,
+ null,
null);
List<DataFileMeta> dataFiles = Collections.singletonList(dataFile);
@@ -563,6 +627,7 @@ public class ManifestCommittableSerializerCompatibilityTest
{
Arrays.asList("field1", "field2", "field3"),
"hdfs://localhost:9000/path/to/file",
null,
+ null,
null);
List<DataFileMeta> dataFiles = Collections.singletonList(dataFile);
@@ -641,6 +706,7 @@ public class ManifestCommittableSerializerCompatibilityTest
{
Arrays.asList("field1", "field2", "field3"),
"hdfs://localhost:9000/path/to/file",
null,
+ null,
null);
List<DataFileMeta> dataFiles = Collections.singletonList(dataFile);
@@ -716,6 +782,7 @@ public class ManifestCommittableSerializerCompatibilityTest
{
Arrays.asList("field1", "field2", "field3"),
"hdfs://localhost:9000/path/to/file",
null,
+ null,
null);
List<DataFileMeta> dataFiles = Collections.singletonList(dataFile);
@@ -791,6 +858,7 @@ public class ManifestCommittableSerializerCompatibilityTest
{
Arrays.asList("field1", "field2", "field3"),
null,
null,
+ null,
null);
List<DataFileMeta> dataFiles = Collections.singletonList(dataFile);
@@ -865,6 +933,7 @@ public class ManifestCommittableSerializerCompatibilityTest
{
Arrays.asList("field1", "field2", "field3"),
null,
null,
+ null,
null);
List<DataFileMeta> dataFiles = Collections.singletonList(dataFile);
@@ -940,6 +1009,7 @@ public class
ManifestCommittableSerializerCompatibilityTest {
null,
null,
null,
+ null,
null);
List<DataFileMeta> dataFiles = Collections.singletonList(dataFile);
@@ -1015,6 +1085,7 @@ public class
ManifestCommittableSerializerCompatibilityTest {
null,
null,
null,
+ null,
null);
List<DataFileMeta> dataFiles = Collections.singletonList(dataFile);
@@ -1090,6 +1161,7 @@ public class
ManifestCommittableSerializerCompatibilityTest {
null,
null,
null,
+ null,
null);
List<DataFileMeta> 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 0547c9acce..5005f04cc9 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.apache.paimon.utils.ObjectSerializerTestBase;
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 @@ public class ManifestEntrySerializerTest extends
ObjectSerializerTestBase<Manife
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<ManifestEntry> 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 6da1b8d2f7..0120471b4d 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 @@ public class ManifestFileMetaTest extends
ManifestFileMetaTestBase {
null,
null,
null,
+ null,
null));
}
@@ -2804,7 +2805,8 @@ public class ManifestFileMetaTest extends
ManifestFileMetaTestBase {
null,
null,
firstRowId,
- Collections.singletonList("f0")));
+ Collections.singletonList("f0"),
+ null));
}
private List<String> readFileNames(
@@ -2887,7 +2889,8 @@ public class ManifestFileMetaTest extends
ManifestFileMetaTestBase {
null,
externalPath,
firstRowId,
- Collections.singletonList("f0")));
+ Collections.singletonList("f0"),
+ null));
}
private List<ManifestEntry> readEntries(List<ManifestFileMeta>
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 649ebb73ec..b7eff0526b 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 @@ public abstract class ManifestFileMetaTestBase {
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 d3ba2527de..b551edc959 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.Path;
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 @@ public class ManifestFileTest {
}
}
+ @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<DataField> 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<InternalRow> 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 52ac56608b..3170c08d02 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 @@ public class NoPartitionManifestFileMetaTest extends
ManifestFileMetaTestBase {
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 02d52df949..d34cc1cfd9 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 @@ public class LevelsOrderingFixTest {
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 464b26b944..ede8b678e9 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 @@ public class IntervalPartitionTest {
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 1c41741b70..859e8bc5b8 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 @@ public class BlobFallbackRecordReaderTest {
null,
null,
firstRowId,
- Arrays.asList(BLOB_FIELD));
+ Arrays.asList(BLOB_FIELD),
+ null);
}
private static List<Range> 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 51a0a93aa2..24d2994957 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 @@ public class DataEvolutionReadTest {
null,
null,
firstRowId,
- Arrays.asList("blob_col"));
+ Arrays.asList("blob_col"),
+ null);
}
/** Creates a blob file with custom write columns. */
@@ -384,7 +385,8 @@ public class DataEvolutionReadTest {
null,
null,
firstRowId,
- writeCols);
+ writeCols,
+ null);
}
private DataFileMeta createVectorFile(
@@ -445,7 +447,8 @@ public class DataEvolutionReadTest {
null,
null,
firstRowId,
- writeCols);
+ writeCols,
+ null);
}
@Test
@@ -523,6 +526,7 @@ public class DataEvolutionReadTest {
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 143f8445b8..46c7626f11 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 class ExpireSnapshotsTest {
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 class ExpireSnapshotsTest {
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 f93d1e7134..2d08a158d5 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 @@ class ManifestEntryRunMergeTest extends
ManifestFileMetaTestBase {
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 d54e29e91f..658507da86 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 @@ class ManifestRewriteCleanupTest extends
ManifestFileMetaTestBase {
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 c4f519e84e..bc36deedca 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 class CommitMessageSerializerTest {
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 a0c76537ab..e3f5403c68 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.FloatType;
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 @@ import static org.assertj.core.api.Assertions.assertThat;
/** 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 class DataSplitCompatibleTest {
DataFileTestDataGenerator gen =
DataFileTestDataGenerator.builder().build();
DataFileTestDataGenerator.Data data = gen.next();
List<DataFileMeta> 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 class DataSplitCompatibleTest {
null,
null,
null,
+ null,
null);
List<DataFileMeta> dataFiles = Collections.singletonList(dataFile);
@@ -336,6 +341,7 @@ public class DataSplitCompatibleTest {
null,
null,
null,
+ null,
null);
List<DataFileMeta> dataFiles = Collections.singletonList(dataFile);
@@ -401,6 +407,7 @@ public class DataSplitCompatibleTest {
Arrays.asList("field1", "field2", "field3"),
null,
null,
+ null,
null);
List<DataFileMeta> dataFiles = Collections.singletonList(dataFile);
@@ -470,6 +477,7 @@ public class DataSplitCompatibleTest {
Arrays.asList("field1", "field2", "field3"),
null,
null,
+ null,
null);
List<DataFileMeta> dataFiles = Collections.singletonList(dataFile);
@@ -539,6 +547,7 @@ public class DataSplitCompatibleTest {
Arrays.asList("field1", "field2", "field3"),
"hdfs:///path/to/warehouse",
null,
+ null,
null);
List<DataFileMeta> dataFiles = Collections.singletonList(dataFile);
@@ -608,6 +617,7 @@ public class DataSplitCompatibleTest {
Arrays.asList("field1", "field2", "field3"),
"hdfs:///path/to/warehouse",
null,
+ null,
null);
List<DataFileMeta> dataFiles = Collections.singletonList(dataFile);
@@ -678,6 +688,7 @@ public class DataSplitCompatibleTest {
Arrays.asList("field1", "field2", "field3"),
"hdfs:///path/to/warehouse",
12L,
+ null,
null);
List<DataFileMeta> dataFiles = Collections.singletonList(dataFile);
@@ -748,7 +759,8 @@ public class DataSplitCompatibleTest {
Arrays.asList("field1", "field2", "field3"),
"hdfs:///path/to/warehouse",
12L,
- Arrays.asList("a", "b", "c", "f"));
+ Arrays.asList("a", "b", "c", "f"),
+ null);
List<DataFileMeta> dataFiles = Collections.singletonList(dataFile);
DeletionFile deletionFile = new DeletionFile("deletion_file", 100, 22,
33L);
@@ -784,6 +796,86 @@ public class DataSplitCompatibleTest {
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<DataFileMeta> dataFiles = Collections.singletonList(dataFile);
+
+ DeletionFile deletionFile = new DeletionFile("deletion_file", 100, 22,
33L);
+ List<DeletionFile> 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 0fd6f982b5..3b6afc3acf 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
@@ -250,6 +250,7 @@ public class SplitSerializerTest {
null,
null,
null,
+ null,
null);
}
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 9ddcf7dfda..9797af2a91 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 @@ public class ChainTableUtilsTest {
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 c499c5f276..99c4508bd6 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 class DataEvolutionUtilsTest {
.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 0000000000..6277c55665
Binary files /dev/null and
b/paimon-core/src/test/resources/compatibility/datasplit-v9 differ
diff --git
a/paimon-core/src/test/resources/compatibility/manifest-committable-v13-v5
b/paimon-core/src/test/resources/compatibility/manifest-committable-v13-v5
new file mode 100644
index 0000000000..d919649c29
Binary files /dev/null and
b/paimon-core/src/test/resources/compatibility/manifest-committable-v13-v5
differ
diff --git a/paimon-core/src/test/resources/compatibility/split-v1-chain
b/paimon-core/src/test/resources/compatibility/split-v1-chain
index 08f61c52fa..d7fa63bf04 100644
Binary files a/paimon-core/src/test/resources/compatibility/split-v1-chain and
b/paimon-core/src/test/resources/compatibility/split-v1-chain differ
diff --git a/paimon-core/src/test/resources/compatibility/split-v1-data
b/paimon-core/src/test/resources/compatibility/split-v1-data
index 6cbade9c58..9d2f6c085f 100644
Binary files a/paimon-core/src/test/resources/compatibility/split-v1-data and
b/paimon-core/src/test/resources/compatibility/split-v1-data differ
diff --git a/paimon-core/src/test/resources/compatibility/split-v1-fallback
b/paimon-core/src/test/resources/compatibility/split-v1-fallback
index f83872aea1..9dbf5ee44f 100644
Binary files a/paimon-core/src/test/resources/compatibility/split-v1-fallback
and b/paimon-core/src/test/resources/compatibility/split-v1-fallback differ
diff --git
a/paimon-core/src/test/resources/compatibility/split-v1-fallback-data
b/paimon-core/src/test/resources/compatibility/split-v1-fallback-data
index 24be9e0ad5..42d174704c 100644
Binary files
a/paimon-core/src/test/resources/compatibility/split-v1-fallback-data and
b/paimon-core/src/test/resources/compatibility/split-v1-fallback-data differ
diff --git a/paimon-core/src/test/resources/compatibility/split-v1-incremental
b/paimon-core/src/test/resources/compatibility/split-v1-incremental
index cb057e42d7..e5368a7ea6 100644
Binary files
a/paimon-core/src/test/resources/compatibility/split-v1-incremental and
b/paimon-core/src/test/resources/compatibility/split-v1-incremental differ
diff --git a/paimon-core/src/test/resources/compatibility/split-v1-indexed
b/paimon-core/src/test/resources/compatibility/split-v1-indexed
index 0d20df1012..eab6fbe8f8 100644
Binary files a/paimon-core/src/test/resources/compatibility/split-v1-indexed
and b/paimon-core/src/test/resources/compatibility/split-v1-indexed differ
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 fce1aa5fb9..5e5acea55d 100644
Binary files a/paimon-core/src/test/resources/compatibility/split-v1-query-auth
and b/paimon-core/src/test/resources/compatibility/split-v1-query-auth differ
diff --git
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/compact/changelog/ChangelogCompactTaskSerializer.java
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/compact/changelog/ChangelogCompactTaskSerializer.java
index c7f56dc5de..e796b92054 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 @@ import static
org.apache.paimon.utils.SerializationUtils.serializeBinaryRow;
public class ChangelogCompactTaskSerializer
implements SimpleVersionedSerializer<ChangelogCompactTask> {
- 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 70530de9dc..bd706147a5 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 @@ public class ChangelogCompactSortOperatorTest {
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 df62115f27..2dba94abd8 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 @@ public class ChangelogCompactTaskSerializerTest {
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 57d58b34c6..343c744194 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 @@ class GenericIndexTopoBuilderTest {
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 2e4b66983e..e89ef16bb2 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 @@ import static org.assertj.core.api.Assertions.assertThat;
/** 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 class PendingSplitsCheckpointSerializerTest {
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<String, String> branchMapping = new LinkedHashMap<>();
+ branchMapping.put(before.fileName(), "snapshot");
+ branchMapping.put(after.fileName(), "delta");
+ Map<String, String> 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<FileStoreSourceSplit> 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<FileStoreSourceSplit> 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 @@ public class PendingSplitsCheckpointSerializerTest {
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<String, String> branchMapping,
+ Map<String, String> 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<DataFileMeta> 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 0000000000..dffb60857f
Binary files /dev/null and
b/paimon-flink/paimon-flink-common/src/test/resources/compatibility/pending-splits-incremental-v1-chain-v2
differ
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 1a1e12efbe..88dd1fb403 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 class CopyFilesUtil {
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 0000000000..2f973ffce3
--- /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 f40c751fe2..f25bde2ba5 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 @@ public class CreateGlobalIndexProcedureTest {
null,
null,
firstRowId,
+ null,
null);
}