From 75e3582435469db21d6355801024aa72c82c355b Mon Sep 17 00:00:00 2001 From: Md Nasim Akhtar Date: Mon, 29 Jun 2026 12:29:13 +0530 Subject: [PATCH] Spark: Fix DeleteOrphanFilesSparkAction sibling path matching bug This commit fixes issue #16493 where file_list_view was scoped using raw string prefix matching, which allowed sibling paths to fall inside orphan cleanup (e.g. s3://bucket/table-backup matched s3://bucket/table). Added a trailing slash to the location before prefix matching in filteredCompareToFileList to ensure precise directory scoping. --- .../spark/actions/DeleteOrphanFilesSparkAction.java | 8 +++++++- .../spark/actions/TestRemoveOrphanFilesAction.java | 5 ++++- .../spark/actions/DeleteOrphanFilesSparkAction.java | 8 +++++++- .../spark/actions/TestRemoveOrphanFilesAction.java | 5 ++++- .../spark/actions/DeleteOrphanFilesSparkAction.java | 8 +++++++- .../spark/actions/TestRemoveOrphanFilesAction.java | 5 ++++- 6 files changed, 33 insertions(+), 6 deletions(-) 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