yyanyy commented on code in PR #58298:
URL: https://github.com/apache/spark/pull/58298#discussion_r4022018264


##########
sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/v2/AnalyzedSchemaProjection.scala:
##########
@@ -0,0 +1,356 @@
+/*
+ * 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.execution.datasources.v2
+
+import org.apache.spark.{SparkException, SparkThrowable}
+import org.apache.spark.sql.catalyst.SQLConfHelper
+import org.apache.spark.sql.catalyst.analysis.Resolver
+import org.apache.spark.sql.catalyst.expressions.{Alias, ArrayTransform, 
AttributeReference, AttributeSeq, CreateNamedStruct, Expression, ExtractValue, 
GetStructField, If, IsNull, KnownNotNull, LambdaFunction, Literal, 
MetadataAttributeWithLogicalName, NamedLambdaVariable, TaggingExpression, 
TransformKeys, TransformValues, UnresolvedNamedLambdaVariable}
+import org.apache.spark.sql.catalyst.plans.logical.{LogicalPlan, Project}
+import org.apache.spark.sql.catalyst.util.MetadataColumnHelper
+import org.apache.spark.sql.types.{ArrayType, DataType, MapType, Metadata, 
StructType}
+import org.apache.spark.sql.util.SchemaUtils
+
+/**
+ * Restores the output a plan was analyzed with on a relation whose table has 
since changed - the
+ * schema `V2TableUtil.validateCapturedColumns` calls the captured one.
+ *
+ * The relation exposes the table's current schema, so its output stays 
aligned with the physical
+ * scan, and a projection on top recreates the columns, types and expression 
IDs the parent plan was
+ * analyzed against.
+ */
+private[sql] object AnalyzedSchemaProjection extends SQLConfHelper {
+
+  /**
+   * Prevents [[CreateNamedStruct]] from inheriting metadata from a field 
value while leaving the
+   * value's type, nullability, evaluation, and code generation unchanged.
+   */
+  private case class MetadataPropagationBarrier(child: Expression) extends 
TaggingExpression {
+    override protected def withNewChildInternal(
+        newChild: Expression): MetadataPropagationBarrier = copy(child = 
newChild)
+  }
+
+  def rebindToAnalyzedSchema(relation: DataSourceV2Relation): LogicalPlan = {
+    // The relation still carries the output captured at analysis time; only 
its table has been
+    // swapped for the current one.
+    val capturedOutput = relation.output
+    val resolver = conf.resolver
+    val caseSensitive = conf.caseSensitiveAnalysis
+    val current = DataSourceV2Relation.create(
+      relation.table,
+      relation.catalog,
+      relation.identifier,
+      relation.options,
+      relation.timeTravelSpec)
+    val currentMetadataOutput = current.metadataOutput
+    val currentMetadata = capturedOutput.filter(_.isMetadataCol).map { 
captured =>
+      val logicalName = metadataLogicalName(captured)
+      matchFoldedName(currentMetadataOutput, logicalName, 
caseSensitive)(metadataLogicalName)
+        .map(pos => currentMetadataOutput(pos))
+        .getOrElse {
+          // The connector still reports this metadata column, so it can only 
be absent here
+          // because a data column has taken its name and the connector 
suppresses rather than
+          // renames the conflict (`canRenameConflictingMetadataColumns`). 
Validation owns
+          // rejecting that.
+          unexpectedSchemaChange(
+            s"captured metadata column $logicalName is missing from the 
current relation")
+        }
+    }
+
+    val currentOutput = current.output ++ currentMetadata
+
+    // Refresh may visit an already rebound relation. Preserve its attributes 
so the projection
+    // above it continues to reference valid expression IDs.
+    //
+    // A further schema change on such a relation adds a second projection 
instead of replacing
+    // the first. Only the cache stores a refreshed plan, so the effect is 
limited to that entry:
+    // it stops matching the single projection a query rebuilds from its own 
captured output, and
+    // is no longer reused. Results stay correct.
+    if (sameOutputShape(capturedOutput, currentOutput)) {
+      return relation
+    }
+
+    val capturedIndex = new AttributeIndex(capturedOutput, resolver, 
caseSensitive)
+    val reboundOutput = currentOutput.map { currentAttr =>
+      capturedIndex.get(currentAttr).filter(canReuse(_, 
currentAttr)).getOrElse(currentAttr)
+    }
+    val reboundRelation = relation.copy(output = reboundOutput)
+
+    val reboundIndex = new AttributeIndex(reboundOutput, resolver, 
caseSensitive)
+    val projectList = capturedOutput.map { capturedAttr =>

Review Comment:
   Good catch. I agree this case is real. It requires an already-analyzed query 
plus an external catalog change that adds a resolver-conflicting nested 
field—Spark DDL rejects this addition—and the affected captured column must be 
wholly unused. The result is a query failure rather than incorrect data.
   I looked at pruning unused outputs here, but doing that safely before 
optimizer pruning becomes plan-wide attribute-liveness analysis rather than a 
local filter. LogicalPlan.references excludes implicitly consumed attributes: 
unions map children by position, DISTINCT/INTERSECT/EXCEPT consume whole rows, 
and CTE or subquery references may use different expression IDs.
   Given how narrow this case is, I would prefer to handle that broader 
analysis in a follow-up rather than expand this PR further. I removed “safely” 
from the migration guide so it no longer makes the broad claim identified here.



-- 
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]

Reply via email to