apache / apache/datafusion

NestedLoopJoin buffers the build side into a single concat_batches allocation: 2x transient peak, invisible to the memory pool

Open
#24,819 0 comments 0 reactions 0 assignees View on GitHub
Dominant language
Rust
Stars
9.3k
Forks
2.4k
Avg merge
3d 7h
Merged PRs (30d)
344

Description

### Is your feature request related to a problem or challenge?

`NestedLoopJoinExec` collects every build-side batch and concatenates them into one `RecordBatch` (`concat_batches` in `collect_left_input`). Three concrete failures follow from that single allocation:

**1. 2x transient peak.** The input batches and the concat output coexist until the copy finishes, so peak memory is twice the build side (the NestedLoopJoin instance of #23076). Reproducer, ~880 MB build side, inequality join so it cannot hash:

```sql
-- nlj_mem.sql
set datafusion.execution.target_partitions = 1;
select count(big.s)
from (select v1, repeat('x', 32) as s from generate_series(1, 20000000) as t(v1)) as big
join (select 2147483647 as v2) as small
on big.v1 > small.v2;
```

```
/usr/bin/time -l datafusion-cli -f nlj_mem.sql # macOS ("maximum resident set size")
/usr/bin/time -v datafusion-cli -f nlj_mem.sql # Linux ("Maximum resident set size")
```

Measured on main (5e168c9d0, release build, 3 runs, stable to 1 MB): peak RSS **1738 MB**, 1.97x the build side.

**2. The peak is invisible to the memory pool.** Each input batch is reserved as it arrives, but the concat output is never reserved. The same query with `datafusion-cli -m 1g` completes successfully while the process peaks at **1739 MB** — 1.7x its configured memory limit. A scheduler provisioning by the pool limit sees almost double the promised footprint.

**3. All-or-nothing allocation.** One giant batch cannot be released or spilled incrementally by the memory-limited mode, and a build side with more than `i32::MAX` bytes in a single string column cannot be represented at all — #23032 reports hitting exactly this offset overflow in production at ~60M rows.

### Describe the solution you'd like

Keep the build side as a `Vec` of target-batch-size chunks produced by arrow's `BatchCoalescer` (the utility #23076 points at), with prefix-sum row offsets so the visited-left bitmap keeps addressing global row numbers, and route the build-side spill through the same coalescer so the memory-limited replay reads back uniform chunks. Batches already at or above half the target size bypass the coalescer and are retained zero-copy.

With that implemented, the reproducer above peaks at **899 MB** (1.02x the build side), and at **900 MB** under `-m 1g`. Probe throughput is unchanged: a build side delivered as 250,000 8-row batches probes at parity with main (7.36 s vs 7.20 s medians over interleaved runs, ~1.6e10 pair evaluations) because the coalescer compacts before probing starts, and normal 8192-row batches are also at parity (7.08 s vs 7.04 s) since they are retained without copying.

### Describe alternatives you've considered

`concat_batches_owned` (#23598): consuming inputs incrementally shrinks the transient overlap but still produces one giant allocation, so the pool blindness improves while the offset overflow and the all-or-nothing release remain; it was closed by its author in favor of avoiding the concat entirely. Retaining the raw input batches (#23032) covers hash join and piecewise merge join in the same change, but leaves probe cost tied to whatever batch sizes the input produced and leaves sliced-batch accounting open (both raised in that PR's review); coalescing to uniform chunks answers both for NestedLoopJoin.

### Additional context

Part of #23076 / #23031, scoped to NestedLoopJoin only; hash join is a natural follow-up. Cross-partition sharing of the visited bitmap stays #22038. Found while running NestedLoopJoin under DataFusion Comet with Spark-sized broadcast build sides.

Contributor guide

Open the contributing guide

Research direction

Locate NestedLoopJoinExec and collect_left_input, then inspect Arrow's BatchCoalescer and the existing build-side spill and replay path. Use nlj_mem.sql with datafusion-cli to establish the memory baseline. Done means chunked build batches use prefix-sum row offsets, retain large batches without copying, and keep memory usage within the configured pool during probing and replay.

Written by the indexing model from the issue text.

Assessment

Tech stack
rust, sql
Domain
databases
Issue type
Feature
Difficulty
5/5
Estimated time
Over a week
Activity status
Active
Clarity
Mostly clear
Newbie friendliness
48/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.