[ 
https://issues.apache.org/jira/browse/SPARK-59009?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
 ]

ASF GitHub Bot updated SPARK-59009:
-----------------------------------
    Labels: pull-request-available regresion  (was: regresion)

> InMemoryRelation.newInstance() does not remap outputOrdering, causing 
> NoSuchElementException during canonicalization
> --------------------------------------------------------------------------------------------------------------------
>
>                 Key: SPARK-59009
>                 URL: https://issues.apache.org/jira/browse/SPARK-59009
>             Project: Spark
>          Issue Type: Bug
>          Components: SQL
>    Affects Versions: 4.1.0, 4.0.2, 3.5.8, 4.2.0, 4.1.1, 4.1.2, 3.5.9, 4.0.3, 
> 4.3.0, 4.1.3, 4.2.1, 4.1.4
>            Reporter: James Willis
>            Priority: Major
>              Labels: pull-request-available, regresion
>
> h2. Symptom
> A query over a cached ({{{}persist(){}}}) DataFrame that has a non-empty 
> {{outputOrdering}} fails at the first action with an internal error, if the 
> cached relation is referenced more than once (for example a CTE that is 
> referenced twice and then self-joined):
> {noformat}
> java.util.NoSuchElementException: key not found: tile#1L
>   at scala.collection.MapOps.default(Map.scala:289)
>   at 
> org.apache.spark.sql.catalyst.expressions.AttributeMap.apply(AttributeMap.scala:41)
>   at 
> org.apache.spark.sql.execution.columnar.InMemoryRelation.$anonfun$withOutput$1(InMemoryRelation.scala:438)
>   at 
> org.apache.spark.sql.execution.columnar.InMemoryRelation.withOutput(InMemoryRelation.scala:438)
>   at 
> org.apache.spark.sql.execution.columnar.InMemoryRelation.doCanonicalize(InMemoryRelation.scala:406)
>   ...
>   at 
> org.apache.spark.sql.catalyst.plans.QueryPlan.canonicalized(QueryPlan.scala:633)
>   at 
> org.apache.spark.sql.execution.adaptive.AdaptiveSparkPlanExec.createNonResultQueryStages(AdaptiveSparkPlanExec.scala:589)
> {noformat}
> The error surfaces at {{.show()}} / {{.collect()}} and gives the user nothing 
> actionable, since the query itself is well formed.
> h2. Reproduction
> Deterministic, no external data:
> {code:python}
> spark.range(0, 20).selectExpr("id", "id % 3 AS 
> tile").createOrReplaceTempView("a")
> # a persisted DataFrame with a non-empty outputOrdering
> b = spark.sql("SELECT id, tile FROM a ORDER BY tile, id")
> b.persist()
> b.count()
> b.createOrReplaceTempView("b")
> # window-rank it, then self-join the ranked result
> spark.sql("""
>   WITH r AS (SELECT *, row_number() OVER (PARTITION BY tile ORDER BY id DESC) 
> AS rn FROM b)
>   SELECT x.id AS p, y.id AS q
>   FROM r x JOIN r y ON x.tile = y.tile AND x.rn = 1 AND y.rn = 2
> """).show()
> {code}
> Verified failing on 4.0.4, 4.1.3 and 4.2.0. Passes on 4.0.1.
> The invariant break can also be shown directly, without running a query:
> {code:scala}
> val df = spark.range(0, 20).selectExpr("id", "id % 3 AS 
> tile").orderBy("tile", "id")
> df.persist(); df.count()
> val ir = 
> spark.sharedState.cacheManager.lookupCachedData(df).get.cachedRepresentation
> ir.output          // id#0L, tile#1L
> ir.outputOrdering  // tile#1L ASC, id#0L ASC        <- consistent
> val ni = ir.newInstance()
> ni.output          // id#68L, tile#69L              <- refreshed
> ni.outputOrdering  // tile#1L ASC, id#0L ASC        <- NOT refreshed
> ni.canonicalized   // throws NoSuchElementException
> {code}
> h2. Root cause
> {{InMemoryRelation}} carries an implicit invariant: {{outputOrdering}} may 
> only reference attributes present in {{{}output{}}}. {{newInstance()}} breaks 
> it, because it gives {{output}} fresh exprIds but passes {{outputOrdering}} 
> through unchanged:
> {code:scala}
> override def newInstance(): this.type = {
>   InMemoryRelation(
>     output.map(_.newInstance()),   // fresh exprIds
>     cacheBuilder,
>     outputOrdering,                // NOT remapped -- still the old exprIds
>     statsOfPlanToCache).asInstanceOf[this.type]
> }
> {code}
> This was harmless until SPARK-53738, which rewired {{doCanonicalize}} through 
> {{withOutput}} and made {{withOutput}} remap the ordering with a strict 
> {{AttributeMap}} lookup:
> {code:scala}
> override def doCanonicalize(): logical.LogicalPlan =
>   withOutput(output.map(QueryPlan.normalizeExpressions(_, output)))
> def withOutput(newOutput: Seq[Attribute]): InMemoryRelation = {
>   val map = AttributeMap(output.zip(newOutput))
>   val newOutputOrdering = outputOrdering
>     .map(_.transform { case a: Attribute => map(a) })   // throws on a 
> missing key
>     .asInstanceOf[Seq[SortOrder]]
>   InMemoryRelation(newOutput, cacheBuilder, newOutputOrdering, 
> statsOfPlanToCache)
> }
> {code}
> So any {{InMemoryRelation}} that has been through {{newInstance()}} now fails 
> the moment anything canonicalizes it.
> How {{newInstance()}} comes to be called on a cached relation: cache 
> substitution ({{{}CacheManager.useCachedData{}}}) runs on the analyzed plan, 
> before the optimizer. {{InlineCTE}} then inlines a CTE that is referenced 
> more than once, and to obtain a fresh-exprId copy it runs 
> {{DeduplicateRelations}} over a synthetic self-join (InlineCTE.scala, the 
> {{case ref: CTERelationRef}} branch). By that point the plan already contains 
> {{{}InMemoryRelation{}}}, which is a {{{}MultiInstanceRelation{}}}, so 
> {{DeduplicateRelations}} calls {{newInstance()}} on it.
> This is also why a plain self-join of a cached temp view does *not* 
> reproduce: that deduplication happens in the analyzer, before caching is 
> substituted, so the {{InMemoryRelation}} is never passed through 
> {{{}newInstance(){}}}.
> h2. Affected versions
> The {{withOutput}} change is present in 3.5.8+, 4.0.2+, and every branch from 
> branch-4.0 through branch-4.3, branch-4.x and master (5.0.0-SNAPSHOT). 3.5.7, 
> 4.0.0 and 4.0.1 are unaffected. {{newInstance()}} is still unpatched 
> everywhere, so no released version since 4.0.1 contains a fix.
> h2. Suggested fix
> Route {{newInstance()}} through {{withOutput}} so the ordering is remapped 
> onto the fresh exprIds. This preserves the ordering rather than dropping it:
> {code:scala}
> override def newInstance(): this.type =
>   withOutput(output.map(_.newInstance())).asInstanceOf[this.type]
> {code}
> And make {{withOutput}} defensive, since {{outputOrdering}} is only an 
> optimization hint and no cache hint should be able to hard-fail a query:
> {code:scala}
> val newOutputOrdering = if 
> (outputOrdering.forall(_.references.forall(map.contains))) {
>   outputOrdering.map(_.transform { case a: Attribute => map(a) 
> }).asInstanceOf[Seq[SortOrder]]
> } else {
>   Nil
> }
> {code}
> Note that {{withOutput}} also silently truncates when {{newOutput}} is 
> shorter than {{{}output{}}}, because of the {{{}output.zip(newOutput){}}}. 
> The guard above covers that case too.
> h2. Workarounds
> {{localCheckpoint()}} or {{checkpoint()}} instead of {{{}persist(){}}}, or 
> breaking the ordering before caching (for example {{repartition()}} after the 
> sort). There is no configuration-level workaround: disabling AQE 
> ({{{}spark.sql.adaptive.enabled=false{}}}), disabling broadcast joins 
> ({{{}spark.sql.autoBroadcastJoinThreshold=-1{}}}), and excluding 
> {{InlineCTE}} via {{spark.sql.optimizer.excludedRules}} all still fail.



--
This message was sent by Atlassian Jira
(v8.20.10#820010)

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

Reply via email to