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

    Projects

    No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions