Skip to content
Merged
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
49 changes: 41 additions & 8 deletions docs/source/user-guide/tuning-guide.md
Original file line number Diff line number Diff line change
Expand Up @@ -145,34 +145,67 @@ 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
full build side must fit in memory per task.
- 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;
```

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
Expand Down