[
https://issues.apache.org/jira/browse/SPARK-59122?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
ASF GitHub Bot updated SPARK-59122:
-----------------------------------
Labels: pull-request-available (was: )
> UnionExec fused into whole-stage codegen can fail with "key not found:
> numOutputRows"
> -------------------------------------------------------------------------------------
>
> Key: SPARK-59122
> URL: https://issues.apache.org/jira/browse/SPARK-59122
> Project: Spark
> Issue Type: Bug
> Components: SQL
> Affects Versions: 5.0.0
> Reporter: Yang Jie
> Priority: Major
> Labels: pull-request-available
>
> A {{UnionExec}} that whole-stage codegen fused can fail at execution with:
> {code}
> java.util.NoSuchElementException: key not found: numOutputRows
> at org.apache.spark.sql.execution.SparkPlan.longMetric(SparkPlan.scala:154)
> at
> org.apache.spark.sql.execution.CodegenSupport.metricTerm(WholeStageCodegenExec.scala:71)
> at
> org.apache.spark.sql.execution.UnionExec.doProduce(basicPhysicalOperators.scala:1148)
> at
> org.apache.spark.sql.execution.WholeStageCodegenExec.doCodeGen(WholeStageCodegenExec.scala:676)
> {code}
> Reproduced on master with default configuration:
> {code}
> spark.range(0, 200, 1, 4).selectExpr("id % 10 AS k", "id AS v")
> .groupBy("k").agg(sum("v").as("s")).createOrReplaceTempView("v")
> spark.catalog.cacheTable("v")
> spark.sql("SELECT k, abs(s) AS s FROM v UNION ALL SELECT k, s FROM
> v").collect()
> {code}
> {{UnionExec}} derives three things from its children's
> {{outputPartitioning}}: whether codegen
> fusion applies, whether {{numOutputRows}} is registered, and which RDD
> {{doExecute}} builds. That
> input is not stable while the plan is prepared -- {{InMemoryTableScanExec}}
> reads
> {{cachedPlan.outputPartitioning}}, and the inner {{AdaptiveSparkPlanExec}}
> answers
> {{UnknownPartitioning}} until its final plan exists.
> In the query above an expression on one branch only keeps a {{ProjectExec}}
> from being collapsed
> away, so the two children are shaped differently, {{comparePartitioning}}
> rejects the pair, and the
> union looks plain and is fused. {{CollapseCodegenStages.insertInputAdapter}}
> then rebuilds the node
> through {{withNewChildren}} and puts that copy inside the
> {{WholeStageCodegenExec}} it just built,
> so the copy evaluates the gate for the first time after the cache stages have
> finalised. By then
> both children report the same {{HashPartitioning}}, the gate answers
> "partitioning-aware",
> {{metrics}} comes back empty, and the generated code still increments the
> metric.
> Registering the metric unconditionally would turn the crash into a wrong
> answer: a fused union
> concatenates its children's partitions, so if it went on claiming their
> {{HashPartitioning}} a
> parent could satisfy a clustered distribution from an RDD that does not have
> it.
> Affects the union codegen fusion added in SPARK-56482, which is present in
> 4.2.0, branch-4.3 and
> master.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]