diff --git a/spark/v3.5/spark/src/main/java/org/apache/iceberg/spark/actions/DeleteOrphanFilesSparkAction.java b/spark/v3.5/spark/src/main/java/org/apache/iceberg/spark/actions/DeleteOrphanFilesSparkAction.java index 92bfc880ad7f..0bf6485b5480 100644 --- a/spark/v3.5/spark/src/main/java/org/apache/iceberg/spark/actions/DeleteOrphanFilesSparkAction.java +++ b/spark/v3.5/spark/src/main/java/org/apache/iceberg/spark/actions/DeleteOrphanFilesSparkAction.java @@ -52,6 +52,7 @@ import org.apache.iceberg.relocated.com.google.common.collect.Maps; import org.apache.iceberg.spark.JobGroupInfo; import org.apache.iceberg.util.FileSystemWalker; +import org.apache.iceberg.util.LocationUtil; import org.apache.iceberg.util.Pair; import org.apache.iceberg.util.PropertyUtil; import org.apache.iceberg.util.Tasks; @@ -225,7 +226,12 @@ public DeleteOrphanFilesSparkAction usePrefixListing(boolean newUsePrefixListing private Dataset filteredCompareToFileList() { Dataset files = compareToFileList; if (location != null) { - files = files.filter(files.col(FILE_PATH).startsWith(location)); + String strippedLocation = LocationUtil.stripTrailingSlash(location); + String locationWithTrailingSlash = + strippedLocation.endsWith(LocationUtil.PATH_SEPARATOR) + ? strippedLocation + : strippedLocation + LocationUtil.PATH_SEPARATOR; + files = files.filter(files.col(FILE_PATH).startsWith(locationWithTrailingSlash)); } return files .filter(files.col(LAST_MODIFIED).lt(new Timestamp(olderThanTimestamp))) diff --git a/spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/actions/TestRemoveOrphanFilesAction.java b/spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/actions/TestRemoveOrphanFilesAction.java index 34a02a93faf1..90165056ecf0 100644 --- a/spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/actions/TestRemoveOrphanFilesAction.java +++ b/spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/actions/TestRemoveOrphanFilesAction.java @@ -1001,7 +1001,10 @@ public void testCompareToFileList() throws IOException { assertThat(actualRecords).as("Rows must match").isEqualTo(expectedRecords); List outsideLocationMockFiles = - Lists.newArrayList(new FilePathLastModifiedRecord("/tmp/mock1", new Timestamp(0L))); + Lists.newArrayList( + new FilePathLastModifiedRecord("/tmp/mock1", new Timestamp(0L)), + new FilePathLastModifiedRecord( + tableLocation + "-backup/mock1.parquet", new Timestamp(0L))); Dataset compareToFileListWithOutsideLocation = spark diff --git a/spark/v4.0/spark/src/main/java/org/apache/iceberg/spark/actions/DeleteOrphanFilesSparkAction.java b/spark/v4.0/spark/src/main/java/org/apache/iceberg/spark/actions/DeleteOrphanFilesSparkAction.java index 92bfc880ad7f..0bf6485b5480 100644 --- a/spark/v4.0/spark/src/main/java/org/apache/iceberg/spark/actions/DeleteOrphanFilesSparkAction.java +++ b/spark/v4.0/spark/src/main/java/org/apache/iceberg/spark/actions/DeleteOrphanFilesSparkAction.java @@ -52,6 +52,7 @@ import org.apache.iceberg.relocated.com.google.common.collect.Maps; import org.apache.iceberg.spark.JobGroupInfo; import org.apache.iceberg.util.FileSystemWalker; +import org.apache.iceberg.util.LocationUtil; import org.apache.iceberg.util.Pair; import org.apache.iceberg.util.PropertyUtil; import org.apache.iceberg.util.Tasks; @@ -225,7 +226,12 @@ public DeleteOrphanFilesSparkAction usePrefixListing(boolean newUsePrefixListing private Dataset filteredCompareToFileList() { Dataset files = compareToFileList; if (location != null) { - files = files.filter(files.col(FILE_PATH).startsWith(location)); + String strippedLocation = LocationUtil.stripTrailingSlash(location); + String locationWithTrailingSlash = + strippedLocation.endsWith(LocationUtil.PATH_SEPARATOR) + ? strippedLocation + : strippedLocation + LocationUtil.PATH_SEPARATOR; + files = files.filter(files.col(FILE_PATH).startsWith(locationWithTrailingSlash)); } return files .filter(files.col(LAST_MODIFIED).lt(new Timestamp(olderThanTimestamp))) diff --git a/spark/v4.0/spark/src/test/java/org/apache/iceberg/spark/actions/TestRemoveOrphanFilesAction.java b/spark/v4.0/spark/src/test/java/org/apache/iceberg/spark/actions/TestRemoveOrphanFilesAction.java index 34a02a93faf1..90165056ecf0 100644 --- a/spark/v4.0/spark/src/test/java/org/apache/iceberg/spark/actions/TestRemoveOrphanFilesAction.java +++ b/spark/v4.0/spark/src/test/java/org/apache/iceberg/spark/actions/TestRemoveOrphanFilesAction.java @@ -1001,7 +1001,10 @@ public void testCompareToFileList() throws IOException { assertThat(actualRecords).as("Rows must match").isEqualTo(expectedRecords); List outsideLocationMockFiles = - Lists.newArrayList(new FilePathLastModifiedRecord("/tmp/mock1", new Timestamp(0L))); + Lists.newArrayList( + new FilePathLastModifiedRecord("/tmp/mock1", new Timestamp(0L)), + new FilePathLastModifiedRecord( + tableLocation + "-backup/mock1.parquet", new Timestamp(0L))); Dataset compareToFileListWithOutsideLocation = spark diff --git a/spark/v4.1/spark/src/main/java/org/apache/iceberg/spark/actions/DeleteOrphanFilesSparkAction.java b/spark/v4.1/spark/src/main/java/org/apache/iceberg/spark/actions/DeleteOrphanFilesSparkAction.java index 92bfc880ad7f..0bf6485b5480 100644 --- a/spark/v4.1/spark/src/main/java/org/apache/iceberg/spark/actions/DeleteOrphanFilesSparkAction.java +++ b/spark/v4.1/spark/src/main/java/org/apache/iceberg/spark/actions/DeleteOrphanFilesSparkAction.java @@ -52,6 +52,7 @@ import org.apache.iceberg.relocated.com.google.common.collect.Maps; import org.apache.iceberg.spark.JobGroupInfo; import org.apache.iceberg.util.FileSystemWalker; +import org.apache.iceberg.util.LocationUtil; import org.apache.iceberg.util.Pair; import org.apache.iceberg.util.PropertyUtil; import org.apache.iceberg.util.Tasks; @@ -225,7 +226,12 @@ public DeleteOrphanFilesSparkAction usePrefixListing(boolean newUsePrefixListing private Dataset filteredCompareToFileList() { Dataset files = compareToFileList; if (location != null) { - files = files.filter(files.col(FILE_PATH).startsWith(location)); + String strippedLocation = LocationUtil.stripTrailingSlash(location); + String locationWithTrailingSlash = + strippedLocation.endsWith(LocationUtil.PATH_SEPARATOR) + ? strippedLocation + : strippedLocation + LocationUtil.PATH_SEPARATOR; + files = files.filter(files.col(FILE_PATH).startsWith(locationWithTrailingSlash)); } return files .filter(files.col(LAST_MODIFIED).lt(new Timestamp(olderThanTimestamp))) diff --git a/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/actions/TestRemoveOrphanFilesAction.java b/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/actions/TestRemoveOrphanFilesAction.java index 78e8a0b000a4..ffc7661edb38 100644 --- a/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/actions/TestRemoveOrphanFilesAction.java +++ b/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/actions/TestRemoveOrphanFilesAction.java @@ -1002,7 +1002,10 @@ public void testCompareToFileList() throws IOException { assertThat(actualRecords).as("Rows must match").isEqualTo(expectedRecords); List outsideLocationMockFiles = - Lists.newArrayList(new FilePathLastModifiedRecord("/tmp/mock1", new Timestamp(0L))); + Lists.newArrayList( + new FilePathLastModifiedRecord("/tmp/mock1", new Timestamp(0L)), + new FilePathLastModifiedRecord( + tableLocation + "-backup/mock1.parquet", new Timestamp(0L))); Dataset compareToFileListWithOutsideLocation = spark