This is an automated email from the ASF dual-hosted git repository.

jackylee-ch pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/gluten.git


The following commit(s) were added to refs/heads/main by this push:
     new 93ed31c8ee [GLUTEN-13106][CORE] Release all task resources even when 
one release fails (#13107)
93ed31c8ee is described below

commit 93ed31c8ee9c0381e462766dfa3160818c2bac4e
Author: YangJie <[email protected]>
AuthorDate: Wed Sep 30 02:37:33 2026 -0400

    [GLUTEN-13106][CORE] Release all task resources even when one release fails 
(#13107)
---
 .../org/apache/spark/task/TaskResources.scala      | 60 ++++++++++++++++++++--
 .../org/apache/gluten/task/TaskResourceSuite.scala | 36 +++++++++++++
 2 files changed, 92 insertions(+), 4 deletions(-)

diff --git 
a/gluten-core/src/main/scala/org/apache/spark/task/TaskResources.scala 
b/gluten-core/src/main/scala/org/apache/spark/task/TaskResources.scala
index f20b0b6525..e0daeb985c 100644
--- a/gluten-core/src/main/scala/org/apache/spark/task/TaskResources.scala
+++ b/gluten-core/src/main/scala/org/apache/spark/task/TaskResources.scala
@@ -30,6 +30,7 @@ import java.util.concurrent.atomic.AtomicLong
 
 import scala.collection.mutable
 import scala.compat.Platform.ConcurrentModificationException
+import scala.util.control.NonFatal
 
 object TaskResources extends TaskListener with Logging {
   // And open java assert mode to get memory stack
@@ -224,9 +225,22 @@ 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 {
+              // Removing the registry must happen even if the metrics update 
throws,
+              // otherwise the registry stays reachable and leaks across 
tasks. The metrics
+              // update itself is best-effort: catch it here so a metrics 
failure cannot
+              // replace (mask) a releaseAll failure propagating from the 
outer try.
+              try {
+                
context.taskMetrics().incPeakExecutionMemory(registry.getSharedUsage().peak())
+              } catch {
+                case NonFatal(e) =>
+                  logWarning("Failed to record peak execution memory", e)
+              } finally {
+                RESOURCE_REGISTRIES.remove(context)
+              }
+            }
           }
         }
       })
@@ -290,12 +304,50 @@ 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]
+    def safeResourceName(resource: TaskResource): String =
+      try resource.resourceName()
+      catch {
+        // Best-effort log label only: catch everything, including fatal 
errors, so a
+        // throwing resourceName() can never abort the release loop before the 
remaining
+        // resources are freed and the maps are cleared. Fatal errors from 
release()
+        // itself are still left to propagate (see the NonFatal handler below).
+        case _: Throwable => s"resource@${System.identityHashCode(resource)}"
+      }
     priorityToResourcesMapping.toSeq.sortBy(-_._1).foreach {
       case (_, resources) =>
-        resources.toSeq.reverse.foreach(release)
+        resources.toSeq.reverse.foreach {
+          resource =>
+            try release(resource)
+            catch {
+              case e: InterruptedException =>
+                // The catch cleared the interrupt status; restore it so task
+                // cancellation still propagates, then record and keep 
releasing so
+                // the remaining resources are freed and the registry is 
cleared.
+                Thread.currentThread().interrupt()
+                failures += e
+                logError(s"Interrupted while releasing resource 
${safeResourceName(resource)}", e)
+              case NonFatal(e) =>
+                // 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. Fatal throwables are 
left to
+                // propagate immediately.
+                failures += e
+                logError(s"Failed to release resource 
${safeResourceName(resource)}", e)
+            }
+        }
     }
     priorityToResourcesMapping.clear()
     resources.clear()
+    failures.headOption.foreach {
+      failure =>
+        // Keep the remaining failures attached; the logs are the only other
+        // record and may be swallowed by the completion-listener machinery.
+        // Skip entries identical to `failure` by reference: addSuppressed 
throws
+        // IllegalArgumentException on self-suppression.
+        failures.tail.filterNot(_ eq failure).foreach(failure.addSuppressed)
+        throw failure
+    }
   }
 
   /** Release single resource by ID */
diff --git 
a/gluten-core/src/test/scala/org/apache/gluten/task/TaskResourceSuite.scala 
b/gluten-core/src/test/scala/org/apache/gluten/task/TaskResourceSuite.scala
index 8e4fd48284..03ba024b4f 100644
--- a/gluten-core/src/test/scala/org/apache/gluten/task/TaskResourceSuite.scala
+++ b/gluten-core/src/test/scala/org/apache/gluten/task/TaskResourceSuite.scala
@@ -83,4 +83,40 @@ class TaskResourceSuite extends AnyFunSuite with SQLHelper {
     }
     assert(unregisteredCount == 2)
   }
+
+  test("Run unsafe - one failing release does not skip the remaining 
resources") {
+    var goodReleased = 0
+    val result = scala.util.Try(TaskResources.runUnsafe {
+      // Higher priority is released first, so the failing resource releases
+      // before the good one.
+      TaskResources.addResource(
+        UUID.randomUUID().toString,
+        new TaskResource {
+          override def priority(): Int = 200
+          override def release(): Unit = throw new RuntimeException("release 
failed")
+          override def resourceName(): String = "failing resource"
+        }
+      )
+      TaskResources.addResource(
+        UUID.randomUUID().toString,
+        new TaskResource {
+          override def priority(): Int = 100
+          override def release(): Unit = goodReleased += 1
+          override def resourceName(): String = "good resource"
+        }
+      )
+    })
+    // The good (lower-priority) resource is still released despite the 
earlier failure.
+    assert(goodReleased == 1)
+    // The release failure is rethrown, not swallowed.
+    assert(result.isFailure)
+    val trace = {
+      val writer = new java.io.StringWriter()
+      result.failed.get.printStackTrace(new java.io.PrintWriter(writer))
+      writer.toString
+    }
+    assert(trace.contains("release failed"))
+    // The registry was cleared on the way out, so a fresh task runs cleanly.
+    TaskResources.runUnsafe(assert(TaskResources.inSparkTask()))
+  }
 }


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to