kinolaev commented on code in PR #15727:
URL: https://github.com/apache/iceberg/pull/15727#discussion_r3979807632
##########
spark/v3.5/spark/src/main/java/org/apache/iceberg/spark/actions/RemoveDanglingDeletesSparkAction.java:
##########
@@ -18,189 +18,123 @@
*/
package org.apache.iceberg.spark.actions;
-import static org.apache.spark.sql.functions.col;
-import static org.apache.spark.sql.functions.min;
-
-import java.util.Collections;
+import java.io.Closeable;
+import java.io.IOException;
+import java.util.Iterator;
import java.util.List;
-import java.util.stream.Collectors;
-import org.apache.iceberg.DataFile;
+import java.util.stream.StreamSupport;
import org.apache.iceberg.DeleteFile;
-import org.apache.iceberg.FileFormat;
-import org.apache.iceberg.MetadataTableType;
-import org.apache.iceberg.Partitioning;
-import org.apache.iceberg.RewriteFiles;
+import org.apache.iceberg.FileScanTask;
+import org.apache.iceberg.ManifestFiles;
+import org.apache.iceberg.ManifestReader;
+import org.apache.iceberg.Snapshot;
import org.apache.iceberg.Table;
-import org.apache.iceberg.actions.ImmutableRemoveDanglingDeleteFiles;
+import org.apache.iceberg.TableScan;
import org.apache.iceberg.actions.RemoveDanglingDeleteFiles;
-import org.apache.iceberg.spark.JobGroupInfo;
-import org.apache.iceberg.spark.SparkDeleteFile;
-import org.apache.iceberg.types.Types;
-import org.apache.iceberg.util.DeleteFileSet;
-import org.apache.spark.sql.Column;
-import org.apache.spark.sql.Dataset;
-import org.apache.spark.sql.Row;
+import org.apache.iceberg.actions.RemoveDanglingDeleteFilesAction;
+import
org.apache.iceberg.actions.RemoveDanglingDeleteFilesAction.DeleteFileKey;
+import org.apache.iceberg.io.CloseableIterable;
+import org.apache.iceberg.io.CloseableIterator;
+import org.apache.iceberg.io.ClosingIterator;
+import org.apache.iceberg.relocated.com.google.common.collect.ImmutableList;
+import org.apache.iceberg.spark.source.SerializableTableWithSize;
+import org.apache.spark.api.java.JavaPairRDD;
+import org.apache.spark.broadcast.Broadcast;
import org.apache.spark.sql.SparkSession;
-import org.apache.spark.sql.types.StructType;
-import org.slf4j.Logger;
-import org.slf4j.LoggerFactory;
+import scala.Tuple2;
/**
* An action that removes dangling delete files from the current snapshot. A
delete file is dangling
* if its deletes no longer applies to any live data files.
- *
- * <p>The following dangling delete files are removed:
- *
- * <ul>
- * <li>Position delete files with a data sequence number less than that of
any data file in the
- * same partition
- * <li>Equality delete files with a data sequence number less than or equal
to that of any data
- * file in the same partition
- * </ul>
*/
class RemoveDanglingDeletesSparkAction
extends BaseSnapshotUpdateSparkAction<RemoveDanglingDeletesSparkAction>
implements RemoveDanglingDeleteFiles {
- private static final Logger LOG =
LoggerFactory.getLogger(RemoveDanglingDeletesSparkAction.class);
private final Table table;
+ private final RemoveDanglingDeleteFilesAction action;
protected RemoveDanglingDeletesSparkAction(SparkSession spark, Table table) {
super(spark);
this.table = table;
+ this.action = new RemoveDanglingDeleteFilesAction(table,
this::findDanglingDeletes);
}
@Override
protected RemoveDanglingDeletesSparkAction self() {
return this;
}
+ public RemoveDanglingDeletesSparkAction toBranch(String targetBranch) {
+ action.toBranch(targetBranch);
+ return this;
+ }
+
@Override
public Result execute() {
- if (table.specs().size() == 1 && table.spec().isUnpartitioned()) {
- // ManifestFilterManager already performs this table-wide delete on each
commit
- return ImmutableRemoveDanglingDeleteFiles.Result.builder()
- .removedDeleteFiles(Collections.emptyList())
- .build();
- }
-
+ commitSummary().forEach(action::set);
String desc = String.format("Removing dangling delete files in %s",
table.name());
- JobGroupInfo info = newJobGroupInfo("REMOVE-DELETES", desc);
- return withJobGroupInfo(info, this::doExecute);
+ return withJobGroupInfo(newJobGroupInfo("REMOVE-DELETES", desc),
action::execute);
}
- Result doExecute() {
- RewriteFiles rewriteFiles = table.newRewrite();
- DeleteFileSet danglingDeletes = DeleteFileSet.create();
- danglingDeletes.addAll(findDanglingDeletes());
- danglingDeletes.addAll(findDanglingDvs());
+ private List<DeleteFile> findDanglingDeletes(Snapshot snapshot) {
+ Broadcast<Table> tableBroadcast =
+ sparkContext().broadcast(SerializableTableWithSize.copyOf(table));
+
+ JavaPairRDD<DeleteFileKey, Void> referencedKeys =
+ sparkContext()
+ .parallelize(ImmutableList.of(snapshot.snapshotId()), 1)
+ .flatMap(
+ snapshotId -> {
+ TableScan scan =
tableBroadcast.value().newScan().useSnapshot(snapshotId);
+ return new ClosingIterator<>(new
DeleteFileKeyIterator(scan.planFiles()));
+ })
+ .mapToPair(key -> new Tuple2<>(key, (Void) null));
+
+ List<ManifestFileBean> deleteManifests =
+
snapshot.deleteManifests(table.io()).stream().map(ManifestFileBean::fromManifest).toList();
+ JavaPairRDD<DeleteFileKey, DeleteFile> allDeletes =
+ sparkContext()
+ .parallelize(deleteManifests, deleteManifests.size())
+ .flatMap(
+ manifest -> {
+ ManifestReader<DeleteFile> reader =
+ ManifestFiles.readDeleteManifest(
+ manifest, tableBroadcast.value().io(),
tableBroadcast.value().specs());
+ return new ClosingIterator<>(reader.iterator());
+ })
+ .mapToPair(file -> new Tuple2<>(new DeleteFileKey(file),
file.copyWithoutStats()));
+
+ return allDeletes.subtractByKey(referencedKeys).values().collect();
+ }
- for (DeleteFile deleteFile : danglingDeletes) {
- LOG.debug("Removing dangling delete file {}", deleteFile.location());
- rewriteFiles.deleteFile(deleteFile);
+ private static class DeleteFileKeyIterator implements
CloseableIterator<DeleteFileKey> {
+ private final Closeable closeable;
+ private final Iterator<DeleteFileKey> iterator;
+
+ DeleteFileKeyIterator(CloseableIterable<FileScanTask> tasks) {
+ this.closeable = tasks;
+ this.iterator =
+ StreamSupport.stream(tasks.spliterator(), false)
+ .flatMap(task -> task.deletes().stream())
+ .map(DeleteFileKey::new)
+ .distinct()
Review Comment:
The full set of referenced delete key is not a problem, it adds only a few
dozens bytes per delete file. The real problem is `DeleteFileIndex`.
`scan.planFiles()` loads all delete manifests eagerly to build
`DeleteFileIndex` before planning any data files, and this index stays alive
until the `tasks` iterable is closed. Consequently, the `String location`,
`Long contentOffset`, and `Long contentSizeInBytes` for all delete files remain
in the heap, even without `.distinct()`.
I haven't found a way to avoid keeping the full `DeleteFileIndex` on the
driver or the executor. I checked
[SparkDistributedDataScan](https://github.com/apache/iceberg/blob/apache-iceberg-1.11.0/spark/v4.1/spark/src/main/java/org/apache/iceberg/SparkDistributedDataScan.java#L157-L175),
but there's no magic there - it collects all delete files from executors to
build the `DeleteFileIndex` on the driver.
@szehon-ho, do you have any ideas on how we could distribute building the
referenced delete file key set?
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]