From c262b5af537ee06c768959dc9f6a88da9c9de88e Mon Sep 17 00:00:00 2001 From: aokolnychyi Date: Tue, 21 Jun 2022 08:58:37 -0700 Subject: [PATCH 1/3] API: Access deleted and added delete files in Snapshot --- .palantir/revapi.yml | 4 +- .../java/org/apache/iceberg/Snapshot.java | 51 ++++++--- .../java/org/apache/iceberg/BaseSnapshot.java | 96 +++++++++++++---- .../apache/iceberg/CherryPickOperation.java | 6 +- .../org/apache/iceberg/util/SnapshotUtil.java | 2 +- .../apache/iceberg/TestRemoveSnapshots.java | 8 +- .../java/org/apache/iceberg/TestSnapshot.java | 100 +++++++++++++++++- .../apache/iceberg/TestSnapshotSelection.java | 2 +- .../org/apache/iceberg/TestWapWorkflow.java | 22 ++-- .../hadoop/TestCatalogUtilDropTable.java | 2 +- .../source/TestIcebergSourceTablesBase.java | 2 +- .../spark/source/SparkMicroBatchStream.java | 2 +- .../actions/TestExpireSnapshotsAction.java | 8 +- .../actions/TestRewriteDataFilesAction.java | 6 +- .../spark/source/TestDataFrameWrites.java | 2 +- .../spark/source/TestDataSourceOptions.java | 2 +- .../source/TestIcebergSourceTablesBase.java | 10 +- .../iceberg/spark/sql/TestRefreshTable.java | 2 +- .../spark/source/SparkMicroBatchStream.java | 2 +- .../actions/TestExpireSnapshotsAction.java | 8 +- .../actions/TestRewriteDataFilesAction.java | 6 +- .../spark/source/TestDataFrameWrites.java | 2 +- .../spark/source/TestDataSourceOptions.java | 2 +- .../source/TestIcebergSourceTablesBase.java | 10 +- .../iceberg/spark/sql/TestRefreshTable.java | 2 +- .../iceberg/spark/extensions/TestMerge.java | 2 +- .../iceberg/spark/extensions/TestUpdate.java | 2 +- ...SourceParquetMultiDeleteFileBenchmark.java | 2 +- ...cebergSourceParquetPosDeleteBenchmark.java | 2 +- ...ceParquetWithUnrelatedDeleteBenchmark.java | 2 +- .../actions/TestExpireSnapshotsAction.java | 8 +- .../actions/TestRewriteDataFilesAction.java | 10 +- .../spark/source/TestDataFrameWrites.java | 2 +- .../spark/source/TestDataSourceOptions.java | 2 +- .../source/TestIcebergSourceTablesBase.java | 12 +-- .../iceberg/spark/sql/TestRefreshTable.java | 2 +- .../actions/TestRewriteDataFilesAction.java | 10 +- .../spark/source/TestDataFrameWrites.java | 2 +- .../source/TestIcebergSourceTablesBase.java | 12 +-- .../iceberg/spark/sql/TestRefreshTable.java | 2 +- 40 files changed, 306 insertions(+), 125 deletions(-) diff --git a/.palantir/revapi.yml b/.palantir/revapi.yml index 49308ebb8075..5785c623b345 100644 --- a/.palantir/revapi.yml +++ b/.palantir/revapi.yml @@ -31,10 +31,10 @@ acceptedBreaks: new: "method boolean org.apache.iceberg.expressions.BoundTerm::isEquivalentTo(org.apache.iceberg.expressions.BoundTerm)" justification: "new API method" - code: "java.method.addedToInterface" - new: "method java.lang.Iterable org.apache.iceberg.Snapshot::addedFiles(org.apache.iceberg.io.FileIO)" + new: "method java.lang.Iterable org.apache.iceberg.Snapshot::addedDataFiles(org.apache.iceberg.io.FileIO)" justification: "Allow adding a new method to the interface - old method is deprecated" - code: "java.method.addedToInterface" - new: "method java.lang.Iterable org.apache.iceberg.Snapshot::deletedFiles(org.apache.iceberg.io.FileIO)" + new: "method java.lang.Iterable org.apache.iceberg.Snapshot::removedDataFiles(org.apache.iceberg.io.FileIO)" justification: "Allow adding a new method to the interface - old method is deprecated" - code: "java.method.addedToInterface" new: "method java.util.List org.apache.iceberg.Snapshot::allManifests(org.apache.iceberg.io.FileIO)" diff --git a/api/src/main/java/org/apache/iceberg/Snapshot.java b/api/src/main/java/org/apache/iceberg/Snapshot.java index bdb3da4160d9..cfaa7f9b24e3 100644 --- a/api/src/main/java/org/apache/iceberg/Snapshot.java +++ b/api/src/main/java/org/apache/iceberg/Snapshot.java @@ -108,7 +108,6 @@ public interface Snapshot extends Serializable { @Deprecated List deleteManifests(); - /** * Return a {@link ManifestFile} for each delete manifest in this snapshot. * @@ -133,50 +132,76 @@ public interface Snapshot extends Serializable { Map summary(); /** - * Return all files added to the table in this snapshot. + * Return all data files added to the table in this snapshot. *

* The files returned include the following columns: file_path, file_format, partition, * record_count, and file_size_in_bytes. Other columns will be null. * - * @return all files added to the table in this snapshot. - * @deprecated since 0.14.0, will be removed in 1.0.0; Use {@link Snapshot#addedFiles(FileIO)} instead. + * @return all data files added to the table in this snapshot. + * @deprecated since 0.14.0, will be removed in 1.0.0; Use {@link Snapshot#addedDataFiles(FileIO)} instead. */ @Deprecated Iterable addedFiles(); /** - * Return all files added to the table in this snapshot. + * Return all data files added to the table in this snapshot. *

* The files returned include the following columns: file_path, file_format, partition, * record_count, and file_size_in_bytes. Other columns will be null. * * @param io a {@link FileIO} instance used for reading files from storage - * @return all files added to the table in this snapshot. + * @return all data files added to the table in this snapshot. */ - Iterable addedFiles(FileIO io); + Iterable addedDataFiles(FileIO io); /** - * Return all files deleted from the table in this snapshot. + * Return all data files deleted from the table in this snapshot. *

* The files returned include the following columns: file_path, file_format, partition, * record_count, and file_size_in_bytes. Other columns will be null. * - * @return all files deleted from the table in this snapshot. - * @deprecated since 0.14.0, will be removed in 1.0.0; Use {@link Snapshot#deletedFiles(FileIO)} instead. + * @return all data files deleted from the table in this snapshot. + * @deprecated since 0.14.0, will be removed in 1.0.0; Use {@link Snapshot#removedDataFiles(FileIO)} instead. */ @Deprecated Iterable deletedFiles(); /** - * Return all files deleted from the table in this snapshot. + * Return all data files removed from the table in this snapshot. + *

+ * The files returned include the following columns: file_path, file_format, partition, + * record_count, and file_size_in_bytes. Other columns will be null. + * + * @param io a {@link FileIO} instance used for reading files from storage + * @return all data files removed from the table in this snapshot. + */ + Iterable removedDataFiles(FileIO io); + + /** + * Return all delete files added to the table in this snapshot. *

* The files returned include the following columns: file_path, file_format, partition, * record_count, and file_size_in_bytes. Other columns will be null. * * @param io a {@link FileIO} instance used for reading files from storage - * @return all files deleted from the table in this snapshot. + * @return all delete files added to the table in this snapshot */ - Iterable deletedFiles(FileIO io); + default Iterable addedDeleteFiles(FileIO io) { + throw new UnsupportedOperationException(this.getClass().getName() + " doesn't implement addedDeleteFiles"); + } + + /** + * Return all delete files removed from the table in this snapshot. + *

+ * The files returned include the following columns: file_path, file_format, partition, + * record_count, and file_size_in_bytes. Other columns will be null. + * + * @param io a {@link FileIO} instance used for reading files from storage + * @return all delete files removed from the table in this snapshot + */ + default Iterable removedDeleteFiles(FileIO io) { + throw new UnsupportedOperationException(this.getClass().getName() + " doesn't implement removedDeleteFiles"); + } /** * Return the location of this snapshot's manifest list, or null if it is not separate. diff --git a/core/src/main/java/org/apache/iceberg/BaseSnapshot.java b/core/src/main/java/org/apache/iceberg/BaseSnapshot.java index 6422a437db06..33d016b8ab8c 100644 --- a/core/src/main/java/org/apache/iceberg/BaseSnapshot.java +++ b/core/src/main/java/org/apache/iceberg/BaseSnapshot.java @@ -20,6 +20,7 @@ package org.apache.iceberg; import java.io.IOException; +import java.io.UncheckedIOException; import java.util.Arrays; import java.util.List; import java.util.Map; @@ -28,6 +29,7 @@ import org.apache.iceberg.io.FileIO; import org.apache.iceberg.relocated.com.google.common.base.MoreObjects; import org.apache.iceberg.relocated.com.google.common.base.Objects; +import org.apache.iceberg.relocated.com.google.common.base.Preconditions; import org.apache.iceberg.relocated.com.google.common.collect.ImmutableList; import org.apache.iceberg.relocated.com.google.common.collect.Iterables; import org.apache.iceberg.relocated.com.google.common.collect.Lists; @@ -54,8 +56,10 @@ class BaseSnapshot implements Snapshot { private transient List allManifests = null; private transient List dataManifests = null; private transient List deleteManifests = null; - private transient List cachedAdds = null; - private transient List cachedDeletes = null; + private transient List addedDataFiles = null; + private transient List removedDataFiles = null; + private transient List addedDeleteFiles = null; + private transient List removedDeleteFiles = null; /** * For testing only. @@ -235,43 +239,59 @@ public List deleteManifests() { } @Override - public List addedFiles(FileIO fileIO) { - if (cachedAdds == null) { - cacheChanges(fileIO); + public List addedDataFiles(FileIO fileIO) { + if (addedDataFiles == null) { + cacheDataFileChanges(fileIO); } - return cachedAdds; + return addedDataFiles; } /** - * @deprecated since 0.14.0, will be removed in 1.0.0; Use {@link Snapshot#addedFiles(FileIO)} instead. + * @deprecated since 0.14.0, will be removed in 1.0.0; Use {@link Snapshot#addedDataFiles(FileIO)} instead. */ @Override @Deprecated public List addedFiles() { - if (cachedAdds == null) { - cacheChanges(io); + if (addedDataFiles == null) { + cacheDataFileChanges(io); } - return cachedAdds; + return addedDataFiles; } @Override - public List deletedFiles(FileIO fileIO) { - if (cachedDeletes == null) { - cacheChanges(fileIO); + public List removedDataFiles(FileIO fileIO) { + if (removedDataFiles == null) { + cacheDataFileChanges(fileIO); } - return cachedDeletes; + return removedDataFiles; } /** - * @deprecated since 0.14.0, will be removed in 1.0.0; Use {@link Snapshot#deletedFiles(FileIO)} instead. + * @deprecated since 0.14.0, will be removed in 1.0.0; Use {@link Snapshot#removedDataFiles(FileIO)} instead. */ @Override @Deprecated public List deletedFiles() { - if (cachedDeletes == null) { - cacheChanges(io); + if (removedDataFiles == null) { + cacheDataFileChanges(io); } - return cachedDeletes; + return removedDataFiles; + } + + @Override + public Iterable addedDeleteFiles(FileIO fileIO) { + if (addedDeleteFiles == null) { + cacheDeleteFileChanges(fileIO); + } + return addedDeleteFiles; + } + + @Override + public Iterable removedDeleteFiles(FileIO fileIO) { + if (removedDeleteFiles == null) { + cacheDeleteFileChanges(fileIO); + } + return removedDeleteFiles; } @Override @@ -279,11 +299,41 @@ public String manifestListLocation() { return manifestListLocation; } - private void cacheChanges(FileIO fileIO) { - if (fileIO == null) { - throw new IllegalArgumentException("Cannot cache changes: FileIO is null"); + private void cacheDeleteFileChanges(FileIO fileIO) { + Preconditions.checkArgument(fileIO != null, "Cannot cache delete file changes: FileIO is null"); + + ImmutableList.Builder adds = ImmutableList.builder(); + ImmutableList.Builder deletes = ImmutableList.builder(); + + Iterable changedManifests = Iterables.filter(deleteManifests(fileIO), + manifest -> Objects.equal(manifest.snapshotId(), snapshotId)); + + for (ManifestFile manifest : changedManifests) { + try (ManifestReader reader = ManifestFiles.readDeleteManifest(manifest, fileIO, null)) { + for (ManifestEntry entry : reader.entries()) { + switch (entry.status()) { + case ADDED: + adds.add(entry.file().copy()); + break; + case DELETED: + deletes.add(entry.file().copyWithoutStats()); + break; + default: + // ignore existing + } + } + } catch (IOException e) { + throw new UncheckedIOException("Failed to close manifest reader", e); + } } + this.addedDeleteFiles = adds.build(); + this.removedDeleteFiles = deletes.build(); + } + + private void cacheDataFileChanges(FileIO fileIO) { + Preconditions.checkArgument(fileIO != null, "Cannot cache data file changes: FileIO is null"); + ImmutableList.Builder adds = ImmutableList.builder(); ImmutableList.Builder deletes = ImmutableList.builder(); @@ -310,8 +360,8 @@ private void cacheChanges(FileIO fileIO) { throw new RuntimeIOException(e, "Failed to close entries while caching changes"); } - this.cachedAdds = adds.build(); - this.cachedDeletes = deletes.build(); + this.addedDataFiles = adds.build(); + this.removedDataFiles = deletes.build(); } @Override diff --git a/core/src/main/java/org/apache/iceberg/CherryPickOperation.java b/core/src/main/java/org/apache/iceberg/CherryPickOperation.java index 1b96743ed34c..3e6978052281 100644 --- a/core/src/main/java/org/apache/iceberg/CherryPickOperation.java +++ b/core/src/main/java/org/apache/iceberg/CherryPickOperation.java @@ -79,7 +79,7 @@ public CherryPickOperation cherrypick(long snapshotId) { set(SnapshotSummary.SOURCE_SNAPSHOT_ID_PROP, String.valueOf(snapshotId)); // Pick modifications from the snapshot - for (DataFile addedFile : cherrypickSnapshot.addedFiles(io)) { + for (DataFile addedFile : cherrypickSnapshot.addedDataFiles(io)) { add(addedFile); } @@ -106,13 +106,13 @@ public CherryPickOperation cherrypick(long snapshotId) { // copy adds from the picked snapshot this.replacedPartitions = PartitionSet.create(specsById); - for (DataFile addedFile : cherrypickSnapshot.addedFiles(io)) { + for (DataFile addedFile : cherrypickSnapshot.addedDataFiles(io)) { add(addedFile); replacedPartitions.add(addedFile.specId(), addedFile.partition()); } // copy deletes from the picked snapshot - for (DataFile deletedFile : cherrypickSnapshot.deletedFiles(io)) { + for (DataFile deletedFile : cherrypickSnapshot.removedDataFiles(io)) { delete(deletedFile); } diff --git a/core/src/main/java/org/apache/iceberg/util/SnapshotUtil.java b/core/src/main/java/org/apache/iceberg/util/SnapshotUtil.java index dd733ddd4ced..b60af75f79c0 100644 --- a/core/src/main/java/org/apache/iceberg/util/SnapshotUtil.java +++ b/core/src/main/java/org/apache/iceberg/util/SnapshotUtil.java @@ -271,7 +271,7 @@ public static List newFiles( return newFiles; } - Iterables.addAll(newFiles, currentSnapshot.addedFiles(io)); + Iterables.addAll(newFiles, currentSnapshot.addedDataFiles(io)); } ValidationException.check(Objects.equals(lastSnapshot.parentId(), baseSnapshotId), diff --git a/core/src/test/java/org/apache/iceberg/TestRemoveSnapshots.java b/core/src/test/java/org/apache/iceberg/TestRemoveSnapshots.java index 5e492e3972f4..9d364bc24225 100644 --- a/core/src/test/java/org/apache/iceberg/TestRemoveSnapshots.java +++ b/core/src/test/java/org/apache/iceberg/TestRemoveSnapshots.java @@ -907,7 +907,7 @@ public void testWithExpiringDanglingStageCommit() { expectedDeletes.add(snapshotA.manifestListLocation()); // Files should be deleted of dangling staged snapshot - snapshotB.addedFiles(table.io()).forEach(i -> { + snapshotB.addedDataFiles(table.io()).forEach(i -> { expectedDeletes.add(i.path().toString()); }); @@ -982,7 +982,7 @@ public void testWithCherryPickTableSnapshot() { // Make sure no dataFiles are deleted for the B, C, D snapshot Lists.newArrayList(snapshotB, snapshotC, snapshotD).forEach(i -> { - i.addedFiles(table.io()).forEach(item -> { + i.addedDataFiles(table.io()).forEach(item -> { Assert.assertFalse(deletedFiles.contains(item.path().toString())); }); }); @@ -1035,7 +1035,7 @@ public void testWithExpiringStagedThenCherrypick() { // Make sure no dataFiles are deleted for the staged snapshot Lists.newArrayList(snapshotB).forEach(i -> { - i.addedFiles(table.io()).forEach(item -> { + i.addedDataFiles(table.io()).forEach(item -> { Assert.assertFalse(deletedFiles.contains(item.path().toString())); }); }); @@ -1048,7 +1048,7 @@ public void testWithExpiringStagedThenCherrypick() { // Make sure no dataFiles are deleted for the staged and cherry-pick Lists.newArrayList(snapshotB, snapshotD).forEach(i -> { - i.addedFiles(table.io()).forEach(item -> { + i.addedDataFiles(table.io()).forEach(item -> { Assert.assertFalse(deletedFiles.contains(item.path().toString())); }); }); diff --git a/core/src/test/java/org/apache/iceberg/TestSnapshot.java b/core/src/test/java/org/apache/iceberg/TestSnapshot.java index 11d8f0ac2f14..60de9652ed0e 100644 --- a/core/src/test/java/org/apache/iceberg/TestSnapshot.java +++ b/core/src/test/java/org/apache/iceberg/TestSnapshot.java @@ -19,6 +19,11 @@ package org.apache.iceberg; +import org.apache.iceberg.expressions.Expressions; +import org.apache.iceberg.relocated.com.google.common.collect.ImmutableSet; +import org.apache.iceberg.relocated.com.google.common.collect.Iterables; +import org.junit.Assert; +import org.junit.Assume; import org.junit.Test; import org.junit.runner.RunWith; import org.junit.runners.Parameterized; @@ -42,7 +47,7 @@ public void testAppendFilesFromTable() { .commit(); // collect data files from deserialization - Iterable filesToAdd = table.currentSnapshot().addedFiles(table.io()); + Iterable filesToAdd = table.currentSnapshot().addedDataFiles(table.io()); table.newDelete().deleteFile(FILE_A).deleteFile(FILE_B).commit(); @@ -82,4 +87,97 @@ public void testAppendFoundFiles() { validateSnapshot(oldSnapshot, newSnapshot, FILE_A, FILE_B); } + @Test + public void testCachedDataFiles() { + table.newFastAppend() + .appendFile(FILE_A) + .appendFile(FILE_B) + .commit(); + + table.updateSpec() + .addField(Expressions.truncate("data", 2)) + .commit(); + + DataFile secondSnapshotDataFile = newDataFile("data_bucket=8/data_trunc_2=aa"); + + table.newFastAppend() + .appendFile(secondSnapshotDataFile) + .commit(); + + DataFile thirdSnapshotDataFile = newDataFile("data_bucket=8/data_trunc_2=bb"); + + table.newOverwrite() + .deleteFile(FILE_A) + .addFile(thirdSnapshotDataFile) + .commit(); + + Snapshot thirdSnapshot = table.currentSnapshot(); + + Iterable deletedDataFiles = thirdSnapshot.removedDataFiles(FILE_IO); + Assert.assertEquals("Must have 1 deleted data file", 1, Iterables.size(deletedDataFiles)); + + DataFile deletedDataFile = Iterables.getOnlyElement(deletedDataFiles); + Assert.assertEquals("Path must match", FILE_A.path(), deletedDataFile.path()); + Assert.assertEquals("Spec ID must match", FILE_A.specId(), deletedDataFile.specId()); + Assert.assertEquals("Partition must match", FILE_A.partition(), deletedDataFile.partition()); + + Iterable addedDataFiles = thirdSnapshot.addedDataFiles(FILE_IO); + Assert.assertEquals("Must have 1 added data file", 1, Iterables.size(addedDataFiles)); + + DataFile addedDataFile = Iterables.getOnlyElement(addedDataFiles); + Assert.assertEquals("Path must match", thirdSnapshotDataFile.path(), addedDataFile.path()); + Assert.assertEquals("Spec ID must match", thirdSnapshotDataFile.specId(), addedDataFile.specId()); + Assert.assertEquals("Partition must match", thirdSnapshotDataFile.partition(), addedDataFile.partition()); + } + + @Test + public void testCachedDeleteFiles() { + Assume.assumeTrue("Delete files only supported in V2", formatVersion >= 2); + + table.newFastAppend() + .appendFile(FILE_A) + .appendFile(FILE_B) + .commit(); + + table.updateSpec() + .addField(Expressions.truncate("data", 2)) + .commit(); + + int specId = table.spec().specId(); + + DataFile secondSnapshotDataFile = newDataFile("data_bucket=8/data_trunc_2=aa"); + DeleteFile secondSnapshotDeleteFile = newDeleteFile(specId, "data_bucket=8/data_trunc_2=aa"); + + table.newRowDelta() + .addRows(secondSnapshotDataFile) + .addDeletes(secondSnapshotDeleteFile) + .commit(); + + DeleteFile thirdSnapshotDeleteFile = newDeleteFile(specId, "data_bucket=8/data_trunc_2=aa"); + + ImmutableSet replacedDeleteFiles = ImmutableSet.of(secondSnapshotDeleteFile); + ImmutableSet newDeleteFiles = ImmutableSet.of(thirdSnapshotDeleteFile); + + table.newRewrite() + .rewriteFiles(ImmutableSet.of(), replacedDeleteFiles, ImmutableSet.of(), newDeleteFiles) + .commit(); + + Snapshot thirdSnapshot = table.currentSnapshot(); + + Iterable deletedDeleteFiles = thirdSnapshot.removedDeleteFiles(FILE_IO); + Assert.assertEquals("Must have 1 deleted delete file", 1, Iterables.size(deletedDeleteFiles)); + + DeleteFile deletedDeleteFile = Iterables.getOnlyElement(deletedDeleteFiles); + Assert.assertEquals("Path must match", secondSnapshotDeleteFile.path(), deletedDeleteFile.path()); + Assert.assertEquals("Spec ID must match", secondSnapshotDeleteFile.specId(), deletedDeleteFile.specId()); + Assert.assertEquals("Partition must match", secondSnapshotDeleteFile.partition(), deletedDeleteFile.partition()); + + Iterable addedDeleteFiles = thirdSnapshot.addedDeleteFiles(FILE_IO); + Assert.assertEquals("Must have 1 added delete file", 1, Iterables.size(addedDeleteFiles)); + + DeleteFile addedDeleteFile = Iterables.getOnlyElement(addedDeleteFiles); + Assert.assertEquals("Path must match", thirdSnapshotDeleteFile.path(), addedDeleteFile.path()); + Assert.assertEquals("Spec ID must match", thirdSnapshotDeleteFile.specId(), addedDeleteFile.specId()); + Assert.assertEquals("Partition must match", thirdSnapshotDeleteFile.partition(), addedDeleteFile.partition()); + } } diff --git a/core/src/test/java/org/apache/iceberg/TestSnapshotSelection.java b/core/src/test/java/org/apache/iceberg/TestSnapshotSelection.java index d828baaea4b6..8002ecb78c66 100644 --- a/core/src/test/java/org/apache/iceberg/TestSnapshotSelection.java +++ b/core/src/test/java/org/apache/iceberg/TestSnapshotSelection.java @@ -79,7 +79,7 @@ public void testSnapshotStatsForAddedFiles() { .commit(); Snapshot snapshot = table.currentSnapshot(); - Iterable addedFiles = snapshot.addedFiles(table.io()); + Iterable addedFiles = snapshot.addedDataFiles(table.io()); Assert.assertEquals(1, Iterables.size(addedFiles)); DataFile dataFile = Iterables.getOnlyElement(addedFiles); Assert.assertNotNull("Value counts should be not null", dataFile.valueCounts()); diff --git a/core/src/test/java/org/apache/iceberg/TestWapWorkflow.java b/core/src/test/java/org/apache/iceberg/TestWapWorkflow.java index c1efec46018f..3d78fc0faf0f 100644 --- a/core/src/test/java/org/apache/iceberg/TestWapWorkflow.java +++ b/core/src/test/java/org/apache/iceberg/TestWapWorkflow.java @@ -143,7 +143,7 @@ public void testCurrentSnapshotOperation() { Assert.assertEquals("Should contain manifests for both files", 2, base.currentSnapshot().allManifests(table.io()).size()); Assert.assertEquals("Should contain append from last commit", 1, - Iterables.size(base.currentSnapshot().addedFiles(table.io()))); + Iterables.size(base.currentSnapshot().addedDataFiles(table.io()))); Assert.assertEquals("Snapshot log should indicate number of snapshots committed", 2, base.snapshotLog().size()); } @@ -172,7 +172,7 @@ public void testSetCurrentSnapshotNoWAP() { Assert.assertEquals("Should contain manifests for both files", 1, base.currentSnapshot().allManifests(table.io()).size()); Assert.assertEquals("Should contain append from last commit", 1, - Iterables.size(base.currentSnapshot().addedFiles(table.io()))); + Iterables.size(base.currentSnapshot().addedDataFiles(table.io()))); Assert.assertEquals("Snapshot log should indicate number of snapshots committed", 3, base.snapshotLog().size()); } @@ -218,7 +218,7 @@ public void testRollbackOnInvalidNonAncestor() { Assert.assertEquals("Should contain manifests for one snapshot", 1, base.currentSnapshot().allManifests(table.io()).size()); Assert.assertEquals("Should contain append from last commit", 1, - Iterables.size(base.currentSnapshot().addedFiles(table.io()))); + Iterables.size(base.currentSnapshot().addedDataFiles(table.io()))); Assert.assertEquals("Snapshot log should indicate number of snapshots committed", 1, base.snapshotLog().size()); } @@ -339,7 +339,7 @@ public void testWithCherryPicking() { Assert.assertEquals("Should contain manifests for both files", 2, base.currentSnapshot().allManifests(table.io()).size()); Assert.assertEquals("Should contain append from last commit", 1, - Iterables.size(base.currentSnapshot().addedFiles(table.io()))); + Iterables.size(base.currentSnapshot().addedDataFiles(table.io()))); Assert.assertEquals("Snapshot log should indicate number of snapshots committed", 2, base.snapshotLog().size()); } @@ -403,7 +403,7 @@ public void testWithTwoPhaseCherryPicking() { Assert.assertEquals("Should contain manifests for both files", 2, base.currentSnapshot().allManifests(table.io()).size()); Assert.assertEquals("Should contain append from last commit", 1, - Iterables.size(base.currentSnapshot().addedFiles(table.io()))); + Iterables.size(base.currentSnapshot().addedDataFiles(table.io()))); Assert.assertEquals("Parent snapshot id should change to latest snapshot before commit", parentSnapshot.snapshotId(), base.currentSnapshot().parentId().longValue()); Assert.assertEquals("Snapshot log should indicate number of snapshots committed", 2, @@ -423,7 +423,7 @@ public void testWithTwoPhaseCherryPicking() { Assert.assertEquals("Should contain manifests for both files", 3, base.currentSnapshot().allManifests(table.io()).size()); Assert.assertEquals("Should contain append from last commit", 1, - Iterables.size(base.currentSnapshot().addedFiles(table.io()))); + Iterables.size(base.currentSnapshot().addedDataFiles(table.io()))); Assert.assertEquals("Parent snapshot id should change to latest snapshot before commit", parentSnapshot.snapshotId(), base.currentSnapshot().parentId().longValue()); Assert.assertEquals("Snapshot log should indicate number of snapshots committed", 3, @@ -504,7 +504,7 @@ public void testWithCommitsBetweenCherryPicking() { Assert.assertEquals("Should contain manifests for three files", 3, base.currentSnapshot().allManifests(table.io()).size()); Assert.assertEquals("Should contain append from last commit", 1, - Iterables.size(base.currentSnapshot().addedFiles(table.io()))); + Iterables.size(base.currentSnapshot().addedDataFiles(table.io()))); Assert.assertEquals("Parent snapshot id should point to same snapshot", parentSnapshot.snapshotId(), base.currentSnapshot().parentId().longValue()); Assert.assertEquals("Snapshot log should indicate number of snapshots committed", 3, @@ -524,7 +524,7 @@ public void testWithCommitsBetweenCherryPicking() { Assert.assertEquals("Should contain manifests for four files", 4, base.currentSnapshot().allManifests(table.io()).size()); Assert.assertEquals("Should contain append from last commit", 1, - Iterables.size(base.currentSnapshot().addedFiles(table.io()))); + Iterables.size(base.currentSnapshot().addedDataFiles(table.io()))); Assert.assertEquals("Parent snapshot id should point to same snapshot", parentSnapshot.snapshotId(), base.currentSnapshot().parentId().longValue()); Assert.assertEquals("Snapshot log should indicate number of snapshots committed", 4, @@ -581,7 +581,7 @@ public void testWithCherryPickingWithCommitRetry() { Assert.assertEquals("Should contain manifests for both files", 2, base.currentSnapshot().allManifests(table.io()).size()); Assert.assertEquals("Should not contain redundant append due to retry", 1, - Iterables.size(base.currentSnapshot().addedFiles(table.io()))); + Iterables.size(base.currentSnapshot().addedDataFiles(table.io()))); Assert.assertEquals("Parent snapshot id should change to latest snapshot before commit", parentSnapshot.snapshotId(), base.currentSnapshot().parentId().longValue()); Assert.assertEquals("Snapshot log should indicate number of snapshots committed", 2, @@ -629,7 +629,7 @@ public void testCherrypickingAncestor() { Assert.assertEquals("Should contain manifests for both files", 2, base.currentSnapshot().allManifests(table.io()).size()); Assert.assertEquals("Should contain append from last commit", 1, - Iterables.size(base.currentSnapshot().addedFiles(table.io()))); + Iterables.size(base.currentSnapshot().addedDataFiles(table.io()))); Assert.assertEquals("Snapshot log should indicate number of snapshots committed", 2, base.snapshotLog().size()); @@ -685,7 +685,7 @@ public void testDuplicateCherrypick() { Assert.assertEquals("Should contain manifests for both files", 2, base.currentSnapshot().allManifests(table.io()).size()); Assert.assertEquals("Should contain append from last commit", 1, - Iterables.size(base.currentSnapshot().addedFiles(table.io()))); + Iterables.size(base.currentSnapshot().addedDataFiles(table.io()))); Assert.assertEquals("Snapshot log should indicate number of snapshots committed", 2, base.snapshotLog().size()); diff --git a/core/src/test/java/org/apache/iceberg/hadoop/TestCatalogUtilDropTable.java b/core/src/test/java/org/apache/iceberg/hadoop/TestCatalogUtilDropTable.java index 4dbe28378887..4c114a280465 100644 --- a/core/src/test/java/org/apache/iceberg/hadoop/TestCatalogUtilDropTable.java +++ b/core/src/test/java/org/apache/iceberg/hadoop/TestCatalogUtilDropTable.java @@ -150,7 +150,7 @@ private Set manifestLocations(Set snapshotSet, FileIO io) { private Set dataLocations(Set snapshotSet, FileIO io) { return snapshotSet.stream() - .flatMap(snapshot -> StreamSupport.stream(snapshot.addedFiles(io).spliterator(), false)) + .flatMap(snapshot -> StreamSupport.stream(snapshot.addedDataFiles(io).spliterator(), false)) .map(dataFile -> dataFile.path().toString()) .collect(Collectors.toSet()); } diff --git a/spark/v2.4/spark/src/test/java/org/apache/iceberg/spark/source/TestIcebergSourceTablesBase.java b/spark/v2.4/spark/src/test/java/org/apache/iceberg/spark/source/TestIcebergSourceTablesBase.java index 4d2cde2f9d4a..6b62d2e9bef1 100644 --- a/spark/v2.4/spark/src/test/java/org/apache/iceberg/spark/source/TestIcebergSourceTablesBase.java +++ b/spark/v2.4/spark/src/test/java/org/apache/iceberg/spark/source/TestIcebergSourceTablesBase.java @@ -1014,7 +1014,7 @@ public void testAllManifestsTable() { .set(TableProperties.FORMAT_VERSION, "2") .commit(); - DataFile dataFile = Iterables.getFirst(table.currentSnapshot().addedFiles(table.io()), null); + DataFile dataFile = Iterables.getFirst(table.currentSnapshot().addedDataFiles(table.io()), null); PartitionSpec dataFileSpec = table.specs().get(dataFile.specId()); StructLike dataFilePartition = dataFile.partition(); diff --git a/spark/v3.0/spark/src/main/java/org/apache/iceberg/spark/source/SparkMicroBatchStream.java b/spark/v3.0/spark/src/main/java/org/apache/iceberg/spark/source/SparkMicroBatchStream.java index d72928b6b75d..794b6a9462e5 100644 --- a/spark/v3.0/spark/src/main/java/org/apache/iceberg/spark/source/SparkMicroBatchStream.java +++ b/spark/v3.0/spark/src/main/java/org/apache/iceberg/spark/source/SparkMicroBatchStream.java @@ -112,7 +112,7 @@ public Offset latestOffset() { Snapshot latestSnapshot = table.currentSnapshot(); return new StreamingOffset( - latestSnapshot.snapshotId(), Iterables.size(latestSnapshot.addedFiles(table.io())), false); + latestSnapshot.snapshotId(), Iterables.size(latestSnapshot.addedDataFiles(table.io())), false); } @Override diff --git a/spark/v3.0/spark/src/test/java/org/apache/iceberg/spark/actions/TestExpireSnapshotsAction.java b/spark/v3.0/spark/src/test/java/org/apache/iceberg/spark/actions/TestExpireSnapshotsAction.java index 9b8c9d2501db..d411abdb8e26 100644 --- a/spark/v3.0/spark/src/test/java/org/apache/iceberg/spark/actions/TestExpireSnapshotsAction.java +++ b/spark/v3.0/spark/src/test/java/org/apache/iceberg/spark/actions/TestExpireSnapshotsAction.java @@ -570,7 +570,7 @@ public void testWithExpiringDanglingStageCommit() { expectedDeletes.add(snapshotA.manifestListLocation()); // Files should be deleted of dangling staged snapshot - snapshotB.addedFiles(table.io()).forEach(i -> { + snapshotB.addedDataFiles(table.io()).forEach(i -> { expectedDeletes.add(i.path().toString()); }); @@ -645,7 +645,7 @@ public void testWithCherryPickTableSnapshot() { // Make sure no dataFiles are deleted for the B, C, D snapshot Lists.newArrayList(snapshotB, snapshotC, snapshotD).forEach(i -> { - i.addedFiles(table.io()).forEach(item -> { + i.addedDataFiles(table.io()).forEach(item -> { Assert.assertFalse(deletedFiles.contains(item.path().toString())); }); }); @@ -700,7 +700,7 @@ public void testWithExpiringStagedThenCherrypick() { // Make sure no dataFiles are deleted for the staged snapshot Lists.newArrayList(snapshotB).forEach(i -> { - i.addedFiles(table.io()).forEach(item -> { + i.addedDataFiles(table.io()).forEach(item -> { Assert.assertFalse(deletedFiles.contains(item.path().toString())); }); }); @@ -714,7 +714,7 @@ public void testWithExpiringStagedThenCherrypick() { // Make sure no dataFiles are deleted for the staged and cherry-pick Lists.newArrayList(snapshotB, snapshotD).forEach(i -> { - i.addedFiles(table.io()).forEach(item -> { + i.addedDataFiles(table.io()).forEach(item -> { Assert.assertFalse(deletedFiles.contains(item.path().toString())); }); }); diff --git a/spark/v3.0/spark/src/test/java/org/apache/iceberg/spark/actions/TestRewriteDataFilesAction.java b/spark/v3.0/spark/src/test/java/org/apache/iceberg/spark/actions/TestRewriteDataFilesAction.java index 9d460276525a..b057b8e62ba7 100644 --- a/spark/v3.0/spark/src/test/java/org/apache/iceberg/spark/actions/TestRewriteDataFilesAction.java +++ b/spark/v3.0/spark/src/test/java/org/apache/iceberg/spark/actions/TestRewriteDataFilesAction.java @@ -44,6 +44,7 @@ import org.apache.iceberg.RewriteJobOrder; import org.apache.iceberg.RowDelta; import org.apache.iceberg.Schema; +import org.apache.iceberg.Snapshot; import org.apache.iceberg.SortOrder; import org.apache.iceberg.StructLike; import org.apache.iceberg.Table; @@ -963,7 +964,7 @@ public void testAutoSortShuffleOutput() { Assert.assertEquals("Should have 1 fileGroups", result.rewriteResults().size(), 1); Assert.assertTrue("Should have written 40+ files", - Iterables.size(table.currentSnapshot().addedFiles(table.io())) >= 40); + Iterables.size(table.currentSnapshot().addedDataFiles(table.io())) >= 40); table.refresh(); @@ -1224,7 +1225,8 @@ private List, Pair>> checkForOverlappingFiles(Table ta NestedField field = table.schema().caseInsensitiveFindField(column); Class javaClass = (Class) field.type().typeId().javaClass(); - Map> filesByPartition = Streams.stream(table.currentSnapshot().addedFiles(table.io())) + Snapshot snapshot = table.currentSnapshot(); + Map> filesByPartition = Streams.stream(snapshot.addedDataFiles(table.io())) .collect(Collectors.groupingBy(DataFile::partition)); Stream, Pair>> overlaps = diff --git a/spark/v3.0/spark/src/test/java/org/apache/iceberg/spark/source/TestDataFrameWrites.java b/spark/v3.0/spark/src/test/java/org/apache/iceberg/spark/source/TestDataFrameWrites.java index f6d292f89f8b..b0a77b72b431 100644 --- a/spark/v3.0/spark/src/test/java/org/apache/iceberg/spark/source/TestDataFrameWrites.java +++ b/spark/v3.0/spark/src/test/java/org/apache/iceberg/spark/source/TestDataFrameWrites.java @@ -181,7 +181,7 @@ private void writeAndValidateWithLocations(Table table, File location, File expe } Assert.assertEquals("Both iterators should be exhausted", expectedIter.hasNext(), actualIter.hasNext()); - table.currentSnapshot().addedFiles(table.io()).forEach(dataFile -> + table.currentSnapshot().addedDataFiles(table.io()).forEach(dataFile -> Assert.assertTrue( String.format( "File should have the parent directory %s, but has: %s.", diff --git a/spark/v3.0/spark/src/test/java/org/apache/iceberg/spark/source/TestDataSourceOptions.java b/spark/v3.0/spark/src/test/java/org/apache/iceberg/spark/source/TestDataSourceOptions.java index 4fa8fdfc6d75..ffcb86052074 100644 --- a/spark/v3.0/spark/src/test/java/org/apache/iceberg/spark/source/TestDataSourceOptions.java +++ b/spark/v3.0/spark/src/test/java/org/apache/iceberg/spark/source/TestDataSourceOptions.java @@ -207,7 +207,7 @@ public void testSplitOptionsOverridesTableProperties() throws IOException { .mode("append") .save(tableLocation); - List files = Lists.newArrayList(icebergTable.currentSnapshot().addedFiles(icebergTable.io())); + List files = Lists.newArrayList(icebergTable.currentSnapshot().addedDataFiles(icebergTable.io())); Assert.assertEquals("Should have written 1 file", 1, files.size()); long fileSize = files.get(0).fileSizeInBytes(); diff --git a/spark/v3.0/spark/src/test/java/org/apache/iceberg/spark/source/TestIcebergSourceTablesBase.java b/spark/v3.0/spark/src/test/java/org/apache/iceberg/spark/source/TestIcebergSourceTablesBase.java index b3ceb192bbb6..7258fbd9690c 100644 --- a/spark/v3.0/spark/src/test/java/org/apache/iceberg/spark/source/TestIcebergSourceTablesBase.java +++ b/spark/v3.0/spark/src/test/java/org/apache/iceberg/spark/source/TestIcebergSourceTablesBase.java @@ -219,7 +219,7 @@ public void testEntriesTableDataFilePrune() throws Exception { .save(loadLocation(tableIdentifier)); table.refresh(); - DataFile file = table.currentSnapshot().addedFiles(table.io()).iterator().next(); + DataFile file = table.currentSnapshot().addedDataFiles(table.io()).iterator().next(); List singleActual = rowsToJava(spark.read() .format("iceberg") @@ -246,7 +246,7 @@ public void testEntriesTableDataFilePruneMulti() throws Exception { .save(loadLocation(tableIdentifier)); table.refresh(); - DataFile file = table.currentSnapshot().addedFiles(table.io()).iterator().next(); + DataFile file = table.currentSnapshot().addedDataFiles(table.io()).iterator().next(); List multiActual = rowsToJava(spark.read() .format("iceberg") @@ -274,7 +274,7 @@ public void testFilesSelectMap() throws Exception { .save(loadLocation(tableIdentifier)); table.refresh(); - DataFile file = table.currentSnapshot().addedFiles(table.io()).iterator().next(); + DataFile file = table.currentSnapshot().addedDataFiles(table.io()).iterator().next(); List multiActual = rowsToJava(spark.read() .format("iceberg") @@ -544,7 +544,7 @@ public void testFilesUnpartitionedTable() throws Exception { .save(loadLocation(tableIdentifier)); table.refresh(); - DataFile toDelete = Iterables.getOnlyElement(table.currentSnapshot().addedFiles(table.io())); + DataFile toDelete = Iterables.getOnlyElement(table.currentSnapshot().addedDataFiles(table.io())); // add a second file df2.select("id", "data").write() @@ -1017,7 +1017,7 @@ public void testAllManifestsTable() { .set(TableProperties.FORMAT_VERSION, "2") .commit(); - DataFile dataFile = Iterables.getFirst(table.currentSnapshot().addedFiles(table.io()), null); + DataFile dataFile = Iterables.getFirst(table.currentSnapshot().addedDataFiles(table.io()), null); PartitionSpec dataFileSpec = table.specs().get(dataFile.specId()); StructLike dataFilePartition = dataFile.partition(); diff --git a/spark/v3.0/spark/src/test/java/org/apache/iceberg/spark/sql/TestRefreshTable.java b/spark/v3.0/spark/src/test/java/org/apache/iceberg/spark/sql/TestRefreshTable.java index 510d7c40eecb..7489c3963d3a 100644 --- a/spark/v3.0/spark/src/test/java/org/apache/iceberg/spark/sql/TestRefreshTable.java +++ b/spark/v3.0/spark/src/test/java/org/apache/iceberg/spark/sql/TestRefreshTable.java @@ -66,7 +66,7 @@ public void testRefreshCommand() { // Modify table outside of spark, it should be cached so Spark should see the same value after mutation Table table = validationCatalog.loadTable(tableIdent); - DataFile file = table.currentSnapshot().addedFiles(table.io()).iterator().next(); + DataFile file = table.currentSnapshot().addedDataFiles(table.io()).iterator().next(); table.newDelete().deleteFile(file).commit(); List cachedActual = sql("SELECT * FROM %s", tableName); diff --git a/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/source/SparkMicroBatchStream.java b/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/source/SparkMicroBatchStream.java index d72928b6b75d..794b6a9462e5 100644 --- a/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/source/SparkMicroBatchStream.java +++ b/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/source/SparkMicroBatchStream.java @@ -112,7 +112,7 @@ public Offset latestOffset() { Snapshot latestSnapshot = table.currentSnapshot(); return new StreamingOffset( - latestSnapshot.snapshotId(), Iterables.size(latestSnapshot.addedFiles(table.io())), false); + latestSnapshot.snapshotId(), Iterables.size(latestSnapshot.addedDataFiles(table.io())), false); } @Override diff --git a/spark/v3.1/spark/src/test/java/org/apache/iceberg/spark/actions/TestExpireSnapshotsAction.java b/spark/v3.1/spark/src/test/java/org/apache/iceberg/spark/actions/TestExpireSnapshotsAction.java index 9b8c9d2501db..d411abdb8e26 100644 --- a/spark/v3.1/spark/src/test/java/org/apache/iceberg/spark/actions/TestExpireSnapshotsAction.java +++ b/spark/v3.1/spark/src/test/java/org/apache/iceberg/spark/actions/TestExpireSnapshotsAction.java @@ -570,7 +570,7 @@ public void testWithExpiringDanglingStageCommit() { expectedDeletes.add(snapshotA.manifestListLocation()); // Files should be deleted of dangling staged snapshot - snapshotB.addedFiles(table.io()).forEach(i -> { + snapshotB.addedDataFiles(table.io()).forEach(i -> { expectedDeletes.add(i.path().toString()); }); @@ -645,7 +645,7 @@ public void testWithCherryPickTableSnapshot() { // Make sure no dataFiles are deleted for the B, C, D snapshot Lists.newArrayList(snapshotB, snapshotC, snapshotD).forEach(i -> { - i.addedFiles(table.io()).forEach(item -> { + i.addedDataFiles(table.io()).forEach(item -> { Assert.assertFalse(deletedFiles.contains(item.path().toString())); }); }); @@ -700,7 +700,7 @@ public void testWithExpiringStagedThenCherrypick() { // Make sure no dataFiles are deleted for the staged snapshot Lists.newArrayList(snapshotB).forEach(i -> { - i.addedFiles(table.io()).forEach(item -> { + i.addedDataFiles(table.io()).forEach(item -> { Assert.assertFalse(deletedFiles.contains(item.path().toString())); }); }); @@ -714,7 +714,7 @@ public void testWithExpiringStagedThenCherrypick() { // Make sure no dataFiles are deleted for the staged and cherry-pick Lists.newArrayList(snapshotB, snapshotD).forEach(i -> { - i.addedFiles(table.io()).forEach(item -> { + i.addedDataFiles(table.io()).forEach(item -> { Assert.assertFalse(deletedFiles.contains(item.path().toString())); }); }); diff --git a/spark/v3.1/spark/src/test/java/org/apache/iceberg/spark/actions/TestRewriteDataFilesAction.java b/spark/v3.1/spark/src/test/java/org/apache/iceberg/spark/actions/TestRewriteDataFilesAction.java index 9d460276525a..b057b8e62ba7 100644 --- a/spark/v3.1/spark/src/test/java/org/apache/iceberg/spark/actions/TestRewriteDataFilesAction.java +++ b/spark/v3.1/spark/src/test/java/org/apache/iceberg/spark/actions/TestRewriteDataFilesAction.java @@ -44,6 +44,7 @@ import org.apache.iceberg.RewriteJobOrder; import org.apache.iceberg.RowDelta; import org.apache.iceberg.Schema; +import org.apache.iceberg.Snapshot; import org.apache.iceberg.SortOrder; import org.apache.iceberg.StructLike; import org.apache.iceberg.Table; @@ -963,7 +964,7 @@ public void testAutoSortShuffleOutput() { Assert.assertEquals("Should have 1 fileGroups", result.rewriteResults().size(), 1); Assert.assertTrue("Should have written 40+ files", - Iterables.size(table.currentSnapshot().addedFiles(table.io())) >= 40); + Iterables.size(table.currentSnapshot().addedDataFiles(table.io())) >= 40); table.refresh(); @@ -1224,7 +1225,8 @@ private List, Pair>> checkForOverlappingFiles(Table ta NestedField field = table.schema().caseInsensitiveFindField(column); Class javaClass = (Class) field.type().typeId().javaClass(); - Map> filesByPartition = Streams.stream(table.currentSnapshot().addedFiles(table.io())) + Snapshot snapshot = table.currentSnapshot(); + Map> filesByPartition = Streams.stream(snapshot.addedDataFiles(table.io())) .collect(Collectors.groupingBy(DataFile::partition)); Stream, Pair>> overlaps = diff --git a/spark/v3.1/spark/src/test/java/org/apache/iceberg/spark/source/TestDataFrameWrites.java b/spark/v3.1/spark/src/test/java/org/apache/iceberg/spark/source/TestDataFrameWrites.java index f6d292f89f8b..b0a77b72b431 100644 --- a/spark/v3.1/spark/src/test/java/org/apache/iceberg/spark/source/TestDataFrameWrites.java +++ b/spark/v3.1/spark/src/test/java/org/apache/iceberg/spark/source/TestDataFrameWrites.java @@ -181,7 +181,7 @@ private void writeAndValidateWithLocations(Table table, File location, File expe } Assert.assertEquals("Both iterators should be exhausted", expectedIter.hasNext(), actualIter.hasNext()); - table.currentSnapshot().addedFiles(table.io()).forEach(dataFile -> + table.currentSnapshot().addedDataFiles(table.io()).forEach(dataFile -> Assert.assertTrue( String.format( "File should have the parent directory %s, but has: %s.", diff --git a/spark/v3.1/spark/src/test/java/org/apache/iceberg/spark/source/TestDataSourceOptions.java b/spark/v3.1/spark/src/test/java/org/apache/iceberg/spark/source/TestDataSourceOptions.java index 4fa8fdfc6d75..ffcb86052074 100644 --- a/spark/v3.1/spark/src/test/java/org/apache/iceberg/spark/source/TestDataSourceOptions.java +++ b/spark/v3.1/spark/src/test/java/org/apache/iceberg/spark/source/TestDataSourceOptions.java @@ -207,7 +207,7 @@ public void testSplitOptionsOverridesTableProperties() throws IOException { .mode("append") .save(tableLocation); - List files = Lists.newArrayList(icebergTable.currentSnapshot().addedFiles(icebergTable.io())); + List files = Lists.newArrayList(icebergTable.currentSnapshot().addedDataFiles(icebergTable.io())); Assert.assertEquals("Should have written 1 file", 1, files.size()); long fileSize = files.get(0).fileSizeInBytes(); diff --git a/spark/v3.1/spark/src/test/java/org/apache/iceberg/spark/source/TestIcebergSourceTablesBase.java b/spark/v3.1/spark/src/test/java/org/apache/iceberg/spark/source/TestIcebergSourceTablesBase.java index 72f8ea38f54c..5e8776c4c68a 100644 --- a/spark/v3.1/spark/src/test/java/org/apache/iceberg/spark/source/TestIcebergSourceTablesBase.java +++ b/spark/v3.1/spark/src/test/java/org/apache/iceberg/spark/source/TestIcebergSourceTablesBase.java @@ -219,7 +219,7 @@ public void testEntriesTableDataFilePrune() throws Exception { .save(loadLocation(tableIdentifier)); table.refresh(); - DataFile file = table.currentSnapshot().addedFiles(table.io()).iterator().next(); + DataFile file = table.currentSnapshot().addedDataFiles(table.io()).iterator().next(); List singleActual = rowsToJava(spark.read() .format("iceberg") @@ -246,7 +246,7 @@ public void testEntriesTableDataFilePruneMulti() throws Exception { .save(loadLocation(tableIdentifier)); table.refresh(); - DataFile file = table.currentSnapshot().addedFiles(table.io()).iterator().next(); + DataFile file = table.currentSnapshot().addedDataFiles(table.io()).iterator().next(); List multiActual = rowsToJava(spark.read() .format("iceberg") @@ -274,7 +274,7 @@ public void testFilesSelectMap() throws Exception { .save(loadLocation(tableIdentifier)); table.refresh(); - DataFile file = table.currentSnapshot().addedFiles(table.io()).iterator().next(); + DataFile file = table.currentSnapshot().addedDataFiles(table.io()).iterator().next(); List multiActual = rowsToJava(spark.read() .format("iceberg") @@ -544,7 +544,7 @@ public void testFilesUnpartitionedTable() throws Exception { .save(loadLocation(tableIdentifier)); table.refresh(); - DataFile toDelete = Iterables.getOnlyElement(table.currentSnapshot().addedFiles(table.io())); + DataFile toDelete = Iterables.getOnlyElement(table.currentSnapshot().addedDataFiles(table.io())); // add a second file df2.select("id", "data").write() @@ -1017,7 +1017,7 @@ public void testAllManifestsTable() { .set(TableProperties.FORMAT_VERSION, "2") .commit(); - DataFile dataFile = Iterables.getFirst(table.currentSnapshot().addedFiles(table.io()), null); + DataFile dataFile = Iterables.getFirst(table.currentSnapshot().addedDataFiles(table.io()), null); PartitionSpec dataFileSpec = table.specs().get(dataFile.specId()); StructLike dataFilePartition = dataFile.partition(); diff --git a/spark/v3.1/spark/src/test/java/org/apache/iceberg/spark/sql/TestRefreshTable.java b/spark/v3.1/spark/src/test/java/org/apache/iceberg/spark/sql/TestRefreshTable.java index 187dd4470a05..a8bdea77e237 100644 --- a/spark/v3.1/spark/src/test/java/org/apache/iceberg/spark/sql/TestRefreshTable.java +++ b/spark/v3.1/spark/src/test/java/org/apache/iceberg/spark/sql/TestRefreshTable.java @@ -61,7 +61,7 @@ public void testRefreshCommand() { // Modify table outside of spark, it should be cached so Spark should see the same value after mutation Table table = validationCatalog.loadTable(tableIdent); - DataFile file = table.currentSnapshot().addedFiles(table.io()).iterator().next(); + DataFile file = table.currentSnapshot().addedDataFiles(table.io()).iterator().next(); table.newDelete().deleteFile(file).commit(); List cachedActual = sql("SELECT * FROM %s", tableName); diff --git a/spark/v3.2/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/TestMerge.java b/spark/v3.2/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/TestMerge.java index 0173f9797d0b..52f7efceb74a 100644 --- a/spark/v3.2/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/TestMerge.java +++ b/spark/v3.2/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/TestMerge.java @@ -101,7 +101,7 @@ public void testMergeWithStaticPredicatePushDown() { "{ \"id\": 2, \"dep\": \"hardware\" }"); // remove the data file from the 'hr' partition to ensure it is not scanned - withUnavailableFiles(snapshot.addedFiles(table.io()), () -> { + withUnavailableFiles(snapshot.addedDataFiles(table.io()), () -> { // disable dynamic pruning and rely only on static predicate pushdown withSQLConf(ImmutableMap.of(SQLConf.DYNAMIC_PARTITION_PRUNING_ENABLED().key(), "false"), () -> { sql("MERGE INTO %s t USING source " + diff --git a/spark/v3.2/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/TestUpdate.java b/spark/v3.2/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/TestUpdate.java index edc3944e69fe..7871e02c5b02 100644 --- a/spark/v3.2/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/TestUpdate.java +++ b/spark/v3.2/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/TestUpdate.java @@ -905,7 +905,7 @@ public void testUpdateWithStaticPredicatePushdown() { Assert.assertEquals("Must have 2 files before UPDATE", "2", dataFilesCount); // remove the data file from the 'hr' partition to ensure it is not scanned - DataFile dataFile = Iterables.getOnlyElement(snapshot.addedFiles(table.io())); + DataFile dataFile = Iterables.getOnlyElement(snapshot.addedDataFiles(table.io())); table.io().deleteFile(dataFile.path().toString()); // disable dynamic pruning and rely only on static predicate pushdown diff --git a/spark/v3.2/spark/src/jmh/java/org/apache/iceberg/spark/source/parquet/IcebergSourceParquetMultiDeleteFileBenchmark.java b/spark/v3.2/spark/src/jmh/java/org/apache/iceberg/spark/source/parquet/IcebergSourceParquetMultiDeleteFileBenchmark.java index 0a77b2889776..3b1b4211174c 100644 --- a/spark/v3.2/spark/src/jmh/java/org/apache/iceberg/spark/source/parquet/IcebergSourceParquetMultiDeleteFileBenchmark.java +++ b/spark/v3.2/spark/src/jmh/java/org/apache/iceberg/spark/source/parquet/IcebergSourceParquetMultiDeleteFileBenchmark.java @@ -47,7 +47,7 @@ protected void appendData() throws IOException { writeData(fileNum); table().refresh(); - for (DataFile file : table().currentSnapshot().addedFiles(table().io())) { + for (DataFile file : table().currentSnapshot().addedDataFiles(table().io())) { writePosDeletes(file.path(), NUM_ROWS, 0.25, numDeleteFile); } } diff --git a/spark/v3.2/spark/src/jmh/java/org/apache/iceberg/spark/source/parquet/IcebergSourceParquetPosDeleteBenchmark.java b/spark/v3.2/spark/src/jmh/java/org/apache/iceberg/spark/source/parquet/IcebergSourceParquetPosDeleteBenchmark.java index 617f8fd069d7..5e36bad11f7a 100644 --- a/spark/v3.2/spark/src/jmh/java/org/apache/iceberg/spark/source/parquet/IcebergSourceParquetPosDeleteBenchmark.java +++ b/spark/v3.2/spark/src/jmh/java/org/apache/iceberg/spark/source/parquet/IcebergSourceParquetPosDeleteBenchmark.java @@ -49,7 +49,7 @@ protected void appendData() throws IOException { if (percentDeleteRow > 0) { // add pos-deletes table().refresh(); - for (DataFile file : table().currentSnapshot().addedFiles(table().io())) { + for (DataFile file : table().currentSnapshot().addedDataFiles(table().io())) { writePosDeletes(file.path(), NUM_ROWS, percentDeleteRow); } } diff --git a/spark/v3.2/spark/src/jmh/java/org/apache/iceberg/spark/source/parquet/IcebergSourceParquetWithUnrelatedDeleteBenchmark.java b/spark/v3.2/spark/src/jmh/java/org/apache/iceberg/spark/source/parquet/IcebergSourceParquetWithUnrelatedDeleteBenchmark.java index 3c57464b4fff..5105c340518d 100644 --- a/spark/v3.2/spark/src/jmh/java/org/apache/iceberg/spark/source/parquet/IcebergSourceParquetWithUnrelatedDeleteBenchmark.java +++ b/spark/v3.2/spark/src/jmh/java/org/apache/iceberg/spark/source/parquet/IcebergSourceParquetWithUnrelatedDeleteBenchmark.java @@ -48,7 +48,7 @@ protected void appendData() throws IOException { writeData(fileNum); table().refresh(); - for (DataFile file : table().currentSnapshot().addedFiles(table().io())) { + for (DataFile file : table().currentSnapshot().addedDataFiles(table().io())) { writePosDeletesWithNoise(file.path(), NUM_ROWS, PERCENT_DELETE_ROW, (int) (percentUnrelatedDeletes / PERCENT_DELETE_ROW), 1); } diff --git a/spark/v3.2/spark/src/test/java/org/apache/iceberg/spark/actions/TestExpireSnapshotsAction.java b/spark/v3.2/spark/src/test/java/org/apache/iceberg/spark/actions/TestExpireSnapshotsAction.java index 20b399e12580..2c324a88f18c 100644 --- a/spark/v3.2/spark/src/test/java/org/apache/iceberg/spark/actions/TestExpireSnapshotsAction.java +++ b/spark/v3.2/spark/src/test/java/org/apache/iceberg/spark/actions/TestExpireSnapshotsAction.java @@ -593,7 +593,7 @@ public void testWithExpiringDanglingStageCommit() { expectedDeletes.add(snapshotA.manifestListLocation()); // Files should be deleted of dangling staged snapshot - snapshotB.addedFiles(table.io()).forEach(i -> { + snapshotB.addedDataFiles(table.io()).forEach(i -> { expectedDeletes.add(i.path().toString()); }); @@ -668,7 +668,7 @@ public void testWithCherryPickTableSnapshot() { // Make sure no dataFiles are deleted for the B, C, D snapshot Lists.newArrayList(snapshotB, snapshotC, snapshotD).forEach(i -> { - i.addedFiles(table.io()).forEach(item -> { + i.addedDataFiles(table.io()).forEach(item -> { Assert.assertFalse(deletedFiles.contains(item.path().toString())); }); }); @@ -723,7 +723,7 @@ public void testWithExpiringStagedThenCherrypick() { // Make sure no dataFiles are deleted for the staged snapshot Lists.newArrayList(snapshotB).forEach(i -> { - i.addedFiles(table.io()).forEach(item -> { + i.addedDataFiles(table.io()).forEach(item -> { Assert.assertFalse(deletedFiles.contains(item.path().toString())); }); }); @@ -737,7 +737,7 @@ public void testWithExpiringStagedThenCherrypick() { // Make sure no dataFiles are deleted for the staged and cherry-pick Lists.newArrayList(snapshotB, snapshotD).forEach(i -> { - i.addedFiles(table.io()).forEach(item -> { + i.addedDataFiles(table.io()).forEach(item -> { Assert.assertFalse(deletedFiles.contains(item.path().toString())); }); }); diff --git a/spark/v3.2/spark/src/test/java/org/apache/iceberg/spark/actions/TestRewriteDataFilesAction.java b/spark/v3.2/spark/src/test/java/org/apache/iceberg/spark/actions/TestRewriteDataFilesAction.java index f7394f645016..ea2a7755f263 100644 --- a/spark/v3.2/spark/src/test/java/org/apache/iceberg/spark/actions/TestRewriteDataFilesAction.java +++ b/spark/v3.2/spark/src/test/java/org/apache/iceberg/spark/actions/TestRewriteDataFilesAction.java @@ -44,6 +44,7 @@ import org.apache.iceberg.RewriteJobOrder; import org.apache.iceberg.RowDelta; import org.apache.iceberg.Schema; +import org.apache.iceberg.Snapshot; import org.apache.iceberg.SortOrder; import org.apache.iceberg.StructLike; import org.apache.iceberg.Table; @@ -996,7 +997,7 @@ public void testAutoSortShuffleOutput() { Assert.assertEquals("Should have 1 fileGroups", result.rewriteResults().size(), 1); Assert.assertTrue("Should have written 40+ files", - Iterables.size(table.currentSnapshot().addedFiles(table.io())) >= 40); + Iterables.size(table.currentSnapshot().addedDataFiles(table.io())) >= 40); table.refresh(); @@ -1063,7 +1064,7 @@ public void testZOrderSort() { .execute(); Assert.assertEquals("Should have 1 fileGroups", 1, result.rewriteResults().size()); - int zOrderedFilesTotal = Iterables.size(table.currentSnapshot().addedFiles(table.io())); + int zOrderedFilesTotal = Iterables.size(table.currentSnapshot().addedDataFiles(table.io())); Assert.assertTrue("Should have written 40+ files", zOrderedFilesTotal >= 40); table.refresh(); @@ -1104,7 +1105,7 @@ public void testZOrderAllTypesSort() { .execute(); Assert.assertEquals("Should have 1 fileGroups", 1, result.rewriteResults().size()); - int zOrderedFilesTotal = Iterables.size(table.currentSnapshot().addedFiles(table.io())); + int zOrderedFilesTotal = Iterables.size(table.currentSnapshot().addedDataFiles(table.io())); Assert.assertEquals("Should have written 1 file", 1, zOrderedFilesTotal); table.refresh(); @@ -1341,7 +1342,8 @@ private List, Pair>> checkForOverlappingFiles(Table ta NestedField field = table.schema().caseInsensitiveFindField(column); Class javaClass = (Class) field.type().typeId().javaClass(); - Map> filesByPartition = Streams.stream(table.currentSnapshot().addedFiles(table.io())) + Snapshot snapshot = table.currentSnapshot(); + Map> filesByPartition = Streams.stream(snapshot.addedDataFiles(table.io())) .collect(Collectors.groupingBy(DataFile::partition)); Stream, Pair>> overlaps = diff --git a/spark/v3.2/spark/src/test/java/org/apache/iceberg/spark/source/TestDataFrameWrites.java b/spark/v3.2/spark/src/test/java/org/apache/iceberg/spark/source/TestDataFrameWrites.java index f6d292f89f8b..b0a77b72b431 100644 --- a/spark/v3.2/spark/src/test/java/org/apache/iceberg/spark/source/TestDataFrameWrites.java +++ b/spark/v3.2/spark/src/test/java/org/apache/iceberg/spark/source/TestDataFrameWrites.java @@ -181,7 +181,7 @@ private void writeAndValidateWithLocations(Table table, File location, File expe } Assert.assertEquals("Both iterators should be exhausted", expectedIter.hasNext(), actualIter.hasNext()); - table.currentSnapshot().addedFiles(table.io()).forEach(dataFile -> + table.currentSnapshot().addedDataFiles(table.io()).forEach(dataFile -> Assert.assertTrue( String.format( "File should have the parent directory %s, but has: %s.", diff --git a/spark/v3.2/spark/src/test/java/org/apache/iceberg/spark/source/TestDataSourceOptions.java b/spark/v3.2/spark/src/test/java/org/apache/iceberg/spark/source/TestDataSourceOptions.java index 00a11c55a9d1..0c56cb328648 100644 --- a/spark/v3.2/spark/src/test/java/org/apache/iceberg/spark/source/TestDataSourceOptions.java +++ b/spark/v3.2/spark/src/test/java/org/apache/iceberg/spark/source/TestDataSourceOptions.java @@ -207,7 +207,7 @@ public void testSplitOptionsOverridesTableProperties() throws IOException { .mode("append") .save(tableLocation); - List files = Lists.newArrayList(icebergTable.currentSnapshot().addedFiles(icebergTable.io())); + List files = Lists.newArrayList(icebergTable.currentSnapshot().addedDataFiles(icebergTable.io())); Assert.assertEquals("Should have written 1 file", 1, files.size()); long fileSize = files.get(0).fileSizeInBytes(); diff --git a/spark/v3.2/spark/src/test/java/org/apache/iceberg/spark/source/TestIcebergSourceTablesBase.java b/spark/v3.2/spark/src/test/java/org/apache/iceberg/spark/source/TestIcebergSourceTablesBase.java index a80e5ee79f9f..ce40f179d649 100644 --- a/spark/v3.2/spark/src/test/java/org/apache/iceberg/spark/source/TestIcebergSourceTablesBase.java +++ b/spark/v3.2/spark/src/test/java/org/apache/iceberg/spark/source/TestIcebergSourceTablesBase.java @@ -219,7 +219,7 @@ public void testEntriesTableDataFilePrune() throws Exception { .save(loadLocation(tableIdentifier)); table.refresh(); - DataFile file = table.currentSnapshot().addedFiles(table.io()).iterator().next(); + DataFile file = table.currentSnapshot().addedDataFiles(table.io()).iterator().next(); List singleActual = rowsToJava(spark.read() .format("iceberg") @@ -246,7 +246,7 @@ public void testEntriesTableDataFilePruneMulti() throws Exception { .save(loadLocation(tableIdentifier)); table.refresh(); - DataFile file = table.currentSnapshot().addedFiles(table.io()).iterator().next(); + DataFile file = table.currentSnapshot().addedDataFiles(table.io()).iterator().next(); List multiActual = rowsToJava(spark.read() .format("iceberg") @@ -274,7 +274,7 @@ public void testFilesSelectMap() throws Exception { .save(loadLocation(tableIdentifier)); table.refresh(); - DataFile file = table.currentSnapshot().addedFiles(table.io()).iterator().next(); + DataFile file = table.currentSnapshot().addedDataFiles(table.io()).iterator().next(); List multiActual = rowsToJava(spark.read() .format("iceberg") @@ -544,7 +544,7 @@ public void testFilesUnpartitionedTable() throws Exception { .save(loadLocation(tableIdentifier)); table.refresh(); - DataFile toDelete = Iterables.getOnlyElement(table.currentSnapshot().addedFiles(table.io())); + DataFile toDelete = Iterables.getOnlyElement(table.currentSnapshot().addedDataFiles(table.io())); // add a second file df2.select("id", "data").write() @@ -907,7 +907,7 @@ public void testManifestsTable() { .set(TableProperties.FORMAT_VERSION, "2") .commit(); - DataFile dataFile = Iterables.getFirst(table.currentSnapshot().addedFiles(table.io()), null); + DataFile dataFile = Iterables.getFirst(table.currentSnapshot().addedDataFiles(table.io()), null); PartitionSpec dataFileSpec = table.specs().get(dataFile.specId()); StructLike dataFilePartition = dataFile.partition(); @@ -1036,7 +1036,7 @@ public void testAllManifestsTable() { .set(TableProperties.FORMAT_VERSION, "2") .commit(); - DataFile dataFile = Iterables.getFirst(table.currentSnapshot().addedFiles(table.io()), null); + DataFile dataFile = Iterables.getFirst(table.currentSnapshot().addedDataFiles(table.io()), null); PartitionSpec dataFileSpec = table.specs().get(dataFile.specId()); StructLike dataFilePartition = dataFile.partition(); diff --git a/spark/v3.2/spark/src/test/java/org/apache/iceberg/spark/sql/TestRefreshTable.java b/spark/v3.2/spark/src/test/java/org/apache/iceberg/spark/sql/TestRefreshTable.java index 187dd4470a05..a8bdea77e237 100644 --- a/spark/v3.2/spark/src/test/java/org/apache/iceberg/spark/sql/TestRefreshTable.java +++ b/spark/v3.2/spark/src/test/java/org/apache/iceberg/spark/sql/TestRefreshTable.java @@ -61,7 +61,7 @@ public void testRefreshCommand() { // Modify table outside of spark, it should be cached so Spark should see the same value after mutation Table table = validationCatalog.loadTable(tableIdent); - DataFile file = table.currentSnapshot().addedFiles(table.io()).iterator().next(); + DataFile file = table.currentSnapshot().addedDataFiles(table.io()).iterator().next(); table.newDelete().deleteFile(file).commit(); List cachedActual = sql("SELECT * FROM %s", tableName); diff --git a/spark/v3.3/spark/src/test/java/org/apache/iceberg/spark/actions/TestRewriteDataFilesAction.java b/spark/v3.3/spark/src/test/java/org/apache/iceberg/spark/actions/TestRewriteDataFilesAction.java index f7394f645016..ea2a7755f263 100644 --- a/spark/v3.3/spark/src/test/java/org/apache/iceberg/spark/actions/TestRewriteDataFilesAction.java +++ b/spark/v3.3/spark/src/test/java/org/apache/iceberg/spark/actions/TestRewriteDataFilesAction.java @@ -44,6 +44,7 @@ import org.apache.iceberg.RewriteJobOrder; import org.apache.iceberg.RowDelta; import org.apache.iceberg.Schema; +import org.apache.iceberg.Snapshot; import org.apache.iceberg.SortOrder; import org.apache.iceberg.StructLike; import org.apache.iceberg.Table; @@ -996,7 +997,7 @@ public void testAutoSortShuffleOutput() { Assert.assertEquals("Should have 1 fileGroups", result.rewriteResults().size(), 1); Assert.assertTrue("Should have written 40+ files", - Iterables.size(table.currentSnapshot().addedFiles(table.io())) >= 40); + Iterables.size(table.currentSnapshot().addedDataFiles(table.io())) >= 40); table.refresh(); @@ -1063,7 +1064,7 @@ public void testZOrderSort() { .execute(); Assert.assertEquals("Should have 1 fileGroups", 1, result.rewriteResults().size()); - int zOrderedFilesTotal = Iterables.size(table.currentSnapshot().addedFiles(table.io())); + int zOrderedFilesTotal = Iterables.size(table.currentSnapshot().addedDataFiles(table.io())); Assert.assertTrue("Should have written 40+ files", zOrderedFilesTotal >= 40); table.refresh(); @@ -1104,7 +1105,7 @@ public void testZOrderAllTypesSort() { .execute(); Assert.assertEquals("Should have 1 fileGroups", 1, result.rewriteResults().size()); - int zOrderedFilesTotal = Iterables.size(table.currentSnapshot().addedFiles(table.io())); + int zOrderedFilesTotal = Iterables.size(table.currentSnapshot().addedDataFiles(table.io())); Assert.assertEquals("Should have written 1 file", 1, zOrderedFilesTotal); table.refresh(); @@ -1341,7 +1342,8 @@ private List, Pair>> checkForOverlappingFiles(Table ta NestedField field = table.schema().caseInsensitiveFindField(column); Class javaClass = (Class) field.type().typeId().javaClass(); - Map> filesByPartition = Streams.stream(table.currentSnapshot().addedFiles(table.io())) + Snapshot snapshot = table.currentSnapshot(); + Map> filesByPartition = Streams.stream(snapshot.addedDataFiles(table.io())) .collect(Collectors.groupingBy(DataFile::partition)); Stream, Pair>> overlaps = diff --git a/spark/v3.3/spark/src/test/java/org/apache/iceberg/spark/source/TestDataFrameWrites.java b/spark/v3.3/spark/src/test/java/org/apache/iceberg/spark/source/TestDataFrameWrites.java index f6d292f89f8b..b0a77b72b431 100644 --- a/spark/v3.3/spark/src/test/java/org/apache/iceberg/spark/source/TestDataFrameWrites.java +++ b/spark/v3.3/spark/src/test/java/org/apache/iceberg/spark/source/TestDataFrameWrites.java @@ -181,7 +181,7 @@ private void writeAndValidateWithLocations(Table table, File location, File expe } Assert.assertEquals("Both iterators should be exhausted", expectedIter.hasNext(), actualIter.hasNext()); - table.currentSnapshot().addedFiles(table.io()).forEach(dataFile -> + table.currentSnapshot().addedDataFiles(table.io()).forEach(dataFile -> Assert.assertTrue( String.format( "File should have the parent directory %s, but has: %s.", diff --git a/spark/v3.3/spark/src/test/java/org/apache/iceberg/spark/source/TestIcebergSourceTablesBase.java b/spark/v3.3/spark/src/test/java/org/apache/iceberg/spark/source/TestIcebergSourceTablesBase.java index a80e5ee79f9f..ce40f179d649 100644 --- a/spark/v3.3/spark/src/test/java/org/apache/iceberg/spark/source/TestIcebergSourceTablesBase.java +++ b/spark/v3.3/spark/src/test/java/org/apache/iceberg/spark/source/TestIcebergSourceTablesBase.java @@ -219,7 +219,7 @@ public void testEntriesTableDataFilePrune() throws Exception { .save(loadLocation(tableIdentifier)); table.refresh(); - DataFile file = table.currentSnapshot().addedFiles(table.io()).iterator().next(); + DataFile file = table.currentSnapshot().addedDataFiles(table.io()).iterator().next(); List singleActual = rowsToJava(spark.read() .format("iceberg") @@ -246,7 +246,7 @@ public void testEntriesTableDataFilePruneMulti() throws Exception { .save(loadLocation(tableIdentifier)); table.refresh(); - DataFile file = table.currentSnapshot().addedFiles(table.io()).iterator().next(); + DataFile file = table.currentSnapshot().addedDataFiles(table.io()).iterator().next(); List multiActual = rowsToJava(spark.read() .format("iceberg") @@ -274,7 +274,7 @@ public void testFilesSelectMap() throws Exception { .save(loadLocation(tableIdentifier)); table.refresh(); - DataFile file = table.currentSnapshot().addedFiles(table.io()).iterator().next(); + DataFile file = table.currentSnapshot().addedDataFiles(table.io()).iterator().next(); List multiActual = rowsToJava(spark.read() .format("iceberg") @@ -544,7 +544,7 @@ public void testFilesUnpartitionedTable() throws Exception { .save(loadLocation(tableIdentifier)); table.refresh(); - DataFile toDelete = Iterables.getOnlyElement(table.currentSnapshot().addedFiles(table.io())); + DataFile toDelete = Iterables.getOnlyElement(table.currentSnapshot().addedDataFiles(table.io())); // add a second file df2.select("id", "data").write() @@ -907,7 +907,7 @@ public void testManifestsTable() { .set(TableProperties.FORMAT_VERSION, "2") .commit(); - DataFile dataFile = Iterables.getFirst(table.currentSnapshot().addedFiles(table.io()), null); + DataFile dataFile = Iterables.getFirst(table.currentSnapshot().addedDataFiles(table.io()), null); PartitionSpec dataFileSpec = table.specs().get(dataFile.specId()); StructLike dataFilePartition = dataFile.partition(); @@ -1036,7 +1036,7 @@ public void testAllManifestsTable() { .set(TableProperties.FORMAT_VERSION, "2") .commit(); - DataFile dataFile = Iterables.getFirst(table.currentSnapshot().addedFiles(table.io()), null); + DataFile dataFile = Iterables.getFirst(table.currentSnapshot().addedDataFiles(table.io()), null); PartitionSpec dataFileSpec = table.specs().get(dataFile.specId()); StructLike dataFilePartition = dataFile.partition(); diff --git a/spark/v3.3/spark/src/test/java/org/apache/iceberg/spark/sql/TestRefreshTable.java b/spark/v3.3/spark/src/test/java/org/apache/iceberg/spark/sql/TestRefreshTable.java index 187dd4470a05..a8bdea77e237 100644 --- a/spark/v3.3/spark/src/test/java/org/apache/iceberg/spark/sql/TestRefreshTable.java +++ b/spark/v3.3/spark/src/test/java/org/apache/iceberg/spark/sql/TestRefreshTable.java @@ -61,7 +61,7 @@ public void testRefreshCommand() { // Modify table outside of spark, it should be cached so Spark should see the same value after mutation Table table = validationCatalog.loadTable(tableIdent); - DataFile file = table.currentSnapshot().addedFiles(table.io()).iterator().next(); + DataFile file = table.currentSnapshot().addedDataFiles(table.io()).iterator().next(); table.newDelete().deleteFile(file).commit(); List cachedActual = sql("SELECT * FROM %s", tableName); From 9b55b4787dc81f7582957028f9de11c1e1bd08a0 Mon Sep 17 00:00:00 2001 From: aokolnychyi Date: Thu, 30 Jun 2022 10:07:35 -0700 Subject: [PATCH 2/3] Rename test vars --- .../java/org/apache/iceberg/TestSnapshot.java | 24 +++++++++---------- 1 file changed, 12 insertions(+), 12 deletions(-) diff --git a/core/src/test/java/org/apache/iceberg/TestSnapshot.java b/core/src/test/java/org/apache/iceberg/TestSnapshot.java index 60de9652ed0e..6bb6384ac4db 100644 --- a/core/src/test/java/org/apache/iceberg/TestSnapshot.java +++ b/core/src/test/java/org/apache/iceberg/TestSnapshot.java @@ -113,13 +113,13 @@ public void testCachedDataFiles() { Snapshot thirdSnapshot = table.currentSnapshot(); - Iterable deletedDataFiles = thirdSnapshot.removedDataFiles(FILE_IO); - Assert.assertEquals("Must have 1 deleted data file", 1, Iterables.size(deletedDataFiles)); + Iterable removedDataFiles = thirdSnapshot.removedDataFiles(FILE_IO); + Assert.assertEquals("Must have 1 removed data file", 1, Iterables.size(removedDataFiles)); - DataFile deletedDataFile = Iterables.getOnlyElement(deletedDataFiles); - Assert.assertEquals("Path must match", FILE_A.path(), deletedDataFile.path()); - Assert.assertEquals("Spec ID must match", FILE_A.specId(), deletedDataFile.specId()); - Assert.assertEquals("Partition must match", FILE_A.partition(), deletedDataFile.partition()); + DataFile removedDataFile = Iterables.getOnlyElement(removedDataFiles); + Assert.assertEquals("Path must match", FILE_A.path(), removedDataFile.path()); + Assert.assertEquals("Spec ID must match", FILE_A.specId(), removedDataFile.specId()); + Assert.assertEquals("Partition must match", FILE_A.partition(), removedDataFile.partition()); Iterable addedDataFiles = thirdSnapshot.addedDataFiles(FILE_IO); Assert.assertEquals("Must have 1 added data file", 1, Iterables.size(addedDataFiles)); @@ -164,13 +164,13 @@ public void testCachedDeleteFiles() { Snapshot thirdSnapshot = table.currentSnapshot(); - Iterable deletedDeleteFiles = thirdSnapshot.removedDeleteFiles(FILE_IO); - Assert.assertEquals("Must have 1 deleted delete file", 1, Iterables.size(deletedDeleteFiles)); + Iterable removedDeleteFiles = thirdSnapshot.removedDeleteFiles(FILE_IO); + Assert.assertEquals("Must have 1 removed delete file", 1, Iterables.size(removedDeleteFiles)); - DeleteFile deletedDeleteFile = Iterables.getOnlyElement(deletedDeleteFiles); - Assert.assertEquals("Path must match", secondSnapshotDeleteFile.path(), deletedDeleteFile.path()); - Assert.assertEquals("Spec ID must match", secondSnapshotDeleteFile.specId(), deletedDeleteFile.specId()); - Assert.assertEquals("Partition must match", secondSnapshotDeleteFile.partition(), deletedDeleteFile.partition()); + DeleteFile removedDeleteFile = Iterables.getOnlyElement(removedDeleteFiles); + Assert.assertEquals("Path must match", secondSnapshotDeleteFile.path(), removedDeleteFile.path()); + Assert.assertEquals("Spec ID must match", secondSnapshotDeleteFile.specId(), removedDeleteFile.specId()); + Assert.assertEquals("Partition must match", secondSnapshotDeleteFile.partition(), removedDeleteFile.partition()); Iterable addedDeleteFiles = thirdSnapshot.addedDeleteFiles(FILE_IO); Assert.assertEquals("Must have 1 added delete file", 1, Iterables.size(addedDeleteFiles)); From 86913c0d7feb5a78de03f581312c7cf0b1b030c5 Mon Sep 17 00:00:00 2001 From: aokolnychyi Date: Thu, 30 Jun 2022 10:15:39 -0700 Subject: [PATCH 3/3] More renames for 3.3 in tests --- .../org/apache/iceberg/spark/extensions/TestMerge.java | 2 +- .../org/apache/iceberg/spark/extensions/TestUpdate.java | 2 +- .../IcebergSourceParquetMultiDeleteFileBenchmark.java | 2 +- .../parquet/IcebergSourceParquetPosDeleteBenchmark.java | 2 +- .../IcebergSourceParquetWithUnrelatedDeleteBenchmark.java | 2 +- .../iceberg/spark/actions/TestExpireSnapshotsAction.java | 8 ++++---- .../iceberg/spark/source/TestDataSourceOptions.java | 2 +- 7 files changed, 10 insertions(+), 10 deletions(-) diff --git a/spark/v3.3/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/TestMerge.java b/spark/v3.3/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/TestMerge.java index 0173f9797d0b..52f7efceb74a 100644 --- a/spark/v3.3/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/TestMerge.java +++ b/spark/v3.3/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/TestMerge.java @@ -101,7 +101,7 @@ public void testMergeWithStaticPredicatePushDown() { "{ \"id\": 2, \"dep\": \"hardware\" }"); // remove the data file from the 'hr' partition to ensure it is not scanned - withUnavailableFiles(snapshot.addedFiles(table.io()), () -> { + withUnavailableFiles(snapshot.addedDataFiles(table.io()), () -> { // disable dynamic pruning and rely only on static predicate pushdown withSQLConf(ImmutableMap.of(SQLConf.DYNAMIC_PARTITION_PRUNING_ENABLED().key(), "false"), () -> { sql("MERGE INTO %s t USING source " + diff --git a/spark/v3.3/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/TestUpdate.java b/spark/v3.3/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/TestUpdate.java index edc3944e69fe..7871e02c5b02 100644 --- a/spark/v3.3/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/TestUpdate.java +++ b/spark/v3.3/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/TestUpdate.java @@ -905,7 +905,7 @@ public void testUpdateWithStaticPredicatePushdown() { Assert.assertEquals("Must have 2 files before UPDATE", "2", dataFilesCount); // remove the data file from the 'hr' partition to ensure it is not scanned - DataFile dataFile = Iterables.getOnlyElement(snapshot.addedFiles(table.io())); + DataFile dataFile = Iterables.getOnlyElement(snapshot.addedDataFiles(table.io())); table.io().deleteFile(dataFile.path().toString()); // disable dynamic pruning and rely only on static predicate pushdown diff --git a/spark/v3.3/spark/src/jmh/java/org/apache/iceberg/spark/source/parquet/IcebergSourceParquetMultiDeleteFileBenchmark.java b/spark/v3.3/spark/src/jmh/java/org/apache/iceberg/spark/source/parquet/IcebergSourceParquetMultiDeleteFileBenchmark.java index 69b5be1d9171..24c2676af24d 100644 --- a/spark/v3.3/spark/src/jmh/java/org/apache/iceberg/spark/source/parquet/IcebergSourceParquetMultiDeleteFileBenchmark.java +++ b/spark/v3.3/spark/src/jmh/java/org/apache/iceberg/spark/source/parquet/IcebergSourceParquetMultiDeleteFileBenchmark.java @@ -47,7 +47,7 @@ protected void appendData() throws IOException { writeData(fileNum); table().refresh(); - for (DataFile file : table().currentSnapshot().addedFiles(table().io())) { + for (DataFile file : table().currentSnapshot().addedDataFiles(table().io())) { writePosDeletes(file.path(), NUM_ROWS, 0.25, numDeleteFile); } } diff --git a/spark/v3.3/spark/src/jmh/java/org/apache/iceberg/spark/source/parquet/IcebergSourceParquetPosDeleteBenchmark.java b/spark/v3.3/spark/src/jmh/java/org/apache/iceberg/spark/source/parquet/IcebergSourceParquetPosDeleteBenchmark.java index 326ecd5fa28b..988eeb751258 100644 --- a/spark/v3.3/spark/src/jmh/java/org/apache/iceberg/spark/source/parquet/IcebergSourceParquetPosDeleteBenchmark.java +++ b/spark/v3.3/spark/src/jmh/java/org/apache/iceberg/spark/source/parquet/IcebergSourceParquetPosDeleteBenchmark.java @@ -49,7 +49,7 @@ protected void appendData() throws IOException { if (percentDeleteRow > 0) { // add pos-deletes table().refresh(); - for (DataFile file : table().currentSnapshot().addedFiles(table().io())) { + for (DataFile file : table().currentSnapshot().addedDataFiles(table().io())) { writePosDeletes(file.path(), NUM_ROWS, percentDeleteRow); } } diff --git a/spark/v3.3/spark/src/jmh/java/org/apache/iceberg/spark/source/parquet/IcebergSourceParquetWithUnrelatedDeleteBenchmark.java b/spark/v3.3/spark/src/jmh/java/org/apache/iceberg/spark/source/parquet/IcebergSourceParquetWithUnrelatedDeleteBenchmark.java index 548c391746b2..5088843ca13e 100644 --- a/spark/v3.3/spark/src/jmh/java/org/apache/iceberg/spark/source/parquet/IcebergSourceParquetWithUnrelatedDeleteBenchmark.java +++ b/spark/v3.3/spark/src/jmh/java/org/apache/iceberg/spark/source/parquet/IcebergSourceParquetWithUnrelatedDeleteBenchmark.java @@ -48,7 +48,7 @@ protected void appendData() throws IOException { writeData(fileNum); table().refresh(); - for (DataFile file : table().currentSnapshot().addedFiles(table().io())) { + for (DataFile file : table().currentSnapshot().addedDataFiles(table().io())) { writePosDeletesWithNoise(file.path(), NUM_ROWS, PERCENT_DELETE_ROW, (int) (percentUnrelatedDeletes / PERCENT_DELETE_ROW), 1); } diff --git a/spark/v3.3/spark/src/test/java/org/apache/iceberg/spark/actions/TestExpireSnapshotsAction.java b/spark/v3.3/spark/src/test/java/org/apache/iceberg/spark/actions/TestExpireSnapshotsAction.java index 20b399e12580..2c324a88f18c 100644 --- a/spark/v3.3/spark/src/test/java/org/apache/iceberg/spark/actions/TestExpireSnapshotsAction.java +++ b/spark/v3.3/spark/src/test/java/org/apache/iceberg/spark/actions/TestExpireSnapshotsAction.java @@ -593,7 +593,7 @@ public void testWithExpiringDanglingStageCommit() { expectedDeletes.add(snapshotA.manifestListLocation()); // Files should be deleted of dangling staged snapshot - snapshotB.addedFiles(table.io()).forEach(i -> { + snapshotB.addedDataFiles(table.io()).forEach(i -> { expectedDeletes.add(i.path().toString()); }); @@ -668,7 +668,7 @@ public void testWithCherryPickTableSnapshot() { // Make sure no dataFiles are deleted for the B, C, D snapshot Lists.newArrayList(snapshotB, snapshotC, snapshotD).forEach(i -> { - i.addedFiles(table.io()).forEach(item -> { + i.addedDataFiles(table.io()).forEach(item -> { Assert.assertFalse(deletedFiles.contains(item.path().toString())); }); }); @@ -723,7 +723,7 @@ public void testWithExpiringStagedThenCherrypick() { // Make sure no dataFiles are deleted for the staged snapshot Lists.newArrayList(snapshotB).forEach(i -> { - i.addedFiles(table.io()).forEach(item -> { + i.addedDataFiles(table.io()).forEach(item -> { Assert.assertFalse(deletedFiles.contains(item.path().toString())); }); }); @@ -737,7 +737,7 @@ public void testWithExpiringStagedThenCherrypick() { // Make sure no dataFiles are deleted for the staged and cherry-pick Lists.newArrayList(snapshotB, snapshotD).forEach(i -> { - i.addedFiles(table.io()).forEach(item -> { + i.addedDataFiles(table.io()).forEach(item -> { Assert.assertFalse(deletedFiles.contains(item.path().toString())); }); }); diff --git a/spark/v3.3/spark/src/test/java/org/apache/iceberg/spark/source/TestDataSourceOptions.java b/spark/v3.3/spark/src/test/java/org/apache/iceberg/spark/source/TestDataSourceOptions.java index 00a11c55a9d1..0c56cb328648 100644 --- a/spark/v3.3/spark/src/test/java/org/apache/iceberg/spark/source/TestDataSourceOptions.java +++ b/spark/v3.3/spark/src/test/java/org/apache/iceberg/spark/source/TestDataSourceOptions.java @@ -207,7 +207,7 @@ public void testSplitOptionsOverridesTableProperties() throws IOException { .mode("append") .save(tableLocation); - List files = Lists.newArrayList(icebergTable.currentSnapshot().addedFiles(icebergTable.io())); + List files = Lists.newArrayList(icebergTable.currentSnapshot().addedDataFiles(icebergTable.io())); Assert.assertEquals("Should have written 1 file", 1, files.size()); long fileSize = files.get(0).fileSizeInBytes();