JingsongLi commented on code in PR #9706:
URL: https://github.com/apache/paimon/pull/9706#discussion_r3986192841


##########
paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/procedure/SparkRemoveUnexistingFiles.scala:
##########
@@ -96,15 +95,20 @@ case class SparkRemoveUnexistingFiles(
               case (_, bytes) => 
messages.add(serializer.deserialize(serializer.getVersion, bytes))
             }
             val commit = table.newCommit(UUID.randomUUID().toString)
-            commit.commit(Long.MaxValue, messages)
+            try {
+              commit.commit(Long.MaxValue, messages)
+            } finally {
+              commit.close()

Review Comment:
   [P2] Use the batch commit lifecycle before closing this committer
   
   The existing call to commit(Long.MaxValue, messages) takes the streaming 
overload and does not set TableCommitImpl.batchCommitted. With 
snapshot.expire.execution-mode=async, maintenance runs on maintainExecutor; the 
newly added close immediately calls shutdownNow(), so expiration/tag 
maintenance can be cancelled or interrupted after the procedure reports a 
successful commit. This is introduced by adding close to that overload. A 
focused probe with the head TableCommitImpl reproduced the interruption.
   
   Use commit(messages), which selects the one-shot batch lifecycle and 
performs maintenance inline, or explicitly await completion and propagate its 
failure before closing. The one-argument overload passed the control probe.



##########
paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/procedure/SparkRemoveUnexistingFiles.scala:
##########
@@ -96,15 +95,20 @@ case class SparkRemoveUnexistingFiles(
               case (_, bytes) => 
messages.add(serializer.deserialize(serializer.getVersion, bytes))
             }
             val commit = table.newCommit(UUID.randomUUID().toString)
-            commit.commit(Long.MaxValue, messages)
+            try {
+              commit.commit(Long.MaxValue, messages)
+            } finally {
+              commit.close()
+            }
           }
       }
     }
 
-    pathAndMessage.mapPartitions(
-      iter => {
-        iter.flatMap { case (paths, _) => paths }
-      })
+    try {

Review Comment:
   [P2] Include the commit action inside the cache cleanup scope
   
   The cached RDD is materialized by foreachPartition at lines 88–104 before 
this try starts. If committing a partition fails after task retries, execution 
never reaches this finally and the persisted blocks remain in the long-lived 
SparkContext. That leaves the resource leak this PR is meant to fix on the 
failure path. Move the try/finally to cover both foreachPartition and collect, 
and verify that persistent RDDs return to their initial count after an injected 
commit failure.



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

Reply via email to