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]