[ 
https://issues.apache.org/jira/browse/SPARK-58120?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
 ]

Cheng Pan reassigned SPARK-58120:
---------------------------------

    Assignee: Yuming Wang

> Fix BatchScanExec.doCanonicalize reorders keyGroupedPartitioning expressions
> ----------------------------------------------------------------------------
>
>                 Key: SPARK-58120
>                 URL: https://issues.apache.org/jira/browse/SPARK-58120
>             Project: Spark
>          Issue Type: Bug
>          Components: SQL
>    Affects Versions: 4.2.0, 5.0.0
>            Reporter: Yuming Wang
>            Assignee: Yuming Wang
>            Priority: Major
>              Labels: pull-request-available
>
> h2. Problem
> `BatchScanExec.doCanonicalize()` uses `QueryPlan.normalizePredicates()` to 
> canonicalize `keyGroupedPartitioning` expressions:
> {code:scala}
> override def doCanonicalize(): BatchScanExec = {
>   this.copy(
>     output = output.map(QueryPlan.normalizeExpressions(_, output)),
>     runtimeFilters = QueryPlan.normalizePredicates(
>       runtimeFilters.filterNot(_ == 
> DynamicPruningExpression(Literal.TrueLiteral)),
>       output),
>     keyGroupedPartitioning = 
> keyGroupedPartitioning.map(QueryPlan.normalizePredicates(_, output)))
>   )
> }
> {code}
> `QueryPlan.normalizePredicates()` combines expressions with `And`, 
> canonicalizes the combined expression (which can reorder children), then 
> splits back. This destroys the original expression order.
> For example, if `keyGroupedPartitioning` is `Seq(id, data)` where `id` is 
> `IntegerType` and `data` is `StringType`, after canonicalization the 
> expressions may be reordered to `Seq(data, id)`. When this canonicalized 
> `BatchScanExec` ends up in the execution plan (e.g., through a fallback or 
> reuse path), `outputPartitioning` creates a `KeyedPartitioning` with 
> expression data types `[StringType, IntegerType]` but partition key row 
> values `(Integer, String)`. This type mismatch causes `ClassCastException` 
> during `TimSort` when sorting partition keys by the expression-derived 
> ordering.
> h2. Impact
> This bug manifests when join keys are a subset of partition keys (e.g., 
> SPARK-44647, SPARK-47094 scenarios). The `ClassCastException` occurs in 
> `DataSourceV2ScanExecBase.outputPartitioning()` when sorting 
> `inputPartitions` by `partitionKey()` using `keyRowOrdering` built from the 
> (reordered) expression data types:
> {noformat}
> java.lang.ClassCastException: class java.lang.Integer cannot be cast to class 
> org.apache.spark.unsafe.types.UTF8String
>   at 
> org.apache.spark.sql.catalyst.expressions.BaseGenericInternalRow.getUTF8String(rows.scala:45)
>   ...
>   at 
> org.apache.spark.sql.execution.datasources.v2.DataSourceV2ScanExecBase.outputPartitioning(DataSourceV2ScanExecBase.scala:99)
> {noformat}
> h2. Reproduction
> Create a table partitioned by two columns with different types (e.g., 
> `identity(id)` as `IntegerType` and `identity(data)` as `StringType`). When 
> the plan is canonicalized (which happens through exchange reuse, AQE, or 
> Gluten's fallback mechanism), the expression order is reversed, causing the 
> `ClassCastException`.
> h2. Fix
> Use `QueryPlan.normalizeExpressions` instead of 
> `QueryPlan.normalizePredicates` for `keyGroupedPartitioning`. 
> `normalizeExpressions` normalizes attribute references without 
> combining/reordering, preserving the expression order:
> {code:scala}
> override def doCanonicalize(): BatchScanExec = {
>   this.copy(
>     output = output.map(QueryPlan.normalizeExpressions(_, output)),
>     runtimeFilters = QueryPlan.normalizePredicates(
>       runtimeFilters.filterNot(_ == 
> DynamicPruningExpression(Literal.TrueLiteral)),
>       output),
>     keyGroupedPartitioning = keyGroupedPartitioning.map(_.map(
>       QueryPlan.normalizeExpressions(_, output))))
>   )
> }
> {code}
> h2. Affected versions
> This affects all Spark versions that have the `keyGroupedPartitioning` field 
> in `BatchScanExec.doCanonicalize()`, including the current master branch.



--
This message was sent by Atlassian Jira
(v8.20.10#820010)

---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to