NoahKusaba opened a new pull request, #2501:
URL: https://github.com/apache/datafusion-ballista/pull/2501
## Which issue does this PR close?
- Closes #2454.
## Rationale for this change
AQE decides to broadcast a join's build side from *estimated* statistics.
When that side then *measures* larger than the probe side, the replan's
`join_selection` swaps the `CollectLeft` join and the broadcast `ExchangeExec`
lands on the probe side. A broadcast reader reports a single partition, and a
`CollectLeft` join has as many partitions as its probe side, so the whole join
stage collapses to one task.
TPC-H q8 at SF1000 hit this in stage 5: one task, 79.9M input rows, `min =
median = max = 7429ms`, feeding a 256-way shuffle write.
```
SortShuffleWriterExec: partitioning=Hash([l_orderkey@0], 256)
HashJoinExec: mode=CollectLeft, join_type=Inner, on=[(s_suppkey@0,
l_suppkey@1)], ...
CoalescePartitionsExec
ShuffleReaderExec: upstream_stage: 3, partitioning:
UnknownPartitioning(256)
ShuffleReaderExec: upstream_stage: 2, broadcast: true,
upstream_partition_count: 256
```
The probe-side reader should be `partitioning: UnknownPartitioning(256)`, so
the stage runs 256 tasks.
This supersedes #2459, which fixed the same bug through a new DataFusion
`ExecutionPlan` hook (apache/datafusion#25354). Both are closed: the fix
belongs in Ballista, next to `DemoteUnsafeBroadcastJoinRule`, and this version
builds against released DataFusion.
## What changes are included in this PR?
- **`PartitionProbeSideBroadcastRule`** runs straight after
`join_selection`, alongside `DemoteUnsafeBroadcastJoinRule`. When a
`CollectLeft` hash join's probe side is a broadcast `ExchangeExec`, it replaces
it with the same stage read partitioned. It leaves alone:
- null-aware joins, which track probe-side NULLs in-process and must run
as one task;
- broadcasts over ordered inputs, which the k-way merge reader handles.
- **`ExchangeExec::to_partitioned`**, the inverse of `to_broadcast`: the
same exchange with `broadcast` off, sharing its stage state and plan id. It
returns `None` when the stage's stored locations don't match its input's
partition count (a stage that ran as a hash shuffle and was broadcast later),
since a partitioned read would then silently skip locations.
- **`take_stage_output_partitions`** stores a broadcast stage's locations
one entry per partition, in partition order, sized by the stage's *input*
partition count. It used to ask the exchange, which reports `1` for a
broadcast, and stored the locations unordered from a `HashMap`. That was fine
while the only reader merged everything into one partition; a partitioned read
needs entry `i` to be partition `i`.
- **`StageOutput::partition_locations_broadcast`** is removed;
`partition_locations(n)` now serves both readers.
## What is the testing strategy for this PR?
End-to-end, through `AdaptivePlanner`
(`scheduler/src/state/aqe/test/probe_side_broadcast.rs`):
- `swapped_broadcast_is_read_partitioned_on_the_probe_side`: the q8 shape.
The swapped join's probe side reads 4 partitions, not 1. Fails with the rule
unregistered.
- `swapped_broadcast_reads_the_locations_its_tasks_reported`: a broadcast
stage whose tasks report in reverse order stores its locations in partition-id
order.
- `null_aware_join_keeps_a_single_task_probe_side`: a null-aware `NOT IN`
join still runs as a single task.
Rule unit tests with full plan snapshots
(`optimizer_rule/partition_probe_side.rs`):
- a probe-side broadcast is read partitioned;
- a null-aware join keeps its probe-side broadcast;
- a broadcast over an ordered input is left to the merge reader;
- a broadcast whose stage was written with a different partition count is
left as a broadcast.
Each positive test was checked to fail with its code path disabled.
## Are there any user-facing changes?
No API or configuration changes. Join stages whose broadcast side gets
swapped onto the probe side now run across all partitions instead of a single
task.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]