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..3ad114cb3ba8 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 @@ -51,12 +51,14 @@ import org.apache.iceberg.relocated.com.google.common.collect.Lists; import org.apache.iceberg.relocated.com.google.common.collect.Maps; import org.apache.iceberg.spark.JobGroupInfo; +import org.apache.iceberg.spark.source.SerializableTableWithSize; import org.apache.iceberg.util.FileSystemWalker; import org.apache.iceberg.util.Pair; import org.apache.iceberg.util.PropertyUtil; import org.apache.iceberg.util.Tasks; import org.apache.spark.api.java.JavaRDD; import org.apache.spark.api.java.function.FlatMapFunction; +import org.apache.spark.api.java.function.ForeachPartitionFunction; import org.apache.spark.api.java.function.MapPartitionsFunction; import org.apache.spark.broadcast.Broadcast; import org.apache.spark.sql.Column; @@ -68,6 +70,8 @@ import org.apache.spark.sql.types.DataTypes; import org.apache.spark.sql.types.StructField; import org.apache.spark.sql.types.StructType; +import org.apache.spark.storage.StorageLevel; +import org.apache.spark.util.LongAccumulator; import org.apache.spark.util.SerializableConfiguration; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -105,6 +109,12 @@ public class DeleteOrphanFilesSparkAction extends BaseSparkAction actualFileIdentDS = actualFileIdentDS(); Dataset validFileIdentDS = validFileIdentDS(); @@ -269,6 +284,10 @@ private DeleteOrphanFiles.Result doExecute() { /** * Deletes orphan files from the cached dataset. * + *

When the table IO supports bulk operations and no custom deleteFunc is provided, deletion + * runs on executors in parallel via {@code foreachPartition}. Otherwise, falls back to the stock + * driver-side path (preserves custom deleteFunc and non-bulk IO behavior). + * * @param orphanFileDS the cached dataset of orphan files * @return result with orphan file paths */ @@ -276,6 +295,75 @@ private DeleteOrphanFiles.Result deleteFiles(Dataset orphanFileDS) { int maxSampleSize = PropertyUtil.propertyAsInt( options(), MAX_ORPHAN_FILE_SAMPLE_SIZE, MAX_ORPHAN_FILE_SAMPLE_SIZE_DEFAULT); + + if (deleteFunc == null && table.io() instanceof SupportsBulkOperations) { + return deleteFilesExecutorSide(orphanFileDS, maxSampleSize); + } + return deleteFilesDriverSide(orphanFileDS, maxSampleSize); + } + + private DeleteOrphanFiles.Result deleteFilesExecutorSide( + Dataset orphanFileDS, int maxSampleSize) { + LOG.warn( + "Starting executor-side orphan file deletion (table={}, partitions={})", + table.name(), + orphanFileDS.rdd().getNumPartitions()); + + // Sample paths for the result (cheap: take from already-persisted dataset) + List samplePaths = Lists.newArrayList(orphanFileDS.takeAsList(maxSampleSize)); + LOG.warn("Collected {} sample paths for result", samplePaths.size()); + + LongAccumulator deletedCount = spark().sparkContext().longAccumulator("orphanFilesDeleted"); + + // Broadcast a serializable copy of the table so each executor has its FileIO + Broadcast tableBroadcast = + sparkContext().broadcast(SerializableTableWithSize.copyOf(table)); + final int deleteGroupSize = DELETE_GROUP_SIZE; + + orphanFileDS.foreachPartition( + (ForeachPartitionFunction) + partition -> { + Table broadcastTable = tableBroadcast.value(); + SupportsBulkOperations bulkIO = (SupportsBulkOperations) broadcastTable.io(); + List batch = Lists.newArrayListWithCapacity(deleteGroupSize); + int batchNum = 0; + while (partition.hasNext()) { + batch.add(partition.next()); + if (batch.size() >= deleteGroupSize) { + deleteBatch(bulkIO, batch, ++batchNum, deletedCount); + batch.clear(); + } + } + if (!batch.isEmpty()) { + deleteBatch(bulkIO, batch, ++batchNum, deletedCount); + } + }); + + long total = deletedCount.value(); + LOG.warn("Executor-side deletion complete: deleted {} orphan files", total); + + return ImmutableDeleteOrphanFiles.Result.builder() + .orphanFileLocations(samplePaths) + .orphanFilesCount(total) + .build(); + } + + private static void deleteBatch( + SupportsBulkOperations bulkIO, List batch, int batchNum, LongAccumulator counter) { + int size = batch.size(); + try { + bulkIO.deleteFiles(batch); + counter.add(size); + LOG.warn("Deleted batch #{} ({} files)", batchNum, size); + } catch (BulkDeletionFailureException e) { + int deleted = size - e.numberFailedObjects(); + counter.add(deleted); + LOG.warn("Deleted batch #{} only {} of {} files", batchNum, deleted, size, e); + } + } + + private DeleteOrphanFiles.Result deleteFilesDriverSide( + Dataset orphanFileDS, int maxSampleSize) { List orphanFileList = Lists.newArrayListWithCapacity(maxSampleSize); long filesCount = 0; @@ -362,8 +450,9 @@ static Dataset findOrphanFiles( .joinWith(validFileIdentDS, joinCond, "leftouter") .mapPartitions(new FindOrphanFiles(prefixMismatchMode, conflicts), Encoders.STRING()); - // Cache and force computation to populate conflicts accumulator - orphanFileDS = orphanFileDS.cache(); + // Persist to disk (survives executor decommission better than memory cache) and force + // computation to populate conflicts accumulator + orphanFileDS = orphanFileDS.persist(StorageLevel.DISK_ONLY()); try { orphanFileDS.count(); @@ -422,17 +511,59 @@ private Dataset listedFileDS() { "Cannot use prefix listing with FileIO {} which does not support prefix operations.", table.io()); - Predicate predicate = - fileInfo -> fileInfo.createdAtMillis() < olderThanTimestamp; - FileSystemWalker.listDirRecursivelyWithFileIO( - (SupportsPrefixOperations) table.io(), + if (!parallelPrefixListing()) { + // Driver-side fallback: listPrefix() recursively enumerates the entire table on a single + // driver thread. Kept as an escape hatch behind parallel-prefix-listing=false. + Predicate predicate = + fileInfo -> fileInfo.createdAtMillis() < olderThanTimestamp; + FileSystemWalker.listDirRecursivelyWithFileIO( + (SupportsPrefixOperations) table.io(), + location, + table.specs(), + predicate, + matchingFiles::add); + + JavaRDD matchingFileRDD = sparkContext().parallelize(matchingFiles, 1); + return spark().createDataset(matchingFileRDD.rdd(), Encoders.STRING()); + } + + // Parallel prefix listing. SupportsPrefixOperations.listPrefix() is recursive-only and + // cannot enumerate a single directory level, so it cannot discover sub-prefixes to fan out + // by itself. We discover the shallow sub-prefixes with a cheap depth-limited Hadoop listing + // on the driver (delimiter-based, returns directories), then dispatch each deferred + // sub-prefix to the executors where listPrefix() runs in parallel. This mirrors the Hadoop + // walker's driver-discovery + executor-dispatch pattern below, but performs the expensive + // deep listing with the S3-native (flat) listPrefix() instead of per-directory Hadoop calls. + Predicate driverPredicate = + file -> file.getModificationTime() < olderThanTimestamp; + FileSystemWalker.listDirRecursivelyWithHadoop( location, table.specs(), - predicate, + driverPredicate, + hadoopConf.value(), + MAX_DRIVER_LISTING_DEPTH, + MAX_DRIVER_LISTING_DIRECT_SUB_DIRS, + subDirs::add, matchingFiles::add); JavaRDD matchingFileRDD = sparkContext().parallelize(matchingFiles, 1); - return spark().createDataset(matchingFileRDD.rdd(), Encoders.STRING()); + + if (subDirs.isEmpty()) { + return spark().createDataset(matchingFileRDD.rdd(), Encoders.STRING()); + } + + int parallelism = Math.min(subDirs.size(), listingParallelism); + JavaRDD subDirRDD = sparkContext().parallelize(subDirs, parallelism); + + // Broadcast a serializable copy of the table so each executor has a usable FileIO. + Broadcast
tableBroadcast = + sparkContext().broadcast(SerializableTableWithSize.copyOf(table)); + ListDirsRecursivelyWithFileIO listDirs = + new ListDirsRecursivelyWithFileIO(tableBroadcast, olderThanTimestamp); + JavaRDD matchingLeafFileRDD = subDirRDD.mapPartitions(listDirs); + + JavaRDD completeMatchingFileRDD = matchingFileRDD.union(matchingLeafFileRDD); + return spark().createDataset(completeMatchingFileRDD.rdd(), Encoders.STRING()); } else { Predicate predicate = file -> file.getModificationTime() < olderThanTimestamp; // list at most MAX_DRIVER_LISTING_DEPTH levels and only dirs that have @@ -523,6 +654,40 @@ public Iterator call(Iterator dirs) throws Exception { } } + /** + * Lists each assigned sub-prefix recursively on an executor using {@link + * SupportsPrefixOperations#listPrefix}. The table (and therefore its FileIO) is obtained from a + * broadcast so the prefix listing runs in parallel across the cluster rather than on the driver. + */ + private static class ListDirsRecursivelyWithFileIO + implements FlatMapFunction, String> { + + private final Broadcast
tableBroadcast; + private final long olderThanTimestamp; + + ListDirsRecursivelyWithFileIO(Broadcast
tableBroadcast, long olderThanTimestamp) { + this.tableBroadcast = tableBroadcast; + this.olderThanTimestamp = olderThanTimestamp; + } + + @Override + public Iterator call(Iterator dirs) throws Exception { + Table broadcastTable = tableBroadcast.value(); + SupportsPrefixOperations io = (SupportsPrefixOperations) broadcastTable.io(); + Map specs = broadcastTable.specs(); + Predicate predicate = + fileInfo -> fileInfo.createdAtMillis() < olderThanTimestamp; + + List files = Lists.newArrayList(); + while (dirs.hasNext()) { + FileSystemWalker.listDirRecursivelyWithFileIO( + io, dirs.next(), specs, predicate, files::add); + } + + return files.iterator(); + } + } + private static class FindOrphanFiles implements MapPartitionsFunction, String> {