diff --git a/spark/v3.5/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/TestRemoveOrphanFilesProcedure.java b/spark/v3.5/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/TestRemoveOrphanFilesProcedure.java index c17dbad71100..d2f825540233 100644 --- a/spark/v3.5/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/TestRemoveOrphanFilesProcedure.java +++ b/spark/v3.5/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/TestRemoveOrphanFilesProcedure.java @@ -750,4 +750,61 @@ public void testRemoveOrphanFilesProcedureWithEqualAuthorities() // Dropping the table here sql("DROP TABLE %s", tableName); } + + @TestTemplate + public void testRemoveOrphanFilesFileListViewDoesNotMatchSiblingPaths() + throws NoSuchTableException, ParseException, IOException { + if (catalogName.equals("testhadoop")) { + sql("CREATE TABLE %s (id bigint NOT NULL, data string) USING iceberg", tableName); + } else { + sql( + "CREATE TABLE %s (id bigint NOT NULL, data string) USING iceberg LOCATION '%s'", + tableName, java.nio.file.Files.createTempDirectory(temp, "junit")); + } + Table table = Spark3Util.loadIcebergTable(spark, tableName); + String tableLocation = table.location(); + // The sibling path has the table location as a raw string prefix, e.g. /tmp/abc-sibling when + // table is at /tmp/abc. Without the trailing-slash fix, startsWith(tableLocation) would + // incorrectly match files under this sibling path. + String siblingPath = tableLocation + "-sibling"; + + Timestamp lastModifiedTimestamp = new Timestamp(10000); + String orphanFile = tableLocation + "/orphan.parquet"; + String siblingFile = siblingPath + "/data.parquet"; + + List allFiles = Lists.newArrayList(); + allFiles.add(new FilePathLastModifiedRecord(orphanFile, lastModifiedTimestamp)); + allFiles.add(new FilePathLastModifiedRecord(siblingFile, lastModifiedTimestamp)); + allFiles.add( + new FilePathLastModifiedRecord( + ReachableFileUtil.versionHintLocation(table), lastModifiedTimestamp)); + for (String metaFile : ReachableFileUtil.metadataFileLocations(table, true)) { + allFiles.add(new FilePathLastModifiedRecord(metaFile, lastModifiedTimestamp)); + } + + Dataset compareToFileList = + spark + .createDataFrame(allFiles, FilePathLastModifiedRecord.class) + .withColumnRenamed("filePath", "file_path") + .withColumnRenamed("lastModified", "last_modified"); + String fileListViewName = "files_view"; + compareToFileList.createOrReplaceTempView(fileListViewName); + + List orphanFiles = + sql( + "CALL %s.system.remove_orphan_files(" + + "table => '%s'," + + "dry_run => true," + + "file_list_view => '%s')", + catalogName, tableIdent, fileListViewName); + + assertThat(orphanFiles) + .as("Files under sibling path must not be identified as orphans for this table") + .extracting(row -> row[0]) + .doesNotContain(siblingFile); + assertThat(orphanFiles) + .as("Orphan file under the table location must be identified") + .extracting(row -> row[0]) + .contains(orphanFile); + } } 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..83235cfd5bbb 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,8 @@ public DeleteOrphanFilesSparkAction usePrefixListing(boolean newUsePrefixListing private Dataset filteredCompareToFileList() { Dataset files = compareToFileList; if (location != null) { - files = files.filter(files.col(FILE_PATH).startsWith(location)); + String locationPrefix = location.endsWith("/") ? location : location + "/"; + files = files.filter(files.col(FILE_PATH).startsWith(locationPrefix)); } return files .filter(files.col(LAST_MODIFIED).lt(new Timestamp(olderThanTimestamp))) diff --git a/spark/v4.0/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/TestRemoveOrphanFilesProcedure.java b/spark/v4.0/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/TestRemoveOrphanFilesProcedure.java index a5ac8a7e01ac..c77632803b48 100644 --- a/spark/v4.0/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/TestRemoveOrphanFilesProcedure.java +++ b/spark/v4.0/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/TestRemoveOrphanFilesProcedure.java @@ -749,4 +749,61 @@ public void testRemoveOrphanFilesProcedureWithEqualAuthorities() // Dropping the table here sql("DROP TABLE %s", tableName); } + + @TestTemplate + public void testRemoveOrphanFilesFileListViewDoesNotMatchSiblingPaths() + throws NoSuchTableException, ParseException, IOException { + if (catalogName.equals("testhadoop")) { + sql("CREATE TABLE %s (id bigint NOT NULL, data string) USING iceberg", tableName); + } else { + sql( + "CREATE TABLE %s (id bigint NOT NULL, data string) USING iceberg LOCATION '%s'", + tableName, java.nio.file.Files.createTempDirectory(temp, "junit")); + } + Table table = Spark3Util.loadIcebergTable(spark, tableName); + String tableLocation = table.location(); + // The sibling path has the table location as a raw string prefix, e.g. /tmp/abc-sibling when + // table is at /tmp/abc. Without the trailing-slash fix, startsWith(tableLocation) would + // incorrectly match files under this sibling path. + String siblingPath = tableLocation + "-sibling"; + + Timestamp lastModifiedTimestamp = new Timestamp(10000); + String orphanFile = tableLocation + "/orphan.parquet"; + String siblingFile = siblingPath + "/data.parquet"; + + List allFiles = Lists.newArrayList(); + allFiles.add(new FilePathLastModifiedRecord(orphanFile, lastModifiedTimestamp)); + allFiles.add(new FilePathLastModifiedRecord(siblingFile, lastModifiedTimestamp)); + allFiles.add( + new FilePathLastModifiedRecord( + ReachableFileUtil.versionHintLocation(table), lastModifiedTimestamp)); + for (String metaFile : ReachableFileUtil.metadataFileLocations(table, true)) { + allFiles.add(new FilePathLastModifiedRecord(metaFile, lastModifiedTimestamp)); + } + + Dataset compareToFileList = + spark + .createDataFrame(allFiles, FilePathLastModifiedRecord.class) + .withColumnRenamed("filePath", "file_path") + .withColumnRenamed("lastModified", "last_modified"); + String fileListViewName = "files_view"; + compareToFileList.createOrReplaceTempView(fileListViewName); + + List orphanFiles = + sql( + "CALL %s.system.remove_orphan_files(" + + "table => '%s'," + + "dry_run => true," + + "file_list_view => '%s')", + catalogName, tableIdent, fileListViewName); + + assertThat(orphanFiles) + .as("Files under sibling path must not be identified as orphans for this table") + .extracting(row -> row[0]) + .doesNotContain(siblingFile); + assertThat(orphanFiles) + .as("Orphan file under the table location must be identified") + .extracting(row -> row[0]) + .contains(orphanFile); + } } 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..83235cfd5bbb 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,8 @@ public DeleteOrphanFilesSparkAction usePrefixListing(boolean newUsePrefixListing private Dataset filteredCompareToFileList() { Dataset files = compareToFileList; if (location != null) { - files = files.filter(files.col(FILE_PATH).startsWith(location)); + String locationPrefix = location.endsWith("/") ? location : location + "/"; + files = files.filter(files.col(FILE_PATH).startsWith(locationPrefix)); } return files .filter(files.col(LAST_MODIFIED).lt(new Timestamp(olderThanTimestamp))) diff --git a/spark/v4.1/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/TestRemoveOrphanFilesProcedure.java b/spark/v4.1/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/TestRemoveOrphanFilesProcedure.java index a5ac8a7e01ac..c77632803b48 100644 --- a/spark/v4.1/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/TestRemoveOrphanFilesProcedure.java +++ b/spark/v4.1/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/TestRemoveOrphanFilesProcedure.java @@ -749,4 +749,61 @@ public void testRemoveOrphanFilesProcedureWithEqualAuthorities() // Dropping the table here sql("DROP TABLE %s", tableName); } + + @TestTemplate + public void testRemoveOrphanFilesFileListViewDoesNotMatchSiblingPaths() + throws NoSuchTableException, ParseException, IOException { + if (catalogName.equals("testhadoop")) { + sql("CREATE TABLE %s (id bigint NOT NULL, data string) USING iceberg", tableName); + } else { + sql( + "CREATE TABLE %s (id bigint NOT NULL, data string) USING iceberg LOCATION '%s'", + tableName, java.nio.file.Files.createTempDirectory(temp, "junit")); + } + Table table = Spark3Util.loadIcebergTable(spark, tableName); + String tableLocation = table.location(); + // The sibling path has the table location as a raw string prefix, e.g. /tmp/abc-sibling when + // table is at /tmp/abc. Without the trailing-slash fix, startsWith(tableLocation) would + // incorrectly match files under this sibling path. + String siblingPath = tableLocation + "-sibling"; + + Timestamp lastModifiedTimestamp = new Timestamp(10000); + String orphanFile = tableLocation + "/orphan.parquet"; + String siblingFile = siblingPath + "/data.parquet"; + + List allFiles = Lists.newArrayList(); + allFiles.add(new FilePathLastModifiedRecord(orphanFile, lastModifiedTimestamp)); + allFiles.add(new FilePathLastModifiedRecord(siblingFile, lastModifiedTimestamp)); + allFiles.add( + new FilePathLastModifiedRecord( + ReachableFileUtil.versionHintLocation(table), lastModifiedTimestamp)); + for (String metaFile : ReachableFileUtil.metadataFileLocations(table, true)) { + allFiles.add(new FilePathLastModifiedRecord(metaFile, lastModifiedTimestamp)); + } + + Dataset compareToFileList = + spark + .createDataFrame(allFiles, FilePathLastModifiedRecord.class) + .withColumnRenamed("filePath", "file_path") + .withColumnRenamed("lastModified", "last_modified"); + String fileListViewName = "files_view"; + compareToFileList.createOrReplaceTempView(fileListViewName); + + List orphanFiles = + sql( + "CALL %s.system.remove_orphan_files(" + + "table => '%s'," + + "dry_run => true," + + "file_list_view => '%s')", + catalogName, tableIdent, fileListViewName); + + assertThat(orphanFiles) + .as("Files under sibling path must not be identified as orphans for this table") + .extracting(row -> row[0]) + .doesNotContain(siblingFile); + assertThat(orphanFiles) + .as("Orphan file under the table location must be identified") + .extracting(row -> row[0]) + .contains(orphanFile); + } } 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..83235cfd5bbb 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,8 @@ public DeleteOrphanFilesSparkAction usePrefixListing(boolean newUsePrefixListing private Dataset filteredCompareToFileList() { Dataset files = compareToFileList; if (location != null) { - files = files.filter(files.col(FILE_PATH).startsWith(location)); + String locationPrefix = location.endsWith("/") ? location : location + "/"; + files = files.filter(files.col(FILE_PATH).startsWith(locationPrefix)); } return files .filter(files.col(LAST_MODIFIED).lt(new Timestamp(olderThanTimestamp)))