[ 
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]

Reply via email to