From 8ad6d3483741a233ff61916116329ba27783acc2 Mon Sep 17 00:00:00 2001 From: Noah Date: Tue, 22 Sep 2026 16:12:51 -0400 Subject: [PATCH 1/2] docs: describe how joins are chosen under AQE The Join Strategy section only described the static planner. Under AQE, the default, datafusion.optimizer.prefer_hash_join is not consulted; the join is chosen at runtime from broadcast_join_threshold_bytes and hash_join_max_build_partition_bytes. Also, join reordering compares byte sizes first and falls back to row counts. Co-Authored-By: Claude Opus 5.5 --- docs/source/user-guide/tuning-guide.md | 54 +++++++++++++++++++++----- 1 file changed, 44 insertions(+), 10 deletions(-) diff --git a/docs/source/user-guide/tuning-guide.md b/docs/source/user-guide/tuning-guide.md index 9c238a503..31e481489 100644 --- a/docs/source/user-guide/tuning-guide.md +++ b/docs/source/user-guide/tuning-guide.md @@ -145,7 +145,38 @@ Because operators spill to disk rather than fail, make sure the executor's ## Join Strategy -Ballista defaults to **sort-merge join** rather than hash join. This is the +How Ballista picks a join depends on whether adaptive query execution is on. + +### With AQE (the default) + +AQE chooses each join at runtime, from measured sizes, after the join's inputs +have run. It turns both hash joins and sort-merge joins into a single runtime +node and decides from that evidence, so `datafusion.optimizer.prefer_hash_join` +has no effect here. + +- **Broadcast.** If the smaller side is under + `ballista.optimizer.broadcast_join_threshold_bytes` (128 MiB by default) and + the join type allows it, that side is broadcast and the join runs as a + `CollectLeft` hash join, without shuffling the larger side. +- **Partitioned hash join.** Otherwise both sides are shuffled on the join key, + and the join runs as a hash join when every build partition is under + `ballista.optimizer.hash_join_max_build_partition_bytes` (64 MiB by default). +- **Sort-merge join.** A build partition over that limit falls back to + sort-merge join, which spills to disk under memory pressure. + +When a build side's size is only an estimate, AQE can shuffle that side on its +own first to measure it, so a join whose build side was over-estimated can still +be broadcast. See [What AQE does today](#what-aqe-does-today). + +Setting `ballista.optimizer.hash_join_max_build_partition_bytes` to `0` disables +the per-partition check, so AQE uses a hash join regardless of build size and +sort-merge join becomes unreachable. + +### With AQE turned off + +The static planner keeps DataFusion's physical planning choice. +`SessionConfig::new_with_ballista()` sets `datafusion.optimizer.prefer_hash_join` +to `false`, so joins are planned as **sort-merge joins** by default. This is the opposite of DataFusion's standalone default and reflects two facts: - DataFusion's hash join implementation does not yet support spilling: the @@ -153,15 +184,16 @@ opposite of DataFusion's standalone default and reflects two facts: - Ballista executors run multiple tasks in parallel per host, so per-task build sides aggregate quickly under load and can OOM the executor. -Sort-merge join spills under memory pressure (via the executor's memory -pool, when configured), making it the safer default for distributed -execution. +Sort-merge join spills under memory pressure via the executor's memory pool, +making it the safer default for distributed execution. -If you know the build side of a particular query fits comfortably in -memory and you want hash-join performance, opt back in at the session -level: +The static planner only promotes hash joins to broadcast, so with this default +Ballista does not broadcast joins. If you know the build side of a particular +query fits comfortably in memory and you want hash-join performance, including +broadcast promotion, opt back in at the session level: ```sql +SET ballista.planner.adaptive.enabled = false; SET datafusion.optimizer.prefer_hash_join = true; ``` @@ -169,10 +201,11 @@ or in code: ```rust let session_config = SessionConfig::new_with_ballista() + .set_bool("ballista.planner.adaptive.enabled", false) .set_bool("datafusion.optimizer.prefer_hash_join", true); ``` -This setting applies per session and does not require restarting the +These settings apply per session and do not require restarting the scheduler or executors. ## Shuffle Implementation @@ -259,8 +292,9 @@ shuffle stage completes, the planner re-optimizes the remaining plan and emits the next set of runnable stages. Two adaptive optimizations are currently implemented: -- **Join reordering.** Uses runtime row counts from completed stages so the - smaller side drives the join. +- **Join reordering.** Uses runtime byte sizes from completed stages, falling + back to row counts when sizes are unavailable, so the smaller side drives the + join. - **Broadcast join selection.** When a join input's runtime size falls under `ballista.optimizer.broadcast_join_threshold_bytes` (or the row-count fallback), the smaller side is broadcast (`CollectLeft`) instead of shuffled. From 044ece6e1cfdf4f51ce4eaf313116d2a2780dc91 Mon Sep 17 00:00:00 2001 From: Noah Date: Tue, 22 Sep 2026 19:17:07 -0400 Subject: [PATCH 2/2] docs: leave the AQE join bullet to #2478 #2478 now merges the join reordering and broadcast join selection bullets into one Join selection bullet that covers the byte-size comparison, so this PR no longer needs to edit it, and would conflict if it did. Co-Authored-By: Claude Opus 5.5 --- docs/source/user-guide/tuning-guide.md | 5 ++--- 1 file changed, 2 insertions(+), 3 deletions(-) diff --git a/docs/source/user-guide/tuning-guide.md b/docs/source/user-guide/tuning-guide.md index 31e481489..d224aba0c 100644 --- a/docs/source/user-guide/tuning-guide.md +++ b/docs/source/user-guide/tuning-guide.md @@ -292,9 +292,8 @@ shuffle stage completes, the planner re-optimizes the remaining plan and emits the next set of runnable stages. Two adaptive optimizations are currently implemented: -- **Join reordering.** Uses runtime byte sizes from completed stages, falling - back to row counts when sizes are unavailable, so the smaller side drives the - join. +- **Join reordering.** Uses runtime row counts from completed stages so the + smaller side drives the join. - **Broadcast join selection.** When a join input's runtime size falls under `ballista.optimizer.broadcast_join_threshold_bytes` (or the row-count fallback), the smaller side is broadcast (`CollectLeft`) instead of shuffled.