Skip to content

Broadcast lowering still hash-repartitions the build and probe sides of a CollectLeft join #1905

Description

@andygrove

Is your feature request related to a problem or challenge?

When the static distributed planner promotes a join to a broadcast
HashJoinExec(CollectLeft) (maybe_promote_to_broadcast + the broadcast lowering
in plan_query_stages_internal), the join's inputs still carry the join-key
RepartitionExec(Hash) that EnforceDistribution inserted for the original
partitioned join. As a result:

  • Build side is hash-repartitioned into a shuffle stage and then written
    again as a broadcast stage — two shuffles where broadcast should need one (just
    replicate the build's natural partitions).
  • Probe side still reshuffles on the join key, even though a CollectLeft
    join replicates the build to every probe task and does not require the probe to
    be partitioned on the join key.

So the broadcast promotion adds a broadcast without removing the reshuffles it
was meant to avoid. On TPC-H SF10 (AQE off) the SMJ-broadcast path (#1904) still
gives ~16% because eliminating the sort + collecting a small side helps, but a
large part of the intended benefit — skipping the join-key shuffles entirely — is
left on the table.

This affects both broadcast paths:

Describe the solution you'd like

During broadcast lowering, strip the redundant join-key RepartitionExec(Hash)
from the converted join's inputs:

  • broadcast the build side from its natural (upstream) partitions instead of
    hash-repartitioning then broadcasting;
  • keep the probe side at its upstream partitioning rather than reshuffling on the
    join key.

Additional context

Correctness to verify: a CollectLeft join's output partitioning follows the
probe side. Removing the probe's join-key repartition changes the probe (and
therefore the join output) partitioning, so downstream operators that assume
hash-on-join-key partitioning must be re-checked / re-satisfied by
EnforceDistribution/EnforceSorting. Build-side broadcast must still replicate
all build partitions to every probe task.

Follow-up to #1904 (#1679). Related: #342, #348.

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Type

    No type

    Projects

    No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions