[
https://issues.apache.org/jira/browse/SPARK-59120?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
Peter Toth updated SPARK-59120:
-------------------------------
Fix Version/s: 4.2.1
> Fix ClassCastException in a storage-partitioned join whose partition keys
> were reduced
> --------------------------------------------------------------------------------------
>
> Key: SPARK-59120
> URL: https://issues.apache.org/jira/browse/SPARK-59120
> Project: Spark
> Issue Type: Bug
> Components: SQL
> Affects Versions: 4.2.0, 4.3.0, 5.0.0, 4.4.0
> Reporter: Peter Toth
> Assignee: Peter Toth
> Priority: Major
> Labels: pull-request-available
> Fix For: 4.3.0, 4.2.1
>
>
> With {{spark.sql.sources.v2.bucketing.allowCompatibleTransforms.enabled}} a
> storage-partitioned join reduces one or both sides' partition keys onto a
> common key space, while the partitioning keeps reporting the partition
> expressions it was built from. The stored key rows are then no longer what
> those expressions describe, and where the reduction also changes the type,
> reading a key row at the expression's type unboxes an Integer as a Long. Two
> queries fail with {{ClassCastException}} because of that.
> *1. At planning.* Two tables partitioned by {{identity(ts)}} and by
> {{years(ts)}}, joined on {{ts}}, with the subset opt-in on. AQE re-runs
> {{createShuffleSpec}} on the already reduced children through
> {{ValidateRequirements}}, which projects and sorts the keys, and the sort
> orders {{IntegerType}} year values with a comparator built for
> {{TimestampType}}:
> {noformat}
> java.lang.ClassCastException: class java.lang.Integer cannot be cast to class
> java.lang.Long
> at KeyedPartitioning.toGrouped
> at KeyedPartitioning.createShuffleSpec
> at ValidateRequirements
> {noformat}
> *2. While the shuffle is prepared.* The same two tables plus a third,
> unpartitioned one joined on the same column, with
> {{spark.sql.sources.v2.bucketing.shuffle.enabled}} on. A reduced partitioning
> is grouped and does satisfy the distribution, so {{EnsureRequirements}}
> offers it as the layout to shuffle the third child onto.
> {{ShuffleExchangeExec}} then builds the shuffle's key map by re-wrapping the
> stored keys at the expressions' types, and the same unboxing throws on the
> driver.
> The reduce, and the reporting of reduced keys under the original expressions,
> arrived in SPARK-47094 (4.0.0). The crashes start in 4.2.0, with the
> {{KeyedPartitioning}} / {{GroupPartitionsExec}} refactor that derives the key
> ordering and types from the reported expressions. On 4.1 and earlier the same
> key-space mismatch is there, but the symptom is wrong results rather than an
> exception, and the code is different enough to need its own ticket.
> Not in scope here: a reduction that keeps the type, {{bucket(12)}} and
> {{bucket(8)}} both reducing onto {{bucket(4)}}, still misroutes rows
> silently. Answering that needs the partitioning to say whether its
> expressions still describe its keys, which is a follow-up.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]