adriangb opened a new pull request, #25910:
URL: https://github.com/apache/datafusion/pull/25910
## Which issue does this PR close?
- Related to #25804 (finding 8: spill-merge admission charges 2 × the
largest batch of each run).
- Related to #25565 and #25853, which bound batches by bytes in the sort
merge output and in aggregate spilling.
## Rationale for this change
An `ORDER BY` over wide rows fails with `ResourcesExhausted` under a memory
limit, even when each input batch fits in the pool and the sort spills:
```sql
-- 512 rows of 64 KiB each (32 MiB), in input batches of 4 rows,
-- FairSpillPool of 12 MiB, target_partitions = 1, batch_size = 512,
-- sort_spill_reservation_bytes = 1 MiB
SELECT id, payload FROM t ORDER BY id DESC;
```
```text
Resources exhausted: Failed to allocate additional 20.0 MB for
ExternalSorterMerge[0] with 1024.0 KB already allocated for this reservation -
11.0 MB remain available for the total memory pool: fair(pool_size: 12.0 MB)
```
The same query with a `Utf8View` payload fails the same way (16.0 MB).
The sort writes its spill files in batches of up to `batch_size` rows,
whatever their size in bytes. With wide rows, one spilled batch is larger than
the pool. The final merge reserves memory for the largest batch of each spill
file (`max_record_batch_memory`), so it cannot start. Smaller input batches do
not help, because the sort puts the rows back together into `batch_size`-row
batches before it writes them.
#25565 bounds the merge output by bytes for `Utf8` and `Binary`. In a local
test with #25565 applied, the `Utf8` case of this query passes, but the
`Utf8View` case still fails. The problem is also not specific to the sort:
every spill file that is written in batches of `batch_size` rows has it.
## What changes are included in this PR?
The bound is in the central spill writer, so each operator can use it:
- `SpillManager::with_max_batch_bytes(Option<usize>)`. When it is set,
`InProgressSpillFile::append_batch` and `append_batch_async` split a larger
batch into row ranges by recursive halving, and write each range as its own IPC
message. Row order is kept. One appended batch can come back as several batches
when the file is read.
- The split uses the bytes that a row range points at, not the size of the
buffers it keeps alive. A slice of a large view array otherwise reports the
whole parent, and no split seems to help. View arrays count 16 bytes per row
plus each value that is not inline. Dictionaries count the keys plus the values
that the keys use.
- Each piece is compacted on its own: view arrays with `gc_view_arrays` as
before, and, when a batch is split, dictionaries drop the values that the piece
does not use. A piece that compacts to more than the estimate (a view builder
allocates its blocks with spare capacity) is halved again.
- The split stops at one row, or when a split does not divide the payload.
That piece is written as it is.
- `append_batch` returns the size of the largest piece, so
`max_record_batch_memory` records the largest batch that a reader decodes. The
reader is unchanged.
- `ExternalSorter` opts in, with a bound of `sort_spill_reservation_bytes`,
which is the headroom that the final merge starts with. No other operator
changes its behavior.
Aggregate spilling can opt in with the same call. I am checking how much of
#25853 that would cover and will report the results here.
## Are these changes tested?
- `datafusion/core/tests/memory_limit/wide_row_sort_spill.rs`: the query
above with a `Utf8` and a `Utf8View` payload. Both fail on `main` with
`ResourcesExhausted` and pass with this change.
- Unit tests in `in_progress_spill_file.rs`:
- slices of one large `Utf8View` array are split, and each piece is within
the bound;
- repeated view values are split, and row order is kept;
- a dictionary with large values is divided between the pieces, and the
file passes the reader's size check;
- one row that is larger than the bound is written as it is, and is
recorded correctly;
- a `SpillManager` without a bound writes each batch as one batch, as
before.
## Are there any user-facing changes?
External sorts over wide rows can finish under memory limits where they
failed before. Sort spill files can have more, smaller batches. There is a new
public method, `SpillManager::with_max_batch_bytes`. There are no configuration
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]