Peter Toth created SPARK-59900:
----------------------------------
Summary: Let a storage-partitioned join compare transforms of
join-key expressions by their shape
Key: SPARK-59900
URL: https://issues.apache.org/jira/browse/SPARK-59900
Project: Spark
Issue Type: Improvement
Components: SQL
Affects Versions: 5.0.0
Reporter: Peter Toth
SPARK-59887 makes a storage-partitioned join refuse any partition transform
whose argument is not a bare column, such as the {{bucket(4, b + 1)}} or
{{bucket(4, s.a)}} that a one-side shuffle builds from a join key. That is
correct, but it costs shuffles where the pairing was sound:
* Two sides of the same shape, e.g. {{bucket(4, b + 1)}} and {{bucket(4, c +
1)}} joined on {{b = c}}, are both shuffled.
* A join output that offers such a layout next to a usable one cannot be
shuffled onto at all, because {{ShuffleSpecCollection.canCreatePartitioning}}
requires every member to be usable. That turns 2 shuffles into 3.
* Present before SPARK-59887: the one-side shuffle's own {{bucket(4, b + 1)}}
partitioning does not satisfy the {{[b + 1]}} distribution it was built for, so
{{ValidateRequirements}} rejects that join.
The way out is to describe the argument relative to its column:
* {{TransformFunctionId}} records each argument with its single reference
replaced by a placeholder, and a nested transform by its own id, canonicalized.
Then {{b + 1}} and {{c + 1}} compare equal, while {{b + 1}} and {{x}} do not.
* {{KeyedShuffleSpec.createPartitioning}} replaces the reference inside the
argument rather than the whole argument, so it builds {{bucket(4, y + 1)}}
instead of {{bucket(4, y)}}.
* {{TransformExpression.reducers}} requires equal normalized arguments.
* The {{canCreatePartitioning}} clause from SPARK-59887 then goes away.
SPARK-50593 generalizes {{TransformFunctionId}} to per-argument shapes, so this
should build on it or be done together with it.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]