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]