[
https://issues.apache.org/jira/browse/SPARK-59015?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
Yan Yan updated SPARK-59015:
----------------------------
Issue Type: Improvement (was: Bug)
> Rebind refreshed DSv2 relations to the current table schema
> -----------------------------------------------------------
>
> Key: SPARK-59015
> URL: https://issues.apache.org/jira/browse/SPARK-59015
> Project: Spark
> Issue Type: Improvement
> Components: SQL
> Affects Versions: 4.2.0
> Reporter: Yan Yan
> Priority: Minor
>
> Spark refreshes versioned DataSourceV2 tables between analysis and
> optimization. {{V2TableRefreshUtil}} does {{r.copy(table = currentTable)}},
> swapping in the current table while keeping the analyzed output. When the
> table gained a column after analysis, the relation output and the current
> scan schema disagree:
> {code}
> relation.output = [id, salary] // captured at analysis
> scan data schema = [id, salary, new_column] // current snapshot
> {code}
> {{PushDownUtils.pruneColumns}} then calls {{toOutputAttrs(scan.readSchema(),
> relation)}}, which builds its name map only from {{relation.output}}, so any
> field the captured output does not know about throws:
> {code}
> java.util.NoSuchElementException: key not found: new_column
> {code}
> The trigger is precisely that {{Scan.readSchema()}} field names are not a
> subset of {{relation.output}} field names. Reporting more than requested is
> allowed: {{SupportsPushDownRequiredColumns.pruneColumns}} permits partial
> pruning (its example is a source that prunes only top-level columns, not
> nested fields). That example reproduces it:
> {code}
> sql("CREATE TABLE t (id INT, person STRUCT<name: STRING>) USING ...")
> val df = spark.table("t") // captures person: STRUCT<name>
> sql("ALTER TABLE t ADD COLUMN person.age INT FIRST")
> df.collect()
> {code}
> {code}
> [PLAN_VALIDATION_FAILED_RULE_IN_BATCH] ... schema has changed from
> STRUCT<id: INT, person: STRUCT<name: STRING>> to
> STRUCT<id: INT, person: STRUCT<age: INT, name: STRING>>
> {code}
> A one-shot query analyzes and optimizes back to back, so this window is
> normally invisible; it becomes deterministic by forcing analysis on a
> DataFrame before changing the table, or via a temp view created from a
> Dataset. Not Connect-specific.
> There is also a second, silent bug: because {{toOutputAttrs}} keys by name, a
> captured metadata column whose name a new data column has taken collapses
> onto one attribute together with that data column, so a query for the
> metadata column silently returned the data column's value. The user-facing
> rejection for that is SPARK-59014 (PR #58295); this issue's fix turns it into
> an internal-error assertion.
> *Fix:* in {{V2TableRefreshUtil}} (ALLOW_NEW_FIELDS mode), rebind the
> refreshed relation to the current table schema and interpose a {{Project}}
> that recreates the captured output (names, expression IDs, types,
> nullability, metadata) via a new {{CapturedSchemaProjection}} that recurses
> through structs, arrays and maps. Unchanged when the table has not changed
> (idempotent under transformDown).
> PR: https://github.com/apache/spark/pull/58298
--
This message was sent by Atlassian Jira
(v8.20.10#820010)
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]