szehon-ho commented on code in PR #15727:
URL: https://github.com/apache/iceberg/pull/15727#discussion_r3974674670


##########
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())

Review Comment:
   Could we handle an empty `deleteManifests` list before creating this RDD, 
ideally before the broadcast and scan? Spark rejects `parallelize(..., 0)` with 
`IllegalArgumentException: Positive number of partitions required`, so this 
fails for a populated table with no delete manifests.



##########
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:
   Could we avoid performing `distinct()` inside this single scan task? It 
retains every referenced delete key in one executor’s heap, while removing it 
may significantly increase shuffle traffic when deletes apply to many file 
tasks. Is there a way to preserve deduplication without concentrating the full 
set in one executor?



-- 
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