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 b47922820d21..3dc42417fa2b 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 @@ -225,7 +225,14 @@ public DeleteOrphanFilesSparkAction usePrefixListing(boolean newUsePrefixListing private Dataset filteredCompareToFileList() { Dataset files = compareToFileList; if (location != null) { - files = files.filter(files.col(FILE_PATH).startsWith(location)); + // Ensure path boundary is respected: s3://bucket/table should not match + // s3://bucket/table-backup/... which is a sibling path, not a subdirectory + String locationPrefix = + location.endsWith("/") ? location : location + "/"; + files = + files.filter( + files.col(FILE_PATH).startsWith(locationPrefix) + .or(files.col(FILE_PATH).equalTo(location))); } 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..c7a36140151b 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 @@ -1027,6 +1027,56 @@ protected long waitUntilAfter(long timestampMillis) { return current; } + @TestTemplate + public void testCompareToFileListDoesNotMatchSiblingPaths() throws IOException { + assumeThat(usePrefixListing) + .as("Should not test both prefix listing and Hadoop file listing (redundant)") + .isEqualTo(false); + Table table = TABLES.create(SCHEMA, PartitionSpec.unpartitioned(), properties, tableLocation); + + List records = + Lists.newArrayList(new ThreeColumnRecord(1, "AAAAAAAAAA", "AAAA")); + + Dataset df = spark.createDataFrame(records, ThreeColumnRecord.class).coalesce(1); + + df.select("c1", "c2", "c3").write().format("iceberg").mode("append").save(tableLocation); + + // sibling paths that share the table location as a raw string prefix but are NOT + // subdirectories of the table location (e.g. table-backup/...). These must NOT be + // treated as in-scope orphan candidates. + String sibling1 = tableLocation + "-backup/data/sibling1.parquet"; + String sibling2 = tableLocation + "_old/data/sibling2.parquet"; + String insideLocation = tableLocation + "/data/inside.parquet"; + + List mockFiles = + Lists.newArrayList( + new FilePathLastModifiedRecord(sibling1, new Timestamp(0L)), + new FilePathLastModifiedRecord(sibling2, new Timestamp(0L)), + new FilePathLastModifiedRecord(insideLocation, new Timestamp(0L))); + + Dataset compareToFileList = + spark + .createDataFrame(mockFiles, FilePathLastModifiedRecord.class) + .withColumnRenamed("filePath", "file_path") + .withColumnRenamed("lastModified", "last_modified"); + + DeleteOrphanFiles.Result result = + SparkActions.get() + .deleteOrphanFiles(table) + .compareToFileList(compareToFileList) + .olderThan(System.currentTimeMillis()) + .deleteWith(s -> {}) + .execute(); + + // Only the file actually inside the table location should be considered in scope. + assertThat(result.orphanFileLocations()) + .as("Only files inside the table location should be in scope") + .containsExactly(insideLocation); + assertThat(result.orphanFilesCount()) + .as("Only 1 file inside the table location should be in scope") + .isEqualTo(1L); + } + @TestTemplate public void testRemoveOrphanFilesWithStatisticFiles() throws Exception { assumeThat(usePrefixListing) 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 b47922820d21..3dc42417fa2b 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 @@ -225,7 +225,14 @@ public DeleteOrphanFilesSparkAction usePrefixListing(boolean newUsePrefixListing private Dataset filteredCompareToFileList() { Dataset files = compareToFileList; if (location != null) { - files = files.filter(files.col(FILE_PATH).startsWith(location)); + // Ensure path boundary is respected: s3://bucket/table should not match + // s3://bucket/table-backup/... which is a sibling path, not a subdirectory + String locationPrefix = + location.endsWith("/") ? location : location + "/"; + files = + files.filter( + files.col(FILE_PATH).startsWith(locationPrefix) + .or(files.col(FILE_PATH).equalTo(location))); } 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..c7a36140151b 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 @@ -1027,6 +1027,56 @@ protected long waitUntilAfter(long timestampMillis) { return current; } + @TestTemplate + public void testCompareToFileListDoesNotMatchSiblingPaths() throws IOException { + assumeThat(usePrefixListing) + .as("Should not test both prefix listing and Hadoop file listing (redundant)") + .isEqualTo(false); + Table table = TABLES.create(SCHEMA, PartitionSpec.unpartitioned(), properties, tableLocation); + + List records = + Lists.newArrayList(new ThreeColumnRecord(1, "AAAAAAAAAA", "AAAA")); + + Dataset df = spark.createDataFrame(records, ThreeColumnRecord.class).coalesce(1); + + df.select("c1", "c2", "c3").write().format("iceberg").mode("append").save(tableLocation); + + // sibling paths that share the table location as a raw string prefix but are NOT + // subdirectories of the table location (e.g. table-backup/...). These must NOT be + // treated as in-scope orphan candidates. + String sibling1 = tableLocation + "-backup/data/sibling1.parquet"; + String sibling2 = tableLocation + "_old/data/sibling2.parquet"; + String insideLocation = tableLocation + "/data/inside.parquet"; + + List mockFiles = + Lists.newArrayList( + new FilePathLastModifiedRecord(sibling1, new Timestamp(0L)), + new FilePathLastModifiedRecord(sibling2, new Timestamp(0L)), + new FilePathLastModifiedRecord(insideLocation, new Timestamp(0L))); + + Dataset compareToFileList = + spark + .createDataFrame(mockFiles, FilePathLastModifiedRecord.class) + .withColumnRenamed("filePath", "file_path") + .withColumnRenamed("lastModified", "last_modified"); + + DeleteOrphanFiles.Result result = + SparkActions.get() + .deleteOrphanFiles(table) + .compareToFileList(compareToFileList) + .olderThan(System.currentTimeMillis()) + .deleteWith(s -> {}) + .execute(); + + // Only the file actually inside the table location should be considered in scope. + assertThat(result.orphanFileLocations()) + .as("Only files inside the table location should be in scope") + .containsExactly(insideLocation); + assertThat(result.orphanFilesCount()) + .as("Only 1 file inside the table location should be in scope") + .isEqualTo(1L); + } + @TestTemplate public void testRemoveOrphanFilesWithStatisticFiles() throws Exception { assumeThat(usePrefixListing) 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 b47922820d21..3dc42417fa2b 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 @@ -225,7 +225,14 @@ public DeleteOrphanFilesSparkAction usePrefixListing(boolean newUsePrefixListing private Dataset filteredCompareToFileList() { Dataset files = compareToFileList; if (location != null) { - files = files.filter(files.col(FILE_PATH).startsWith(location)); + // Ensure path boundary is respected: s3://bucket/table should not match + // s3://bucket/table-backup/... which is a sibling path, not a subdirectory + String locationPrefix = + location.endsWith("/") ? location : location + "/"; + files = + files.filter( + files.col(FILE_PATH).startsWith(locationPrefix) + .or(files.col(FILE_PATH).equalTo(location))); } 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..73f1d4056fb8 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 @@ -1028,6 +1028,56 @@ protected long waitUntilAfter(long timestampMillis) { return current; } + @TestTemplate + public void testCompareToFileListDoesNotMatchSiblingPaths() throws IOException { + assumeThat(usePrefixListing) + .as("Should not test both prefix listing and Hadoop file listing (redundant)") + .isEqualTo(false); + Table table = TABLES.create(SCHEMA, PartitionSpec.unpartitioned(), properties, tableLocation); + + List records = + Lists.newArrayList(new ThreeColumnRecord(1, "AAAAAAAAAA", "AAAA")); + + Dataset df = spark.createDataFrame(records, ThreeColumnRecord.class).coalesce(1); + + df.select("c1", "c2", "c3").write().format("iceberg").mode("append").save(tableLocation); + + // sibling paths that share the table location as a raw string prefix but are NOT + // subdirectories of the table location (e.g. table-backup/...). These must NOT be + // treated as in-scope orphan candidates. + String sibling1 = tableLocation + "-backup/data/sibling1.parquet"; + String sibling2 = tableLocation + "_old/data/sibling2.parquet"; + String insideLocation = tableLocation + "/data/inside.parquet"; + + List mockFiles = + Lists.newArrayList( + new FilePathLastModifiedRecord(sibling1, new Timestamp(0L)), + new FilePathLastModifiedRecord(sibling2, new Timestamp(0L)), + new FilePathLastModifiedRecord(insideLocation, new Timestamp(0L))); + + Dataset compareToFileList = + spark + .createDataFrame(mockFiles, FilePathLastModifiedRecord.class) + .withColumnRenamed("filePath", "file_path") + .withColumnRenamed("lastModified", "last_modified"); + + DeleteOrphanFiles.Result result = + SparkActions.get() + .deleteOrphanFiles(table) + .compareToFileList(compareToFileList) + .olderThan(System.currentTimeMillis()) + .deleteWith(s -> {}) + .execute(); + + // Only the file actually inside the table location should be considered in scope. + assertThat(result.orphanFileLocations()) + .as("Only files inside the table location should be in scope") + .containsExactly(insideLocation); + assertThat(result.orphanFilesCount()) + .as("Only 1 file inside the table location should be in scope") + .isEqualTo(1L); + } + @TestTemplate public void testRemoveOrphanFilesWithStatisticFiles() throws Exception { assumeThat(usePrefixListing)