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]

Reply via email to