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]