Copilot commented on code in PR #13107:
URL: https://github.com/apache/gluten/pull/13107#discussion_r4082226511
##########
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:
`RESOURCE_REGISTRIES.remove(context)` is intended to run even when
cleanup/metrics code fails, but it currently won’t execute if
`incPeakExecutionMemory(...)` (or `registry.getSharedUsage().peak()`) throws
inside the `finally`. To guarantee registry removal (and avoid leaving the
registry reachable/leaking across tasks), wrap metrics update in an inner `try`
and put `RESOURCE_REGISTRIES.remove(context)` in an inner `finally` (or perform
removal first, then best-effort metrics).
##########
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:
Catching `Throwable` will also catch fatal errors (e.g.,
`VirtualMachineError`) and `InterruptedException`, and then continue releasing
other resources while holding the registry lock. That can be unsafe for truly
fatal conditions, and it can also clear the thread’s interrupt status. Consider
catching `NonFatal` for the ‘record and continue’ path, and handling
`InterruptedException` explicitly by re-setting the interrupt flag (and still
recording/rethrowing). For truly fatal errors, prefer rethrowing immediately
rather than continuing.
--
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]