Peter Toth created SPARK-59971:
----------------------------------

             Summary: 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: 5.0.0
            Reporter: Peter Toth


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