diff --git a/spark/src/main/scala/org/apache/comet/ExtendedExplainInfo.scala b/spark/src/main/scala/org/apache/comet/ExtendedExplainInfo.scala index 755c345717f..b5b8a53029b 100644 --- a/spark/src/main/scala/org/apache/comet/ExtendedExplainInfo.scala +++ b/spark/src/main/scala/org/apache/comet/ExtendedExplainInfo.scala @@ -25,6 +25,7 @@ import org.apache.spark.sql.ExtendedExplainGenerator import org.apache.spark.sql.catalyst.trees.{TreeNode, TreeNodeTag} import org.apache.spark.sql.execution.{InputAdapter, SparkPlan, WholeStageCodegenExec} import org.apache.spark.sql.execution.adaptive.{AdaptiveSparkPlanExec, QueryStageExec} +import org.apache.spark.sql.execution.exchange.ReusedExchangeExec import org.apache.comet.CometExplainInfo.getActualPlan @@ -158,6 +159,7 @@ object CometExplainInfo { case p: InputAdapter => getActualPlan(p.child) case p: QueryStageExec => getActualPlan(p.plan) case p: WholeStageCodegenExec => getActualPlan(p.child) + case p: ReusedExchangeExec => getActualPlan(p.child) case p => p } } diff --git a/spark/src/main/scala/org/apache/comet/rules/CometExecRule.scala b/spark/src/main/scala/org/apache/comet/rules/CometExecRule.scala index 091f70fdc29..3527b48c308 100644 --- a/spark/src/main/scala/org/apache/comet/rules/CometExecRule.scala +++ b/spark/src/main/scala/org/apache/comet/rules/CometExecRule.scala @@ -535,8 +535,7 @@ case class CometExecRule(session: SparkSession) extends Rule[SparkPlan] { case op => op match { - case _: CometExec | _: AQEShuffleReadExec | _: BroadcastExchangeExec | - _: CometBroadcastExchangeExec | _: CometShuffleExchangeExec | + case _: CometPlan | _: AQEShuffleReadExec | _: BroadcastExchangeExec | _: BroadcastQueryStageExec | _: AdaptiveSparkPlanExec => // Some execs should never be replaced. We include // these cases specially here so we do not add a misleading 'info' message