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..92b1ec80e89e 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 @@ -32,6 +32,7 @@ import java.util.concurrent.TimeUnit; import java.util.function.Consumer; import java.util.function.Predicate; +import javax.annotation.Nullable; import org.apache.hadoop.fs.FileStatus; import org.apache.hadoop.fs.FileSystem; import org.apache.hadoop.fs.Path; @@ -42,6 +43,7 @@ import org.apache.iceberg.actions.ImmutableDeleteOrphanFiles; import org.apache.iceberg.exceptions.ValidationException; import org.apache.iceberg.io.BulkDeletionFailureException; +import org.apache.iceberg.io.FileIO; import org.apache.iceberg.io.SupportsBulkOperations; import org.apache.iceberg.io.SupportsPrefixOperations; import org.apache.iceberg.relocated.com.google.common.annotations.VisibleForTesting; @@ -131,6 +133,7 @@ public class DeleteOrphanFilesSparkAction extends BaseSparkAction deleteFunc = null; private ExecutorService deleteExecutorService = null; private boolean usePrefixListing = false; + private boolean useParallelDeletes = false; private static final Encoder FILE_URI_ENCODER = Encoders.bean(FileURI.class); DeleteOrphanFilesSparkAction(SparkSession spark, Table table) { @@ -222,6 +225,11 @@ public DeleteOrphanFilesSparkAction usePrefixListing(boolean newUsePrefixListing return this; } + public DeleteOrphanFilesSparkAction useParallelDeletes(boolean newUseParallelDeletes) { + this.useParallelDeletes = newUseParallelDeletes; + return this; + } + private Dataset filteredCompareToFileList() { Dataset files = compareToFileList; if (location != null) { @@ -279,9 +287,12 @@ private DeleteOrphanFiles.Result deleteFiles(Dataset orphanFileDS) { List orphanFileList = Lists.newArrayListWithCapacity(maxSampleSize); long filesCount = 0; + if (useParallelDeletes) { + orphanFileDS = orphanFileDS.mapPartitions(new DeleteFiles(table.io()), Encoders.STRING()); + } + Iterator orphanFiles = streamResults() ? orphanFileDS.toLocalIterator() : orphanFileDS.collectAsList().iterator(); - Iterator> fileGroups = Iterators.partition(orphanFiles, DELETE_GROUP_SIZE); while (fileGroups.hasNext()) { @@ -289,10 +300,12 @@ private DeleteOrphanFiles.Result deleteFiles(Dataset orphanFileDS) { collectPathsForOutput(fileGroup, orphanFileList, maxSampleSize); - if (deleteFunc == null && table.io() instanceof SupportsBulkOperations) { + if (!useParallelDeletes + && deleteFunc == null + && table.io() instanceof SupportsBulkOperations) { deleteBulk((SupportsBulkOperations) table.io(), fileGroup); - } else { - deleteNonBulk(fileGroup); + } else if (!useParallelDeletes) { + deleteNonBulk(table.io(), fileGroup, deleteFunc, deleteExecutorService); } filesCount += fileGroup.size(); @@ -316,7 +329,7 @@ private void collectPathsForOutput( } } - private void deleteBulk(SupportsBulkOperations io, List paths) { + private static void deleteBulk(SupportsBulkOperations io, List paths) { try { io.deleteFiles(paths); LOG.info("Deleted {} files using bulk deletes", paths.size()); @@ -327,7 +340,11 @@ private void deleteBulk(SupportsBulkOperations io, List paths) { } } - private void deleteNonBulk(List paths) { + private static void deleteNonBulk( + FileIO io, + List paths, + @Nullable Consumer deleteFunc, + @Nullable ExecutorService deleteExecutorService) { Tasks.Builder deleteTasks = Tasks.foreach(paths) .noRetry() @@ -338,8 +355,8 @@ private void deleteNonBulk(List paths) { if (deleteFunc == null) { LOG.info( "Table IO {} does not support bulk operations. Using non-bulk deletes.", - table.io().getClass().getName()); - deleteTasks.run(table.io()::deleteFile); + io.getClass().getName()); + deleteTasks.run(io::deleteFile); } else { LOG.info("Custom delete function provided. Using non-bulk deletes"); deleteTasks.run(deleteFunc::accept); @@ -567,6 +584,26 @@ private String toOrphanFile(Tuple2 row) { } } + private static final class DeleteFiles implements MapPartitionsFunction { + + private final FileIO io; + + public DeleteFiles(FileIO io) { + this.io = io; + } + + @Override + public Iterator call(Iterator input) { + List paths = Lists.newArrayList(input); + if (io instanceof SupportsBulkOperations) { + deleteBulk((SupportsBulkOperations) io, paths); + } else { + deleteNonBulk(io, paths, null, null); + } + return paths.iterator(); + } + } + @VisibleForTesting static class StringToFileURI extends ToFileURI { StringToFileURI(Map equalSchemes, Map equalAuthorities) { diff --git a/spark/v4.1/spark/src/main/java/org/apache/iceberg/spark/procedures/RemoveOrphanFilesProcedure.java b/spark/v4.1/spark/src/main/java/org/apache/iceberg/spark/procedures/RemoveOrphanFilesProcedure.java index 89e218c10360..35812a27da55 100644 --- a/spark/v4.1/spark/src/main/java/org/apache/iceberg/spark/procedures/RemoveOrphanFilesProcedure.java +++ b/spark/v4.1/spark/src/main/java/org/apache/iceberg/spark/procedures/RemoveOrphanFilesProcedure.java @@ -80,6 +80,9 @@ public class RemoveOrphanFilesProcedure extends BaseProcedure { // Stream results to avoid loading all orphan files in driver memory. Default is false. private static final ProcedureParameter STREAM_RESULTS_PARAM = optionalInParameter("stream_results", DataTypes.BooleanType); + // Delete files in parallel using Spark executors. Default is false. + private static final ProcedureParameter PARALLEL_DELETES_PARAM = + optionalInParameter("parallel_deletes", DataTypes.BooleanType); private static final ProcedureParameter[] PARAMETERS = new ProcedureParameter[] { @@ -93,7 +96,8 @@ public class RemoveOrphanFilesProcedure extends BaseProcedure { EQUAL_AUTHORITIES_PARAM, PREFIX_MISMATCH_MODE_PARAM, PREFIX_LISTING_PARAM, - STREAM_RESULTS_PARAM + STREAM_RESULTS_PARAM, + PARALLEL_DELETES_PARAM }; private static final StructType OUTPUT_TYPE = @@ -149,6 +153,7 @@ public Iterator call(InternalRow args) { boolean prefixListing = input.asBoolean(PREFIX_LISTING_PARAM, false); boolean streamResults = input.asBoolean(STREAM_RESULTS_PARAM, false); + boolean parallelDeletes = input.asBoolean(PARALLEL_DELETES_PARAM, false); return withIcebergTable( tableIdent, @@ -197,6 +202,7 @@ public Iterator call(InternalRow args) { } action.usePrefixListing(prefixListing); + action.useParallelDeletes(parallelDeletes); if (streamResults) { action.option("stream-results", "true");