[
https://issues.apache.org/jira/browse/SPARK-59971?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
Peter Toth updated SPARK-59971:
-------------------------------
Affects Version/s: 4.2.0
4.1.3
4.0.4
(was: 5.0.0)
> The one-side shuffle of a storage-partitioned join ignores
> spark.sql.requireAllClusterKeysForCoPartition
> --------------------------------------------------------------------------------------------------------
>
> Key: SPARK-59971
> URL: https://issues.apache.org/jira/browse/SPARK-59971
> Project: Spark
> Issue Type: Bug
> Components: SQL
> Affects Versions: 4.2.0, 4.1.3, 4.0.4
> Reporter: Peter Toth
> Priority: Major
>
> With spark.sql.sources.v2.bucketing.shuffle.enabled on, a join can
> co-partition on only some of its join keys, although
> spark.sql.requireAllClusterKeysForCoPartition (true by default) is meant to
> prevent that. {{k1(a, b)}} and {{k2(a, b)}} are identity-partitioned on
> {{a}}, and {{plain(a, b)}} is not partitioned:
> {code:sql}
> SELECT * FROM k1 x JOIN k2 y ON x.a = y.a AND x.b = y.b;
> SELECT * FROM k1 x JOIN plain y ON x.a = y.a AND x.b = y.b;
> {code}
> * The first query takes no shuffle. The storage-partitioned join declines,
> since {{a}} does not cover the join key {{b}}. Then the shuffle path keeps
> both sides as they are, so the join still runs on {{k1}}'s partitions.
> * The second query shuffles {{plain}} onto {{k1}}'s layout, so the join runs
> with {{k1}}'s partition count, on {{a}} only.
> With the shuffle conf off, both queries shuffle both sides on {{(a, b)}}. The
> results are correct either way. What is lost is the protection against skew
> and low parallelism. Measured on master and on branch-4.0.
> {{HashShuffleSpec.canCreatePartitioning}} checks the conf, but
> {{KeyedShuffleSpec.canCreatePartitioning}} does not. So a keyed spec can be
> the layout that the other children are shuffled onto or paired with, even
> when its keys cover only some of the join keys.
> This came with SPARK-41471 in 4.0.0, which let a keyed spec create
> partitionings. cloud-fan raised the skew concern in
> https://github.com/apache/spark/pull/42194 after the merge. SPARK-58558 made
> the storage-partitioned join path check that the partition keys cover the
> join keys, and kept a partial cover gated by the conf. It did not change the
> one-side shuffle path.
> The SPARK-54439 and SPARK-52246 tests assert a one-side shuffle onto a keyed
> layout that covers only some of the join keys, with the conf at its default.
> So a fix has to decide whether that shuffle is intended, or update those
> tests.
> https://github.com/apache/spark/pull/59165 lists another case among its
> costs. A join on two keys whose sides pair only through a transform of an
> expression gets one more shuffle with that PR. Both plans co-partition on one
> key only.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]