From 0deedf1400193da7ea7d1f1bb024af5185958871 Mon Sep 17 00:00:00 2001 From: James Xu Date: Sun, 23 Aug 2026 19:55:30 +0800 Subject: [PATCH] [SPARK-58950][SQL] Materialize repeated view references as one CTE Introduce a ConvertViewToMaterializedCTE optimizer rule that rewrites multiple references to the same view into a single forceSkipInline CTERelationDef with a reference per occurrence, so the view's underlying plan is computed once (through exchange reuse) instead of once per reference. The rule runs in FinishAnalysis immediately before EliminateView and is gated behind spark.sql.optimizer.convertViewToMaterializedCTE (off by default). Only deterministic, non-streaming views with matching view SQL configs and aligned output schemas convert. Occurrences are grouped by catalog identity plus canonicalized body, but the bottom-up rewrite matches by identifier alone: by the time an outer view is visited, nested views inside its body have already been rewritten into CTERelationRefs, so its body no longer canonicalizes to the key computed up front. A view nested inside another view's body therefore converts too, and definitions are appended in topological order (referenced definitions first), which ReplaceCTERefWithRepartition relies on when it iterates the definitions in order and resolves references from a map it builds incrementally. A view body may contain correlated subqueries whose outer references resolve to relations inside the body (e.g. t WHERE x IN (SELECT y FROM s WHERE s.k = t.k)). The view is analyzed standalone when it is created, so an outer reference that does not resolve inside the body fails view analysis and can never escape to the outer query. The converted definition contains the whole body, so internal correlations resolve within it and these bodies are safe to convert; InlineCTE's rejection of boundary-crossing outer references is only a generic safety net for non-view forceSkipInline producers. Add catalyst rule tests proving that an internally correlated body and a view-over-view body convert (the latter with two definitions in topological order, surviving the Inline CTE batch), and QuerySuite tests proving such views execute correctly end to end with the conversion enabled. --- .../ConvertViewToMaterializedCTE.scala | 190 ++++++++ .../sql/catalyst/optimizer/Optimizer.scala | 1 + .../apache/spark/sql/internal/SQLConf.scala | 15 + .../ConvertViewToMaterializedCTESuite.scala | 428 ++++++++++++++++++ ...nvertViewToMaterializedCTEQuerySuite.scala | 123 +++++ 5 files changed, 757 insertions(+) create mode 100644 sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/optimizer/ConvertViewToMaterializedCTE.scala create mode 100644 sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/optimizer/ConvertViewToMaterializedCTESuite.scala create mode 100644 sql/core/src/test/scala/org/apache/spark/sql/ConvertViewToMaterializedCTEQuerySuite.scala diff --git a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/optimizer/ConvertViewToMaterializedCTE.scala b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/optimizer/ConvertViewToMaterializedCTE.scala new file mode 100644 index 0000000000000..32e390d9e075e --- /dev/null +++ b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/optimizer/ConvertViewToMaterializedCTE.scala @@ -0,0 +1,190 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.spark.sql.catalyst.optimizer + +import scala.collection.mutable + +import org.apache.spark.sql.catalyst.TableIdentifier +import org.apache.spark.sql.catalyst.expressions.{Alias, Attribute} +import org.apache.spark.sql.catalyst.plans.logical._ +import org.apache.spark.sql.catalyst.rules.Rule +import org.apache.spark.sql.internal.SQLConf + +/** + * Rewrites multiple references to the same view into a single `CTERelationDef` with multiple + * `CTERelationRef`s, so that the view's underlying plan is computed once (through exchange + * reuse at the physical layer) instead of once per reference. + * + * The rule runs in `FinishAnalysis`, immediately before `EliminateView`: after `EliminateView` + * no `View` nodes remain and every reference site holds an independent copy of the view's plan. + * + * A converted definition always sets `forceSkipInline = true`; otherwise `InlineCTE` would + * immediately flatten it back into duplicated subtrees (the definition body is deterministic + * in every case we convert), making the rule a no-op. + * + * Only deterministic, batch views are eligible: a multi-reference CTE guarantees that its + * definition is evaluated exactly once (even for non-deterministic definitions), while + * multiple references to a non-deterministic view are evaluated independently today. + * Converting such views would change query results. + * + * A view body may contain correlated subqueries whose outer references resolve to relations + * inside the same body (e.g. `t WHERE x IN (SELECT y FROM s WHERE s.k = t.k)`). The view is + * analyzed standalone when it is created, so an outer reference that does not resolve inside + * the body fails view analysis and can never escape to the outer query. The converted + * definition contains the whole body, so internal correlations resolve within it and these + * bodies are safe to convert. `InlineCTE`'s rejection of boundary-crossing outer references + * is only a generic safety net for non-view `forceSkipInline` producers. + * + * Each reference site gains a shuffle boundary added by `ReplaceCTERefWithRepartition` and + * deduplicated by exchange reuse, so the conversion trades recomputation for a + * shuffle plus reuse; it is therefore gated behind + * [[SQLConf.CONVERT_VIEW_TO_MATERIALIZED_CTE]] and off by default. + */ +object ConvertViewToMaterializedCTE extends Rule[LogicalPlan] { + + override def apply(plan: LogicalPlan): LogicalPlan = { + if (!SQLConf.get.getConf(SQLConf.CONVERT_VIEW_TO_MATERIALIZED_CTE)) return plan + val occurrences = plan.collectWithSubqueries { case v: View => v } + if (occurrences.length < 2) return plan + + // Group occurrences of the same view by their canonicalized body and captured SQL + // configs. Occurrences of one view differ only in renewed expression ids, which + // canonicalization normalizes away. + val qualifiedGroups = occurrences.groupBy(groupKey).values.filter(qualifies) + if (qualifiedGroups.isEmpty) return plan + // Match by identifier during the transform, not by the full group key: the key embeds + // the canonicalized body, and by the time an outer view is visited in the bottom-up + // traversal, nested views inside its body have already been rewritten into + // `CTERelationRef`s, so the body no longer canonicalizes to the key computed here. + // An identifier mapping to more than one qualified group would mean occurrences whose + // canonicalized bodies diverge; refuse conversion rather than rewrite all of them + // against whichever definition happens to be visited first. + val qualifiedIdentifiers = qualifiedGroups + .groupBy(g => groupKeyOf(g)._1) + .values + .filter(_.size == 1) + .map(_.head) + .map(groupKeyOf(_)._1) + .toSet + if (qualifiedIdentifiers.isEmpty) return plan + + // Bottom-up rewrite: nested views are visited (and their definitions appended) before + // the views containing them, so a referenced definition always precedes its referrer + // in `cteDefs`. The first occurrence of a group creates the definition, replaced by a + // bare reference; later occurrences are wrapped in a `Project` re-minting the + // occurrence's ids from the definition output, so consumers above need no rewriting. + val cteDefs = mutable.ArrayBuffer.empty[CTERelationDef] + val defByGroup = mutable.HashMap.empty[TableIdentifier, CTERelationDef] + + val rewritten = plan.transformUpWithSubqueries { + case v: View if qualifiedIdentifiers.contains(v.desc.identifier) => + defByGroup.get(v.desc.identifier) match { + case Some(cteDef) => + // Later occurrence: re-bind the reference output to this occurrence's + // attributes positionally. The group qualification has already asserted that + // name, type and nullability align element-wise. + val ref = CTERelationRef( + cteDef.id, + _resolved = true, + output = cteDef.output, + isStreaming = v.child.isStreaming) + Project(rebindingProjectList(v.output, cteDef.output), ref) + + case None => + // First occurrence: consumers above already reference this occurrence's + // expression ids, which are exactly the definition output, so the bare + // reference is output-compatible. + val cteDef = CTERelationDef(v.child, forceSkipInline = true) + defByGroup.put(v.desc.identifier, cteDef) + cteDefs += cteDef + CTERelationRef( + cteDef.id, + _resolved = true, + output = v.child.output, + isStreaming = v.child.isStreaming) + } + } + + if (cteDefs.isEmpty) { + plan + } else { + attachDefs(rewritten, cteDefs.toSeq) + } + } + + private type GroupKey = (TableIdentifier, LogicalPlan) + + // Group by view identity, not by structural body fingerprint: the rule dedupes references + // of the SAME view, not distinct views with coincidentally equal bodies. Because the + // analyzer resolves each occurrence of one view through the same deterministic path, all + // occurrences share a canonically equal body and (db-qualified) identifier, so identity + // alone is sufficient and unambiguous. The canonicalized body is kept in the key as a + // precondition guard: if a future change ever made two occurrences of one view diverge + // structurally, we skip conversion instead of building a wrong shared definition. + private def groupKey(v: View): GroupKey = + (v.desc.identifier, v.child.canonicalized) + + private def groupKeyOf(occs: Seq[View]): GroupKey = groupKey(occs.head) + + private def qualifies(occs: Seq[View]): Boolean = { + val first = occs.head + occs.length >= 2 && occs.forall { v => + v.resolved && + v.child.deterministic && + !v.child.isStreaming && + v.desc.viewSQLConfigs == first.desc.viewSQLConfigs && + schemasAlign(first, v) + } + } + + /** + * Occurrences of the same view are produced by resolution-time attribute renewal, which + * preserves column order, name, type and nullability. Assert the invariant element-wise + * before zipping positions: silent wrong results are the failure mode we cannot tolerate. + */ + private def schemasAlign(first: View, other: View): Boolean = + first.output.length == other.output.length && + first.output.zip(other.output).forall { case (l, r) => + l.name == r.name && l.dataType == r.dataType && l.nullable == r.nullable + } + + /** + * Builds a project list that re-mints `target`'s attributes from the definition output + * positionally, preserving names, expression ids, qualifiers and metadata so that + * consumers referencing `target` resolve unchanged. + */ + private def rebindingProjectList(target: Seq[Attribute], source: Seq[Attribute]): Seq[Alias] = + target.zip(source).map { case (t, s) => + Alias(s, t.name)(exprId = t.exprId, qualifier = t.qualifier, + explicitMetadata = Some(t.metadata)) + } + + /** + * Attaches the new definitions at the scope root, mirroring how `CTESubstitution` groups + * definitions: merged into an existing top-level `WithCTE`, spread onto command children + * for plans implementing `CTEInChildren`, or wrapped around the plan otherwise. References + * inside subquery expressions resolve against the top-level scope, as they do for regular + * user-written CTEs. + */ + private def attachDefs(plan: LogicalPlan, newDefs: Seq[CTERelationDef]): LogicalPlan = + plan match { + case WithCTE(child, cteDefs) => WithCTE(child, newDefs ++ cteDefs) + case cmd: LogicalPlan with CTEInChildren => cmd.withCTEDefs(newDefs) + case other => WithCTE(other, newDefs) + } +} diff --git a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/optimizer/Optimizer.scala b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/optimizer/Optimizer.scala index 58b6479ddee30..e9d80e437fb26 100644 --- a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/optimizer/Optimizer.scala +++ b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/optimizer/Optimizer.scala @@ -349,6 +349,7 @@ abstract class Optimizer(catalogManager: CatalogManager) EliminateResolvedHint, EliminateSubqueryAliases, EliminatePipeOperators, + ConvertViewToMaterializedCTE, EliminateView, EliminateSQLFunctionNode, ReplaceExpressions, diff --git a/sql/catalyst/src/main/scala/org/apache/spark/sql/internal/SQLConf.scala b/sql/catalyst/src/main/scala/org/apache/spark/sql/internal/SQLConf.scala index 2a8d23eaceba6..1fb4605e11f3f 100644 --- a/sql/catalyst/src/main/scala/org/apache/spark/sql/internal/SQLConf.scala +++ b/sql/catalyst/src/main/scala/org/apache/spark/sql/internal/SQLConf.scala @@ -983,6 +983,21 @@ object SQLConf { .booleanConf .createWithDefault(false) + val CONVERT_VIEW_TO_MATERIALIZED_CTE = + buildConf("spark.sql.optimizer.convertViewToMaterializedCTE") + .doc("When true, the optimizer rewrites multiple references to the same view " + + "into a single materialized CTE definition with multiple references, so the " + + "view's underlying plan is computed once (through exchange reuse) instead of " + + "once per reference. Views that contain non-deterministic expressions or " + + "streaming sources are never converted, as the conversion would change their " + + "semantics. Note that each reference site gains a shuffle boundary, so this " + + "helps only when the view subtree is expensive relative to a shuffle of its " + + "output.") + .version("4.4.0") + .withBindingPolicy(ConfigBindingPolicy.NOT_APPLICABLE) + .booleanConf + .createWithDefault(false) + val COMBINE_APPROXIMATE_PERCENTILES_ENABLED = buildConf("spark.sql.optimizer.combineApproximatePercentiles.enabled") .doc("When true, combines compatible scalar approximate percentile aggregates into a " + diff --git a/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/optimizer/ConvertViewToMaterializedCTESuite.scala b/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/optimizer/ConvertViewToMaterializedCTESuite.scala new file mode 100644 index 0000000000000..82fc825cb4a3a --- /dev/null +++ b/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/optimizer/ConvertViewToMaterializedCTESuite.scala @@ -0,0 +1,428 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.spark.sql.catalyst.optimizer + +import org.apache.spark.sql.catalyst.TableIdentifier +import org.apache.spark.sql.catalyst.catalog.{CatalogStorageFormat, CatalogTable, CatalogTableType} +import org.apache.spark.sql.catalyst.dsl.expressions._ +import org.apache.spark.sql.catalyst.dsl.plans._ +import org.apache.spark.sql.catalyst.expressions.{AttributeReference, ExprId, In, ListQuery, Literal, NamedExpression, OuterReference, ScalarSubquery} +import org.apache.spark.sql.catalyst.plans.Inner +import org.apache.spark.sql.catalyst.plans.PlanTest +import org.apache.spark.sql.catalyst.plans.logical._ +import org.apache.spark.sql.catalyst.rules.RuleExecutor +import org.apache.spark.sql.internal.SQLConf +import org.apache.spark.sql.types.{IntegerType, StructType} + +class ConvertViewToMaterializedCTESuite extends PlanTest { + + object Optimize extends RuleExecutor[LogicalPlan] { + val batches = + Batch("Convert View to Materialized CTE", FixedPoint(1), ConvertViewToMaterializedCTE) :: Nil + } + + object OptimizeWithInlineCTE extends RuleExecutor[LogicalPlan] { + val batches = + Batch("Convert View to Materialized CTE", FixedPoint(1), ConvertViewToMaterializedCTE) :: + Batch("Inline CTE", FixedPoint(1), InlineCTE()) :: Nil + } + + private def attr(name: String, id: Long): AttributeReference = + AttributeReference(name, IntegerType, nullable = true)(exprId = ExprId(id)) + + private def viewDesc(name: String, schema: StructType): CatalogTable = + CatalogTable( + identifier = TableIdentifier(name), + tableType = CatalogTableType.VIEW, + storage = CatalogStorageFormat.empty, + schema = schema) + + private def tempView(name: String, child: LogicalPlan): View = + View(viewDesc(name, child.schema), isTempView = true, child) + + /** + * One view `name` referenced twice, as the analyzer produces it: two `View` occurrences + * whose bodies are built from the same base attributes, except that the second + * occurrence's attributes carry renewed expression ids. The renewed attributes are + * returned as well for tests that reference them (e.g. in a join condition). Calling + * this instead of writing `tempView("v", ...)` twice makes it unambiguous that the pair + * stands for the same view used twice, not two views that happen to share a name. + */ + private def sameViewTwice( + name: String, + base: AttributeReference, + make: AttributeReference => LogicalPlan): (View, View, AttributeReference) = { + val renewed = base.withExprId(NamedExpression.newExprId) + (tempView(name, make(base)), tempView(name, make(renewed)), renewed) + } + + private def sameViewTwice( + name: String, + base: Seq[AttributeReference], + make: Seq[AttributeReference] => LogicalPlan): (View, View, Seq[AttributeReference]) = { + val renewed = base.map(_.withExprId(NamedExpression.newExprId)) + (tempView(name, make(base)), tempView(name, make(renewed)), renewed) + } + + // Same as above, for bodies that mint their own occurrence-specific ids per call + // (e.g. through `NamedExpression.newExprId`), so there is no base attribute to renew. + private def sameViewTwice(name: String, make: () => LogicalPlan): (View, View) = { + (tempView(name, make()), tempView(name, make())) + } + + // A simple deterministic view body: LocalRelation [a, b] filtered on a. + private def simpleBody(a: AttributeReference, b: AttributeReference): LogicalPlan = + LocalRelation(Seq(a, b)).where(a > 10) + + // A view body that contains a surviving (multi-ref) inner CTE. The references carry + // distinct expression ids, mirroring analyzer output. The inner CTE survives inlining + // either because it is non-deterministic or because it sets forceSkipInline. + private def nestedCteBody(defId: Long, deterministic: Boolean): LogicalPlan = { + val innerProject = + if (deterministic) OneRowRelation().select(Literal(1).as("r")) + else OneRowRelation().select(rand(0).as("r")) + val innerDef = CTERelationDef( + innerProject, + id = defId, + forceSkipInline = deterministic) + def mkRef(): CTERelationRef = CTERelationRef( + defId, + _resolved = true, + output = innerDef.output.map(_.withExprId(NamedExpression.newExprId)), + isStreaming = false) + WithCTE(Except(mkRef(), mkRef(), isAll = true), Seq(innerDef)) + } + + test("converts a self-joined view into one CTE definition and two references") { + withSQLConf(SQLConf.CONVERT_VIEW_TO_MATERIALIZED_CTE.key -> "true") { + val base = Seq(attr("a", 100), attr("b", 101)) + val (v1, v2, renewed) = sameViewTwice( + "v", base, (as: Seq[AttributeReference]) => simpleBody(as(0), as(1))) + val query = Join(v1, v2, Inner, Some(base(0) === renewed(0)), JoinHint(None, None)) + + val optimized = Optimize.execute(query) + + // The root is wrapped in a WithCTE carrying exactly one definition. + val WithCTE(mainPlan, cteDefs) = optimized + assert(cteDefs.length == 1) + val cteDef = cteDefs.head + assert(cteDef.forceSkipInline, + "converted CTE must set forceSkipInline so InlineCTE keeps it materialized") + assert(cteDef.child.canonicalized == simpleBody(base(0), base(1)).canonicalized) + + // Exactly two references in the main plan. + val refs = mainPlan.collect { case r: CTERelationRef => r } + assert(refs.length == 2) + + // The first reference adopts the definition output directly. + assert(refs.exists(_.output.map(_.exprId) == cteDef.output.map(_.exprId))) + + // The second reference is re-bound through an aliasing Project that re-mints the + // original expression ids of that occurrence, so consumers need no rewriting. + val projects = mainPlan.collect { + case p @ Project(_, _: CTERelationRef) => p + } + assert(projects.length == 1) + assert(projects.head.output.map(_.name) == Seq("a", "b")) + assert(projects.head.output.map(_.exprId) == renewed.map(_.exprId)) + + // The join condition still references the original attributes of both occurrences. + val join = mainPlan.collect { case j: Join => j }.head + assert(join.condition.get == (base(0) === renewed(0))) + + // The query output schema is unchanged. + assert(optimized.output == query.output) + } + } + + test("leaves a single-reference view unchanged") { + withSQLConf(SQLConf.CONVERT_VIEW_TO_MATERIALIZED_CTE.key -> "true") { + val a1 = attr("a", 100) + val b1 = attr("b", 101) + val query = Filter(a1 > 5, tempView("v", simpleBody(a1, b1))) + comparePlans(Optimize.execute(query), query) + } + } + + test("does not convert non-deterministic views") { + withSQLConf(SQLConf.CONVERT_VIEW_TO_MATERIALIZED_CTE.key -> "true") { + def randBody(r: AttributeReference): LogicalPlan = + Project(Seq(rand(0).as("r")), OneRowRelation()) + val (v1, v2, _) = sameViewTwice("v", attr("r", 100), r => randBody(r)) + val query = Join(v1, v2, Inner, None, JoinHint(None, None)) + comparePlans(Optimize.execute(query), query) + } + } + + test("does not convert streaming views") { + withSQLConf(SQLConf.CONVERT_VIEW_TO_MATERIALIZED_CTE.key -> "true") { + val (v1, v2, _) = sameViewTwice( + "v", attr("a", 100), + (a: AttributeReference) => LocalRelation(Seq(a), Nil, isStreaming = true)) + val query = Join(v1, v2, Inner, None, JoinHint(None, None)) + comparePlans(Optimize.execute(query), query) + } + } + + test("does not convert views with mixed effective SQL configs") { + withSQLConf(SQLConf.CONVERT_VIEW_TO_MATERIALIZED_CTE.key -> "true") { + val a1 = attr("a", 100) + val b1 = attr("b", 101) + val a2 = attr("a", 200) + val b2 = attr("b", 201) + val confKey = s"${CatalogTable.VIEW_SQL_CONFIG_PREFIX}spark.sql.foo" + val descWithConf = viewDesc("v", simpleBody(a1, b1).schema).copy( + properties = Map(confKey -> "bar")) + val v1 = View(descWithConf, isTempView = true, simpleBody(a1, b1)) + val v2 = tempView("v", simpleBody(a2, b2)) + val query = Join(v1, v2, Inner, Some(a1 === a2), JoinHint(None, None)) + comparePlans(Optimize.execute(query), query) + } + } + + test("bails out when occurrence schemas mismatch") { + withSQLConf(SQLConf.CONVERT_VIEW_TO_MATERIALIZED_CTE.key -> "true") { + // The two bodies are canonically equal, but the output column names differ + // ("x" vs "y"), so positional re-binding is unsafe and the group must be skipped. + val a1 = attr("a", 100) + val a2 = attr("a", 200) + val v1 = tempView("v", Project(Seq(a1.as("x")), LocalRelation(Seq(a1)))) + val v2 = tempView("v", Project(Seq(a2.as("y")), LocalRelation(Seq(a2)))) + val query = Join(v1, v2, Inner, None, JoinHint(None, None)) + comparePlans(Optimize.execute(query), query) + } + } + + test("refuses conversion when occurrences of the same view diverge into multiple groups") { + withSQLConf(SQLConf.CONVERT_VIEW_TO_MATERIALIZED_CTE.key -> "true") { + // Degenerate scenario the analyzer cannot produce today: two pairs of occurrences + // of the same view whose canonicalized bodies diverge. Each pair qualifies on its + // own, but rewriting all four occurrences against whichever definition happens to + // be visited first would re-bind the divergent pair positionally against the wrong + // schema, so the identifier must be skipped entirely. + val a1 = attr("a", 100) + val a2 = attr("a", 200) + val a3 = attr("a", 300) + val a4 = attr("a", 400) + def bodyWithOne(a: AttributeReference): LogicalPlan = + Project(Seq(a.as("x")), LocalRelation(Seq(a))) + def bodyWithTwo(a: AttributeReference): LogicalPlan = + Project(Seq(a.as("x"), a.as("y")), LocalRelation(Seq(a))) + val left = Join( + tempView("v", bodyWithOne(a1)), tempView("v", bodyWithOne(a2)), + Inner, None, JoinHint(None, None)) + val right = Join( + tempView("v", bodyWithTwo(a3)), tempView("v", bodyWithTwo(a4)), + Inner, None, JoinHint(None, None)) + val query = Join(left, right, Inner, None, JoinHint(None, None)) + comparePlans(Optimize.execute(query), query) + } + } + + test("converts views referenced inside scalar subqueries") { + withSQLConf(SQLConf.CONVERT_VIEW_TO_MATERIALIZED_CTE.key -> "true") { + val (v1, v2, _) = sameViewTwice( + "v", Seq(attr("a", 100), attr("b", 101)), + (as: Seq[AttributeReference]) => simpleBody(as(0), as(1))) + val sq1 = ScalarSubquery(v1) + val sq2 = ScalarSubquery(v2) + val query = OneRowRelation().select(sq1.as("s1"), sq2.as("s2")) + + val optimized = Optimize.execute(query) + + // The definitions are attached at the top-level scope while the references live + // inside the scalar subqueries. + val WithCTE(mainPlan, cteDefs) = optimized + assert(cteDefs.length == 1) + assert(cteDefs.head.forceSkipInline) + val refs = optimized.collectWithSubqueries { case r: CTERelationRef => r } + assert(refs.length == 2) + assert(mainPlan.collect { case r: CTERelationRef => r }.isEmpty, + "references must live inside the subqueries, not in the main plan") + assert(optimized.output == query.output) + } + } + + test("converted definition survives the Inline CTE batch") { + withSQLConf(SQLConf.CONVERT_VIEW_TO_MATERIALIZED_CTE.key -> "true") { + val base = Seq(attr("a", 100), attr("b", 101)) + val (v1, v2, renewed) = sameViewTwice( + "v", base, (as: Seq[AttributeReference]) => simpleBody(as(0), as(1))) + val query = Join(v1, v2, Inner, Some(base(0) === renewed(0)), JoinHint(None, None)) + + val optimized = OptimizeWithInlineCTE.execute(query) + + val withCTEs = optimized.collect { case w: WithCTE => w } + assert(withCTEs.nonEmpty, + "deterministic single-ref CTEs get inlined, but the converted def must survive") + val defs = optimized.collect { case d: CTERelationDef => d } + assert(defs.length == 1) + assert(optimized.output == query.output) + } + } + + test("does not convert views whose body contains non-deterministic inner CTEs") { + withSQLConf(SQLConf.CONVERT_VIEW_TO_MATERIALIZED_CTE.key -> "true") { + // The view body itself evaluates a multi-ref non-deterministic inner CTE, so the view + // yields different results per evaluation. Converting it to a compute-once CTE would + // change query results, hence it must stay unconverted. + val (v1, v2) = sameViewTwice( + "v", () => nestedCteBody(987654321L, deterministic = false)) + val query = Join(v1, v2, Inner, None, JoinHint(None, None)) + comparePlans(Optimize.execute(query), query) + } + } + + test("converts a view whose body contains a surviving deterministic inner CTE") { + withSQLConf(SQLConf.CONVERT_VIEW_TO_MATERIALIZED_CTE.key -> "true") { + // Both occurrences share the same inner CTE definition id, mirroring how the + // analyzer duplicates a view body while keeping its inner CTE ids intact. + val (v1, v2) = sameViewTwice( + "v", () => nestedCteBody(987654321L, deterministic = true)) + val query = Join(v1, v2, Inner, None, JoinHint(None, None)) + + val optimized = Optimize.execute(query) + + val WithCTE(mainPlan, cteDefs) = optimized + assert(cteDefs.length == 1) + assert(cteDefs.head.forceSkipInline) + // The converted definition wraps the inner WithCTE of the view body. + assert(cteDefs.head.child.isInstanceOf[WithCTE]) + val refs = mainPlan.collect { case r: CTERelationRef => r } + assert(refs.length == 2) + assert(optimized.output == query.output) + } + } + + test("converts a view whose body references another view") { + withSQLConf(SQLConf.CONVERT_VIEW_TO_MATERIALIZED_CTE.key -> "true") { + // v2's body is the v1 view referenced in the query twice, so both views are + // referenced twice (v1 once inside each v2 body) and both groups qualify. + val (v2a, v2b, _) = sameViewTwice( + "v2", Seq(attr("a", 100), attr("b", 101)), + (as: Seq[AttributeReference]) => tempView("v1", simpleBody(as(0), as(1)))) + val query = Join(v2a, v2b, Inner, None, JoinHint(None, None)) + + val optimized = Optimize.execute(query) + + val WithCTE(mainPlan, cteDefs) = optimized + assert(cteDefs.length == 2) + assert(cteDefs.forall(_.forceSkipInline)) + // Bottom-up creation order puts the referenced (inner) definition before the + // definition that references it, so the outer body resolves its inner reference + // without needing the inner definition to be attached later. + val innerDef = cteDefs.head + val outerDef = cteDefs.last + assert(innerDef.child.collect { case _: CTERelationRef => true }.isEmpty) + val innerRefs = outerDef.child.collect { case r: CTERelationRef => r } + assert(innerRefs.map(_.cteId) == Seq(innerDef.id)) + + // Only the outer view is referenced from the main plan, once bare and once + // re-bound through a rebinding Project. + val refs = mainPlan.collect { case r: CTERelationRef => r } + assert(refs.length == 2) + assert(refs.forall(_.cteId == outerDef.id)) + assert(optimized.output == query.output) + } + } + + test("nested view definitions survive the Inline CTE batch") { + withSQLConf(SQLConf.CONVERT_VIEW_TO_MATERIALIZED_CTE.key -> "true") { + val (v2a, v2b, _) = sameViewTwice( + "v2", Seq(attr("a", 100), attr("b", 101)), + (as: Seq[AttributeReference]) => tempView("v1", simpleBody(as(0), as(1)))) + val query = Join(v2a, v2b, Inner, None, JoinHint(None, None)) + + val optimized = OptimizeWithInlineCTE.execute(query) + + val defs = optimized.collect { case d: CTERelationDef => d } + assert(defs.length == 2, "both nested definitions must survive inlining") + // The order follows the bottom-up creation order: the referenced inner view's + // definition (a bare body with no references) comes first, and the outer view's + // definition wraps a reference to it. + assert(defs.head.child.collect { case _: CTERelationRef => true }.isEmpty) + val innerRefs = defs.last.child.collect { case r: CTERelationRef => r } + assert(innerRefs.map(_.cteId) == Seq(defs.head.id)) + assert(optimized.output == query.output) + } + } + + test("is disabled by default") { + val (v1, v2, _) = sameViewTwice( + "v", Seq(attr("a", 100), attr("b", 101)), + (as: Seq[AttributeReference]) => simpleBody(as(0), as(1))) + val query = Join(v1, v2, Inner, None, JoinHint(None, None)) + comparePlans(Optimize.execute(query), query) + } + + test("does not convert distinct views even with identical bodies") { + withSQLConf(SQLConf.CONVERT_VIEW_TO_MATERIALIZED_CTE.key -> "true") { + // The rule dedupes references of the SAME view, identified by its catalog identifier, + // not views that happen to share a body. Two differently-named views with identical + // bodies must not be merged into one CTE definition. + val a1 = attr("a", 100) + val b1 = attr("b", 101) + val a2 = attr("a", 200) + val b2 = attr("b", 201) + val v1 = tempView("v1", simpleBody(a1, b1)) + val v2 = tempView("v2", simpleBody(a2, b2)) + val query = Join(v1, v2, Inner, Some(a1 === a2), JoinHint(None, None)) + comparePlans(Optimize.execute(query), query) + } + } + + test("is idempotent") { + withSQLConf(SQLConf.CONVERT_VIEW_TO_MATERIALIZED_CTE.key -> "true") { + val base = Seq(attr("a", 100), attr("b", 101)) + val (v1, v2, renewed) = sameViewTwice( + "v", base, (as: Seq[AttributeReference]) => simpleBody(as(0), as(1))) + val query = Join(v1, v2, Inner, Some(base(0) === renewed(0)), JoinHint(None, None)) + val once = Optimize.execute(query) + val twice = Optimize.execute(once) + comparePlans(twice, once) + } + } + + test("converts a view whose body contains an internally correlated subquery") { + withSQLConf(SQLConf.CONVERT_VIEW_TO_MATERIALIZED_CTE.key -> "true") { + // The body mirrors an analyzed `t WHERE x IN (SELECT y FROM s WHERE s.k = t.k)`: + // the correlation references `t.k`, an attribute produced by a relation inside the + // body itself, so it resolves within the body and does not escape it. The correlation + // survives as an `OuterReference` inside the subquery plan; the body must not be + // rejected for carrying it. + def body(x: AttributeReference, tk: AttributeReference, + s: AttributeReference, y: AttributeReference): LogicalPlan = { + val correlated = LocalRelation(Seq(s, y)).where(OuterReference(tk) === s).select(y) + LocalRelation(Seq(x, tk)).where( + In(x, Seq(ListQuery(correlated, outerAttrs = Seq(tk), numCols = 1)))) + } + val v1 = tempView("v", body(attr("x", 100), attr("t", 101), attr("s", 102), attr("y", 103))) + val v2 = tempView("v", body(attr("x", 200), attr("t", 201), attr("s", 202), attr("y", 203))) + val query = Join(v1, v2, Inner, None, JoinHint(None, None)) + + val optimized = Optimize.execute(query) + + val WithCTE(mainPlan, cteDefs) = optimized + assert(cteDefs.length == 1) + assert(cteDefs.head.forceSkipInline) + val refs = mainPlan.collect { case r: CTERelationRef => r } + assert(refs.length == 2) + assert(optimized.output == query.output) + } + } +} diff --git a/sql/core/src/test/scala/org/apache/spark/sql/ConvertViewToMaterializedCTEQuerySuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/ConvertViewToMaterializedCTEQuerySuite.scala new file mode 100644 index 0000000000000..70a021e0114a9 --- /dev/null +++ b/sql/core/src/test/scala/org/apache/spark/sql/ConvertViewToMaterializedCTEQuerySuite.scala @@ -0,0 +1,123 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.spark.sql + +import org.apache.spark.sql.catalyst.plans.logical.RepartitionByExpression +import org.apache.spark.sql.functions.rand +import org.apache.spark.sql.internal.SQLConf +import org.apache.spark.sql.test.SharedSparkSession + +/** + * Integration tests for the `ConvertViewToMaterializedCTE` optimizer rule: repeated + * references to the same view are rewritten into one CTE definition with multiple + * references when `spark.sql.optimizer.convertViewToMaterializedCTE` is enabled. After + * the final `Replace CTE with Repartition` batch, a converted view shows up as one + * repartition node per reference site; exchange reuse deduplicates them at execution + * time. + */ +class ConvertViewToMaterializedCTEQuerySuite extends QueryTest with SharedSparkSession { + import testImplicits._ + + private val selfJoinQuery = + "SELECT t1.id, t2.k FROM v t1 JOIN v t2 ON t1.id = t2.id WHERE t1.id < 10" + + private def withSelfJoinedView(f: => Unit): Unit = { + withTempView("v") { + spark.range(0, 100).select($"id", ($"id" % 10).as("k")).createOrReplaceTempView("v") + f + } + } + + private def countRepartitions(query: String): Int = + spark.sql(query).queryExecution.optimizedPlan.collect { + case _: RepartitionByExpression => true + }.length + + test("self-joined view returns identical results with conversion enabled") { + withSelfJoinedView { + val expected = spark.sql(selfJoinQuery).collect() + withSQLConf(SQLConf.CONVERT_VIEW_TO_MATERIALIZED_CTE.key -> "true") { + checkAnswer(spark.sql(selfJoinQuery), expected) + } + } + } + + test("conversion adds one repartition per reference site") { + withSelfJoinedView { + assert(countRepartitions(selfJoinQuery) == 0) + withSQLConf(SQLConf.CONVERT_VIEW_TO_MATERIALIZED_CTE.key -> "true") { + // One shuffle boundary per reference site; identical shuffles are then reused. + assert(countRepartitions(selfJoinQuery) == 2) + } + } + } + + test("non-deterministic views are not converted") { + withTempView("rand_view") { + spark.range(0, 10).select($"id", rand(0).as("r")) + .createOrReplaceTempView("rand_view") + val query = "SELECT t1.r FROM rand_view t1 JOIN rand_view t2 ON t1.id = t2.id" + withSQLConf(SQLConf.CONVERT_VIEW_TO_MATERIALIZED_CTE.key -> "true") { + assert(countRepartitions(query) == 0) + } + } + } + + test("insert into a table selecting from a repeatedly referenced view") { + withTable("dest") { + withSelfJoinedView { + sql("CREATE TABLE dest (id BIGINT, k BIGINT) USING parquet") + withSQLConf(SQLConf.CONVERT_VIEW_TO_MATERIALIZED_CTE.key -> "true") { + sql(s"INSERT INTO dest $selfJoinQuery") + } + checkAnswer(spark.table("dest"), spark.range(0, 10).select($"id", ($"id" % 10))) + } + } + } + + test("internally correlated view executes correctly with conversion enabled") { + withTempView("t", "s", "v") { + spark.range(5).selectExpr("id AS k", "id AS x").createOrReplaceTempView("t") + spark.range(5).selectExpr("id AS k", "id AS y").createOrReplaceTempView("s") + sql("CREATE OR REPLACE TEMP VIEW v AS " + + "SELECT * FROM t WHERE x IN (SELECT y FROM s WHERE s.k = t.k)") + val query = "SELECT * FROM v, v" + val expected = spark.sql(query).collect() + withSQLConf(SQLConf.CONVERT_VIEW_TO_MATERIALIZED_CTE.key -> "true") { + checkAnswer(spark.sql(query), expected) + } + } + } + + test("nested view is converted together with the view it references") { + withTempView("v1", "v2") { + spark.range(0, 100).select($"id", ($"id" % 10).as("k")).createOrReplaceTempView("v1") + sql("CREATE OR REPLACE TEMP VIEW v2 AS SELECT id, k FROM v1 WHERE k < 5") + val query = "SELECT t1.id FROM v2 t1 JOIN v2 t2 ON t1.id = t2.id" + val expected = spark.sql(query).collect() + assert(countRepartitions(query) == 0) + withSQLConf(SQLConf.CONVERT_VIEW_TO_MATERIALIZED_CTE.key -> "true") { + checkAnswer(spark.sql(query), expected) + // Both views convert: each reference site of v2 renders its own shuffle + // boundary with v1's boundary nested inside, so 2 sites x 2 boundaries = 4. + // (v1 alone would yield 2, no conversion 0.) + assert(countRepartitions(query) == 4) + } + } + } +}