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]