andygrove opened a new pull request, #25820:
URL: https://github.com/apache/datafusion/pull/25820

   ## Which issue does this PR close?
   
   - Part of #25804 (finding 5).
   - Related to #22526.
   
   ## Rationale for this change
   
   A grouped aggregate emits its output as zero-copy slices of one batch. 
`SortExec` reserves every buffered batch's full buffer capacity, so each slice 
is charged for the whole aggregate output. A sort after a `GROUP BY` therefore 
reserves several times the memory it holds: it spills when it doesn't need to, 
and with spilling disabled it fails with `ResourcesExhausted`. For example, 
`select group_key, sum(payload) as s from t group by group_key order by 
group_key` over 100,000 groups needs a 24 MB memory limit on `main` and runs in 
6 MB with this change.
   
   ## What changes are included in this PR?
   
   - `ExternalSorter` counts its buffered batches with a 
`RecordBatchMemoryCounter`, as #22862 did for the hash join build side, so a 
buffer shared between batches is reserved once. Each batch still reserves its 
sliced size for sorting it.
   - When the buffered batches are sorted as separate runs and the runs share 
buffers, the reservation for the runs' input is held until every run is sorted, 
because a run still waiting to be sorted keeps the shared buffers alive. Runs 
that share nothing release their input as soon as they are sorted, as before.
   - Coalescing batches into runs for single-column sorts sizes the batches the 
same way.
   - Schemas with view, dictionary or list view types keep reserving each batch 
on its own. Sorting those types keeps referencing the input's data buffers or 
values, and the sorted runs and the merge charge those buffers once per run. 
Counting them once at insert would let a spill take more runs than the merge 
can then hold (`sort_spill_is_compacted_for_view_arrays` fails that way). 
Counting them once needs the merge side to change too; see #25800 and #25791.
   
   ## What is the testing strategy for this PR?
   
   - `sort_after_group_by_reserves_shared_output_once` in `memory_limit` runs 
`GROUP BY ... ORDER BY` with one sort key and with two, under a 12 MB limit 
with spilling disabled. It fails on `main` with `Memory Exhausted while 
Sorting`.
   - `test_sliced_runs_reserve_shared_buffers_once` checks that slices of one 
batch reserve its buffers once, and that the buffers stay reserved until the 
last run holding them is sorted.
   - The `spill_tests` fixture now builds each batch separately. Its batches 
were slices of one batch, and the tests relied on each slice being charged for 
the whole batch to create memory pressure.
   
   Benchmarks: `cargo bench -p datafusion --bench sort` with the 
`release-nonlto` profile, running `main` and this PR alternately, median of 
three runs:
   
   | benchmark | `main` | this PR | change |
   |---|---|---|---|
   | sort i64 1M | 24.11 ms | 23.94 ms | -0.7% |
   | sort partitioned i64 1M | 3.05 ms | 3.21 ms | +5.3% |
   | sort utf8 high cardinality 1M | 59.54 ms | 60.72 ms | +2.0% |
   | sort partitioned utf8 high cardinality 1M | 4.84 ms | 4.93 ms | +2.0% |
   | sort utf8 dictionary tuple 1M | 510.35 ms | 511.02 ms | +0.1% |
   | sort partitioned utf8 dictionary tuple 1M | 10.44 ms | 10.65 ms | +2.0% |
   | sort mixed tuple 1M | 106.05 ms | 106.43 ms | +0.4% |
   | sort partitioned mixed tuple 1M | 11.51 ms | 11.12 ms | -3.4% |
   
   `sort partitioned i64 1M` was 5-12% slower across longer runs. The 
difference comes from moving the shared reservation, empty in this benchmark, 
into each run's sort future: a build that doesn't move it matches `main`. The 
counting itself costs about 2 µs per partition. The dictionary cases don't take 
the new path, which puts the noise at about 2%.
   
   ## Are there any user-facing changes?
   
   Sorts over sliced input without view or dictionary columns, such as the 
output of a grouped aggregate, reserve less memory and spill less. No API 
changes.
   


-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to