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]

Reply via email to