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

philo-he 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 83eaedc5d9 [GLUTEN-13091][CORE] Support resource profile adjustment 
for final stage (#13092)
83eaedc5d9 is described below

commit 83eaedc5d955edbc268b3726a554909972fb1247
Author: Wechar Yu <[email protected]>
AuthorDate: Wed Sep 30 15:44:27 2026 +0800

    [GLUTEN-13091][CORE] Support resource profile adjustment for final stage 
(#13092)
---
 .../AutoAdjustStageResourceProfileSuite.scala      | 72 ++++++++++++++++++----
 .../sql/execution/ApplyResourceProfileExec.scala   |  8 +++
 .../GlutenAutoAdjustStageResourceProfile.scala     | 45 ++++++++------
 3 files changed, 93 insertions(+), 32 deletions(-)

diff --git 
a/backends-velox/src/test/scala/org/apache/gluten/execution/AutoAdjustStageResourceProfileSuite.scala
 
b/backends-velox/src/test/scala/org/apache/gluten/execution/AutoAdjustStageResourceProfileSuite.scala
index 2ef614e415..583d11bbe1 100644
--- 
a/backends-velox/src/test/scala/org/apache/gluten/execution/AutoAdjustStageResourceProfileSuite.scala
+++ 
b/backends-velox/src/test/scala/org/apache/gluten/execution/AutoAdjustStageResourceProfileSuite.scala
@@ -20,7 +20,7 @@ import org.apache.gluten.config.GlutenConfig
 
 import org.apache.spark.SparkConf
 import org.apache.spark.annotation.Experimental
-import org.apache.spark.sql.execution.{ApplyResourceProfileExec, 
ColumnarShuffleExchangeExec, SparkPlan}
+import org.apache.spark.sql.execution.{ApplyResourceProfileExec, 
ColumnarShuffleExchangeExec, CommandResultExec, SparkPlan}
 import org.apache.spark.sql.execution.adaptive.AdaptiveSparkPlanHelper
 import org.apache.spark.sql.execution.exchange.ShuffleExchangeExec
 import org.apache.spark.sql.internal.SQLConf
@@ -87,7 +87,12 @@ class AutoAdjustStageResourceProfileSuite
   }
 
   private def collectApplyResourceProfileExec(plan: SparkPlan): Int = {
-    collect(plan) { case c: ApplyResourceProfileExec => c }.size
+    plan match {
+      case command: CommandResultExec =>
+        collectApplyResourceProfileExec(command.commandPhysicalPlan)
+      case _ =>
+        collect(plan) { case c: ApplyResourceProfileExec => c }.size
+    }
   }
 
   test("stage contains fallback nodes and apply new resource profile") {
@@ -179,20 +184,61 @@ class AutoAdjustStageResourceProfileSuite
         // scalastyle:off
         // format: off
         /*
-         DeserializeToObject createexternalrow(java_method(java.lang.Integer, 
signum, c1)#35.toString, count(1)#36L, 
StructField(java_method(java.lang.Integer, signum, c1),StringType,true), 
StructField(count(1),LongType,false)), obj#42: org.apache.spark.sql.Row
-         +- *(3) HashAggregate(keys=[_nondeterministic#37], 
functions=[count(1)], output=[java_method(java.lang.Integer, signum, c1)#35, 
count(1)#36L])
-            +- AQEShuffleRead coalesced
-               +- ShuffleQueryStage 0
-                  +- Exchange hashpartitioning(_nondeterministic#37, 5), 
ENSURE_REQUIREMENTS, [plan_id=607]
-                     +- ApplyResourceProfile Profile: id = 0, executor 
resources: cores -> name: cores, amount: 1, script: , vendor: ,memory -> name: 
memory, amount: 1024, script: , vendor: ,offHeap -> name: offHeap, amount: 
2048, script: , vendor: , task resources: cpus -> name: cpus, amount: 1.0
-                        +- *(2) HashAggregate(keys=[_nondeterministic#37], 
functions=[partial_count(1)], output=[_nondeterministic#37, count#41L])
-                           +- Project [java_method(java.lang.Integer, signum, 
c1#22) AS _nondeterministic#37]
-                              +- *(1) ColumnarToRow
-                                 +- FileScan parquet default.tmp1[c1#22] 
Batched: true, DataFilters: [], Format: Parquet
+        ResultQueryStage 1
+        +- ApplyResourceProfile Profile: id = 0,
+            +- *(3) HashAggregate(keys=[_nondeterministic#25], 
functions=[count(1)], output=[java_method(java.lang.Integer, signum, c1)#23, 
count(1)#24L])
+              +- AQEShuffleRead coalesced
+                  +- ShuffleQueryStage 0
+                    +- Exchange hashpartitioning(_nondeterministic#25, 5), 
ENSURE_REQUIREMENTS, [plan_id=1275]
+                        +- ApplyResourceProfile Profile: id = 0,
+                          +- *(2) HashAggregate(keys=[_nondeterministic#25], 
functions=[partial_count(1)], output=[_nondeterministic#25, count#27L])
+                              +- Project [java_method(java.lang.Integer, 
signum, c1#9, true) AS _nondeterministic#25]
+                                +- *(1) ColumnarToRow
+                                    +- FileScan parquet 
spark_catalog.default.tmp1[c1#9] Batched: true, DataFilters: [],
          */
         // format: on
         // scalastyle:on
-        df => 
assert(collectApplyResourceProfileExec(df.queryExecution.executedPlan) == 1)
+        df => 
assert(collectApplyResourceProfileExec(df.queryExecution.executedPlan) == 2)
+      }
+    }
+  }
+
+  test("Apply new resource profile to a direct data-writing stage") {
+    withSQLConf(
+      GlutenConfig.NATIVE_WRITER_ENABLED.key -> "false",
+      GlutenConfig.COLUMNAR_FALLBACK_PREFER_COLUMNAR.key -> "false",
+      GlutenConfig.COLUMNAR_FALLBACK_IGNORE_ROW_TO_COLUMNAR.key -> "false",
+      GlutenConfig.AUTO_ADJUST_STAGE_RESOURCES_FALLEN_NODE_RATIO_THRESHOLD.key 
-> "0.1"
+    ) {
+      withTable("t") {
+        spark.sql("CREATE TABLE t (c1 STRING) USING parquet")
+        runQueryAndCompare(s"""
+                              |INSERT OVERWRITE TABLE t
+                              |SELECT java_method('java.lang.Integer', 
'signum', c1)
+                              |FROM tmp1
+                              |""".stripMargin) {
+          df =>
+            val plan = df.queryExecution.executedPlan
+            // scalastyle:off
+            // format: off
+            /*
+            CommandResult <empty>
+              +- Execute InsertIntoHadoopFsRelationCommand
+                  +- ApplyResourceProfile Profile: id = 0,
+                    +- WriteFiles
+                        +- VeloxColumnarToRow
+                          +- ^(2) ProjectExecTransformer 
[java_method(java.lang.Integer, signum, c1)#17 AS c1#18]
+                              +- ^(2) 
InputIteratorTransformer[java_method(java.lang.Integer, signum, c1)#17]
+                                +- RowToVeloxColumnar
+                                    +- Project [java_method(java.lang.Integer, 
signum, c1#10, true) AS java_method(java.lang.Integer, signum, c1)#17]
+                                      +- VeloxColumnarToRow
+                                          +- ^(1) 
FileFileSourceScanExecTransformer parquet spark_catalog.default.tmp1[c1#10] 
Batched: true, DataFilters: [],
+             */
+            // scalastyle:on
+            // format: on
+            assert(collectShuffleExchange(plan) == 0)
+            assert(collectApplyResourceProfileExec(plan) == 1)
+        }
       }
     }
   }
diff --git 
a/gluten-substrait/src/main/scala/org/apache/spark/sql/execution/ApplyResourceProfileExec.scala
 
b/gluten-substrait/src/main/scala/org/apache/spark/sql/execution/ApplyResourceProfileExec.scala
index b175cb5a05..c1ba887712 100644
--- 
a/gluten-substrait/src/main/scala/org/apache/spark/sql/execution/ApplyResourceProfileExec.scala
+++ 
b/gluten-substrait/src/main/scala/org/apache/spark/sql/execution/ApplyResourceProfileExec.scala
@@ -25,6 +25,8 @@ import org.apache.spark.resource.ResourceProfile
 import org.apache.spark.sql.catalyst.InternalRow
 import org.apache.spark.sql.catalyst.expressions.{Attribute, SortOrder}
 import org.apache.spark.sql.catalyst.plans.physical.{Distribution, 
Partitioning}
+import org.apache.spark.sql.connector.write.WriterCommitMessage
+import org.apache.spark.sql.execution.datasources.WriteFilesSpec
 import org.apache.spark.sql.vectorized.ColumnarBatch
 
 /**
@@ -72,6 +74,12 @@ case class ApplyResourceProfileExec(child: SparkPlan, 
resourceProfile: ResourceP
     child.executeColumnar.withResources(resourceProfile)
   }
 
+  override protected def doExecuteWrite(writeFilesSpec: WriteFilesSpec)
+      : RDD[WriterCommitMessage] = {
+    log.info(s"Apply $resourceProfile for write plan ${child.nodeName}")
+    child.executeWrite(writeFilesSpec).withResources(resourceProfile)
+  }
+
   override def output: scala.Seq[Attribute] = child.output
 
   override protected def withNewChildInternal(newChild: SparkPlan): 
ApplyResourceProfileExec =
diff --git 
a/gluten-substrait/src/main/scala/org/apache/spark/sql/execution/GlutenAutoAdjustStageResourceProfile.scala
 
b/gluten-substrait/src/main/scala/org/apache/spark/sql/execution/GlutenAutoAdjustStageResourceProfile.scala
index 14669c51ac..48154b931b 100644
--- 
a/gluten-substrait/src/main/scala/org/apache/spark/sql/execution/GlutenAutoAdjustStageResourceProfile.scala
+++ 
b/gluten-substrait/src/main/scala/org/apache/spark/sql/execution/GlutenAutoAdjustStageResourceProfile.scala
@@ -32,6 +32,7 @@ import 
org.apache.spark.sql.execution.{GlutenAutoAdjustStageResourceProfile => G
 import org.apache.spark.sql.execution.adaptive.QueryStageExec
 import org.apache.spark.sql.execution.columnar.InMemoryTableScanExec
 import org.apache.spark.sql.execution.command.{DataWritingCommandExec, 
ExecutedCommandExec}
+import org.apache.spark.sql.execution.datasources.v2.{V2CommandExec, 
V2TableWriteExec}
 import org.apache.spark.sql.execution.exchange.Exchange
 import org.apache.spark.sql.internal.SQLConf
 import org.apache.spark.util.{SparkResourceUtil, SparkTestUtil}
@@ -52,8 +53,6 @@ import scala.collection.mutable.ArrayBuffer
  *   3. Partial fallback: if the ratio of fallen (non-Gluten) nodes in a stage 
exceeds
  *      `spark.gluten.auto.adjustStageResources.fallenNode.ratio.threshold`, 
heap memory is
  *      increased and off-heap memory is decreased proportionally.
- *
- * * Note: Case 2 and 3 are not applied to final (non-Exchange) stages yet.
  */
 @Experimental
 case class GlutenAutoAdjustStageResourceProfile(glutenConf: GlutenConfig, 
spark: SparkSession)
@@ -105,10 +104,6 @@ case class 
GlutenAutoAdjustStageResourceProfile(glutenConf: GlutenConfig, spark:
       }
     }
 
-    if (!plan.isInstanceOf[Exchange]) {
-      // todo: support set resource profile for final stage
-      return plan
-    }
     val planNodes = GlutenResourceProfile.collectStagePlan(plan)
     if (planNodes.isEmpty) {
       return plan
@@ -120,11 +115,16 @@ case class 
GlutenAutoAdjustStageResourceProfile(glutenConf: GlutenConfig, spark:
     logInfo(s"default memory request $memoryRequest")
     logInfo(s"default offheap request $offheapRequest")
 
+    val countedPlanNodes = 
planNodes.filterNot(GlutenExplainUtils.shouldIgnoreInFallbackStats)
+    if (countedPlanNodes.isEmpty) {
+      return plan
+    }
+
     // case 1: whole stage fallback to vanilla spark in such case we increase 
the heap
     //
     // one stage is considered as fallback if all node is not GlutenPlan
     // or all GlutenPlan node is C2R node.
-    val wholeStageFallback = planNodes
+    val wholeStageFallback = countedPlanNodes
       .filter(_.isInstanceOf[GlutenPlan])
       .count(!_.isInstanceOf[ColumnarToRowExecBase]) == 0
     if (wholeStageFallback) {
@@ -147,7 +147,6 @@ case class GlutenAutoAdjustStageResourceProfile(glutenConf: 
GlutenConfig, spark:
 
     // case 2: check whether fallback exists and decide whether increase heap 
memory
     // and decrease offheap memory.
-    val countedPlanNodes = 
planNodes.filterNot(GlutenExplainUtils.shouldIgnoreInFallbackStats)
     val fallenNodeCnt = countedPlanNodes.count {
       case _: GlutenPlan => false
       case i: InMemoryTableScanExec => !PlanUtil.isGlutenTableCache(i)
@@ -184,9 +183,14 @@ object GlutenAutoAdjustStageResourceProfile extends 
Logging {
   def collectStagePlan(plan: SparkPlan): ArrayBuffer[SparkPlan] = {
 
     def collectStagePlan(plan: SparkPlan, planNodes: ArrayBuffer[SparkPlan]): 
Unit = {
-      if (plan.isInstanceOf[DataWritingCommandExec] || 
plan.isInstanceOf[ExecutedCommandExec]) {
-        // todo: support set final stage's resource profile
-        return
+      plan match {
+        // V1/V2 writes have a physical computation child and must remain 
eligible for profiling.
+        case _: DataWritingCommandExec | _: V2TableWriteExec =>
+        case _: CommandResultExec | _: ExecutedCommandExec | _: V2CommandExec 
=>
+          // Limitation: RunnableCommand exposes no physical child, so this 
collector cannot attach
+          // a profile to worker RDDs created internally (e.g. by 
InsertIntoDataSourceDirCommand).
+          return
+        case _ =>
       }
       planNodes += plan
       if (plan.isInstanceOf[QueryStageExec]) {
@@ -269,16 +273,19 @@ object GlutenAutoAdjustStageResourceProfile extends 
Logging {
       taskResource: mutable.Map[String, TaskResourceRequest],
       rpManager: ResourceProfileManager,
       sparkConf: SparkConf): SparkPlan = {
-    val rp = new ResourceProfile(executorResource.toMap, taskResource.toMap)
-    val finalRP = getFinalResourceProfile(rpManager, rp)
-    updateResourceSetting(finalRP, sparkConf)
+    lazy val finalRP = {
+      val rp = new ResourceProfile(executorResource.toMap, taskResource.toMap)
+      val profile = getFinalResourceProfile(rpManager, rp)
+      updateResourceSetting(profile, sparkConf)
+      profile
+    }
 
     plan match {
-      case shuffle: Exchange =>
-        logInfo(s"Apply resource profile $finalRP for plan 
${shuffle.child.nodeName}")
-        // Wrap the plan with ApplyResourceProfileExec so that we can apply 
new ResourceProfile
-        val wrapperPlan = ApplyResourceProfileExec(shuffle.child, finalRP)
-        shuffle.withNewChildren(Seq(wrapperPlan))
+      case _: Exchange | _: DataWritingCommandExec | _: V2TableWriteExec =>
+        val child = plan.children.head
+        logInfo(s"Apply resource profile $finalRP for child ${child.nodeName}")
+        // Wrap the child with ApplyResourceProfileExec so that we can apply 
new ResourceProfile
+        plan.withNewChildren(Seq(ApplyResourceProfileExec(child, finalRP)))
       case other =>
         logInfo(s"Apply resource profile $finalRP for plan ${other.nodeName}")
         ApplyResourceProfileExec(other, finalRP)


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

Reply via email to