Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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;
Expand Down Expand Up @@ -131,6 +133,7 @@ public class DeleteOrphanFilesSparkAction extends BaseSparkAction<DeleteOrphanFi
private Consumer<String> deleteFunc = null;
private ExecutorService deleteExecutorService = null;
private boolean usePrefixListing = false;
private boolean useParallelDeletes = false;
private static final Encoder<FileURI> FILE_URI_ENCODER = Encoders.bean(FileURI.class);

DeleteOrphanFilesSparkAction(SparkSession spark, Table table) {
Expand Down Expand Up @@ -222,6 +225,11 @@ public DeleteOrphanFilesSparkAction usePrefixListing(boolean newUsePrefixListing
return this;
}

public DeleteOrphanFilesSparkAction useParallelDeletes(boolean newUseParallelDeletes) {
this.useParallelDeletes = newUseParallelDeletes;
return this;
}

private Dataset<String> filteredCompareToFileList() {
Dataset<Row> files = compareToFileList;
if (location != null) {
Expand Down Expand Up @@ -279,20 +287,25 @@ private DeleteOrphanFiles.Result deleteFiles(Dataset<String> orphanFileDS) {
List<String> orphanFileList = Lists.newArrayListWithCapacity(maxSampleSize);
long filesCount = 0;

if (useParallelDeletes) {
orphanFileDS = orphanFileDS.mapPartitions(new DeleteFiles(table.io()), Encoders.STRING());
}

Iterator<String> orphanFiles =
streamResults() ? orphanFileDS.toLocalIterator() : orphanFileDS.collectAsList().iterator();

Iterator<List<String>> fileGroups = Iterators.partition(orphanFiles, DELETE_GROUP_SIZE);

while (fileGroups.hasNext()) {
List<String> fileGroup = fileGroups.next();

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();
Expand All @@ -316,7 +329,7 @@ private void collectPathsForOutput(
}
}

private void deleteBulk(SupportsBulkOperations io, List<String> paths) {
private static void deleteBulk(SupportsBulkOperations io, List<String> paths) {
try {
io.deleteFiles(paths);
LOG.info("Deleted {} files using bulk deletes", paths.size());
Expand All @@ -327,7 +340,11 @@ private void deleteBulk(SupportsBulkOperations io, List<String> paths) {
}
}

private void deleteNonBulk(List<String> paths) {
private static void deleteNonBulk(
FileIO io,
List<String> paths,
@Nullable Consumer<String> deleteFunc,
@Nullable ExecutorService deleteExecutorService) {
Tasks.Builder<String> deleteTasks =
Tasks.foreach(paths)
.noRetry()
Expand All @@ -338,8 +355,8 @@ private void deleteNonBulk(List<String> 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);
Expand Down Expand Up @@ -567,6 +584,26 @@ private String toOrphanFile(Tuple2<FileURI, FileURI> row) {
}
}

private static final class DeleteFiles implements MapPartitionsFunction<String, String> {

private final FileIO io;

public DeleteFiles(FileIO io) {
this.io = io;
}

@Override
public Iterator<String> call(Iterator<String> input) {
List<String> 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<String> {
StringToFileURI(Map<String, String> equalSchemes, Map<String, String> equalAuthorities) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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[] {
Expand All @@ -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 =
Expand Down Expand Up @@ -149,6 +153,7 @@ public Iterator<Scan> 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,
Expand Down Expand Up @@ -197,6 +202,7 @@ public Iterator<Scan> call(InternalRow args) {
}

action.usePrefixListing(prefixListing);
action.useParallelDeletes(parallelDeletes);

if (streamResults) {
action.option("stream-results", "true");
Expand Down
Loading