diff --git a/core/src/main/java/org/apache/iceberg/util/PartitionUtil.java b/core/src/main/java/org/apache/iceberg/util/PartitionUtil.java index ad6ef605420a..c03a43a3dfdf 100644 --- a/core/src/main/java/org/apache/iceberg/util/PartitionUtil.java +++ b/core/src/main/java/org/apache/iceberg/util/PartitionUtil.java @@ -62,7 +62,7 @@ private PartitionUtil() {} idToConstant.put( MetadataColumns.LAST_UPDATED_SEQUENCE_NUMBER.fieldId(), - convertConstant.apply(Types.LongType.get(), task.file().fileSequenceNumber())); + convertConstant.apply(Types.LongType.get(), task.file().dataSequenceNumber())); // add _file idToConstant.put( diff --git a/core/src/test/java/org/apache/iceberg/TestRowLineageAssignment.java b/core/src/test/java/org/apache/iceberg/TestRowLineageAssignment.java index 91aec244d426..51a608ed44d8 100644 --- a/core/src/test/java/org/apache/iceberg/TestRowLineageAssignment.java +++ b/core/src/test/java/org/apache/iceberg/TestRowLineageAssignment.java @@ -26,12 +26,14 @@ import java.io.UncheckedIOException; import java.util.List; import java.util.Map; +import java.util.Set; import org.apache.iceberg.data.Record; import org.apache.iceberg.io.CloseableIterable; import org.apache.iceberg.io.InputFile; import org.apache.iceberg.relocated.com.google.common.collect.Iterables; import org.apache.iceberg.types.Types; import org.apache.iceberg.types.Types.NestedField; +import org.apache.iceberg.util.PartitionUtil; import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; @@ -722,6 +724,45 @@ public void testRowDeltaAssignmentAfterUpgrade(@TempDir File altLocation) { assertThat(manifests.get(1).path()).isEqualTo(existingManifests.get(1).path()); } + @Test + public void lastUpdatedAfterUpgrade(@TempDir File altLocation) throws IOException { + BaseTable upgradeTable = + TestTables.create(altLocation, "test_upgrade", SCHEMA, PartitionSpec.unpartitioned(), 2); + + upgradeTable.newAppend().appendFile(FILE_A).commit(); + Snapshot originalSnapshot = upgradeTable.currentSnapshot(); + long originalSequenceNumber = originalSnapshot.sequenceNumber(); + + upgradeTable + .newRewrite() + .validateFromSnapshot(originalSnapshot.snapshotId()) + .rewriteFiles(Set.of(FILE_A), Set.of(FILE_B), originalSequenceNumber) + .commit(); + + TestTables.upgrade(altLocation, "test_upgrade", 3); + upgradeTable.refresh(); + + // Assign row IDs to the upgraded metadata tree without rewriting the data file. + upgradeTable.newFastAppend().commit(); + + try (CloseableIterable tasks = upgradeTable.newScan().planFiles()) { + FileScanTask task = Iterables.getOnlyElement(tasks); + + assertThat(task.file().location()).isEqualTo(FILE_B.location()); + assertThat(task.file().dataSequenceNumber()).isEqualTo(originalSequenceNumber); + assertThat(task.file().fileSequenceNumber()).isGreaterThan(originalSequenceNumber); + assertThat(task.file().firstRowId()).isNotNull(); + + // Scan planning projects row lineage metadata columns through constantsMap. + // The upgraded data file is not rewritten, so validate the value readers see. + assertThat( + PartitionUtil.constantsMap(task) + .get(MetadataColumns.LAST_UPDATED_SEQUENCE_NUMBER.fieldId())) + .as("Last updated should preserve the original data sequence after upgrade") + .isEqualTo(originalSequenceNumber); + } + } + @Test public void testUpgradeAssignmentWithManifestCompaction(@TempDir File altLocation) { // create a non-empty upgrade table with FILE_A diff --git a/spark/v3.4/spark/src/test/java/org/apache/iceberg/spark/actions/TestRewriteDataFilesAction.java b/spark/v3.4/spark/src/test/java/org/apache/iceberg/spark/actions/TestRewriteDataFilesAction.java index 50e42dcbb5cb..eac8336253a8 100644 --- a/spark/v3.4/spark/src/test/java/org/apache/iceberg/spark/actions/TestRewriteDataFilesAction.java +++ b/spark/v3.4/spark/src/test/java/org/apache/iceberg/spark/actions/TestRewriteDataFilesAction.java @@ -2065,6 +2065,49 @@ public void testUnpartitionedRewriteDataFilesPreservesLineage() throws NoSuchTab assertEquals("Rows must match", expectedRecordsWithLineage, actualRecordsWithLineage); } + @TestTemplate + public void testUpgradePreservesDataSequence() throws NoSuchTableException { + assumeThat(formatVersion).isEqualTo(2); + + Table table = createTable(); + writeRecords(2, 4); + table.refresh(); + shouldHaveFiles(table, 2); + long committedDataSequence = table.currentSnapshot().sequenceNumber(); + + Result result = basicRewrite(table).execute(); + assertThat(result.rewrittenDataFilesCount()).isEqualTo(2); + assertThat(result.addedDataFilesCount()).isOne(); + table.refresh(); + shouldHaveFiles(table, 1); + + DataFile compactedFile = Iterables.getOnlyElement(currentDataFiles(table)); + long dataSequenceNumber = compactedFile.dataSequenceNumber(); + assertThat(dataSequenceNumber) + .as("Compaction must preserve the original data sequence number") + .isEqualTo(committedDataSequence); + assertThat(compactedFile.fileSequenceNumber()) + .as("Compaction must bump the file sequence above the preserved data sequence") + .isGreaterThan(dataSequenceNumber); + + table.updateProperties().set(TableProperties.FORMAT_VERSION, "3").commit(); + table.rewriteManifests().rewriteIf(manifest -> true).commit(); + table.refresh(); + + List expectedLineage = + Lists.newArrayList( + row(0L, committedDataSequence, ANY, ANY, ANY), + row(1L, committedDataSequence, ANY, ANY, ANY), + row(2L, committedDataSequence, ANY, ANY, ANY), + row(3L, committedDataSequence, ANY, ANY, ANY)); + + assertEquals( + "First snapshot after upgrade to v3 assigns row IDs and inherits the committed data" + + " sequence as _last_updated_sequence_number", + expectedLineage, + currentDataWithLineage()); + } + @TestTemplate public void testRewriteDataFilesPreservesLineage() throws NoSuchTableException { assumeThat(formatVersion).isGreaterThan(2); diff --git a/spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/actions/TestRewriteDataFilesAction.java b/spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/actions/TestRewriteDataFilesAction.java index 50e42dcbb5cb..eac8336253a8 100644 --- a/spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/actions/TestRewriteDataFilesAction.java +++ b/spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/actions/TestRewriteDataFilesAction.java @@ -2065,6 +2065,49 @@ public void testUnpartitionedRewriteDataFilesPreservesLineage() throws NoSuchTab assertEquals("Rows must match", expectedRecordsWithLineage, actualRecordsWithLineage); } + @TestTemplate + public void testUpgradePreservesDataSequence() throws NoSuchTableException { + assumeThat(formatVersion).isEqualTo(2); + + Table table = createTable(); + writeRecords(2, 4); + table.refresh(); + shouldHaveFiles(table, 2); + long committedDataSequence = table.currentSnapshot().sequenceNumber(); + + Result result = basicRewrite(table).execute(); + assertThat(result.rewrittenDataFilesCount()).isEqualTo(2); + assertThat(result.addedDataFilesCount()).isOne(); + table.refresh(); + shouldHaveFiles(table, 1); + + DataFile compactedFile = Iterables.getOnlyElement(currentDataFiles(table)); + long dataSequenceNumber = compactedFile.dataSequenceNumber(); + assertThat(dataSequenceNumber) + .as("Compaction must preserve the original data sequence number") + .isEqualTo(committedDataSequence); + assertThat(compactedFile.fileSequenceNumber()) + .as("Compaction must bump the file sequence above the preserved data sequence") + .isGreaterThan(dataSequenceNumber); + + table.updateProperties().set(TableProperties.FORMAT_VERSION, "3").commit(); + table.rewriteManifests().rewriteIf(manifest -> true).commit(); + table.refresh(); + + List expectedLineage = + Lists.newArrayList( + row(0L, committedDataSequence, ANY, ANY, ANY), + row(1L, committedDataSequence, ANY, ANY, ANY), + row(2L, committedDataSequence, ANY, ANY, ANY), + row(3L, committedDataSequence, ANY, ANY, ANY)); + + assertEquals( + "First snapshot after upgrade to v3 assigns row IDs and inherits the committed data" + + " sequence as _last_updated_sequence_number", + expectedLineage, + currentDataWithLineage()); + } + @TestTemplate public void testRewriteDataFilesPreservesLineage() throws NoSuchTableException { assumeThat(formatVersion).isGreaterThan(2); diff --git a/spark/v4.0/spark/src/test/java/org/apache/iceberg/spark/actions/TestRewriteDataFilesAction.java b/spark/v4.0/spark/src/test/java/org/apache/iceberg/spark/actions/TestRewriteDataFilesAction.java index 255938eb3f3c..1407a5009b52 100644 --- a/spark/v4.0/spark/src/test/java/org/apache/iceberg/spark/actions/TestRewriteDataFilesAction.java +++ b/spark/v4.0/spark/src/test/java/org/apache/iceberg/spark/actions/TestRewriteDataFilesAction.java @@ -2069,6 +2069,49 @@ public void testUnpartitionedRewriteDataFilesPreservesLineage() throws NoSuchTab assertEquals("Rows must match", expectedRecordsWithLineage, actualRecordsWithLineage); } + @TestTemplate + public void testUpgradePreservesDataSequence() throws NoSuchTableException { + assumeThat(formatVersion).isEqualTo(2); + + Table table = createTable(); + writeRecords(2, 4); + table.refresh(); + shouldHaveFiles(table, 2); + long committedDataSequence = table.currentSnapshot().sequenceNumber(); + + Result result = basicRewrite(table).execute(); + assertThat(result.rewrittenDataFilesCount()).isEqualTo(2); + assertThat(result.addedDataFilesCount()).isOne(); + table.refresh(); + shouldHaveFiles(table, 1); + + DataFile compactedFile = Iterables.getOnlyElement(currentDataFiles(table)); + long dataSequenceNumber = compactedFile.dataSequenceNumber(); + assertThat(dataSequenceNumber) + .as("Compaction must preserve the original data sequence number") + .isEqualTo(committedDataSequence); + assertThat(compactedFile.fileSequenceNumber()) + .as("Compaction must bump the file sequence above the preserved data sequence") + .isGreaterThan(dataSequenceNumber); + + table.updateProperties().set(TableProperties.FORMAT_VERSION, "3").commit(); + table.rewriteManifests().rewriteIf(manifest -> true).commit(); + table.refresh(); + + List expectedLineage = + Lists.newArrayList( + row(0L, committedDataSequence, ANY, ANY, ANY), + row(1L, committedDataSequence, ANY, ANY, ANY), + row(2L, committedDataSequence, ANY, ANY, ANY), + row(3L, committedDataSequence, ANY, ANY, ANY)); + + assertEquals( + "First snapshot after upgrade to v3 assigns row IDs and inherits the committed data" + + " sequence as _last_updated_sequence_number", + expectedLineage, + currentDataWithLineage()); + } + @TestTemplate public void testRewriteDataFilesPreservesLineage() throws NoSuchTableException { assumeThat(formatVersion).isGreaterThan(2);