Skip to content

Commit 5f15471

Browse files
authored
Core, Spark: Fix row lineage last updated sequence inheritance (#17039)
Core, Spark: Fix row lineage last updated sequence inheritance Restore the row lineage spec wording and update Java scan constants to inherit _last_updated_sequence_number from a data file's data sequence number. This preserves the original row update sequence when a v2 rewrite carries an older data sequence number and the table is later upgraded to v3.
1 parent 30f87ad commit 5f15471

5 files changed

Lines changed: 171 additions & 1 deletion

File tree

core/src/main/java/org/apache/iceberg/util/PartitionUtil.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -62,7 +62,7 @@ private PartitionUtil() {}
6262

6363
idToConstant.put(
6464
MetadataColumns.LAST_UPDATED_SEQUENCE_NUMBER.fieldId(),
65-
convertConstant.apply(Types.LongType.get(), task.file().fileSequenceNumber()));
65+
convertConstant.apply(Types.LongType.get(), task.file().dataSequenceNumber()));
6666

6767
// add _file
6868
idToConstant.put(

core/src/test/java/org/apache/iceberg/TestRowLineageAssignment.java

Lines changed: 41 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -26,12 +26,14 @@
2626
import java.io.UncheckedIOException;
2727
import java.util.List;
2828
import java.util.Map;
29+
import java.util.Set;
2930
import org.apache.iceberg.data.Record;
3031
import org.apache.iceberg.io.CloseableIterable;
3132
import org.apache.iceberg.io.InputFile;
3233
import org.apache.iceberg.relocated.com.google.common.collect.Iterables;
3334
import org.apache.iceberg.types.Types;
3435
import org.apache.iceberg.types.Types.NestedField;
36+
import org.apache.iceberg.util.PartitionUtil;
3537
import org.junit.jupiter.api.AfterEach;
3638
import org.junit.jupiter.api.BeforeEach;
3739
import org.junit.jupiter.api.Test;
@@ -722,6 +724,45 @@ public void testRowDeltaAssignmentAfterUpgrade(@TempDir File altLocation) {
722724
assertThat(manifests.get(1).path()).isEqualTo(existingManifests.get(1).path());
723725
}
724726

727+
@Test
728+
public void lastUpdatedAfterUpgrade(@TempDir File altLocation) throws IOException {
729+
BaseTable upgradeTable =
730+
TestTables.create(altLocation, "test_upgrade", SCHEMA, PartitionSpec.unpartitioned(), 2);
731+
732+
upgradeTable.newAppend().appendFile(FILE_A).commit();
733+
Snapshot originalSnapshot = upgradeTable.currentSnapshot();
734+
long originalSequenceNumber = originalSnapshot.sequenceNumber();
735+
736+
upgradeTable
737+
.newRewrite()
738+
.validateFromSnapshot(originalSnapshot.snapshotId())
739+
.rewriteFiles(Set.of(FILE_A), Set.of(FILE_B), originalSequenceNumber)
740+
.commit();
741+
742+
TestTables.upgrade(altLocation, "test_upgrade", 3);
743+
upgradeTable.refresh();
744+
745+
// Assign row IDs to the upgraded metadata tree without rewriting the data file.
746+
upgradeTable.newFastAppend().commit();
747+
748+
try (CloseableIterable<FileScanTask> tasks = upgradeTable.newScan().planFiles()) {
749+
FileScanTask task = Iterables.getOnlyElement(tasks);
750+
751+
assertThat(task.file().location()).isEqualTo(FILE_B.location());
752+
assertThat(task.file().dataSequenceNumber()).isEqualTo(originalSequenceNumber);
753+
assertThat(task.file().fileSequenceNumber()).isGreaterThan(originalSequenceNumber);
754+
assertThat(task.file().firstRowId()).isNotNull();
755+
756+
// Scan planning projects row lineage metadata columns through constantsMap.
757+
// The upgraded data file is not rewritten, so validate the value readers see.
758+
assertThat(
759+
PartitionUtil.constantsMap(task)
760+
.get(MetadataColumns.LAST_UPDATED_SEQUENCE_NUMBER.fieldId()))
761+
.as("Last updated should preserve the original data sequence after upgrade")
762+
.isEqualTo(originalSequenceNumber);
763+
}
764+
}
765+
725766
@Test
726767
public void testUpgradeAssignmentWithManifestCompaction(@TempDir File altLocation) {
727768
// create a non-empty upgrade table with FILE_A

spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/actions/TestRewriteDataFilesAction.java

Lines changed: 43 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2067,6 +2067,49 @@ public void testUnpartitionedRewriteDataFilesPreservesLineage() throws NoSuchTab
20672067
assertEquals("Rows must match", expectedRecordsWithLineage, actualRecordsWithLineage);
20682068
}
20692069

2070+
@TestTemplate
2071+
public void testUpgradePreservesDataSequence() throws NoSuchTableException {
2072+
assumeThat(formatVersion).isEqualTo(2);
2073+
2074+
Table table = createTable();
2075+
writeRecords(2, 4);
2076+
table.refresh();
2077+
shouldHaveFiles(table, 2);
2078+
long committedDataSequence = table.currentSnapshot().sequenceNumber();
2079+
2080+
Result result = basicRewrite(table).execute();
2081+
assertThat(result.rewrittenDataFilesCount()).isEqualTo(2);
2082+
assertThat(result.addedDataFilesCount()).isOne();
2083+
table.refresh();
2084+
shouldHaveFiles(table, 1);
2085+
2086+
DataFile compactedFile = Iterables.getOnlyElement(currentDataFiles(table));
2087+
long dataSequenceNumber = compactedFile.dataSequenceNumber();
2088+
assertThat(dataSequenceNumber)
2089+
.as("Compaction must preserve the original data sequence number")
2090+
.isEqualTo(committedDataSequence);
2091+
assertThat(compactedFile.fileSequenceNumber())
2092+
.as("Compaction must bump the file sequence above the preserved data sequence")
2093+
.isGreaterThan(dataSequenceNumber);
2094+
2095+
table.updateProperties().set(TableProperties.FORMAT_VERSION, "3").commit();
2096+
table.rewriteManifests().rewriteIf(manifest -> true).commit();
2097+
table.refresh();
2098+
2099+
List<Object[]> expectedLineage =
2100+
Lists.newArrayList(
2101+
row(0L, committedDataSequence, ANY, ANY, ANY),
2102+
row(1L, committedDataSequence, ANY, ANY, ANY),
2103+
row(2L, committedDataSequence, ANY, ANY, ANY),
2104+
row(3L, committedDataSequence, ANY, ANY, ANY));
2105+
2106+
assertEquals(
2107+
"First snapshot after upgrade to v3 assigns row IDs and inherits the committed data"
2108+
+ " sequence as _last_updated_sequence_number",
2109+
expectedLineage,
2110+
currentDataWithLineage());
2111+
}
2112+
20702113
@TestTemplate
20712114
public void testRewriteDataFilesPreservesLineage() throws NoSuchTableException {
20722115
assumeThat(formatVersion).isGreaterThan(2);

spark/v4.0/spark/src/test/java/org/apache/iceberg/spark/actions/TestRewriteDataFilesAction.java

Lines changed: 43 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2069,6 +2069,49 @@ public void testUnpartitionedRewriteDataFilesPreservesLineage() throws NoSuchTab
20692069
assertEquals("Rows must match", expectedRecordsWithLineage, actualRecordsWithLineage);
20702070
}
20712071

2072+
@TestTemplate
2073+
public void testUpgradePreservesDataSequence() throws NoSuchTableException {
2074+
assumeThat(formatVersion).isEqualTo(2);
2075+
2076+
Table table = createTable();
2077+
writeRecords(2, 4);
2078+
table.refresh();
2079+
shouldHaveFiles(table, 2);
2080+
long committedDataSequence = table.currentSnapshot().sequenceNumber();
2081+
2082+
Result result = basicRewrite(table).execute();
2083+
assertThat(result.rewrittenDataFilesCount()).isEqualTo(2);
2084+
assertThat(result.addedDataFilesCount()).isOne();
2085+
table.refresh();
2086+
shouldHaveFiles(table, 1);
2087+
2088+
DataFile compactedFile = Iterables.getOnlyElement(currentDataFiles(table));
2089+
long dataSequenceNumber = compactedFile.dataSequenceNumber();
2090+
assertThat(dataSequenceNumber)
2091+
.as("Compaction must preserve the original data sequence number")
2092+
.isEqualTo(committedDataSequence);
2093+
assertThat(compactedFile.fileSequenceNumber())
2094+
.as("Compaction must bump the file sequence above the preserved data sequence")
2095+
.isGreaterThan(dataSequenceNumber);
2096+
2097+
table.updateProperties().set(TableProperties.FORMAT_VERSION, "3").commit();
2098+
table.rewriteManifests().rewriteIf(manifest -> true).commit();
2099+
table.refresh();
2100+
2101+
List<Object[]> expectedLineage =
2102+
Lists.newArrayList(
2103+
row(0L, committedDataSequence, ANY, ANY, ANY),
2104+
row(1L, committedDataSequence, ANY, ANY, ANY),
2105+
row(2L, committedDataSequence, ANY, ANY, ANY),
2106+
row(3L, committedDataSequence, ANY, ANY, ANY));
2107+
2108+
assertEquals(
2109+
"First snapshot after upgrade to v3 assigns row IDs and inherits the committed data"
2110+
+ " sequence as _last_updated_sequence_number",
2111+
expectedLineage,
2112+
currentDataWithLineage());
2113+
}
2114+
20722115
@TestTemplate
20732116
public void testRewriteDataFilesPreservesLineage() throws NoSuchTableException {
20742117
assumeThat(formatVersion).isGreaterThan(2);

spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/actions/TestRewriteDataFilesAction.java

Lines changed: 43 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2069,6 +2069,49 @@ public void testUnpartitionedRewriteDataFilesPreservesLineage() throws NoSuchTab
20692069
assertEquals("Rows must match", expectedRecordsWithLineage, actualRecordsWithLineage);
20702070
}
20712071

2072+
@TestTemplate
2073+
public void testUpgradePreservesDataSequence() throws NoSuchTableException {
2074+
assumeThat(formatVersion).isEqualTo(2);
2075+
2076+
Table table = createTable();
2077+
writeRecords(2, 4);
2078+
table.refresh();
2079+
shouldHaveFiles(table, 2);
2080+
long committedDataSequence = table.currentSnapshot().sequenceNumber();
2081+
2082+
Result result = basicRewrite(table).execute();
2083+
assertThat(result.rewrittenDataFilesCount()).isEqualTo(2);
2084+
assertThat(result.addedDataFilesCount()).isOne();
2085+
table.refresh();
2086+
shouldHaveFiles(table, 1);
2087+
2088+
DataFile compactedFile = Iterables.getOnlyElement(currentDataFiles(table));
2089+
long dataSequenceNumber = compactedFile.dataSequenceNumber();
2090+
assertThat(dataSequenceNumber)
2091+
.as("Compaction must preserve the original data sequence number")
2092+
.isEqualTo(committedDataSequence);
2093+
assertThat(compactedFile.fileSequenceNumber())
2094+
.as("Compaction must bump the file sequence above the preserved data sequence")
2095+
.isGreaterThan(dataSequenceNumber);
2096+
2097+
table.updateProperties().set(TableProperties.FORMAT_VERSION, "3").commit();
2098+
table.rewriteManifests().rewriteIf(manifest -> true).commit();
2099+
table.refresh();
2100+
2101+
List<Object[]> expectedLineage =
2102+
Lists.newArrayList(
2103+
row(0L, committedDataSequence, ANY, ANY, ANY),
2104+
row(1L, committedDataSequence, ANY, ANY, ANY),
2105+
row(2L, committedDataSequence, ANY, ANY, ANY),
2106+
row(3L, committedDataSequence, ANY, ANY, ANY));
2107+
2108+
assertEquals(
2109+
"First snapshot after upgrade to v3 assigns row IDs and inherits the committed data"
2110+
+ " sequence as _last_updated_sequence_number",
2111+
expectedLineage,
2112+
currentDataWithLineage());
2113+
}
2114+
20722115
@TestTemplate
20732116
public void testRewriteDataFilesPreservesLineage() throws NoSuchTableException {
20742117
assumeThat(formatVersion).isGreaterThan(2);

0 commit comments

Comments
 (0)