jayzhan211 opened a new issue, #26140:
URL: https://github.com/apache/datafusion/issues/26140

   ### Is your feature request related to a problem or challenge?
   
   Operators that keep `RecordBatch`es in memory across polls (sort, window, 
joins, ...) report that memory to the memory pool in inconsistent ways. Two 
examples on current `main`:
   
   - **Over-counting.** A sort over aggregate output charges every slice the 
size of its whole parent buffer:
   
     ```sql
     -- datafusion-cli -m 700M
     SELECT k, c FROM (
       SELECT value % 3000000 AS k, count(*) AS c
       FROM generate_series(1, 6000000) GROUP BY k
     ) ORDER BY c, k;
     -- Resources exhausted: Additional allocation failed for 
ExternalSorterMerge[..]
     ```
   
     The query needs ~346 MB without a limit. `EXPLAIN ANALYZE` shows the final 
`AggregateExec` reporting `output_bytes=1488.0 MB` against `45.8 MB` for the 
`SortExec` above it, for the same rows (#22526).
   
   - **Not counted.** Window operators hold no reservation at all:
   
     ```sql
     -- datafusion-cli -m 200M
     SELECT count(*), max(s) FROM (
       SELECT sum(value) OVER () AS s FROM generate_series(1, 50000000)
     );
     -- succeeds with ~1.27 GB peak RSS
     ```
   
   `RecordBatchMemoryCounter` (#22862) already solves the "count each shared 
buffer once" part, but it can only add: there is no way to stop counting a 
batch. So it only fits operators that build once and keep everything until the 
end (hash join build side, `AsofJoin`). Operators that retain batches 
incrementally and drop them later (sort, window, sort-merge join, TopK) can't 
use it, and keep their own estimates instead.
   
   ### Describe the solution you'd like
   
   Let `RecordBatchMemoryCounter` release batches, as Samyak2 suggested in the 
#22862 review: track a reference count per buffer instead of a set of seen 
buffers.
   
   - Counting a batch increments the count of each of its buffers; a buffer's 
size is added only when its count goes from 0 to 1 (today's behaviour).
   - New `uncount_batch(&mut self, batch: &RecordBatch) -> usize` (plus 
`uncount_array`, and `uncount_batch_with_array_overhead` as the inverse of 
`count_batch_with_array_overhead`) decrements the counts and returns the bytes 
released: the sizes of buffers whose count reaches 0.
   - Document the contract: the counter identifies buffers by address and does 
not keep them alive, so a batch must be uncounted before its buffers are 
dropped.
   
   Example: two zero-copy slices of one 4 MB batch.
   
   | step | `memory_usage()` |
   |---|---|
   | `count_batch(slice1)` | 4 MB (parent buffer counted) |
   | `count_batch(slice2)` | 4 MB (already counted) |
   | `uncount_batch(slice1)` | 4 MB (slice2 still uses it) |
   | `uncount_batch(slice2)` | 0 |
   
   Implementation notes:
   
   - This should be purely additive: existing `count_*` methods and 
`get_record_batch_memory_size` must return exactly what they return today, and 
no caller changes are needed.
   - The per-type buffer walk (null buffers, offsets, view buffers including 
variadic data buffers, dictionaries, nested children, the `ArrayData` fallback) 
should be shared by counting and uncounting rather than duplicated.
   - Keep the allocation-free fast path for typical batches (see #24310 / 
#24319); the `record_batch_memory` benchmark should show no regression for 
counting.
   
   Acceptance criteria:
   
   - Unit tests: count/uncount round trip; the two-slice example above; view 
arrays sharing data buffers; dictionaries sharing values; nested types; more 
than 16 distinct buffers (the inline-to-hash-set promotion) with removals; a 
randomized count/uncount sequence checked against a simple reference model.
   - `record_batch_memory` benchmark before/after in the PR description, plus a 
case for count + uncount.
   
   ### Describe alternatives you've considered
   
   Arrow's `claim()` API (#22898) tracks buffers inside Arrow itself, but it 
doesn't fit per-operator budgets: reservations cannot fail, the last claimer is 
charged, and enabling the `pool` feature adds a mutex to every buffer in the 
process.
   
   ### Additional context
   
   Part of #22758. This is the first step towards a per-operator container that 
owns the batches an operator retains and keeps its `MemoryReservation` equal to 
the unique buffers it holds; window operators and `ExternalSorter` would be the 
first users.
   


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