From 0932ee6ea31ad2b048255d7c89597b8cf68a714f Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E4=BB=9F=E5=BC=8B?= Date: Thu, 6 Aug 2026 16:44:37 +0800 Subject: [PATCH 1/4] [core] Persist column sequences in data file metadata --- .../DataEvolutionCompactTaskSerializer.java | 2 +- .../DataEvolutionNormalCompactTask.java | 38 ++++++ ...ataEvolutionGlobalIndexRefreshPlanner.java | 41 ++++-- .../apache/paimon/io/BinaryDataFileMeta.java | 28 ++++ .../org/apache/paimon/io/DataFileMeta.java | 70 +++++++++- ...ataFileMetaFirstRowIdLegacySerializer.java | 2 +- .../paimon/io/DataFileMetaSerializer.java | 24 +++- ...DataFileMetaWriteColsLegacySerializer.java | 115 ++++++++++++++++ .../apache/paimon/io/PojoDataFileMeta.java | 127 ++++++++++++++++-- .../sink/AppendCompactTaskSerializer.java | 2 +- .../table/sink/CommitMessageSerializer.java | 11 +- .../MultiTableCompactionTaskSerializer.java | 2 +- .../paimon/table/source/ChainSplit.java | 9 +- .../apache/paimon/table/source/DataSplit.java | 7 +- .../paimon/table/source/IncrementalSplit.java | 11 +- .../paimon/table/source/SplitSerializer.java | 23 ++-- .../paimon/utils/DataEvolutionUtils.java | 10 ++ ...volutionGlobalIndexRefreshPlannerTest.java | 45 +++++++ .../sorted/SortedGlobalIndexScannerTest.java | 125 +++++++++++++++++ .../paimon/io/BinaryDataFileMetaTest.java | 33 +++-- .../paimon/io/DataFileMetaSerializerTest.java | 35 ++++- .../sink/CommitMessageSerializerTest.java | 9 ++ .../table/source/DataSplitCompatibleTest.java | 1 + .../table/source/SplitSerializerTest.java | 17 +-- .../ChangelogCompactTaskSerializer.java | 2 +- .../ChangelogCompactTaskSerializerTest.java | 36 ++--- .../paimon/spark/copy/CopyFilesUtil.java | 3 +- 27 files changed, 735 insertions(+), 93 deletions(-) create mode 100644 paimon-core/src/main/java/org/apache/paimon/io/DataFileMetaWriteColsLegacySerializer.java diff --git a/paimon-core/src/main/java/org/apache/paimon/append/dataevolution/DataEvolutionCompactTaskSerializer.java b/paimon-core/src/main/java/org/apache/paimon/append/dataevolution/DataEvolutionCompactTaskSerializer.java index 019cf3458c66..40ae2073dd17 100644 --- a/paimon-core/src/main/java/org/apache/paimon/append/dataevolution/DataEvolutionCompactTaskSerializer.java +++ b/paimon-core/src/main/java/org/apache/paimon/append/dataevolution/DataEvolutionCompactTaskSerializer.java @@ -40,7 +40,7 @@ public class DataEvolutionCompactTaskSerializer implements VersionedSerializer { - private static final int CURRENT_VERSION = 2; + private static final int CURRENT_VERSION = 3; private final DataFileMetaSerializer dataFileSerializer; diff --git a/paimon-core/src/main/java/org/apache/paimon/append/dataevolution/DataEvolutionNormalCompactTask.java b/paimon-core/src/main/java/org/apache/paimon/append/dataevolution/DataEvolutionNormalCompactTask.java index 1db0786d27cc..916a99ce03b2 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 @@ -36,13 +36,17 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import java.util.LinkedHashMap; import java.util.List; +import java.util.Map; import java.util.Set; import java.util.stream.Collectors; import static org.apache.paimon.types.BlobType.fieldNamesInBlobFile; import static org.apache.paimon.types.VectorType.fieldNamesInVectorFile; import static org.apache.paimon.types.VectorType.isVectorStoreFile; +import static org.apache.paimon.utils.DataEvolutionUtils.fieldMaxSequenceNumber; +import static org.apache.paimon.utils.DataEvolutionUtils.fileFieldIds; import static org.apache.paimon.utils.Preconditions.checkArgument; /** Compacts normal structured files of a data evolution table. */ @@ -81,6 +85,8 @@ public CommitMessage doCompact(FileStoreTable table, String commitUser) throws E table.rowType().getFields().stream() .filter(f -> !fieldsInDedicatedFile.contains(f.name())) .collect(Collectors.toList())); + Map columnMaxSequenceNumbers = + compactedColumnMaxSequenceNumbers(table, readWriteType); FileStorePathFactory pathFactory = table.store().pathFactory(); AppendOnlyFileStore store = (AppendOnlyFileStore) table.store(); @@ -122,8 +128,40 @@ public CommitMessage doCompact(FileStoreTable table, String commitUser) throws E dataFileMeta = dataFileMeta.assignSequenceNumber( minSequenceId(compactBefore), maxSequenceId(compactBefore)); + dataFileMeta = dataFileMeta.withColumnMaxSequenceNumbers(columnMaxSequenceNumbers); compactAfter.add(dataFileMeta); return commitMessage(compactBefore, compactAfter); } + + private Map compactedColumnMaxSequenceNumbers( + FileStoreTable table, RowType outputType) { + Map> inputFieldIds = new LinkedHashMap<>(); + for (DataFileMeta input : compactBefore) { + inputFieldIds.put(input, fileFieldIds(table.schemaManager()::schema, input)); + } + + long fallbackSequence = maxSequenceId(compactBefore); + Map result = new LinkedHashMap<>(); + outputType + .getFields() + .forEach( + field -> { + long fieldSequence = Long.MIN_VALUE; + for (DataFileMeta input : compactBefore) { + if (inputFieldIds.get(input).contains(field.id())) { + fieldSequence = + Math.max( + fieldSequence, + fieldMaxSequenceNumber(input, field.id())); + } + } + result.put( + field.id(), + fieldSequence == Long.MIN_VALUE + ? fallbackSequence + : fieldSequence); + }); + return result; + } } diff --git a/paimon-core/src/main/java/org/apache/paimon/globalindex/DataEvolutionGlobalIndexRefreshPlanner.java b/paimon-core/src/main/java/org/apache/paimon/globalindex/DataEvolutionGlobalIndexRefreshPlanner.java index 44ed39685d24..201849edcc30 100644 --- a/paimon-core/src/main/java/org/apache/paimon/globalindex/DataEvolutionGlobalIndexRefreshPlanner.java +++ b/paimon-core/src/main/java/org/apache/paimon/globalindex/DataEvolutionGlobalIndexRefreshPlanner.java @@ -41,6 +41,7 @@ import java.util.Set; import java.util.TreeMap; +import static org.apache.paimon.utils.DataEvolutionUtils.fieldMaxSequenceNumber; import static org.apache.paimon.utils.DataEvolutionUtils.fileFieldIds; /** Plans existing global index files which need refresh after data-evolution updates. */ @@ -98,7 +99,14 @@ public static List findIndexesToRefresh( Pair.of(file.schemaId(), file.writeCols()), key -> fileFieldIds(schemaManager::schema, file)); if (!disjoint(indexedFieldIds, physicalFieldIds)) { - group.addDataFile(file); + long indexedMaxSequence = Long.MIN_VALUE; + for (Integer fieldId : indexedFieldIds) { + if (physicalFieldIds.contains(fieldId)) { + indexedMaxSequence = + Math.max(indexedMaxSequence, fieldMaxSequenceNumber(file, fieldId)); + } + } + group.addDataFile(file, indexedMaxSequence); } } @@ -119,7 +127,7 @@ public static List findIndexesToRefresh( private static final class RefreshGroup { private final List indexes = new ArrayList<>(); - private final List dataFiles = new ArrayList<>(); + private final List dataUpdates = new ArrayList<>(); private final MergedRanges indexedRanges = new MergedRanges(); private long minScanSnapshotId = Long.MAX_VALUE; @@ -134,21 +142,23 @@ private boolean mayContainUpdate(DataFileMeta file) { && indexedRanges.intersects(file.nonNullRowIdRange()); } - private void addDataFile(DataFileMeta file) { - dataFiles.add(file); + private void addDataFile(DataFileMeta file, long maxSequenceNumber) { + if (maxSequenceNumber > minScanSnapshotId) { + dataUpdates.add(new DataUpdate(file.nonNullRowIdRange(), maxSequenceNumber)); + } } private void markIndexesToRefresh(boolean[] result) { // As scan watermarks decrease, eligible data files only grow. - dataFiles.sort(Comparator.comparingLong(DataFileMeta::maxSequenceNumber).reversed()); + dataUpdates.sort(Comparator.comparingLong(DataUpdate::maxSequenceNumber).reversed()); indexes.sort((left, right) -> Long.compare(right.scanSnapshotId, left.scanSnapshotId)); MergedRanges updatedRanges = new MergedRanges(); int nextFile = 0; for (IndexQuery index : indexes) { - while (nextFile < dataFiles.size() - && dataFiles.get(nextFile).maxSequenceNumber() > index.scanSnapshotId) { - updatedRanges.add(dataFiles.get(nextFile).nonNullRowIdRange()); + while (nextFile < dataUpdates.size() + && dataUpdates.get(nextFile).maxSequenceNumber > index.scanSnapshotId) { + updatedRanges.add(dataUpdates.get(nextFile).rowRange); nextFile++; } if (updatedRanges.intersects(index.rowRange)) { @@ -158,6 +168,21 @@ private void markIndexesToRefresh(boolean[] result) { } } + private static final class DataUpdate { + + private final Range rowRange; + private final long maxSequenceNumber; + + private DataUpdate(Range rowRange, long maxSequenceNumber) { + this.rowRange = rowRange; + this.maxSequenceNumber = maxSequenceNumber; + } + + private long maxSequenceNumber() { + return maxSequenceNumber; + } + } + private static final class IndexQuery { private final int ordinal; diff --git a/paimon-core/src/main/java/org/apache/paimon/io/BinaryDataFileMeta.java b/paimon-core/src/main/java/org/apache/paimon/io/BinaryDataFileMeta.java index be56deb20727..0b49cc71a245 100644 --- a/paimon-core/src/main/java/org/apache/paimon/io/BinaryDataFileMeta.java +++ b/paimon-core/src/main/java/org/apache/paimon/io/BinaryDataFileMeta.java @@ -21,6 +21,7 @@ import org.apache.paimon.data.BinaryRow; import org.apache.paimon.data.BinaryString; import org.apache.paimon.data.InternalArray; +import org.apache.paimon.data.InternalMap; import org.apache.paimon.data.InternalRow; import org.apache.paimon.data.Timestamp; import org.apache.paimon.fs.Path; @@ -33,7 +34,9 @@ import java.util.Arrays; import java.util.Collections; +import java.util.LinkedHashMap; import java.util.List; +import java.util.Map; import java.util.Optional; import static org.apache.paimon.utils.InternalRowUtils.fromStringArrayData; @@ -244,6 +247,24 @@ public List writeCols() { return nullableStringArray(Fields.WRITE_COLS); } + @Nullable + @Override + public Map columnMaxSequenceNumbers() { + int position = requiredPosition(Fields.COLUMN_MAX_SEQUENCE_NUMBERS); + InternalRow row = currentRow(); + if (row.isNullAt(position)) { + return null; + } + InternalMap map = row.getMap(position); + InternalArray keys = map.keyArray(); + InternalArray values = map.valueArray(); + Map result = new LinkedHashMap<>(); + for (int i = 0; i < map.size(); i++) { + result.put(keys.getInt(i), values.getLong(i)); + } + return Collections.unmodifiableMap(result); + } + public boolean containsWriteColumn(BinaryString fieldName) { int position = requiredPosition(Fields.WRITE_COLS); InternalRow row = currentRow(); @@ -281,6 +302,11 @@ public DataFileMeta assignSequenceNumber(long minSequenceNumber, long maxSequenc throw unsupportedOperation("assignSequenceNumber(long, long)"); } + @Override + public DataFileMeta withColumnMaxSequenceNumbers(Map columnMaxSequenceNumbers) { + throw unsupportedOperation("withColumnMaxSequenceNumbers(Map)"); + } + @Override public DataFileMeta assignFirstRowId(long firstRowId) { throw unsupportedOperation("assignFirstRowId(long)"); @@ -378,6 +404,8 @@ private static class Fields { private static final int EXTERNAL_PATH = fieldIndex(DataFileMeta.EXTERNAL_PATH); private static final int FIRST_ROW_ID = fieldIndex(DataFileMeta.FIRST_ROW_ID); private static final int WRITE_COLS = fieldIndex(DataFileMeta.WRITE_COLS); + private static final int COLUMN_MAX_SEQUENCE_NUMBERS = + fieldIndex(DataFileMeta.COLUMN_MAX_SEQUENCE_NUMBERS); } /** Projected data-file schema together with its bound binary field layout. */ diff --git a/paimon-core/src/main/java/org/apache/paimon/io/DataFileMeta.java b/paimon-core/src/main/java/org/apache/paimon/io/DataFileMeta.java index f8b4e6aaf5b7..48c5dcbd7a88 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 @@ -29,6 +29,7 @@ import org.apache.paimon.types.DataField; import org.apache.paimon.types.DataTypes; import org.apache.paimon.types.IntType; +import org.apache.paimon.types.MapType; import org.apache.paimon.types.RowType; import org.apache.paimon.types.TinyIntType; import org.apache.paimon.utils.Range; @@ -80,6 +81,7 @@ public interface DataFileMeta { String EXTERNAL_PATH = "_EXTERNAL_PATH"; String FIRST_ROW_ID = "_FIRST_ROW_ID"; String WRITE_COLS = "_WRITE_COLS"; + String COLUMN_MAX_SEQUENCE_NUMBERS = "_COLUMN_MAX_SEQUENCE_NUMBERS"; RowType SCHEMA = new RowType( @@ -109,7 +111,11 @@ public interface DataFileMeta { new DataField(17, EXTERNAL_PATH, newStringType(true)), new DataField(18, FIRST_ROW_ID, new BigIntType(true)), new DataField( - 19, WRITE_COLS, new ArrayType(true, newStringType(false))))); + 19, WRITE_COLS, new ArrayType(true, newStringType(false))), + new DataField( + 20, + COLUMN_MAX_SEQUENCE_NUMBERS, + new MapType(true, new IntType(false), new BigIntType(false))))); BinaryRow EMPTY_MIN_KEY = EMPTY_ROW; BinaryRow EMPTY_MAX_KEY = EMPTY_ROW; @@ -173,7 +179,7 @@ static DataFileMeta create( @Nullable String externalPath, @Nullable Long firstRowId, @Nullable List writeCols) { - return new PojoDataFileMeta( + return create( fileName, fileSize, rowCount, @@ -193,7 +199,54 @@ static DataFileMeta create( valueStatsCols, externalPath, firstRowId, - writeCols); + writeCols, + null); + } + + static DataFileMeta create( + String fileName, + long fileSize, + long rowCount, + BinaryRow minKey, + BinaryRow maxKey, + SimpleStats keyStats, + SimpleStats valueStats, + long minSequenceNumber, + long maxSequenceNumber, + long schemaId, + int level, + List extraFiles, + Timestamp creationTime, + @Nullable Long deleteRowCount, + @Nullable byte[] embeddedIndex, + @Nullable FileSource fileSource, + @Nullable List valueStatsCols, + @Nullable String externalPath, + @Nullable Long firstRowId, + @Nullable List writeCols, + @Nullable Map columnMaxSequenceNumbers) { + return new PojoDataFileMeta( + fileName, + fileSize, + rowCount, + minKey, + maxKey, + keyStats, + valueStats, + minSequenceNumber, + maxSequenceNumber, + schemaId, + level, + extraFiles, + creationTime, + deleteRowCount, + embeddedIndex, + fileSource, + valueStatsCols, + externalPath, + firstRowId, + writeCols, + columnMaxSequenceNumbers); } static DataFileMeta create( @@ -354,6 +407,15 @@ default Range nonNullRowIdRange() { @Nullable List writeCols(); + /** + * Maximum sequence number per physical table field after data-evolution compaction. + * + *

A null value means that only the file-level sequence range is available. Field ids which + * are absent from a non-null map must also fall back to {@link #maxSequenceNumber()}. + */ + @Nullable + Map columnMaxSequenceNumbers(); + DataFileMeta upgrade(int newLevel); DataFileMeta rename(String newFileName); @@ -362,6 +424,8 @@ default Range nonNullRowIdRange() { DataFileMeta assignSequenceNumber(long minSequenceNumber, long maxSequenceNumber); + DataFileMeta withColumnMaxSequenceNumbers(Map columnMaxSequenceNumbers); + DataFileMeta assignFirstRowId(long firstRowId); DataFileMeta newFirstRowId(@Nullable Long newFirstRowId); diff --git a/paimon-core/src/main/java/org/apache/paimon/io/DataFileMetaFirstRowIdLegacySerializer.java b/paimon-core/src/main/java/org/apache/paimon/io/DataFileMetaFirstRowIdLegacySerializer.java index 59abcc730d38..34ad0dc28d99 100644 --- a/paimon-core/src/main/java/org/apache/paimon/io/DataFileMetaFirstRowIdLegacySerializer.java +++ b/paimon-core/src/main/java/org/apache/paimon/io/DataFileMetaFirstRowIdLegacySerializer.java @@ -36,7 +36,7 @@ public class DataFileMetaFirstRowIdLegacySerializer extends ObjectSerializer fromIntLongMap(InternalMap map) { + InternalArray keys = map.keyArray(); + InternalArray values = map.valueArray(); + Map result = new LinkedHashMap<>(); + for (int i = 0; i < map.size(); i++) { + result.put(keys.getInt(i), values.getLong(i)); + } + return result; } } diff --git a/paimon-core/src/main/java/org/apache/paimon/io/DataFileMetaWriteColsLegacySerializer.java b/paimon-core/src/main/java/org/apache/paimon/io/DataFileMetaWriteColsLegacySerializer.java new file mode 100644 index 000000000000..dd78d47d1d2a --- /dev/null +++ b/paimon-core/src/main/java/org/apache/paimon/io/DataFileMetaWriteColsLegacySerializer.java @@ -0,0 +1,115 @@ +/* + * 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.io; + +import org.apache.paimon.data.BinaryString; +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; +import static org.apache.paimon.utils.InternalRowUtils.toStringArrayData; +import static org.apache.paimon.utils.SerializationUtils.deserializeBinaryRow; +import static org.apache.paimon.utils.SerializationUtils.serializeBinaryRow; + +/** Legacy serializer for {@link DataFileMeta} before column sequence numbers were introduced. */ +public class DataFileMetaWriteColsLegacySerializer extends ObjectSerializer { + + private static final long serialVersionUID = 1L; + + static final RowType SCHEMA = + DataFileMeta.SCHEMA.project( + DataFileMeta.FILE_NAME, + DataFileMeta.FILE_SIZE, + DataFileMeta.ROW_COUNT, + DataFileMeta.MIN_KEY, + DataFileMeta.MAX_KEY, + DataFileMeta.KEY_STATS, + DataFileMeta.VALUE_STATS, + DataFileMeta.MIN_SEQUENCE_NUMBER, + DataFileMeta.MAX_SEQUENCE_NUMBER, + DataFileMeta.SCHEMA_ID, + DataFileMeta.LEVEL, + DataFileMeta.EXTRA_FILES, + DataFileMeta.CREATION_TIME, + DataFileMeta.DELETE_ROW_COUNT, + DataFileMeta.EMBEDDED_FILE_INDEX, + DataFileMeta.FILE_SOURCE, + DataFileMeta.VALUE_STATS_COLS, + DataFileMeta.EXTERNAL_PATH, + DataFileMeta.FIRST_ROW_ID, + DataFileMeta.WRITE_COLS); + + public DataFileMetaWriteColsLegacySerializer() { + super(SCHEMA); + } + + @Override + public InternalRow toRow(DataFileMeta meta) { + return GenericRow.of( + BinaryString.fromString(meta.fileName()), + meta.fileSize(), + meta.rowCount(), + serializeBinaryRow(meta.minKey()), + serializeBinaryRow(meta.maxKey()), + meta.keyStats().toRow(), + meta.valueStats().toRow(), + meta.minSequenceNumber(), + meta.maxSequenceNumber(), + meta.schemaId(), + meta.level(), + toStringArrayData(meta.extraFiles()), + meta.creationTime(), + meta.deleteRowCount().orElse(null), + meta.embeddedIndex(), + meta.fileSource().map(FileSource::toByteValue).orElse(null), + toStringArrayData(meta.valueStatsCols()), + meta.externalPath().map(BinaryString::fromString).orElse(null), + meta.firstRowId(), + meta.writeCols() == null ? null : toStringArrayData(meta.writeCols())); + } + + @Override + public DataFileMeta fromRow(InternalRow row) { + return DataFileMeta.create( + row.getString(0).toString(), + row.getLong(1), + row.getLong(2), + deserializeBinaryRow(row.getBinary(3)), + deserializeBinaryRow(row.getBinary(4)), + SimpleStats.fromRow(row.getRow(5, 3)), + SimpleStats.fromRow(row.getRow(6, 3)), + row.getLong(7), + row.getLong(8), + row.getLong(9), + row.getInt(10), + fromStringArrayData(row.getArray(11)), + row.getTimestamp(12, 3), + row.isNullAt(13) ? null : row.getLong(13), + row.isNullAt(14) ? null : row.getBinary(14), + row.isNullAt(15) ? null : FileSource.fromByteValue(row.getByte(15)), + row.isNullAt(16) ? null : fromStringArrayData(row.getArray(16)), + row.isNullAt(17) ? null : row.getString(17).toString(), + row.isNullAt(18) ? null : row.getLong(18), + row.isNullAt(19) ? null : fromStringArrayData(row.getArray(19))); + } +} diff --git a/paimon-core/src/main/java/org/apache/paimon/io/PojoDataFileMeta.java b/paimon-core/src/main/java/org/apache/paimon/io/PojoDataFileMeta.java index 9b288d1d5f7f..76cf448536c3 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 @@ -28,7 +28,9 @@ import java.util.Arrays; import java.util.Collections; +import java.util.LinkedHashMap; import java.util.List; +import java.util.Map; import java.util.Objects; import java.util.Optional; @@ -78,6 +80,8 @@ public class PojoDataFileMeta implements DataFileMeta { private final @Nullable List writeCols; + private final @Nullable Map columnMaxSequenceNumbers; + public PojoDataFileMeta( String fileName, long fileSize, @@ -99,6 +103,52 @@ public PojoDataFileMeta( @Nullable String externalPath, @Nullable Long firstRowId, @Nullable List writeCols) { + this( + fileName, + fileSize, + rowCount, + minKey, + maxKey, + keyStats, + valueStats, + minSequenceNumber, + maxSequenceNumber, + schemaId, + level, + extraFiles, + creationTime, + deleteRowCount, + embeddedIndex, + fileSource, + valueStatsCols, + externalPath, + firstRowId, + writeCols, + null); + } + + public PojoDataFileMeta( + String fileName, + long fileSize, + long rowCount, + BinaryRow minKey, + BinaryRow maxKey, + SimpleStats keyStats, + SimpleStats valueStats, + long minSequenceNumber, + long maxSequenceNumber, + long schemaId, + int level, + List extraFiles, + Timestamp creationTime, + @Nullable Long deleteRowCount, + @Nullable byte[] embeddedIndex, + @Nullable FileSource fileSource, + @Nullable List valueStatsCols, + @Nullable String externalPath, + @Nullable Long firstRowId, + @Nullable List writeCols, + @Nullable Map columnMaxSequenceNumbers) { this.fileName = fileName; this.fileSize = fileSize; @@ -123,6 +173,11 @@ public PojoDataFileMeta( this.externalPath = externalPath; this.firstRowId = firstRowId; this.writeCols = writeCols; + this.columnMaxSequenceNumbers = + columnMaxSequenceNumbers == null + ? null + : Collections.unmodifiableMap( + new LinkedHashMap<>(columnMaxSequenceNumbers)); } @Override @@ -239,6 +294,12 @@ public List writeCols() { return writeCols; } + @Nullable + @Override + public Map columnMaxSequenceNumbers() { + return columnMaxSequenceNumbers; + } + @Override public PojoDataFileMeta upgrade(int newLevel) { checkArgument(newLevel > this.level); @@ -262,7 +323,8 @@ public PojoDataFileMeta upgrade(int newLevel) { valueStatsCols, externalPath, firstRowId, - writeCols); + writeCols, + columnMaxSequenceNumbers); } @Override @@ -288,7 +350,8 @@ public PojoDataFileMeta rename(String newFileName) { valueStatsCols, newExternalPath, firstRowId, - writeCols); + writeCols, + columnMaxSequenceNumbers); } @Override @@ -313,7 +376,8 @@ public PojoDataFileMeta copyWithoutStats() { Collections.emptyList(), externalPath, firstRowId, - writeCols); + writeCols, + columnMaxSequenceNumbers); } @Override @@ -338,7 +402,35 @@ public PojoDataFileMeta assignSequenceNumber(long minSequenceNumber, long maxSeq valueStatsCols, externalPath, firstRowId, - writeCols); + writeCols, + columnMaxSequenceNumbers); + } + + @Override + public PojoDataFileMeta withColumnMaxSequenceNumbers( + Map 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 +455,8 @@ public PojoDataFileMeta assignFirstRowId(long firstRowId) { valueStatsCols, externalPath, firstRowId, - writeCols); + writeCols, + columnMaxSequenceNumbers); } @Override @@ -388,7 +481,8 @@ public PojoDataFileMeta newFirstRowId(@Nullable Long newFirstRowId) { valueStatsCols, externalPath, newFirstRowId, - writeCols); + writeCols, + columnMaxSequenceNumbers); } @Override @@ -413,7 +507,8 @@ public PojoDataFileMeta copy(List newExtraFiles) { valueStatsCols, externalPath, firstRowId, - writeCols); + writeCols, + columnMaxSequenceNumbers); } @Override @@ -438,7 +533,8 @@ public PojoDataFileMeta newExternalPath(String newExternalPath) { valueStatsCols, newExternalPath, firstRowId, - writeCols); + writeCols, + columnMaxSequenceNumbers); } @Override @@ -463,7 +559,8 @@ public PojoDataFileMeta copy(byte[] newEmbeddedIndex) { valueStatsCols, externalPath, firstRowId, - writeCols); + writeCols, + columnMaxSequenceNumbers); } @Override @@ -494,7 +591,8 @@ public boolean equals(Object o) { && Objects.equals(valueStatsCols, that.valueStatsCols()) && Objects.equals(externalPath, that.externalPath().orElse(null)) && Objects.equals(firstRowId, that.firstRowId()) - && Objects.equals(writeCols, that.writeCols()); + && Objects.equals(writeCols, that.writeCols()) + && Objects.equals(columnMaxSequenceNumbers, that.columnMaxSequenceNumbers()); } @Override @@ -519,7 +617,8 @@ public int hashCode() { valueStatsCols, externalPath, firstRowId, - writeCols); + writeCols, + columnMaxSequenceNumbers); } @Override @@ -529,7 +628,8 @@ public String toString() { + "minKey: %s, maxKey: %s, keyStats: %s, valueStats: %s, " + "minSequenceNumber: %d, maxSequenceNumber: %d, " + "schemaId: %d, level: %d, extraFiles: %s, creationTime: %s, " - + "deleteRowCount: %d, fileSource: %s, valueStatsCols: %s, externalPath: %s, firstRowId: %s, writeCols: %s}", + + "deleteRowCount: %d, fileSource: %s, valueStatsCols: %s, externalPath: %s, " + + "firstRowId: %s, writeCols: %s, columnMaxSequenceNumbers: %s}", fileName, fileSize, rowCount, @@ -549,6 +649,7 @@ public String toString() { valueStatsCols, externalPath, firstRowId, - writeCols); + writeCols, + columnMaxSequenceNumbers); } } diff --git a/paimon-core/src/main/java/org/apache/paimon/table/sink/AppendCompactTaskSerializer.java b/paimon-core/src/main/java/org/apache/paimon/table/sink/AppendCompactTaskSerializer.java index ed78ec2f7f53..ec6df446a869 100644 --- a/paimon-core/src/main/java/org/apache/paimon/table/sink/AppendCompactTaskSerializer.java +++ b/paimon-core/src/main/java/org/apache/paimon/table/sink/AppendCompactTaskSerializer.java @@ -37,7 +37,7 @@ /** Serializer for {@link AppendCompactTask}. */ public class AppendCompactTaskSerializer implements VersionedSerializer { - private static final int CURRENT_VERSION = 2; + private static final int CURRENT_VERSION = 3; private final DataFileMetaSerializer dataFileSerializer; diff --git a/paimon-core/src/main/java/org/apache/paimon/table/sink/CommitMessageSerializer.java b/paimon-core/src/main/java/org/apache/paimon/table/sink/CommitMessageSerializer.java index 2593a2a68ac2..8222b07c8c8c 100644 --- a/paimon-core/src/main/java/org/apache/paimon/table/sink/CommitMessageSerializer.java +++ b/paimon-core/src/main/java/org/apache/paimon/table/sink/CommitMessageSerializer.java @@ -34,6 +34,7 @@ import org.apache.paimon.io.DataFileMeta12LegacySerializer; import org.apache.paimon.io.DataFileMetaFirstRowIdLegacySerializer; import org.apache.paimon.io.DataFileMetaSerializer; +import org.apache.paimon.io.DataFileMetaWriteColsLegacySerializer; import org.apache.paimon.io.DataIncrement; import org.apache.paimon.io.DataInputDeserializer; import org.apache.paimon.io.DataInputView; @@ -52,12 +53,13 @@ /** {@link VersionedSerializer} for {@link CommitMessage}. */ public class CommitMessageSerializer implements VersionedSerializer { - public static final int CURRENT_VERSION = 12; + public static final int CURRENT_VERSION = 13; private final DataFileMetaSerializer dataFileSerializer; private final IndexFileMetaSerializer indexEntrySerializer; private DataFileMetaFirstRowIdLegacySerializer dataFileMetaFirstRowIdLegacySerializer; + private DataFileMetaWriteColsLegacySerializer dataFileMetaWriteColsLegacySerializer; private DataFileMeta12LegacySerializer dataFileMeta12LegacySerializer; private DataFileMeta10LegacySerializer dataFileMeta10LegacySerializer; private DataFileMeta09Serializer dataFile09Serializer; @@ -186,8 +188,13 @@ private CommitMessage deserialize(int version, DataInputView view) throws IOExce private IOExceptionSupplier> fileDeserializer( int version, DataInputView view) { - if (version >= 9) { + if (version >= 13) { return () -> dataFileSerializer.deserializeList(view); + } else if (version >= 9) { + if (dataFileMetaWriteColsLegacySerializer == null) { + dataFileMetaWriteColsLegacySerializer = new DataFileMetaWriteColsLegacySerializer(); + } + return () -> dataFileMetaWriteColsLegacySerializer.deserializeList(view); } else if (version == 8) { if (dataFileMetaFirstRowIdLegacySerializer == null) { dataFileMetaFirstRowIdLegacySerializer = diff --git a/paimon-core/src/main/java/org/apache/paimon/table/sink/MultiTableCompactionTaskSerializer.java b/paimon-core/src/main/java/org/apache/paimon/table/sink/MultiTableCompactionTaskSerializer.java index 8478d04ea3b2..d1212db99f60 100644 --- a/paimon-core/src/main/java/org/apache/paimon/table/sink/MultiTableCompactionTaskSerializer.java +++ b/paimon-core/src/main/java/org/apache/paimon/table/sink/MultiTableCompactionTaskSerializer.java @@ -39,7 +39,7 @@ public class MultiTableCompactionTaskSerializer implements VersionedSerializer { - private static final int CURRENT_VERSION = 1; + private static final int CURRENT_VERSION = 2; private final DataFileMetaSerializer dataFileSerializer; diff --git a/paimon-core/src/main/java/org/apache/paimon/table/source/ChainSplit.java b/paimon-core/src/main/java/org/apache/paimon/table/source/ChainSplit.java index bad2dc1c7bf8..39edeb6c233b 100644 --- a/paimon-core/src/main/java/org/apache/paimon/table/source/ChainSplit.java +++ b/paimon-core/src/main/java/org/apache/paimon/table/source/ChainSplit.java @@ -21,10 +21,12 @@ import org.apache.paimon.data.BinaryRow; import org.apache.paimon.io.DataFileMeta; import org.apache.paimon.io.DataFileMetaSerializer; +import org.apache.paimon.io.DataFileMetaWriteColsLegacySerializer; import org.apache.paimon.io.DataInputView; import org.apache.paimon.io.DataInputViewStreamWrapper; import org.apache.paimon.io.DataOutputView; import org.apache.paimon.io.DataOutputViewStreamWrapper; +import org.apache.paimon.utils.ObjectSerializer; import org.apache.paimon.utils.SerializationUtils; import javax.annotation.Nullable; @@ -49,7 +51,7 @@ public class ChainSplit implements Split { private static final long serialVersionUID = 1L; private static final int VERSION_1 = 1; - private static final int VERSION = 2; + private static final int VERSION = 3; private BinaryRow logicalPartition; private List dataFiles; @@ -217,7 +219,10 @@ public static ChainSplit deserialize(DataInputView in) throws IOException { int n = in.readInt(); List dataFiles = new ArrayList<>(n); - DataFileMetaSerializer dataFileSer = new DataFileMetaSerializer(); + ObjectSerializer dataFileSer = + version <= 2 + ? new DataFileMetaWriteColsLegacySerializer() + : new DataFileMetaSerializer(); for (int i = 0; i < n; i++) { dataFiles.add(dataFileSer.deserialize(in)); } diff --git a/paimon-core/src/main/java/org/apache/paimon/table/source/DataSplit.java b/paimon-core/src/main/java/org/apache/paimon/table/source/DataSplit.java index df31763a3001..88bf60f019c8 100644 --- a/paimon-core/src/main/java/org/apache/paimon/table/source/DataSplit.java +++ b/paimon-core/src/main/java/org/apache/paimon/table/source/DataSplit.java @@ -28,6 +28,7 @@ import org.apache.paimon.io.DataFileMeta12LegacySerializer; import org.apache.paimon.io.DataFileMetaFirstRowIdLegacySerializer; import org.apache.paimon.io.DataFileMetaSerializer; +import org.apache.paimon.io.DataFileMetaWriteColsLegacySerializer; import org.apache.paimon.io.DataInputView; import org.apache.paimon.io.DataInputViewStreamWrapper; import org.apache.paimon.io.DataOutputView; @@ -63,7 +64,7 @@ public class DataSplit implements Split { private static final long serialVersionUID = 7L; private static final long MAGIC = -2394839472490812314L; - private static final int VERSION = 8; + private static final int VERSION = 9; private long snapshotId = 0; private BinaryRow partition; @@ -509,6 +510,10 @@ private static FunctionWithIOException getFileMetaS new DataFileMetaFirstRowIdLegacySerializer(); return serializer::deserialize; } else if (version == 8) { + DataFileMetaWriteColsLegacySerializer serializer = + new DataFileMetaWriteColsLegacySerializer(); + return serializer::deserialize; + } else if (version == 9) { DataFileMetaSerializer serializer = new DataFileMetaSerializer(); return serializer::deserialize; } else { diff --git a/paimon-core/src/main/java/org/apache/paimon/table/source/IncrementalSplit.java b/paimon-core/src/main/java/org/apache/paimon/table/source/IncrementalSplit.java index 5b672a4fefd5..2672bc6d0fbe 100644 --- a/paimon-core/src/main/java/org/apache/paimon/table/source/IncrementalSplit.java +++ b/paimon-core/src/main/java/org/apache/paimon/table/source/IncrementalSplit.java @@ -21,10 +21,12 @@ import org.apache.paimon.data.BinaryRow; import org.apache.paimon.io.DataFileMeta; import org.apache.paimon.io.DataFileMetaSerializer; +import org.apache.paimon.io.DataFileMetaWriteColsLegacySerializer; import org.apache.paimon.io.DataInputView; import org.apache.paimon.io.DataInputViewStreamWrapper; import org.apache.paimon.io.DataOutputViewStreamWrapper; import org.apache.paimon.utils.FunctionWithIOException; +import org.apache.paimon.utils.ObjectSerializer; import javax.annotation.Nullable; @@ -44,7 +46,7 @@ public class IncrementalSplit implements Split { private static final long serialVersionUID = 1L; - private static final int VERSION = 1; + private static final int VERSION = 2; private long snapshotId; private BinaryRow partition; @@ -220,7 +222,7 @@ private void readObject(ObjectInputStream objectInputStream) throws IOException, ClassNotFoundException { DataInputViewStreamWrapper in = new DataInputViewStreamWrapper(objectInputStream); int version = in.readInt(); - if (version != VERSION) { + if (version < 1 || version > VERSION) { throw new UnsupportedOperationException("Unsupported version: " + version); } @@ -229,7 +231,10 @@ private void readObject(ObjectInputStream objectInputStream) bucket = in.readInt(); totalBuckets = in.readInt(); - DataFileMetaSerializer dataFileMetaSerializer = new DataFileMetaSerializer(); + ObjectSerializer dataFileMetaSerializer = + version == 1 + ? new DataFileMetaWriteColsLegacySerializer() + : new DataFileMetaSerializer(); FunctionWithIOException deletionFileSerializer = DeletionFile::deserialize; diff --git a/paimon-core/src/main/java/org/apache/paimon/table/source/SplitSerializer.java b/paimon-core/src/main/java/org/apache/paimon/table/source/SplitSerializer.java index 070956b811d3..fabfb5604e1d 100644 --- a/paimon-core/src/main/java/org/apache/paimon/table/source/SplitSerializer.java +++ b/paimon-core/src/main/java/org/apache/paimon/table/source/SplitSerializer.java @@ -23,12 +23,14 @@ import org.apache.paimon.globalindex.IndexedSplit; import org.apache.paimon.io.DataFileMeta; import org.apache.paimon.io.DataFileMetaSerializer; +import org.apache.paimon.io.DataFileMetaWriteColsLegacySerializer; import org.apache.paimon.io.DataInputDeserializer; import org.apache.paimon.io.DataInputView; import org.apache.paimon.io.DataOutputView; import org.apache.paimon.io.DataOutputViewStreamWrapper; import org.apache.paimon.table.FallbackReadFileStoreTable; import org.apache.paimon.utils.FunctionWithIOException; +import org.apache.paimon.utils.ObjectSerializer; import javax.annotation.Nullable; @@ -54,7 +56,7 @@ public class SplitSerializer { private static final long MAGIC = 0x53504C49545F5631L; // "SPLIT_V1" - private static final int VERSION = 1; + private static final int VERSION = 2; private static final int DATA_SPLIT = 1; private static final int INCREMENTAL_SPLIT = 2; @@ -113,7 +115,7 @@ public static Split deserialize(DataInputView in) throws IOException { } int version = in.readInt(); - if (version != VERSION) { + if (version < 1 || version > VERSION) { throw new IOException("Unsupported split serializer version: " + version); } @@ -122,7 +124,7 @@ public static Split deserialize(DataInputView in) throws IOException { case DATA_SPLIT: return DataSplit.deserialize(in); case INCREMENTAL_SPLIT: - return readIncrementalSplit(in); + return readIncrementalSplit(in, version); case INDEXED_SPLIT: return IndexedSplit.deserialize(in); case CHAIN_SPLIT: @@ -151,17 +153,18 @@ private static void writeIncrementalSplit(IncrementalSplit split, DataOutputView out.writeBoolean(split.isStreaming()); } - private static IncrementalSplit readIncrementalSplit(DataInputView in) throws IOException { + private static IncrementalSplit readIncrementalSplit(DataInputView in, int version) + throws IOException { long snapshotId = in.readLong(); BinaryRow partition = deserializeBinaryRow(in); int bucket = in.readInt(); int totalBuckets = in.readInt(); - List beforeFiles = readDataFiles(in); + List beforeFiles = readDataFiles(in, version); FunctionWithIOException deletionFileSerializer = DeletionFile::deserialize; List beforeDeletionFiles = DeletionFile.deserializeList(in, deletionFileSerializer); - List afterFiles = readDataFiles(in); + List afterFiles = readDataFiles(in, version); List afterDeletionFiles = DeletionFile.deserializeList(in, deletionFileSerializer); boolean isStreaming = in.readBoolean(); @@ -232,10 +235,14 @@ private static void writeDataFiles(List files, DataOutputView out) } } - private static List readDataFiles(DataInputView in) throws IOException { + private static List readDataFiles(DataInputView in, int version) + throws IOException { int size = in.readInt(); List files = new ArrayList<>(size); - DataFileMetaSerializer serializer = new DataFileMetaSerializer(); + ObjectSerializer serializer = + version == 1 + ? new DataFileMetaWriteColsLegacySerializer() + : new DataFileMetaSerializer(); for (int i = 0; i < size; i++) { files.add(serializer.deserialize(in)); } diff --git a/paimon-core/src/main/java/org/apache/paimon/utils/DataEvolutionUtils.java b/paimon-core/src/main/java/org/apache/paimon/utils/DataEvolutionUtils.java index c1294f39462f..dbb6c328a2f3 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 @@ -26,6 +26,7 @@ import java.util.Comparator; import java.util.HashSet; import java.util.List; +import java.util.Map; import java.util.Set; import java.util.function.Function; import java.util.stream.Collectors; @@ -56,6 +57,15 @@ public static Set fileFieldIds( return ids; } + /** Returns the latest sequence known for a field physically present in the file. */ + public static long fieldMaxSequenceNumber(DataFileMeta file, int fieldId) { + Map columnSequences = file.columnMaxSequenceNumbers(); + if (columnSequences == null) { + return file.maxSequenceNumber(); + } + return columnSequences.getOrDefault(fieldId, file.maxSequenceNumber()); + } + /** * 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/globalindex/DataEvolutionGlobalIndexRefreshPlannerTest.java b/paimon-core/src/test/java/org/apache/paimon/globalindex/DataEvolutionGlobalIndexRefreshPlannerTest.java index 68d1a26218de..17423ae36a99 100644 --- a/paimon-core/src/test/java/org/apache/paimon/globalindex/DataEvolutionGlobalIndexRefreshPlannerTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/globalindex/DataEvolutionGlobalIndexRefreshPlannerTest.java @@ -152,6 +152,38 @@ void testUsesStableFieldIdAcrossRenameAndFullWrites() { .containsExactly(index); } + @Test + void testUsesColumnSequenceNumbersForCompactedFullFile() { + IndexManifestEntry index = index("index", 0, 99, 5L, BinaryRow.EMPTY_ROW, 0); + + assertThat( + plan( + Collections.singletonList( + dataWithColumnSequences( + "unrelated-compact", + 0, + 100, + 10, + Collections.singletonMap(VECTOR_FIELD.id(), 5L))), + index)) + .isEmpty(); + assertThat( + plan( + Collections.singletonList( + dataWithColumnSequences( + "index-compact", + 0, + 100, + 10, + Collections.singletonMap(VECTOR_FIELD.id(), 6L))), + index)) + .containsExactly(index); + + // Legacy compacted files have no column metadata and remain conservative. + assertThat(plan(Collections.singletonList(data("legacy", 0, 100, 10, 1)), index)) + .containsExactly(index); + } + @Test void testRefreshesFromUpdateLayerOverBaseSchemaWithoutIndexColumn() { IndexManifestEntry index = index("index", 0, 99, 5L, BinaryRow.EMPTY_ROW, 0); @@ -491,4 +523,17 @@ private ManifestEntry dataWithWriteCols( writeCols); return ManifestEntry.create(FileKind.ADD, BinaryRow.EMPTY_ROW, 0, 1, file); } + + private ManifestEntry dataWithColumnSequences( + String fileName, + long firstRowId, + long rowCount, + long maxSequenceNumber, + java.util.Map columnSequences) { + DataFileMeta file = + data(fileName, firstRowId, rowCount, maxSequenceNumber, 1) + .file() + .withColumnMaxSequenceNumbers(columnSequences); + return ManifestEntry.create(FileKind.ADD, BinaryRow.EMPTY_ROW, 0, 1, file); + } } diff --git a/paimon-core/src/test/java/org/apache/paimon/globalindex/sorted/SortedGlobalIndexScannerTest.java b/paimon-core/src/test/java/org/apache/paimon/globalindex/sorted/SortedGlobalIndexScannerTest.java index f918b1728981..8d7b1de0233e 100644 --- a/paimon-core/src/test/java/org/apache/paimon/globalindex/sorted/SortedGlobalIndexScannerTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/globalindex/sorted/SortedGlobalIndexScannerTest.java @@ -20,6 +20,9 @@ import org.apache.paimon.CoreOptions; import org.apache.paimon.Snapshot; +import org.apache.paimon.append.dataevolution.DataEvolutionCompactCoordinator; +import org.apache.paimon.append.dataevolution.DataEvolutionCompactTask; +import org.apache.paimon.append.dataevolution.DataEvolutionCompactionCommitPreparation; import org.apache.paimon.data.BinaryRow; import org.apache.paimon.data.BinaryString; import org.apache.paimon.data.BlobData; @@ -35,6 +38,8 @@ import org.apache.paimon.manifest.IndexManifestEntry; import org.apache.paimon.manifest.ManifestEntry; import org.apache.paimon.memory.MemorySlice; +import org.apache.paimon.options.ExpireConfig; +import org.apache.paimon.options.Options; import org.apache.paimon.partition.PartitionPredicate; import org.apache.paimon.predicate.Predicate; import org.apache.paimon.schema.Schema; @@ -44,6 +49,7 @@ import org.apache.paimon.table.sink.BatchTableWrite; import org.apache.paimon.table.sink.BatchWriteBuilder; import org.apache.paimon.table.sink.CommitMessage; +import org.apache.paimon.table.sink.CommitMessageImpl; import org.apache.paimon.table.source.DataSplit; import org.apache.paimon.types.DataTypes; import org.apache.paimon.types.RowType; @@ -52,6 +58,7 @@ import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.Test; +import java.time.Duration; import java.util.ArrayList; import java.util.Collections; import java.util.Comparator; @@ -59,6 +66,7 @@ import java.util.List; import java.util.Map; import java.util.Optional; +import java.util.stream.Collectors; import static org.assertj.core.api.Assertions.assertThat; @@ -277,6 +285,123 @@ public void testIncrementalScanWithNewData() throws Exception { 500, totalRowCount, "incrementalScan should only return the newly written rows"); } + @Test + public void testIncrementalScanIgnoresNonIndexColumnCompactionAfterSnapshotExpiration() + throws Exception { + write(); + createIndex(null); + + DataFileMeta firstCompact = updateColumnAndCompact("f1", 1); + int f0Id = getTableDefault().rowType().getField("f0").id(); + long f0Sequence = firstCompact.columnMaxSequenceNumbers().get(f0Id); + assertThat(f0Sequence).isLessThan(firstCompact.maxSequenceNumber()); + + DataFileMeta secondCompact = updateColumnAndCompact("f1", 2); + assertThat(secondCompact.columnMaxSequenceNumbers()).containsEntry(f0Id, f0Sequence); + + FileStoreTable table = getTableDefault(); + table.newExpireSnapshots() + .config( + ExpireConfig.builder() + .snapshotRetainMax(1) + .snapshotRetainMin(1) + .snapshotTimeRetain(Duration.ZERO) + .build()) + .expire(); + assertThat(table.snapshotManager().earliestSnapshotId()) + .isEqualTo(table.snapshotManager().latestSnapshotId()); + + assertThat(dataEvolutionScanner(table).withIndexField("f0").incrementalScan()).isEmpty(); + } + + @Test + public void testIncrementalScanRefreshesIndexColumnCompaction() throws Exception { + write(); + createIndex(null); + DataFileMeta compacted = updateColumnAndCompact("f0", 1); + + int f0Id = getTableDefault().rowType().getField("f0").id(); + assertThat(compacted.columnMaxSequenceNumbers().get(f0Id)) + .isEqualTo(compacted.maxSequenceNumber()); + + Optional> scanResult = + dataEvolutionScanner(getTableDefault()).withIndexField("f0").incrementalScan(); + assertThat(scanResult).isPresent(); + assertThat(scanResult.get().deletedIndexEntries()).isNotEmpty(); + } + + private SortedGlobalIndexScanner dataEvolutionScanner(FileStoreTable table) { + Options options = new Options(); + options.set( + CoreOptions.GLOBAL_INDEX_COLUMN_UPDATE_ACTION, + CoreOptions.GlobalIndexColumnUpdateAction.IGNORE); + return new SortedGlobalIndexScanner(table, "btree", options); + } + + private DataFileMeta updateColumnAndCompact(String column, int updateRound) throws Exception { + Map writeOptions = new HashMap<>(); + writeOptions.put( + CoreOptions.GLOBAL_INDEX_COLUMN_UPDATE_ACTION.key(), + CoreOptions.GlobalIndexColumnUpdateAction.IGNORE.toString()); + FileStoreTable table = getTableDefault().copy(writeOptions); + RowType writeType = table.rowType().project("dt", column); + BatchWriteBuilder writeBuilder = table.newBatchWriteBuilder(); + try (BatchTableWrite batchWrite = writeBuilder.newWrite().withWriteType(writeType)) { + for (int i = 0; i < PART_ROW_NUM; i++) { + Object value = + "f0".equals(column) + ? i + updateRound * (int) PART_ROW_NUM + : BinaryString.fromString("updated_" + updateRound + "_" + i); + batchWrite.write(GenericRow.of(BinaryString.fromString("p0"), value)); + } + List messages = batchWrite.prepareCommit(); + assignFirstRowId(messages, 0L); + try (BatchTableCommit commit = writeBuilder.newCommit()) { + commit.commit(messages); + } + } + + writeOptions.put(CoreOptions.COMPACTION_MIN_FILE_NUM.key(), "2"); + table = getTableDefault().copy(writeOptions); + Snapshot compactSnapshot = table.snapshotManager().latestSnapshot(); + DataEvolutionCompactCoordinator coordinator = + new DataEvolutionCompactCoordinator(table, false, false, compactSnapshot); + List compactMessages = new ArrayList<>(); + for (DataEvolutionCompactTask task : coordinator.plan()) { + compactMessages.add(task.doCompact(table, "test-compact")); + } + assertThat(compactMessages).isNotEmpty(); + compactMessages.addAll( + new DataEvolutionCompactionCommitPreparation(table, compactSnapshot) + .prepare(compactMessages)); + try (BatchTableCommit commit = table.newBatchWriteBuilder().newCommit()) { + commit.commit(compactMessages); + } + + List rowRangeFiles = + getTableDefault().store().newScan().plan().files().stream() + .map(ManifestEntry::file) + .filter(file -> file.firstRowId() != null && file.firstRowId() == 0L) + .collect(Collectors.toList()); + assertThat(rowRangeFiles).hasSize(1); + assertThat(rowRangeFiles.get(0).columnMaxSequenceNumbers()).isNotNull(); + return rowRangeFiles.get(0); + } + + private void assignFirstRowId(List messages, long firstRowId) { + for (CommitMessage message : messages) { + CommitMessageImpl impl = (CommitMessageImpl) message; + List files = new ArrayList<>(impl.newFilesIncrement().newFiles()); + impl.newFilesIncrement().newFiles().clear(); + impl.newFilesIncrement() + .newFiles() + .addAll( + files.stream() + .map(file -> file.assignFirstRowId(firstRowId)) + .collect(Collectors.toList())); + } + } + @Test public void testIncrementalScanWithPartitionPredicate() throws Exception { write(); diff --git a/paimon-core/src/test/java/org/apache/paimon/io/BinaryDataFileMetaTest.java b/paimon-core/src/test/java/org/apache/paimon/io/BinaryDataFileMetaTest.java index df7c7a8f17ef..320e0d924cd1 100644 --- a/paimon-core/src/test/java/org/apache/paimon/io/BinaryDataFileMetaTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/io/BinaryDataFileMetaTest.java @@ -42,20 +42,21 @@ public class BinaryDataFileMetaTest { void testImplementsProjectedDataFileMeta() { DataFileMeta expected = DataFileMeta.forAppend( - "data.parquet", - 123L, - 5L, - SimpleStats.EMPTY_STATS, - 2L, - 3L, - 4L, - Arrays.asList("extra-1", "extra-2"), - new byte[] {1, 2}, - FileSource.COMPACT, - Collections.singletonList("value_col"), - "external/dir/data.parquet", - 10L, - Collections.singletonList("write_col")); + "data.parquet", + 123L, + 5L, + SimpleStats.EMPTY_STATS, + 2L, + 3L, + 4L, + Arrays.asList("extra-1", "extra-2"), + new byte[] {1, 2}, + FileSource.COMPACT, + Collections.singletonList("value_col"), + "external/dir/data.parquet", + 10L, + Collections.singletonList("write_col")) + .withColumnMaxSequenceNumbers(Collections.singletonMap(7, 11L)); BinaryDataFileMeta actual = BinaryDataFileMeta.Projection.create(DataFileMeta.SCHEMA) .createDataFile() @@ -91,6 +92,7 @@ void testImplementsProjectedDataFileMeta() { assertThat(actual.firstRowId()).isEqualTo(10L); assertThat(actual.nonNullFirstRowId()).isEqualTo(10L); assertThat(actual.writeCols()).containsExactly("write_col"); + assertThat(actual.columnMaxSequenceNumbers()).containsEntry(7, 11L); assertThat(actual.containsWriteColumn(BinaryString.fromString("write_col"))).isTrue(); assertThat(actual.containsWriteColumn(BinaryString.fromString("other"))).isFalse(); assertThat(actual.toFileSelection(Collections.singletonList(new Range(11L, 12L)))) @@ -101,6 +103,9 @@ void testImplementsProjectedDataFileMeta() { assertUnsupported(actual::copyWithoutStats, "copyWithoutStats()"); assertUnsupported( () -> actual.assignSequenceNumber(4L, 5L), "assignSequenceNumber(long, long)"); + assertUnsupported( + () -> actual.withColumnMaxSequenceNumbers(Collections.singletonMap(1, 2L)), + "withColumnMaxSequenceNumbers(Map)"); assertUnsupported(() -> actual.assignFirstRowId(20L), "assignFirstRowId(long)"); assertUnsupported(() -> actual.newFirstRowId(20L), "newFirstRowId(Long)"); assertUnsupported(() -> actual.copy(Collections.emptyList()), "copy(List)"); diff --git a/paimon-core/src/test/java/org/apache/paimon/io/DataFileMetaSerializerTest.java b/paimon-core/src/test/java/org/apache/paimon/io/DataFileMetaSerializerTest.java index 5074ccf1cf7d..37be1ffcde50 100644 --- a/paimon-core/src/test/java/org/apache/paimon/io/DataFileMetaSerializerTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/io/DataFileMetaSerializerTest.java @@ -20,7 +20,12 @@ import org.apache.paimon.utils.ObjectSerializerTestBase; +import org.junit.jupiter.api.Test; + import java.util.Arrays; +import java.util.Collections; + +import static org.assertj.core.api.Assertions.assertThat; /** Tests for {@link DataFileMetaSerializer}. */ public class DataFileMetaSerializerTest extends ObjectSerializerTestBase { @@ -34,6 +39,34 @@ protected DataFileMetaSerializer serializer() { @Override protected DataFileMeta object() { - return gen.next().meta.copy(Arrays.asList("extra1", "extra2")); + return gen.next() + .meta + .copy(Arrays.asList("extra1", "extra2")) + .withColumnMaxSequenceNumbers(Collections.singletonMap(3, 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(); + } + + private void assertColumnSequences(DataFileMeta file) { + assertThat(file.columnMaxSequenceNumbers()).containsEntry(3, 42L); } } diff --git a/paimon-core/src/test/java/org/apache/paimon/table/sink/CommitMessageSerializerTest.java b/paimon-core/src/test/java/org/apache/paimon/table/sink/CommitMessageSerializerTest.java index c4f519e84e52..bddcea35bba5 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 @@ -25,6 +25,7 @@ import java.io.IOException; import java.util.Arrays; +import java.util.Collections; import static org.apache.paimon.index.IndexFileMetaSerializerTest.randomIndexFile; import static org.apache.paimon.manifest.ManifestCommittableSerializerTest.randomCompactIncrement; @@ -40,6 +41,14 @@ public void test() throws IOException { CommitMessageSerializer serializer = new CommitMessageSerializer(); DataIncrement dataIncrement = randomNewFilesIncrement(); + dataIncrement + .newFiles() + .set( + 0, + dataIncrement + .newFiles() + .get(0) + .withColumnMaxSequenceNumbers(Collections.singletonMap(3, 42L))); dataIncrement.newIndexFiles().addAll(Arrays.asList(randomIndexFile(), randomIndexFile())); dataIncrement .deletedIndexFiles() diff --git a/paimon-core/src/test/java/org/apache/paimon/table/source/DataSplitCompatibleTest.java b/paimon-core/src/test/java/org/apache/paimon/table/source/DataSplitCompatibleTest.java index a0c76537ab10..66a3c86d9e84 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 @@ -218,6 +218,7 @@ public void testSerializer() throws IOException { DataFileTestDataGenerator gen = DataFileTestDataGenerator.builder().build(); DataFileTestDataGenerator.Data data = gen.next(); List files = new ArrayList<>(); + files.add(gen.next().meta.withColumnMaxSequenceNumbers(Collections.singletonMap(3, 42L))); for (int i = 0; i < ThreadLocalRandom.current().nextInt(10); i++) { files.add(gen.next().meta); } diff --git a/paimon-core/src/test/java/org/apache/paimon/table/source/SplitSerializerTest.java b/paimon-core/src/test/java/org/apache/paimon/table/source/SplitSerializerTest.java index 0fd6f982b527..0a5f7f38a275 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 @@ -28,7 +28,6 @@ import org.apache.paimon.manifest.FileSource; import org.apache.paimon.stats.SimpleStats; import org.apache.paimon.table.FallbackReadFileStoreTable; -import org.apache.paimon.utils.CompatibilityUtils; import org.apache.paimon.utils.IOUtils; import org.apache.paimon.utils.InstantiationUtil; import org.apache.paimon.utils.Range; @@ -51,7 +50,6 @@ /** Test for {@link SplitSerializer}. */ public class SplitSerializerTest { - private static final String GENERATE_GOLDEN_FILES_PROPERTY = "generateSplitGoldenFiles"; private static final String RESOURCE_PREFIX = "compatibility/"; @Test @@ -63,18 +61,11 @@ public void testRoundTrip() throws IOException { } @Test - public void testGoldenFiles() throws IOException { - boolean generateGoldenFiles = - Boolean.parseBoolean( - System.getProperties().getProperty(GENERATE_GOLDEN_FILES_PROPERTY)); + public void testVersion1GoldenFiles() throws IOException { for (GoldenCase goldenCase : goldenCases()) { - byte[] actual = SplitSerializer.serialize(goldenCase.split); - if (generateGoldenFiles) { - CompatibilityUtils.writeCompatibilityFile(goldenCase.fileName, actual); - } else { - assertThat(actual).isEqualTo(readGoldenFile(goldenCase.fileName)); - } - assertSplitEquals(goldenCase.split, SplitSerializer.deserialize(actual)); + assertSplitEquals( + goldenCase.split, + SplitSerializer.deserialize(readGoldenFile(goldenCase.fileName))); } } diff --git a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/compact/changelog/ChangelogCompactTaskSerializer.java b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/compact/changelog/ChangelogCompactTaskSerializer.java index c7f56dc5de15..e796b92054ba 100644 --- a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/compact/changelog/ChangelogCompactTaskSerializer.java +++ b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/compact/changelog/ChangelogCompactTaskSerializer.java @@ -40,7 +40,7 @@ public class ChangelogCompactTaskSerializer implements SimpleVersionedSerializer { - private static final int CURRENT_VERSION = 2; + private static final int CURRENT_VERSION = 3; private final DataFileMetaSerializer dataFileSerializer; diff --git a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/compact/changelog/ChangelogCompactTaskSerializerTest.java b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/compact/changelog/ChangelogCompactTaskSerializerTest.java index df62115f2750..6a02c91ec510 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 @@ -26,6 +26,7 @@ import org.junit.jupiter.api.Test; import java.util.ArrayList; +import java.util.Collections; import java.util.HashMap; import java.util.List; import java.util.UUID; @@ -87,22 +88,23 @@ private List newFiles(int num) { private DataFileMeta newFile() { return DataFileMeta.create( - UUID.randomUUID().toString(), - 0, - 1, - row(0), - row(0), - newSimpleStats(0, 1), - newSimpleStats(0, 1), - 0, - 1, - 0, - 0, - 0L, - null, - FileSource.APPEND, - null, - null, - null); + UUID.randomUUID().toString(), + 0, + 1, + row(0), + row(0), + newSimpleStats(0, 1), + newSimpleStats(0, 1), + 0, + 1, + 0, + 0, + 0L, + null, + FileSource.APPEND, + null, + null, + null) + .withColumnMaxSequenceNumbers(Collections.singletonMap(1, 1L)); } } diff --git a/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/copy/CopyFilesUtil.java b/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/copy/CopyFilesUtil.java index 1a1e12efbe33..7cc896c0b048 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,8 @@ public static DataFileMeta toNewDataFileMeta( oldFileMeta.valueStatsCols(), newExternalPath, oldFileMeta.firstRowId(), - oldFileMeta.writeCols()); + oldFileMeta.writeCols(), + oldFileMeta.columnMaxSequenceNumbers()); } public static IndexFileMeta toNewIndexFileMeta(IndexFileMeta oldFileMeta, String newFileName) { From 8e343ab01f1b379156f2453339d6de04b00970ef Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E4=BB=9F=E5=BC=8B?= Date: Thu, 6 Aug 2026 17:40:39 +0800 Subject: [PATCH 2/4] [core] Use positional arrays for column sequences --- .../DataEvolutionNormalCompactTask.java | 58 ++++++++++--------- ...ataEvolutionGlobalIndexRefreshPlanner.java | 35 +++++------ .../apache/paimon/io/BinaryDataFileMeta.java | 21 +++---- .../org/apache/paimon/io/DataFileMeta.java | 19 +++--- .../paimon/io/DataFileMetaSerializer.java | 20 +++---- .../apache/paimon/io/PojoDataFileMeta.java | 22 +++---- .../paimon/utils/DataEvolutionUtils.java | 40 +++++++++---- ...volutionGlobalIndexRefreshPlannerTest.java | 47 +++++++++++++-- .../sorted/SortedGlobalIndexScannerTest.java | 21 +++++-- .../paimon/io/BinaryDataFileMetaTest.java | 8 +-- .../paimon/io/DataFileMetaSerializerTest.java | 4 +- .../sink/CommitMessageSerializerTest.java | 3 +- .../table/source/DataSplitCompatibleTest.java | 2 +- .../paimon/utils/DataEvolutionUtilsTest.java | 42 ++++++++++++++ .../ChangelogCompactTaskSerializerTest.java | 3 +- 15 files changed, 219 insertions(+), 126 deletions(-) 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 916a99ce03b2..75b999771ce0 100644 --- a/paimon-core/src/main/java/org/apache/paimon/append/dataevolution/DataEvolutionNormalCompactTask.java +++ b/paimon-core/src/main/java/org/apache/paimon/append/dataevolution/DataEvolutionNormalCompactTask.java @@ -28,6 +28,7 @@ import org.apache.paimon.table.FileStoreTable; import org.apache.paimon.table.sink.CommitMessage; import org.apache.paimon.table.source.DataSplit; +import org.apache.paimon.types.DataField; import org.apache.paimon.types.RowType; import org.apache.paimon.utils.FileStorePathFactory; import org.apache.paimon.utils.RecordWriter; @@ -46,7 +47,7 @@ import static org.apache.paimon.types.VectorType.fieldNamesInVectorFile; import static org.apache.paimon.types.VectorType.isVectorStoreFile; import static org.apache.paimon.utils.DataEvolutionUtils.fieldMaxSequenceNumber; -import static org.apache.paimon.utils.DataEvolutionUtils.fileFieldIds; +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. */ @@ -85,8 +86,6 @@ public CommitMessage doCompact(FileStoreTable table, String commitUser) throws E table.rowType().getFields().stream() .filter(f -> !fieldsInDedicatedFile.contains(f.name())) .collect(Collectors.toList())); - Map columnMaxSequenceNumbers = - compactedColumnMaxSequenceNumbers(table, readWriteType); FileStorePathFactory pathFactory = table.store().pathFactory(); AppendOnlyFileStore store = (AppendOnlyFileStore) table.store(); @@ -128,40 +127,43 @@ public CommitMessage doCompact(FileStoreTable table, String commitUser) throws E dataFileMeta = dataFileMeta.assignSequenceNumber( minSequenceId(compactBefore), maxSequenceId(compactBefore)); - dataFileMeta = dataFileMeta.withColumnMaxSequenceNumbers(columnMaxSequenceNumbers); + dataFileMeta = + dataFileMeta.withColumnMaxSequenceNumbers( + compactedColumnMaxSequenceNumbers(table, dataFileMeta)); compactAfter.add(dataFileMeta); return commitMessage(compactBefore, compactAfter); } - private Map compactedColumnMaxSequenceNumbers( - FileStoreTable table, RowType outputType) { - Map> inputFieldIds = new LinkedHashMap<>(); + private long[] compactedColumnMaxSequenceNumbers( + FileStoreTable table, DataFileMeta outputFile) { + Map> inputFields = new LinkedHashMap<>(); for (DataFileMeta input : compactBefore) { - inputFieldIds.put(input, fileFieldIds(table.schemaManager()::schema, input)); + inputFields.put(input, fileFields(table.schemaManager()::schema, input)); } long fallbackSequence = maxSequenceId(compactBefore); - Map result = new LinkedHashMap<>(); - outputType - .getFields() - .forEach( - field -> { - long fieldSequence = Long.MIN_VALUE; - for (DataFileMeta input : compactBefore) { - if (inputFieldIds.get(input).contains(field.id())) { - fieldSequence = - Math.max( - fieldSequence, - fieldMaxSequenceNumber(input, field.id())); - } - } - result.put( - field.id(), - fieldSequence == Long.MIN_VALUE - ? fallbackSequence - : fieldSequence); - }); + List outputFields = fileFields(table.schemaManager()::schema, outputFile); + long[] result = new long[outputFields.size()]; + for (int outputPosition = 0; outputPosition < outputFields.size(); outputPosition++) { + DataField outputField = outputFields.get(outputPosition); + long fieldSequence = Long.MIN_VALUE; + for (DataFileMeta input : compactBefore) { + List fields = inputFields.get(input); + for (int inputPosition = 0; inputPosition < fields.size(); inputPosition++) { + if (fields.get(inputPosition).id() == outputField.id()) { + fieldSequence = + Math.max( + fieldSequence, + fieldMaxSequenceNumber( + input, inputPosition, fields.size())); + break; + } + } + } + result[outputPosition] = + fieldSequence == Long.MIN_VALUE ? fallbackSequence : fieldSequence; + } return result; } } diff --git a/paimon-core/src/main/java/org/apache/paimon/globalindex/DataEvolutionGlobalIndexRefreshPlanner.java b/paimon-core/src/main/java/org/apache/paimon/globalindex/DataEvolutionGlobalIndexRefreshPlanner.java index 201849edcc30..4f0b14f46c34 100644 --- a/paimon-core/src/main/java/org/apache/paimon/globalindex/DataEvolutionGlobalIndexRefreshPlanner.java +++ b/paimon-core/src/main/java/org/apache/paimon/globalindex/DataEvolutionGlobalIndexRefreshPlanner.java @@ -42,7 +42,7 @@ import java.util.TreeMap; import static org.apache.paimon.utils.DataEvolutionUtils.fieldMaxSequenceNumber; -import static org.apache.paimon.utils.DataEvolutionUtils.fileFieldIds; +import static org.apache.paimon.utils.DataEvolutionUtils.fileFields; /** Plans existing global index files which need refresh after data-evolution updates. */ public final class DataEvolutionGlobalIndexRefreshPlanner { @@ -82,7 +82,7 @@ public static List findIndexesToRefresh( .addIndex(i, indexMeta.rowRange(), scanSnapshotId); } - Map>, Set> fileFieldIdsCache = new HashMap<>(); + Map>, List> fileFieldsCache = new HashMap<>(); for (ManifestEntry dataEntry : dataEntries) { DataFileMeta file = dataEntry.file(); if (dataEntry.kind() != FileKind.ADD || file.firstRowId() == null) { @@ -94,18 +94,20 @@ public static List findIndexesToRefresh( continue; } - Set physicalFieldIds = - fileFieldIdsCache.computeIfAbsent( + List physicalFields = + fileFieldsCache.computeIfAbsent( Pair.of(file.schemaId(), file.writeCols()), - key -> fileFieldIds(schemaManager::schema, file)); - if (!disjoint(indexedFieldIds, physicalFieldIds)) { - long indexedMaxSequence = Long.MIN_VALUE; - for (Integer fieldId : indexedFieldIds) { - if (physicalFieldIds.contains(fieldId)) { - indexedMaxSequence = - Math.max(indexedMaxSequence, fieldMaxSequenceNumber(file, fieldId)); - } + key -> fileFields(schemaManager::schema, file)); + long indexedMaxSequence = Long.MIN_VALUE; + for (int position = 0; position < physicalFields.size(); position++) { + if (indexedFieldIds.contains(physicalFields.get(position).id())) { + indexedMaxSequence = + Math.max( + indexedMaxSequence, + fieldMaxSequenceNumber(file, position, physicalFields.size())); } + } + if (indexedMaxSequence != Long.MIN_VALUE) { group.addDataFile(file, indexedMaxSequence); } } @@ -243,13 +245,4 @@ private static boolean matchesFields(GlobalIndexMeta meta, List field } return expectedExtraFields != null && Arrays.equals(actualExtraFields, expectedExtraFields); } - - private static boolean disjoint(Set left, Set right) { - for (Integer value : left) { - if (right.contains(value)) { - return false; - } - } - return true; - } } diff --git a/paimon-core/src/main/java/org/apache/paimon/io/BinaryDataFileMeta.java b/paimon-core/src/main/java/org/apache/paimon/io/BinaryDataFileMeta.java index 0b49cc71a245..f48653571e36 100644 --- a/paimon-core/src/main/java/org/apache/paimon/io/BinaryDataFileMeta.java +++ b/paimon-core/src/main/java/org/apache/paimon/io/BinaryDataFileMeta.java @@ -21,7 +21,6 @@ import org.apache.paimon.data.BinaryRow; import org.apache.paimon.data.BinaryString; import org.apache.paimon.data.InternalArray; -import org.apache.paimon.data.InternalMap; import org.apache.paimon.data.InternalRow; import org.apache.paimon.data.Timestamp; import org.apache.paimon.fs.Path; @@ -34,9 +33,7 @@ import java.util.Arrays; import java.util.Collections; -import java.util.LinkedHashMap; import java.util.List; -import java.util.Map; import java.util.Optional; import static org.apache.paimon.utils.InternalRowUtils.fromStringArrayData; @@ -249,20 +246,18 @@ public List writeCols() { @Nullable @Override - public Map columnMaxSequenceNumbers() { + public long[] columnMaxSequenceNumbers() { int position = requiredPosition(Fields.COLUMN_MAX_SEQUENCE_NUMBERS); InternalRow row = currentRow(); if (row.isNullAt(position)) { return null; } - InternalMap map = row.getMap(position); - InternalArray keys = map.keyArray(); - InternalArray values = map.valueArray(); - Map result = new LinkedHashMap<>(); - for (int i = 0; i < map.size(); i++) { - result.put(keys.getInt(i), values.getLong(i)); + 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 Collections.unmodifiableMap(result); + return result; } public boolean containsWriteColumn(BinaryString fieldName) { @@ -303,8 +298,8 @@ public DataFileMeta assignSequenceNumber(long minSequenceNumber, long maxSequenc } @Override - public DataFileMeta withColumnMaxSequenceNumbers(Map columnMaxSequenceNumbers) { - throw unsupportedOperation("withColumnMaxSequenceNumbers(Map)"); + public DataFileMeta withColumnMaxSequenceNumbers(long[] columnMaxSequenceNumbers) { + throw unsupportedOperation("withColumnMaxSequenceNumbers(long[])"); } @Override 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 48c5dcbd7a88..aa944ae2edbd 100644 --- a/paimon-core/src/main/java/org/apache/paimon/io/DataFileMeta.java +++ b/paimon-core/src/main/java/org/apache/paimon/io/DataFileMeta.java @@ -29,7 +29,6 @@ import org.apache.paimon.types.DataField; import org.apache.paimon.types.DataTypes; import org.apache.paimon.types.IntType; -import org.apache.paimon.types.MapType; import org.apache.paimon.types.RowType; import org.apache.paimon.types.TinyIntType; import org.apache.paimon.utils.Range; @@ -115,7 +114,7 @@ public interface DataFileMeta { new DataField( 20, COLUMN_MAX_SEQUENCE_NUMBERS, - new MapType(true, new IntType(false), new BigIntType(false))))); + new ArrayType(true, new BigIntType(false))))); BinaryRow EMPTY_MIN_KEY = EMPTY_ROW; BinaryRow EMPTY_MAX_KEY = EMPTY_ROW; @@ -224,7 +223,7 @@ static DataFileMeta create( @Nullable String externalPath, @Nullable Long firstRowId, @Nullable List writeCols, - @Nullable Map columnMaxSequenceNumbers) { + @Nullable long[] columnMaxSequenceNumbers) { return new PojoDataFileMeta( fileName, fileSize, @@ -410,11 +409,14 @@ default Range nonNullRowIdRange() { /** * Maximum sequence number per physical table field after data-evolution compaction. * - *

A null value means that only the file-level sequence range is available. Field ids which - * are absent from a non-null map must also fall back to {@link #maxSequenceNumber()}. + *

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 - Map columnMaxSequenceNumbers(); + default long[] columnMaxSequenceNumbers() { + return null; + } DataFileMeta upgrade(int newLevel); @@ -424,7 +426,10 @@ default Range nonNullRowIdRange() { DataFileMeta assignSequenceNumber(long minSequenceNumber, long maxSequenceNumber); - DataFileMeta withColumnMaxSequenceNumbers(Map columnMaxSequenceNumbers); + default DataFileMeta withColumnMaxSequenceNumbers(long[] columnMaxSequenceNumbers) { + throw new UnsupportedOperationException( + "This DataFileMeta implementation does not support column sequence numbers."); + } DataFileMeta assignFirstRowId(long firstRowId); 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 995741dea16e..6e5c6b2bee27 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,18 +19,14 @@ package org.apache.paimon.io; import org.apache.paimon.data.BinaryString; -import org.apache.paimon.data.GenericMap; +import org.apache.paimon.data.GenericArray; import org.apache.paimon.data.GenericRow; import org.apache.paimon.data.InternalArray; -import org.apache.paimon.data.InternalMap; import org.apache.paimon.data.InternalRow; import org.apache.paimon.manifest.FileSource; import org.apache.paimon.stats.SimpleStats; import org.apache.paimon.utils.ObjectSerializer; -import java.util.LinkedHashMap; -import java.util.Map; - import static org.apache.paimon.utils.InternalRowUtils.fromStringArrayData; import static org.apache.paimon.utils.InternalRowUtils.toStringArrayData; import static org.apache.paimon.utils.SerializationUtils.deserializeBinaryRow; @@ -70,7 +66,7 @@ public InternalRow toRow(DataFileMeta meta) { meta.writeCols() == null ? null : toStringArrayData(meta.writeCols()), meta.columnMaxSequenceNumbers() == null ? null - : new GenericMap(meta.columnMaxSequenceNumbers())); + : new GenericArray(meta.columnMaxSequenceNumbers())); } @Override @@ -96,15 +92,13 @@ public DataFileMeta fromRow(InternalRow row) { 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(20) ? null : fromIntLongMap(row.getMap(20))); + row.isNullAt(20) ? null : fromLongArray(row.getArray(20))); } - private static Map fromIntLongMap(InternalMap map) { - InternalArray keys = map.keyArray(); - InternalArray values = map.valueArray(); - Map result = new LinkedHashMap<>(); - for (int i = 0; i < map.size(); i++) { - result.put(keys.getInt(i), values.getLong(i)); + 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/PojoDataFileMeta.java b/paimon-core/src/main/java/org/apache/paimon/io/PojoDataFileMeta.java index 76cf448536c3..5e219f012374 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 @@ -28,9 +28,7 @@ import java.util.Arrays; import java.util.Collections; -import java.util.LinkedHashMap; import java.util.List; -import java.util.Map; import java.util.Objects; import java.util.Optional; @@ -80,7 +78,7 @@ public class PojoDataFileMeta implements DataFileMeta { private final @Nullable List writeCols; - private final @Nullable Map columnMaxSequenceNumbers; + private final @Nullable long[] columnMaxSequenceNumbers; public PojoDataFileMeta( String fileName, @@ -148,7 +146,7 @@ public PojoDataFileMeta( @Nullable String externalPath, @Nullable Long firstRowId, @Nullable List writeCols, - @Nullable Map columnMaxSequenceNumbers) { + @Nullable long[] columnMaxSequenceNumbers) { this.fileName = fileName; this.fileSize = fileSize; @@ -174,10 +172,7 @@ public PojoDataFileMeta( this.firstRowId = firstRowId; this.writeCols = writeCols; this.columnMaxSequenceNumbers = - columnMaxSequenceNumbers == null - ? null - : Collections.unmodifiableMap( - new LinkedHashMap<>(columnMaxSequenceNumbers)); + columnMaxSequenceNumbers == null ? null : columnMaxSequenceNumbers.clone(); } @Override @@ -296,7 +291,7 @@ public List writeCols() { @Nullable @Override - public Map columnMaxSequenceNumbers() { + public long[] columnMaxSequenceNumbers() { return columnMaxSequenceNumbers; } @@ -407,8 +402,7 @@ public PojoDataFileMeta assignSequenceNumber(long minSequenceNumber, long maxSeq } @Override - public PojoDataFileMeta withColumnMaxSequenceNumbers( - Map columnMaxSequenceNumbers) { + public PojoDataFileMeta withColumnMaxSequenceNumbers(long[] columnMaxSequenceNumbers) { return new PojoDataFileMeta( fileName, fileSize, @@ -592,7 +586,7 @@ public boolean equals(Object o) { && Objects.equals(externalPath, that.externalPath().orElse(null)) && Objects.equals(firstRowId, that.firstRowId()) && Objects.equals(writeCols, that.writeCols()) - && Objects.equals(columnMaxSequenceNumbers, that.columnMaxSequenceNumbers()); + && Arrays.equals(columnMaxSequenceNumbers, that.columnMaxSequenceNumbers()); } @Override @@ -618,7 +612,7 @@ public int hashCode() { externalPath, firstRowId, writeCols, - columnMaxSequenceNumbers); + Arrays.hashCode(columnMaxSequenceNumbers)); } @Override @@ -650,6 +644,6 @@ public String toString() { externalPath, firstRowId, writeCols, - columnMaxSequenceNumbers); + Arrays.toString(columnMaxSequenceNumbers)); } } 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 dbb6c328a2f3..96c25850e411 100644 --- a/paimon-core/src/main/java/org/apache/paimon/utils/DataEvolutionUtils.java +++ b/paimon-core/src/main/java/org/apache/paimon/utils/DataEvolutionUtils.java @@ -22,9 +22,10 @@ import org.apache.paimon.schema.TableSchema; import org.apache.paimon.types.DataField; +import java.util.ArrayList; import java.util.Collection; import java.util.Comparator; -import java.util.HashSet; +import java.util.HashMap; import java.util.List; import java.util.Map; import java.util.Set; @@ -44,26 +45,43 @@ public class DataEvolutionUtils { */ public static Set fileFieldIds( Function scanTableSchema, DataFileMeta file) { + return fileFields(scanTableSchema, file).stream() + .map(DataField::id) + .collect(Collectors.toSet()); + } + + /** Table fields physically present in a file, in their physical write order. */ + public static List fileFields( + Function scanTableSchema, DataFileMeta file) { TableSchema schema = scanTableSchema.apply(file.schemaId()); List writeCols = file.writeCols(); - Set writeColNames = writeCols == null ? null : new HashSet<>(writeCols); - Set ids = new HashSet<>(); + if (writeCols == null) { + return schema.fields(); + } + + Map fieldsByName = new HashMap<>(); for (DataField field : schema.fields()) { + fieldsByName.put(field.name(), field); + } + List fields = new ArrayList<>(); + for (String writeCol : writeCols) { // writeCols may also contain physical row-tracking fields outside the table schema. - if (writeColNames == null || writeColNames.contains(field.name())) { - ids.add(field.id()); + DataField field = fieldsByName.get(writeCol); + if (field != null) { + fields.add(field); } } - return ids; + return fields; } - /** Returns the latest sequence known for a field physically present in the file. */ - public static long fieldMaxSequenceNumber(DataFileMeta file, int fieldId) { - Map columnSequences = file.columnMaxSequenceNumbers(); - if (columnSequences == null) { + /** Returns the latest sequence known for a physical field position in the file. */ + public static long fieldMaxSequenceNumber( + DataFileMeta file, int fieldPosition, int physicalFieldCount) { + long[] columnSequences = file.columnMaxSequenceNumbers(); + if (columnSequences == null || columnSequences.length != physicalFieldCount) { return file.maxSequenceNumber(); } - return columnSequences.getOrDefault(fieldId, file.maxSequenceNumber()); + return columnSequences[fieldPosition]; } /** diff --git a/paimon-core/src/test/java/org/apache/paimon/globalindex/DataEvolutionGlobalIndexRefreshPlannerTest.java b/paimon-core/src/test/java/org/apache/paimon/globalindex/DataEvolutionGlobalIndexRefreshPlannerTest.java index 17423ae36a99..3fa8491ff640 100644 --- a/paimon-core/src/test/java/org/apache/paimon/globalindex/DataEvolutionGlobalIndexRefreshPlannerTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/globalindex/DataEvolutionGlobalIndexRefreshPlannerTest.java @@ -164,7 +164,7 @@ void testUsesColumnSequenceNumbersForCompactedFullFile() { 0, 100, 10, - Collections.singletonMap(VECTOR_FIELD.id(), 5L))), + new long[] {5L, 10L, 10L})), index)) .isEmpty(); assertThat( @@ -175,7 +175,7 @@ void testUsesColumnSequenceNumbersForCompactedFullFile() { 0, 100, 10, - Collections.singletonMap(VECTOR_FIELD.id(), 6L))), + new long[] {6L, 10L, 10L})), index)) .containsExactly(index); @@ -184,6 +184,44 @@ void testUsesColumnSequenceNumbersForCompactedFullFile() { .containsExactly(index); } + @Test + void testColumnSequenceNumbersFollowWriteColsOrder() { + IndexManifestEntry index = index("index", 0, 99, 5L, BinaryRow.EMPTY_ROW, 0); + + assertThat( + plan( + Collections.singletonList( + dataWithColumnSequences( + "reordered-compact", + 0, + 100, + 10, + new long[] {10L, 5L}, + "unrelated", + "vector")), + index)) + .isEmpty(); + } + + @Test + void testMalformedColumnSequenceNumbersFallBackToFileSequence() { + IndexManifestEntry index = index("index", 0, 99, 5L, BinaryRow.EMPTY_ROW, 0); + + assertThat( + plan( + Collections.singletonList( + dataWithColumnSequences( + "malformed-compact", + 0, + 100, + 10, + new long[] {5L}, + "vector", + "other")), + index)) + .containsExactly(index); + } + @Test void testRefreshesFromUpdateLayerOverBaseSchemaWithoutIndexColumn() { IndexManifestEntry index = index("index", 0, 99, 5L, BinaryRow.EMPTY_ROW, 0); @@ -529,9 +567,10 @@ private ManifestEntry dataWithColumnSequences( long firstRowId, long rowCount, long maxSequenceNumber, - java.util.Map columnSequences) { + long[] columnSequences, + String... writeCols) { DataFileMeta file = - data(fileName, firstRowId, rowCount, maxSequenceNumber, 1) + data(fileName, firstRowId, rowCount, maxSequenceNumber, 1, writeCols) .file() .withColumnMaxSequenceNumbers(columnSequences); return ManifestEntry.create(FileKind.ADD, BinaryRow.EMPTY_ROW, 0, 1, file); diff --git a/paimon-core/src/test/java/org/apache/paimon/globalindex/sorted/SortedGlobalIndexScannerTest.java b/paimon-core/src/test/java/org/apache/paimon/globalindex/sorted/SortedGlobalIndexScannerTest.java index 8d7b1de0233e..91373f4b0b50 100644 --- a/paimon-core/src/test/java/org/apache/paimon/globalindex/sorted/SortedGlobalIndexScannerTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/globalindex/sorted/SortedGlobalIndexScannerTest.java @@ -51,6 +51,7 @@ import org.apache.paimon.table.sink.CommitMessage; import org.apache.paimon.table.sink.CommitMessageImpl; import org.apache.paimon.table.source.DataSplit; +import org.apache.paimon.types.DataField; import org.apache.paimon.types.DataTypes; import org.apache.paimon.types.RowType; import org.apache.paimon.utils.Pair; @@ -68,6 +69,7 @@ import java.util.Optional; import java.util.stream.Collectors; +import static org.apache.paimon.utils.DataEvolutionUtils.fileFields; import static org.assertj.core.api.Assertions.assertThat; /** Test class for {@link SortedGlobalIndexScanner}. */ @@ -293,11 +295,11 @@ public void testIncrementalScanIgnoresNonIndexColumnCompactionAfterSnapshotExpir DataFileMeta firstCompact = updateColumnAndCompact("f1", 1); int f0Id = getTableDefault().rowType().getField("f0").id(); - long f0Sequence = firstCompact.columnMaxSequenceNumbers().get(f0Id); + long f0Sequence = columnSequence(firstCompact, f0Id); assertThat(f0Sequence).isLessThan(firstCompact.maxSequenceNumber()); DataFileMeta secondCompact = updateColumnAndCompact("f1", 2); - assertThat(secondCompact.columnMaxSequenceNumbers()).containsEntry(f0Id, f0Sequence); + assertThat(columnSequence(secondCompact, f0Id)).isEqualTo(f0Sequence); FileStoreTable table = getTableDefault(); table.newExpireSnapshots() @@ -321,8 +323,7 @@ public void testIncrementalScanRefreshesIndexColumnCompaction() throws Exception DataFileMeta compacted = updateColumnAndCompact("f0", 1); int f0Id = getTableDefault().rowType().getField("f0").id(); - assertThat(compacted.columnMaxSequenceNumbers().get(f0Id)) - .isEqualTo(compacted.maxSequenceNumber()); + assertThat(columnSequence(compacted, f0Id)).isEqualTo(compacted.maxSequenceNumber()); Optional> scanResult = dataEvolutionScanner(getTableDefault()).withIndexField("f0").incrementalScan(); @@ -330,6 +331,18 @@ public void testIncrementalScanRefreshesIndexColumnCompaction() throws Exception assertThat(scanResult.get().deletedIndexEntries()).isNotEmpty(); } + private long columnSequence(DataFileMeta file, int fieldId) throws Exception { + List fields = fileFields(getTableDefault().schemaManager()::schema, file); + long[] sequences = file.columnMaxSequenceNumbers(); + assertThat(sequences).hasSize(fields.size()); + for (int i = 0; i < fields.size(); i++) { + if (fields.get(i).id() == fieldId) { + return sequences[i]; + } + } + throw new IllegalArgumentException("Field not found in data file: " + fieldId); + } + private SortedGlobalIndexScanner dataEvolutionScanner(FileStoreTable table) { Options options = new Options(); options.set( diff --git a/paimon-core/src/test/java/org/apache/paimon/io/BinaryDataFileMetaTest.java b/paimon-core/src/test/java/org/apache/paimon/io/BinaryDataFileMetaTest.java index 320e0d924cd1..2518fdfe1d79 100644 --- a/paimon-core/src/test/java/org/apache/paimon/io/BinaryDataFileMetaTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/io/BinaryDataFileMetaTest.java @@ -56,7 +56,7 @@ void testImplementsProjectedDataFileMeta() { "external/dir/data.parquet", 10L, Collections.singletonList("write_col")) - .withColumnMaxSequenceNumbers(Collections.singletonMap(7, 11L)); + .withColumnMaxSequenceNumbers(new long[] {11L}); BinaryDataFileMeta actual = BinaryDataFileMeta.Projection.create(DataFileMeta.SCHEMA) .createDataFile() @@ -92,7 +92,7 @@ void testImplementsProjectedDataFileMeta() { assertThat(actual.firstRowId()).isEqualTo(10L); assertThat(actual.nonNullFirstRowId()).isEqualTo(10L); assertThat(actual.writeCols()).containsExactly("write_col"); - assertThat(actual.columnMaxSequenceNumbers()).containsEntry(7, 11L); + 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)))) @@ -104,8 +104,8 @@ void testImplementsProjectedDataFileMeta() { assertUnsupported( () -> actual.assignSequenceNumber(4L, 5L), "assignSequenceNumber(long, long)"); assertUnsupported( - () -> actual.withColumnMaxSequenceNumbers(Collections.singletonMap(1, 2L)), - "withColumnMaxSequenceNumbers(Map)"); + () -> actual.withColumnMaxSequenceNumbers(new long[] {2L}), + "withColumnMaxSequenceNumbers(long[])"); assertUnsupported(() -> actual.assignFirstRowId(20L), "assignFirstRowId(long)"); assertUnsupported(() -> actual.newFirstRowId(20L), "newFirstRowId(Long)"); assertUnsupported(() -> actual.copy(Collections.emptyList()), "copy(List)"); diff --git a/paimon-core/src/test/java/org/apache/paimon/io/DataFileMetaSerializerTest.java b/paimon-core/src/test/java/org/apache/paimon/io/DataFileMetaSerializerTest.java index 37be1ffcde50..3bfc25030031 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 @@ -42,7 +42,7 @@ protected DataFileMeta object() { return gen.next() .meta .copy(Arrays.asList("extra1", "extra2")) - .withColumnMaxSequenceNumbers(Collections.singletonMap(3, 42L)); + .withColumnMaxSequenceNumbers(new long[] {3L, 42L}); } @Test @@ -67,6 +67,6 @@ void testLegacySerializerDropsColumnSequences() { } private void assertColumnSequences(DataFileMeta file) { - assertThat(file.columnMaxSequenceNumbers()).containsEntry(3, 42L); + assertThat(file.columnMaxSequenceNumbers()).containsExactly(3L, 42L); } } 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 bddcea35bba5..bc36deedca9a 100644 --- a/paimon-core/src/test/java/org/apache/paimon/table/sink/CommitMessageSerializerTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/table/sink/CommitMessageSerializerTest.java @@ -25,7 +25,6 @@ import java.io.IOException; import java.util.Arrays; -import java.util.Collections; import static org.apache.paimon.index.IndexFileMetaSerializerTest.randomIndexFile; import static org.apache.paimon.manifest.ManifestCommittableSerializerTest.randomCompactIncrement; @@ -48,7 +47,7 @@ public void test() throws IOException { dataIncrement .newFiles() .get(0) - .withColumnMaxSequenceNumbers(Collections.singletonMap(3, 42L))); + .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 66a3c86d9e84..140b964106d7 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 @@ -218,7 +218,7 @@ public void testSerializer() throws IOException { DataFileTestDataGenerator gen = DataFileTestDataGenerator.builder().build(); DataFileTestDataGenerator.Data data = gen.next(); List files = new ArrayList<>(); - files.add(gen.next().meta.withColumnMaxSequenceNumbers(Collections.singletonMap(3, 42L))); + 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); } diff --git a/paimon-core/src/test/java/org/apache/paimon/utils/DataEvolutionUtilsTest.java b/paimon-core/src/test/java/org/apache/paimon/utils/DataEvolutionUtilsTest.java index 33feb9d850e1..b05907306214 100644 --- a/paimon-core/src/test/java/org/apache/paimon/utils/DataEvolutionUtilsTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/utils/DataEvolutionUtilsTest.java @@ -109,6 +109,48 @@ public void testFileFieldIdsHandlesFullEmptyAndUnrelatedWrites() { .containsExactly(2); } + @Test + public void testFileFieldsFollowWriteColsOrderAndIgnoreSystemFields() { + TableSchema schema = + new TableSchema( + 1L, + Arrays.asList( + new DataField(1, "indexed", new IntType()), + new DataField(2, "other", new IntType())), + 2, + Collections.emptyList(), + Collections.emptyList(), + new HashMap<>(), + ""); + + assertThat( + DataEvolutionUtils.fileFields( + ignored -> schema, + dataFile( + "reordered.parquet", + 1, + Arrays.asList( + "other", SpecialFields.ROW_ID.name(), "indexed")))) + .extracting(DataField::id) + .containsExactly(2, 1); + } + + @Test + public void testFieldMaxSequenceNumberFallsBackForMissingOrMalformedArray() { + DataFileMeta legacy = dataFile("legacy.parquet", 10, null); + DataFileMeta malformed = + dataFile("malformed.parquet", 10, null) + .withColumnMaxSequenceNumbers(new long[] {5L}); + DataFileMeta valid = + dataFile("valid.parquet", 10, null) + .withColumnMaxSequenceNumbers(new long[] {5L, 8L}); + + assertThat(DataEvolutionUtils.fieldMaxSequenceNumber(legacy, 0, 2)).isEqualTo(10L); + assertThat(DataEvolutionUtils.fieldMaxSequenceNumber(malformed, 0, 2)).isEqualTo(10L); + assertThat(DataEvolutionUtils.fieldMaxSequenceNumber(valid, 0, 2)).isEqualTo(5L); + assertThat(DataEvolutionUtils.fieldMaxSequenceNumber(valid, 1, 2)).isEqualTo(8L); + } + @Test public void testRetrieveAnchorFileSkipsSpecialFiles() { DataFileMeta blobFile = dataFile("blob-file.blob", 1); diff --git a/paimon-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 6a02c91ec510..2dba94abd828 100644 --- a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/compact/changelog/ChangelogCompactTaskSerializerTest.java +++ b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/compact/changelog/ChangelogCompactTaskSerializerTest.java @@ -26,7 +26,6 @@ import org.junit.jupiter.api.Test; import java.util.ArrayList; -import java.util.Collections; import java.util.HashMap; import java.util.List; import java.util.UUID; @@ -105,6 +104,6 @@ private DataFileMeta newFile() { null, null, null) - .withColumnMaxSequenceNumbers(Collections.singletonMap(1, 1L)); + .withColumnMaxSequenceNumbers(new long[] {1L}); } } From f98f9dceb870bb64434eb6afc23b8aaf461a793f Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E4=BB=9F=E5=BC=8B?= Date: Fri, 7 Aug 2026 06:46:23 +0800 Subject: [PATCH 3/4] [core][spark] Fix column sequence propagation --- .../DataEvolutionNormalCompactTask.java | 30 ++++----- .../paimon/spark/copy/CopyFilesUtil.java | 4 +- .../paimon/spark/copy/CopyFilesUtilTest.java | 61 +++++++++++++++++++ 3 files changed, 75 insertions(+), 20 deletions(-) create mode 100644 paimon-spark/paimon-spark-common/src/test/java/org/apache/paimon/spark/copy/CopyFilesUtilTest.java 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 75b999771ce0..aebb1c9004e0 100644 --- a/paimon-core/src/main/java/org/apache/paimon/append/dataevolution/DataEvolutionNormalCompactTask.java +++ b/paimon-core/src/main/java/org/apache/paimon/append/dataevolution/DataEvolutionNormalCompactTask.java @@ -37,7 +37,7 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; -import java.util.LinkedHashMap; +import java.util.HashMap; import java.util.List; import java.util.Map; import java.util.Set; @@ -137,32 +137,24 @@ public CommitMessage doCompact(FileStoreTable table, String commitUser) throws E private long[] compactedColumnMaxSequenceNumbers( FileStoreTable table, DataFileMeta outputFile) { - Map> inputFields = new LinkedHashMap<>(); + Map fieldMaxSequences = new HashMap<>(); for (DataFileMeta input : compactBefore) { - inputFields.put(input, fileFields(table.schemaManager()::schema, input)); + List inputFields = fileFields(table.schemaManager()::schema, input); + for (int inputPosition = 0; inputPosition < inputFields.size(); inputPosition++) { + fieldMaxSequences.merge( + inputFields.get(inputPosition).id(), + fieldMaxSequenceNumber(input, inputPosition, inputFields.size()), + (left, right) -> Math.max(left, right)); + } } long fallbackSequence = maxSequenceId(compactBefore); List outputFields = fileFields(table.schemaManager()::schema, outputFile); long[] result = new long[outputFields.size()]; for (int outputPosition = 0; outputPosition < outputFields.size(); outputPosition++) { - DataField outputField = outputFields.get(outputPosition); - long fieldSequence = Long.MIN_VALUE; - for (DataFileMeta input : compactBefore) { - List fields = inputFields.get(input); - for (int inputPosition = 0; inputPosition < fields.size(); inputPosition++) { - if (fields.get(inputPosition).id() == outputField.id()) { - fieldSequence = - Math.max( - fieldSequence, - fieldMaxSequenceNumber( - input, inputPosition, fields.size())); - break; - } - } - } result[outputPosition] = - fieldSequence == Long.MIN_VALUE ? fallbackSequence : fieldSequence; + fieldMaxSequences.getOrDefault( + outputFields.get(outputPosition).id(), fallbackSequence); } return result; } 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 7cc896c0b048..88dd1fb40311 100644 --- a/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/copy/CopyFilesUtil.java +++ b/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/copy/CopyFilesUtil.java @@ -74,7 +74,9 @@ public static DataFileMeta toNewDataFileMeta( newExternalPath, oldFileMeta.firstRowId(), oldFileMeta.writeCols(), - oldFileMeta.columnMaxSequenceNumbers()); + // Column sequence numbers are positional and cannot be safely reused after + // changing the schema id. A null value makes readers fall back conservatively. + null); } public static IndexFileMeta toNewIndexFileMeta(IndexFileMeta oldFileMeta, String newFileName) { diff --git a/paimon-spark/paimon-spark-common/src/test/java/org/apache/paimon/spark/copy/CopyFilesUtilTest.java b/paimon-spark/paimon-spark-common/src/test/java/org/apache/paimon/spark/copy/CopyFilesUtilTest.java new file mode 100644 index 000000000000..2f973ffce3ad --- /dev/null +++ b/paimon-spark/paimon-spark-common/src/test/java/org/apache/paimon/spark/copy/CopyFilesUtilTest.java @@ -0,0 +1,61 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.paimon.spark.copy; + +import org.apache.paimon.io.DataFileMeta; +import org.apache.paimon.stats.SimpleStats; + +import org.junit.jupiter.api.Test; + +import java.util.Arrays; +import java.util.Collections; + +import static org.assertj.core.api.Assertions.assertThat; + +/** Tests for {@link CopyFilesUtil}. */ +public class CopyFilesUtilTest { + + @Test + void testClearColumnSequencesWhenChangingSchemaId() { + DataFileMeta source = + DataFileMeta.forAppend( + "source.parquet", + 10L, + 2L, + SimpleStats.EMPTY_STATS, + 1L, + 3L, + 5L, + Collections.emptyList(), + null, + null, + null, + null, + null, + Arrays.asList("a", "b")) + .withColumnMaxSequenceNumbers(new long[] {2L, 3L}); + + DataFileMeta copied = CopyFilesUtil.toNewDataFileMeta(source, "copied.parquet", 6L); + + assertThat(copied.fileName()).isEqualTo("copied.parquet"); + assertThat(copied.schemaId()).isEqualTo(6L); + assertThat(copied.writeCols()).containsExactly("a", "b"); + assertThat(copied.columnMaxSequenceNumbers()).isNull(); + } +} From c454d0c0fdd0244610b738b4a90d5b978c160b32 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E4=BB=9F=E5=BC=8B?= Date: Fri, 7 Aug 2026 16:00:12 +0800 Subject: [PATCH 4/4] [core] Harden column sequence compatibility --- .../paimon/io/DataFileMetaSerializer.java | 5 +- .../apache/paimon/io/PojoDataFileMeta.java | 2 +- .../paimon/io/DataFileMetaSerializerTest.java | 13 ++ ...ommittableSerializerCompatibilityTest.java | 57 +++++++ .../table/source/DataSplitCompatibleTest.java | 81 ++++++++++ .../table/source/SplitSerializerTest.java | 143 ++++++++++++------ .../test/resources/compatibility/datasplit-v9 | Bin 0 -> 1034 bytes .../compatibility/manifest-committable-v13-v5 | Bin 0 -> 3146 bytes .../resources/compatibility/split-v2-chain | Bin 0 -> 1491 bytes .../resources/compatibility/split-v2-data | Bin 0 -> 944 bytes .../resources/compatibility/split-v2-fallback | Bin 0 -> 1508 bytes .../compatibility/split-v2-fallback-data | Bin 0 -> 945 bytes .../compatibility/split-v2-incremental | Bin 0 -> 978 bytes .../resources/compatibility/split-v2-indexed | Bin 0 -> 1009 bytes .../compatibility/split-v2-query-auth | Bin 0 -> 1059 bytes 15 files changed, 249 insertions(+), 52 deletions(-) create mode 100644 paimon-core/src/test/resources/compatibility/datasplit-v9 create mode 100644 paimon-core/src/test/resources/compatibility/manifest-committable-v13-v5 create mode 100644 paimon-core/src/test/resources/compatibility/split-v2-chain create mode 100644 paimon-core/src/test/resources/compatibility/split-v2-data create mode 100644 paimon-core/src/test/resources/compatibility/split-v2-fallback create mode 100644 paimon-core/src/test/resources/compatibility/split-v2-fallback-data create mode 100644 paimon-core/src/test/resources/compatibility/split-v2-incremental create mode 100644 paimon-core/src/test/resources/compatibility/split-v2-indexed create mode 100644 paimon-core/src/test/resources/compatibility/split-v2-query-auth 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 6e5c6b2bee27..f148bb397e85 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 @@ -43,6 +43,7 @@ public DataFileMetaSerializer() { @Override public InternalRow toRow(DataFileMeta meta) { + long[] columnMaxSequenceNumbers = meta.columnMaxSequenceNumbers(); return GenericRow.of( BinaryString.fromString(meta.fileName()), meta.fileSize(), @@ -64,9 +65,9 @@ public InternalRow toRow(DataFileMeta meta) { meta.externalPath().map(BinaryString::fromString).orElse(null), meta.firstRowId(), meta.writeCols() == null ? null : toStringArrayData(meta.writeCols()), - meta.columnMaxSequenceNumbers() == null + columnMaxSequenceNumbers == null ? null - : new GenericArray(meta.columnMaxSequenceNumbers())); + : new GenericArray(columnMaxSequenceNumbers)); } @Override 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 5e219f012374..dcbd650c8b1c 100644 --- a/paimon-core/src/main/java/org/apache/paimon/io/PojoDataFileMeta.java +++ b/paimon-core/src/main/java/org/apache/paimon/io/PojoDataFileMeta.java @@ -292,7 +292,7 @@ public List writeCols() { @Nullable @Override public long[] columnMaxSequenceNumbers() { - return columnMaxSequenceNumbers; + return columnMaxSequenceNumbers == null ? null : columnMaxSequenceNumbers.clone(); } @Override 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 3bfc25030031..be87c063aa8b 100644 --- a/paimon-core/src/test/java/org/apache/paimon/io/DataFileMetaSerializerTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/io/DataFileMetaSerializerTest.java @@ -66,6 +66,19 @@ void testLegacySerializerDropsColumnSequences() { assertThat(file.columnMaxSequenceNumbers()).isNull(); } + @Test + void testColumnSequencesAreDefensivelyCopied() { + long[] sequences = {3L, 42L}; + DataFileMeta file = gen.next().meta.withColumnMaxSequenceNumbers(sequences); + + sequences[0] = 100L; + assertColumnSequences(file); + + long[] returned = file.columnMaxSequenceNumbers(); + returned[1] = 100L; + assertColumnSequences(file); + } + private void assertColumnSequences(DataFileMeta file) { assertThat(file.columnMaxSequenceNumbers()).containsExactly(3L, 42L); } diff --git a/paimon-core/src/test/java/org/apache/paimon/manifest/ManifestCommittableSerializerCompatibilityTest.java b/paimon-core/src/test/java/org/apache/paimon/manifest/ManifestCommittableSerializerCompatibilityTest.java index 20eac11da7e8..b10b6f7dcb94 100644 --- a/paimon-core/src/test/java/org/apache/paimon/manifest/ManifestCommittableSerializerCompatibilityTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/manifest/ManifestCommittableSerializerCompatibilityTest.java @@ -27,6 +27,7 @@ import org.apache.paimon.io.DataIncrement; import org.apache.paimon.stats.SimpleStats; import org.apache.paimon.table.sink.CommitMessageImpl; +import org.apache.paimon.utils.CompatibilityUtils; import org.apache.paimon.utils.IOUtils; import org.junit.jupiter.api.Test; @@ -45,6 +46,62 @@ /** Compatibility Test for {@link ManifestCommittableSerializer}. */ public class ManifestCommittableSerializerCompatibilityTest { + private static final String GENERATE_GOLDEN_FILES_PROPERTY = + "generateManifestCommittableGoldenFiles"; + + @Test + public void testCompatibilityToV5CommitV13() throws IOException { + DataFileMeta dataFile = + DataFileMeta.create( + "column-sequence-file", + 1024L, + 10L, + singleColumn("min_key"), + singleColumn("max_key"), + SimpleStats.EMPTY_STATS, + SimpleStats.EMPTY_STATS, + 1L, + 5L, + 1L, + 0, + Collections.emptyList(), + Timestamp.fromLocalDateTime( + LocalDateTime.parse("2026-08-07T00:00:00")), + 0L, + null, + FileSource.COMPACT, + null, + null, + 1L, + Arrays.asList("a", "b")) + .withColumnMaxSequenceNumbers(new long[] {3L, 5L}); + IndexFileMeta indexFile = + new IndexFileMeta( + "index-type", "index-file", 100L, 10L, (GlobalIndexMeta) null, null); + ManifestCommittable committable = + createManifestCommittable( + Collections.singletonList(dataFile), indexFile, indexFile); + + ManifestCommittableSerializer serializer = new ManifestCommittableSerializer(); + byte[] current = serializer.serialize(committable); + byte[] serialized; + if (Boolean.parseBoolean( + System.getProperties().getProperty(GENERATE_GOLDEN_FILES_PROPERTY))) { + CompatibilityUtils.writeCompatibilityFile("manifest-committable-v13-v5", current); + serialized = current; + } else { + serialized = + IOUtils.readFully( + ManifestCommittableSerializerCompatibilityTest.class + .getClassLoader() + .getResourceAsStream( + "compatibility/manifest-committable-v13-v5"), + true); + } + + assertThat(serializer.deserialize(5, serialized)).isEqualTo(committable); + } + @Test public void testCompatibilityToV5CommitV11() throws IOException { String fileName = "manifest-committable-v11-v5"; diff --git a/paimon-core/src/test/java/org/apache/paimon/table/source/DataSplitCompatibleTest.java b/paimon-core/src/test/java/org/apache/paimon/table/source/DataSplitCompatibleTest.java index 140b964106d7..b37f25946a5b 100644 --- a/paimon-core/src/test/java/org/apache/paimon/table/source/DataSplitCompatibleTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/table/source/DataSplitCompatibleTest.java @@ -36,6 +36,7 @@ import org.apache.paimon.types.IntType; import org.apache.paimon.types.SmallIntType; import org.apache.paimon.types.TimestampType; +import org.apache.paimon.utils.CompatibilityUtils; import org.apache.paimon.utils.IOUtils; import org.apache.paimon.utils.InstantiationUtil; @@ -61,6 +62,8 @@ /** Test for {@link DataSplit}. */ public class DataSplitCompatibleTest { + private static final String GENERATE_GOLDEN_FILES_PROPERTY = "generateDataSplitGoldenFiles"; + @Test public void testSplitMergedRowCount() { // not rawConvertible @@ -785,6 +788,84 @@ public void testSerializerCompatibleV8() throws Exception { assertThat(actual).isEqualTo(split); } + @Test + public void testSerializerCompatibleV9() throws Exception { + SimpleStats keyStats = + new SimpleStats( + singleColumn("min_key"), + singleColumn("max_key"), + fromLongArray(new Long[] {0L})); + SimpleStats valueStats = + new SimpleStats( + singleColumn("min_value"), + singleColumn("max_value"), + fromLongArray(new Long[] {0L})); + + DataFileMeta dataFile = + DataFileMeta.create( + "my_file", + 1024 * 1024, + 1024, + singleColumn("min_key"), + singleColumn("max_key"), + keyStats, + valueStats, + 15, + 200, + 5, + 3, + Arrays.asList("extra1", "extra2"), + Timestamp.fromLocalDateTime( + LocalDateTime.parse("2022-03-02T20:20:12")), + 11L, + new byte[] {1, 2, 4}, + FileSource.COMPACT, + Arrays.asList("field1", "field2", "field3"), + "hdfs:///path/to/warehouse", + 12L, + Arrays.asList("a", "b", "c", "f")) + .withColumnMaxSequenceNumbers(new long[] {15L, 100L, 150L, 200L}); + List dataFiles = Collections.singletonList(dataFile); + + DeletionFile deletionFile = new DeletionFile("deletion_file", 100, 22, 33L); + List deletionFiles = Collections.singletonList(deletionFile); + + BinaryRow partition = new BinaryRow(1); + BinaryRowWriter binaryRowWriter = new BinaryRowWriter(partition); + binaryRowWriter.writeString(0, BinaryString.fromString("aaaaa")); + binaryRowWriter.complete(); + + DataSplit split = + DataSplit.builder() + .withSnapshot(18) + .withPartition(partition) + .withBucket(20) + .withTotalBuckets(32) + .withDataFiles(dataFiles) + .withDataDeletionFiles(deletionFiles) + .withBucketPath("my path") + .build(); + + byte[] current = InstantiationUtil.serializeObject(split); + byte[] serialized; + if (Boolean.parseBoolean( + System.getProperties().getProperty(GENERATE_GOLDEN_FILES_PROPERTY))) { + CompatibilityUtils.writeCompatibilityFile("datasplit-v9", current); + serialized = current; + } else { + serialized = + IOUtils.readFully( + DataSplitCompatibleTest.class + .getClassLoader() + .getResourceAsStream("compatibility/datasplit-v9"), + true); + } + + DataSplit actual = + InstantiationUtil.deserializeObject(serialized, DataSplit.class.getClassLoader()); + assertThat(actual).isEqualTo(split); + } + private DataFileMeta newDataFile(long rowCount) { return newDataFile(rowCount, null, null); } diff --git a/paimon-core/src/test/java/org/apache/paimon/table/source/SplitSerializerTest.java b/paimon-core/src/test/java/org/apache/paimon/table/source/SplitSerializerTest.java index 0a5f7f38a275..eea7a817f2cf 100644 --- a/paimon-core/src/test/java/org/apache/paimon/table/source/SplitSerializerTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/table/source/SplitSerializerTest.java @@ -28,6 +28,7 @@ import org.apache.paimon.manifest.FileSource; import org.apache.paimon.stats.SimpleStats; import org.apache.paimon.table.FallbackReadFileStoreTable; +import org.apache.paimon.utils.CompatibilityUtils; import org.apache.paimon.utils.IOUtils; import org.apache.paimon.utils.InstantiationUtil; import org.apache.paimon.utils.Range; @@ -50,11 +51,12 @@ /** Test for {@link SplitSerializer}. */ public class SplitSerializerTest { + private static final String GENERATE_GOLDEN_FILES_PROPERTY = "generateSplitGoldenFiles"; private static final String RESOURCE_PREFIX = "compatibility/"; @Test public void testRoundTrip() throws IOException { - for (GoldenCase goldenCase : goldenCases()) { + for (GoldenCase goldenCase : goldenCases(2)) { Split actual = SplitSerializer.deserialize(SplitSerializer.serialize(goldenCase.split)); assertSplitEquals(goldenCase.split, actual); } @@ -62,16 +64,34 @@ public void testRoundTrip() throws IOException { @Test public void testVersion1GoldenFiles() throws IOException { - for (GoldenCase goldenCase : goldenCases()) { + for (GoldenCase goldenCase : goldenCases(1)) { assertSplitEquals( goldenCase.split, SplitSerializer.deserialize(readGoldenFile(goldenCase.fileName))); } } + @Test + public void testVersion2GoldenFiles() throws IOException { + boolean generateGoldenFiles = + Boolean.parseBoolean( + System.getProperties().getProperty(GENERATE_GOLDEN_FILES_PROPERTY)); + for (GoldenCase goldenCase : goldenCases(2)) { + byte[] current = SplitSerializer.serialize(goldenCase.split); + byte[] serialized; + if (generateGoldenFiles) { + CompatibilityUtils.writeCompatibilityFile(goldenCase.fileName, current); + serialized = current; + } else { + serialized = readGoldenFile(goldenCase.fileName); + } + assertSplitEquals(goldenCase.split, SplitSerializer.deserialize(serialized)); + } + } + @Test public void testFallbackSplitImplSerializeAndDeserialize() throws IOException { - FallbackReadFileStoreTable.FallbackSplitImpl split = fallbackSplit(); + FallbackReadFileStoreTable.FallbackSplitImpl split = fallbackSplit(true); ByteArrayOutputStream out = new ByteArrayOutputStream(); split.serialize(new DataOutputViewStreamWrapper(out)); @@ -84,7 +104,7 @@ public void testFallbackSplitImplSerializeAndDeserialize() throws IOException { @Test public void testFallbackSplitImplJavaSerializeAndDeserialize() throws Exception { - FallbackReadFileStoreTable.FallbackSplitImpl split = fallbackSplit(); + FallbackReadFileStoreTable.FallbackSplitImpl split = fallbackSplit(true); byte[] bytes = InstantiationUtil.serializeObject(split); FallbackReadFileStoreTable.FallbackSplitImpl deserialized = @@ -102,32 +122,36 @@ private static byte[] readGoldenFile(String fileName) throws IOException { } } - private static List goldenCases() { + private static List goldenCases(int version) { + boolean withColumnSequences = version >= 2; List cases = new ArrayList<>(); - DataSplit dataSplit = dataSplit(); - IncrementalSplit incrementalSplit = incrementalSplit(); + DataSplit dataSplit = dataSplit(withColumnSequences); + IncrementalSplit incrementalSplit = incrementalSplit(withColumnSequences); IndexedSplit indexedSplit = indexedSplit(dataSplit); - ChainSplit chainSplit = chainSplit(); + ChainSplit chainSplit = chainSplit(withColumnSequences); QueryAuthSplit queryAuthSplit = queryAuthSplit(dataSplit); FallbackReadFileStoreTable.FallbackDataSplit fallbackDataSplit = (FallbackReadFileStoreTable.FallbackDataSplit) FallbackReadFileStoreTable.toFallbackSplit(dataSplit, true); - FallbackReadFileStoreTable.FallbackSplitImpl fallbackSplit = fallbackSplit(); - - cases.add(new GoldenCase("split-v1-data", dataSplit)); - cases.add(new GoldenCase("split-v1-incremental", incrementalSplit)); - cases.add(new GoldenCase("split-v1-indexed", indexedSplit)); - cases.add(new GoldenCase("split-v1-chain", chainSplit)); - cases.add(new GoldenCase("split-v1-query-auth", queryAuthSplit)); - cases.add(new GoldenCase("split-v1-fallback-data", fallbackDataSplit)); - cases.add(new GoldenCase("split-v1-fallback", fallbackSplit)); + FallbackReadFileStoreTable.FallbackSplitImpl fallbackSplit = + fallbackSplit(withColumnSequences); + + String prefix = "split-v" + version + "-"; + cases.add(new GoldenCase(prefix + "data", dataSplit)); + cases.add(new GoldenCase(prefix + "incremental", incrementalSplit)); + cases.add(new GoldenCase(prefix + "indexed", indexedSplit)); + cases.add(new GoldenCase(prefix + "chain", chainSplit)); + cases.add(new GoldenCase(prefix + "query-auth", queryAuthSplit)); + cases.add(new GoldenCase(prefix + "fallback-data", fallbackDataSplit)); + cases.add(new GoldenCase(prefix + "fallback", fallbackSplit)); return cases; } - private static DataSplit dataSplit() { + private static DataSplit dataSplit(boolean withColumnSequences) { List files = Arrays.asList( - dataFile("file-a", 0, 1, 10, 100L), dataFile("file-b", 1, 11, 20, 200L)); + dataFile("file-a", 0, 1, 10, 100L, withColumnSequences), + dataFile("file-b", 1, 11, 20, 200L, withColumnSequences)); List deletionFiles = Arrays.asList(null, new DeletionFile("dv/file-b", 2L, 10L, 3L)); return DataSplit.builder() @@ -142,10 +166,13 @@ private static DataSplit dataSplit() { .build(); } - private static IncrementalSplit incrementalSplit() { + private static IncrementalSplit incrementalSplit(boolean withColumnSequences) { List before = - Collections.singletonList(dataFile("before-file", 0, 1, 5, 10L)); - List after = Collections.singletonList(dataFile("after-file", 0, 6, 12, 20L)); + Collections.singletonList( + dataFile("before-file", 0, 1, 5, 10L, withColumnSequences)); + List after = + Collections.singletonList( + dataFile("after-file", 0, 6, 12, 20L, withColumnSequences)); return new IncrementalSplit( 43L, DataFileTestUtils.row(2026, 8), @@ -165,8 +192,8 @@ private static IndexedSplit indexedSplit(DataSplit dataSplit) { new float[] {0.5f, 0.25f, 0.125f}); } - private static ChainSplit chainSplit() { - DataSplit left = dataSplit(); + private static ChainSplit chainSplit(boolean withColumnSequences) { + DataSplit left = dataSplit(withColumnSequences); DataSplit right = DataSplit.builder() .withSnapshot(44L) @@ -175,7 +202,14 @@ private static ChainSplit chainSplit() { .withTotalBuckets(8) .withBucketPath("dt=20260707/bucket-5") .withDataFiles( - Collections.singletonList(dataFile("chain-file", 0, 21, 30, 300L))) + Collections.singletonList( + dataFile( + "chain-file", + 0, + 21, + 30, + 300L, + withColumnSequences))) .withDataDeletionFiles( Collections.singletonList( new DeletionFile("deletion_file", 100, 22, null))) @@ -215,33 +249,44 @@ private static QueryAuthSplit queryAuthSplit(Split split) { Arrays.asList("filter-json-1", "filter-json-2"), columnMasking)); } - private static FallbackReadFileStoreTable.FallbackSplitImpl fallbackSplit() { - return new FallbackReadFileStoreTable.FallbackSplitImpl(chainSplit(), true); + private static FallbackReadFileStoreTable.FallbackSplitImpl fallbackSplit( + boolean withColumnSequences) { + return new FallbackReadFileStoreTable.FallbackSplitImpl( + chainSplit(withColumnSequences), true); } private static DataFileMeta dataFile( - String name, int level, int minKey, int maxKey, long maxSequence) { - return DataFileMeta.create( - name, - maxKey - minKey + 1, - maxKey - minKey + 1, - DataFileTestUtils.row(minKey), - DataFileTestUtils.row(maxKey), - SimpleStats.EMPTY_STATS, - SimpleStats.EMPTY_STATS, - 0L, - maxSequence, - 0L, - level, - Collections.emptyList(), - Timestamp.fromEpochMillis(100), - 0L, - null, - FileSource.APPEND, - null, - null, - null, - null); + String name, + int level, + int minKey, + int maxKey, + long maxSequence, + boolean withColumnSequences) { + DataFileMeta file = + DataFileMeta.create( + name, + maxKey - minKey + 1, + maxKey - minKey + 1, + DataFileTestUtils.row(minKey), + DataFileTestUtils.row(maxKey), + SimpleStats.EMPTY_STATS, + SimpleStats.EMPTY_STATS, + 0L, + maxSequence, + 0L, + level, + Collections.emptyList(), + Timestamp.fromEpochMillis(100), + 0L, + null, + FileSource.APPEND, + null, + null, + null, + null); + return withColumnSequences + ? file.withColumnMaxSequenceNumbers(new long[] {maxSequence}) + : file; } private static void assertSplitEquals(Split expected, Split actual) { diff --git a/paimon-core/src/test/resources/compatibility/datasplit-v9 b/paimon-core/src/test/resources/compatibility/datasplit-v9 new file mode 100644 index 0000000000000000000000000000000000000000..6277c55665ef8369e01e4458cf2ca4f1739174d9 GIT binary patch literal 1034 zcma)5yKWOf6uou`k&^&PKq%=9qDX-}KpI4vC`1THTSTNwlU#ep-XZJVWq16d03w9> z1HOWSPe7ufrlLZK8hT33Jl2k&;Yw$Z=gfW1jNkvF`68#yH19Sz<8~w)8LM8JG&Hwj z*(lO}-jAwybz7Low4*=LK?mzNP zon7qodkeqnfZ@l0$)5qRX^6Qiqno}59QUQ!o!P&BnB#B1sgsX$M?xe=I_JAIv3!pv z<@tA%j6>*_pPYBR3^|uk+AqqU?CqciyEJIewRxX1=DJ_sAR7Gv5m>cf literal 0 HcmV?d00001 diff --git a/paimon-core/src/test/resources/compatibility/manifest-committable-v13-v5 b/paimon-core/src/test/resources/compatibility/manifest-committable-v13-v5 new file mode 100644 index 0000000000000000000000000000000000000000..d919649c29e31701f71f2510d46b07ac53c7b638 GIT binary patch literal 3146 zcmZQz00UMC#Vo1Wg9Xul>u?K5s)?litz$5BM^h+K{Ocf07(HL&aI3uNGvMJ zEXmBzgUB#&gA{?}fHdO)1_lRakc0>jZvZh^Km-E=7lZX$O!A zKsE?KoeAedL>U-3fQ$~XqSl-9VTxgrKpKk>SQaKO0pv{pazOM7Af23_Q<|HnTbx>0 znwpoKs+*RXlM2>=WK?cuUVL_HWjls=Vg<4|SQ->i3P22^4S*QrUI!p16;R6@CaA@U z5V|o5O2fp#kp&ckg*T8143*J%q*nM-%N$sIjm9Ie5E^m$ks*Ph0FrS*<){M?gXjPt zP61*NfR(Qh9wwQYmy%kcTT)p7E!`k|q|y#j-qFdP(ei~_M$7U|<4b79eH@Vi6!_1Y$6F#SUYEM$7U|<4b#(Re^8R$s9NdvMuffxj|fLI2ILGobmiXB2jDG-O5K_sQb z*2uuf%)s2hOh2hKIXksP*O(C~#Q{+a*3a0%z~I0S5=hI;N!3kcXyby)!e|j7e*+MM z=o3Kf0mKiWv~Dj@d3ZG|%+y5L-x+mQ|90df?8 z7(^QYG00sGKujv2mN_uTBEkhNE|SQP3s5{E1DIX3j0{zq?S1_ f$HKx5#AgC}fitB{ADl0epgw@*I9LK>W?%#W8yP0? literal 0 HcmV?d00001 diff --git a/paimon-core/src/test/resources/compatibility/split-v2-fallback b/paimon-core/src/test/resources/compatibility/split-v2-fallback new file mode 100644 index 0000000000000000000000000000000000000000..c92ca5a3eba6f550771799d10f3009938fa3609a GIT binary patch literal 1508 zcmWFz@bL_Z4>M$7U|<4bcE(^-0T!SjGZ2daF(VLz!7Fwc3na(b!NB0a4-!es%t_Tv zWN71pO2cT7<_$m$qE7&^2M|Ai(i(7685mN4V%Pu&P_O`~4wpPOJ=nwqfPxhe`{1_1 z84wL{F3j!7=I{VH3P22^4S*QrE(ahc6;R6@m}3#)f)*D^Q0+g0O=^vJ+U^K{I2Y?tv-vD9{ApQWQEno%$`8r5! zAQuimanS&jL&QH^iX;XoUvvP?NzO>j%+m$sVz_p=&7yE2Fas_whbj){Q7cTTWe&`- zNa4uLzy>X?L{ds@jSP&;49pG8^pi@Hvr|iSjiIth4A^D4kU|Pcg*i;WDKHoiHgFW@ zB^DHCN6aWC{U}ReW literal 0 HcmV?d00001 diff --git a/paimon-core/src/test/resources/compatibility/split-v2-fallback-data b/paimon-core/src/test/resources/compatibility/split-v2-fallback-data new file mode 100644 index 0000000000000000000000000000000000000000..eab303006fcbae62473c737f4affb04f15915895 GIT binary patch literal 945 zcmWFz@bL_Z4>M$7U|<4bwtI&!8R$s9NdvMuffxj|fLI2ILGobmiXB2jDG-O5K_sQb z*2uuf%)s2hOh2hKIXksP*O(C~#Q{+a*3a0%z~I0S5=hI;N!3kcXyby)!e|j7e*+MM z=o3Kf0mKiWv~Dj@d3ZG|%+y5L-x+mQ|90df?8 z7(^QYG00sGKujv2mN_uTBEkhNE|SQP3s5{E1DIX3j0M$7U|@n`AjO~!#4<>HhF9z$VFm^c2n8Zppj<`<2F4i-3=aGtL2e-4 z0K}{y4iLa-5g_{j5QFF&KM$7U|<4b=1)H#ggB|_f%uH~4qr0Rk$jT|WOD*B2xtMZ3=o6Vg25|x z2o0q`9A*ZQloDGb10yp7a|1K|q|)T<)Dm4|MxYc2L@`)DV+R9+13yS0Ei)%oH<6)@ z3n~kvMS%PbKn$W!0I>%UKY-F2AmgAQ1;#+5LADhD)!~vyQ;w4>0F2&VcBG zb75{rHjD?zQ2=5PZ2-g|cR2ttseoGMz#NMR7qqxYB0nxb@q`RucF{5}xREUcrdt>n zly1-gwZf2E=D-{a3pWs-3FrmRlrnvAzDR=l0G8um35*%+Pnamo7#65%Sdj37B$(}i Kgk1v=18D$pvn)0M literal 0 HcmV?d00001 diff --git a/paimon-core/src/test/resources/compatibility/split-v2-query-auth b/paimon-core/src/test/resources/compatibility/split-v2-query-auth new file mode 100644 index 0000000000000000000000000000000000000000..f8a06462a45f641c4314b1f6bfe9ff949291963f GIT binary patch literal 1059 zcmWFz@bL_Z4>M$7U|<4b)?idV%ifx&?vB#@Swld7A@(8dLoh0!8F z{stfh(I=#+pwS@P3V`Zx$)hR9$rS*~R6y*5+X`nubiuhWw<8w-979Y7J*f{gq;kQ6^y3d{q`0RX0%LOuWh literal 0 HcmV?d00001