LuciferYang commented on code in PR #13107:
URL: https://github.com/apache/gluten/pull/13107#discussion_r4110042524


##########
gluten-core/src/main/scala/org/apache/spark/task/TaskResources.scala:
##########
@@ -224,9 +224,12 @@ object TaskResources extends TaskListener with Logging {
             }
             // We should first call `releaseAll` then remove the registries, 
because
             // the functions inside registries may register new resource to 
registries.
-            currentTaskRegistries.releaseAll()
-            
context.taskMetrics().incPeakExecutionMemory(registry.getSharedUsage().peak())
-            RESOURCE_REGISTRIES.remove(context)
+            try {
+              currentTaskRegistries.releaseAll()
+            } finally {
+              
context.taskMetrics().incPeakExecutionMemory(registry.getSharedUsage().peak())
+              RESOURCE_REGISTRIES.remove(context)

Review Comment:
   Done in c2a8dc37f. Wrapped the metrics update in an inner `try` and moved 
`RESOURCE_REGISTRIES.remove(context)` into an inner `finally`, so the registry 
is removed even if `incPeakExecutionMemory` / `getSharedUsage().peak()` throws.



##########
gluten-core/src/main/scala/org/apache/spark/task/TaskResources.scala:
##########
@@ -290,12 +293,41 @@ class TaskResourceRegistry extends Logging {
 
   /** Release all managed resources according to priority and reversed order */
   private[task] def releaseAll(): Unit = lock {
+    val failures = mutable.ArrayBuffer.empty[Throwable]
     priorityToResourcesMapping.toSeq.sortBy(-_._1).foreach {
       case (_, resources) =>
-        resources.toSeq.reverse.foreach(release)
+        resources.toSeq.reverse.foreach {
+          resource =>
+            try release(resource)
+            catch {
+              case e: Throwable =>
+                // One failing release must not skip the remaining ones or 
leave the
+                // registry uncleared; record the failure and rethrow it after 
the
+                // loop so callers still see the error.
+                failures += e
+                // resourceName() is user code too, so guard it with the same
+                // throwable range as release(): a failure building the log 
label
+                // must not abort the loop either.
+                val name =
+                  try resource.resourceName()
+                  catch {
+                    case _: Throwable => 
s"resource@${System.identityHashCode(resource)}"
+                  }
+                logError(s"Failed to release resource $name", e)
+            }

Review Comment:
   Done in c2a8dc37f. The record-and-continue branch now catches `NonFatal`, so 
`VirtualMachineError` and other fatal errors propagate immediately instead of 
being recorded. `InterruptedException` is handled in its own branch: it 
restores the interrupt flag (the catch clears it) so task cancellation still 
propagates, records the failure so the loop finishes freeing the remaining 
resources and clearing the registry, and is rethrown after the loop like any 
other release failure. I also narrowed the `resourceName()` label guard to 
`NonFatal` to match.



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