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


##########
sql/catalyst/src/main/scala/org/apache/spark/sql/util/SchemaUtils.scala:
##########
@@ -529,12 +529,23 @@ private[spark] object SchemaUtils {
     if (field.nullable) s"$name $dataType" else s"$name $dataType NOT NULL"
   }
 
+  /**
+   * Folds a name to the key used to decide whether two names refer to the 
same column or field.
+   *
+   * This is the identity rule name resolution is built on: `AttributeSeq` 
looks attributes up by
+   * this key and only then filters the candidates with the resolver, and the 
duplicate-name checks
+   * above reject a schema that holds two names folding to one key. Matching 
by folded name
+   * therefore finds exactly the field resolution would, while comparing with 
the resolver alone
+   * can match several fields a schema is allowed to keep apart 
(`equalsIgnoreCase` equates U+017F
+   * LONG S with `s`, which this fold does not).
+   */
+  def foldName(name: String, caseSensitiveAnalysis: Boolean): String = {
+    if (caseSensitiveAnalysis) name else name.toLowerCase(Locale.ROOT)

Review Comment:
   [P1] Keep duplicate validation and rebinding on one unique matching rule
   
   The uniqueness invariant claimed here does not currently hold. This helper 
folds with `Locale.ROOT`, while `checkColumnNameDuplication` still uses 
parameterless `toLowerCase` and therefore the JVM default locale. Under Turkish 
locale, retain captured `i\u0307` (ID 1) and prepend U+0130 LATIN CAPITAL I 
WITH DOT ABOVE (ID 2). Duplicate validation lowercases those to distinct 
strings and accepts the schema, but ROOT folding maps both to `i\u0307`. 
Compatibility indexing keeps the later retained field and validates ID 1, while 
`CapturedSchemaProjection.matchName` uses `indexWhere` and selects the first 
newly added field, returning its values under the captured identity. A fresh 
`AttributeSeq` lookup still filters the ROOT-keyed candidates with the resolver 
and selects the retained field. Please use one shared unique matching/indexing 
routine across duplicate validation, compatibility validation, and rebinding, 
and add a Turkish-locale end-to-end test.



##########
sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/v2/CapturedSchemaProjection.scala:
##########
@@ -0,0 +1,304 @@
+/*
+ * 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
+import org.apache.spark.sql.catalyst.SQLConfHelper
+import org.apache.spark.sql.catalyst.expressions.{Alias, ArrayTransform, 
AttributeReference, CreateNamedStruct, Expression, 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
+
+/**
+ * Rebinds a relation that reads a current table schema to output attributes 
captured from an
+ * earlier compatible schema. The current schema is exposed by the relation so 
its output remains
+ * aligned with the physical scan, while a projection recreates the captured 
output for the
+ * already-analyzed parent plan.
+ */
+private[sql] object CapturedSchemaProjection 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 rebindToCapturedSchema(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 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)
+      matchName(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, 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, caseSensitive)
+    val projectList = capturedOutput.map { capturedAttr =>
+      val currentAttr = reboundIndex.get(capturedAttr).getOrElse {
+        unexpectedSchemaChange(
+          s"captured column ${capturedAttr.name} is missing from current table 
${relation.name}")
+      }
+      if (currentAttr.exprId == capturedAttr.exprId &&
+        sameAttributeShape(currentAttr, capturedAttr)) {
+        currentAttr
+      } else {
+        if (currentAttr.nullable != capturedAttr.nullable) {
+          unexpectedSchemaChange(
+            s"nullability changed for captured column ${capturedAttr.name} in 
${relation.name}")
+        }
+        val projected = projectToType(
+          currentAttr, currentAttr.dataType, capturedAttr.dataType, 
caseSensitive)
+        if (projected.dataType != capturedAttr.dataType ||
+          projected.nullable != capturedAttr.nullable) {
+          unexpectedSchemaChange(
+            s"failed to recreate captured column ${capturedAttr.name} in 
${relation.name}")
+        }
+        Alias(projected, capturedAttr.name)(
+          exprId = capturedAttr.exprId,
+          qualifier = capturedAttr.qualifier,
+          explicitMetadata = Some(capturedAttr.metadata))
+      }
+    }
+
+    Project(projectList, reboundRelation)
+  }
+
+  private[v2] def projectToType(
+      input: Expression,
+      from: DataType,
+      to: DataType,
+      caseSensitive: Boolean): Expression = {
+    if (from == to) {
+      return input
+    }
+
+    val projected = (from, to) match {
+      case (fromStruct: StructType, toStruct: StructType) =>
+        val structInput = if (input.nullable) KnownNotNull(input) else input
+        val fields = toStruct.fields.iterator.flatMap { targetField =>
+          val index = matchName(fromStruct, targetField.name, 
caseSensitive)(_.name).getOrElse {
+            unexpectedSchemaChange(
+              s"captured struct field ${targetField.name} is missing from 
$fromStruct")
+          }
+          val sourceField = fromStruct.fields(index)
+          val value = projectToType(
+            GetStructField(structInput, index, Some(sourceField.name)),
+            sourceField.dataType,
+            targetField.dataType,
+            caseSensitive)
+          val namedValue = if (targetField.metadata == Metadata.empty) {
+            // An empty captured value is still an explicit instruction not to 
inherit metadata
+            // from the current GetStructField. CleanupAliases removes an 
empty-metadata Alias, so
+            // use a local passthrough barrier instead of changing shared 
alias cleanup behavior.
+            MetadataPropagationBarrier(value)
+          } else {
+            Alias(value, targetField.name)(explicitMetadata = 
Some(targetField.metadata))
+          }
+          Iterator(Literal(targetField.name), namedValue)
+        }.toSeq
+        val rebuilt = CreateNamedStruct(fields)
+        if (input.nullable) {
+          // The null literal takes the rebuilt type rather than `toStruct` so 
that the type check
+          // below still sees any mismatch: `If` merges its branch types and 
only requires them to
+          // match up to `sameType`, which ignores nullability and metadata.
+          If(IsNull(input), Literal.create(null, rebuilt.dataType), rebuilt)
+        } else {
+          rebuilt
+        }
+
+      case (ArrayType(fromElement, fromContainsNull), ArrayType(toElement, 
toContainsNull)) =>
+        if (fromContainsNull != toContainsNull) {
+          unexpectedSchemaChange(s"array element nullability changed from 
$from to $to")
+        }
+        val element = NamedLambdaVariable(
+          UnresolvedNamedLambdaVariable.freshVarName("element"),
+          fromElement,
+          fromContainsNull)
+        ArrayTransform(
+          input,
+          LambdaFunction(
+            projectToType(element, fromElement, toElement, caseSensitive), 
Seq(element)))
+
+      case (
+            MapType(fromKey, fromValue, fromValueContainsNull),
+            MapType(toKey, toValue, toValueContainsNull)) =>
+        if (fromValueContainsNull != toValueContainsNull) {
+          unexpectedSchemaChange(s"map value nullability changed from $from to 
$to")
+        }
+
+        val withProjectedKeys = if (fromKey != toKey) {

Review Comment:
   [P1] Reject additions inside map keys before rebinding
   
   Adding a field inside a map-key struct is accepted by `ALLOW_NEW_FIELDS`, 
but narrowing that key is not a projection. Two current keys that differ only 
in the added field collapse to the same captured key. With the default policy 
this makes a supposedly compatible refresh fail based on row contents; with 
`LAST_WIN` it silently drops entries. The end-to-end test explicitly blesses a 
current two-entry map becoming size one. A schema change classified as 
compatible must not change stale-query cardinality according to 
`spark.sql.mapKeyDedupPolicy`. Please reject additions beneath map keys during 
compatibility validation with a user-facing schema-change error.



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