Skip to content

AQE regresses join-heavy TPC-H queries at SF1000 instead of promoting broadcasts #2024

Description

@andygrove

Describe the bug

At TPC-H SF1000 with sort-merge joins, enabling the adaptive planner (ballista.planner.adaptive.enabled=true) makes the join-heavy queries slower, not faster, and costs two queries entirely.

These are exactly the queries that dynamic join selection (#1752 / DynamicJoinSelectionExec) is supposed to rescue by promoting the small side to a broadcast. Instead they regress 10–36%, which is what you would expect if the machinery is adding runtime overhead without actually firing the broadcast promotion.

Query AQE off (s) AQE on (s) Delta
3 189.3 248.3 +31%
5 446.0 605.6 +36%
7 481.5 545.7 +13%
8 647.4 733.6 +13%
9 900.2 994.7 +10%
19 28.4 81.9 +188%
17 670.9 348.9 −48%
18 587.9 deadlock (#2023)
21 1012.9 execution error

Totals over the 20 queries the adaptive path completed: AQE off 4173.9s vs AQE on 4344.0s — about 4% slower overall, plus 2 failures. The static path completed 22/22.

Q17 (−48%) shows the promotion machinery can fire. It just doesn't fire on most of the joins that dominate runtime at this scale.

Context: this is where the gap to Spark lives

Same hardware (2 executors × 8 cores / 56Gi), same SF1000 node-local parquet, single iteration:

Engine / config Total, 22 queries
Spark 3.5 (AQE on) 4621.5s
Ballista, sort-merge joins (default: prefer_hash_join=false), AQE off 5774.7s

Spark is 1.25× faster than stock Ballista, and the deficit is concentrated almost entirely in Q7/Q9/Q21. Inspecting the Spark plans, Spark's AQE emits 36 BroadcastHashJoins alongside 93 SortMergeJoins — i.e. Spark's "sort-merge default" is not really all sort-merge; roughly a quarter of its joins get promoted to broadcasts at runtime.

Ballista's static planner promotes zero, because broadcast promotion only fires for HashJoinExec and prefer_hash_join=false makes every join an SMJ (see #1922). The adaptive planner is the mechanism that is supposed to recover those broadcasts, and at SF1000 it is not doing so.

For comparison, forcing hash joins (prefer_hash_join=true) on the same cluster is 1.28× faster than Spark over Q1–17 (2683.6s vs 3441.4s), versus 4039.1s for the SMJ default — so the join strategy alone swings Ballista by ~1.5×. (That config currently cannot complete the suite; see the CollectLeft OOM issue filed alongside this one.)

To Reproduce

Cluster: 2 executors, 8 CPU / 56Gi each, --memory-pool-size=48GB, 1 scheduler. TPC-H SF1000 parquet on node-local disk. Ballista 54.0.0-rc2. Each query run as its own job, 1 iteration.

cargo run --release --bin tpch -- benchmark ballista \
  --host <scheduler> --port 50050 \
  --query 9 --path /mnt/bigdata/tpch/sf1000 --format parquet \
  --partitions 32 --iterations 1 \
  -c datafusion.optimizer.prefer_hash_join=false \
  -c datafusion.optimizer.enable_dynamic_filter_pushdown=false \
  -c ballista.planner.adaptive.enabled=true    # vs =false

Expected behavior

With the adaptive planner enabled, join-heavy queries whose build side fits the broadcast threshold should be promoted to a broadcast and get faster — approaching what Spark's AQE achieves — rather than regressing.

Additional context

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Type

    No type

    Fields

    No fields configured for issues without a type.

    Projects

    No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions