yyanyy commented on code in PR #58298: URL: https://github.com/apache/spark/pull/58298#discussion_r4009344011
########## sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/v2/CapturedSchemaProjection.scala: ########## @@ -0,0 +1,310 @@ +/* + * 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.analysis.Resolver +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} + +/** + * 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 resolver = conf.resolver + 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, resolver)(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 Review Comment: Your point about the test is right and I've fixed it. `numCachedEntries == 1` held whether the entry was usable or not, so it recorded nothing; the test now calls `assertNotCached` after the second refresh. That states the actual limitation, and it will fail once this is fixed, prompting a flip back to `assertCached`. Note the same expression passes `assertCached` after the first refresh in that test, so the new assertion does discriminate. On the fix itself I'd still like to keep it out of this PR, but I don't think my earlier reasoning got a response, so let me name the three parts I'd want you to disagree with specifically: 1. Reaching this needs a cached plan, two external schema changes with a refresh in between, and a Dataset derived from the pre-change one. Results stay correct, and the entry returns as soon as it is re-cached. 2. This area is best-effort by design — "inability to refresh cache shouldn't fail operations" drops the entry outright on any refresh failure. 3. Replacing rather than nesting requires recognizing, from inside the rule, that the projection above a relation is one we generated. I looked at doing that with a `TreeNodeTag`. It would probably work, but making plan identity depend on a mutable annotation that has to survive every rewrite and cache normalization is a fragile mechanism and I'd rather not introduce it here. Without it the fix needs either a first-class representation of the captured output on the relation, or a change to which plan `CacheManager` stores as the key — which is SPARK-54424's own design question. Both are larger than this change. -- 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]
