philo-he commented on code in PR #13092:
URL: https://github.com/apache/gluten/pull/13092#discussion_r4125855769
##########
gluten-substrait/src/main/scala/org/apache/spark/sql/execution/GlutenAutoAdjustStageResourceProfile.scala:
##########
@@ -269,16 +273,25 @@ 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 command: DataWritingCommandExec =>
+ logInfo(s"Apply resource profile $finalRP for write input
${command.child.nodeName}")
+ command.withNewChildren(Seq(ApplyResourceProfileExec(command.child,
finalRP)))
+ case write: V2TableWriteExec =>
+ logInfo(s"Apply resource profile $finalRP for V2 write input
${write.child.nodeName}")
+ write.withNewChildren(Seq(ApplyResourceProfileExec(write.child,
finalRP)))
Review Comment:
Can we consolidate the above 3 cases?
##########
gluten-substrait/src/main/scala/org/apache/spark/sql/execution/GlutenAutoAdjustStageResourceProfile.scala:
##########
@@ -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
=>
+ // Most commands are DDL and expose no physical input plan. Only few
commands
+ // (e.g. InsertIntoDataSourceDirCommand) have a physical child plan.
Review Comment:
For those cases that wrap physical query trees executed on workers, the
above comment doesn't explain why we can still ignore them.
If handling them are out of scope for this PR, let's update the comment to
note this as a known limitation/follow-up.
--
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]