kinolaev commented on code in PR #15727:
URL: https://github.com/apache/iceberg/pull/15727#discussion_r3979421721


##########
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:
   Thanks for catching this! I moved the check to core, before calling 
`findDanglingDeletes` (1e536022da6464f7751e12e01262b9c8).



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