andygrove commented on code in PR #25800:
URL: https://github.com/apache/datafusion/pull/25800#discussion_r4116873947
##########
datafusion/physical-plan/src/sorts/sort.rs:
##########
@@ -790,50 +791,29 @@ impl ExternalSorter {
let schema = batch.schema();
let expressions = self.expr.clone();
let batch_size = self.batch_size;
- let merge_pool = Arc::clone(&self.merge_pool);
let stream = futures::stream::once(async move {
let schema = batch.schema();
// Sort the batch immediately and get all output batches
let sorted_batches = sort_batch_chunked(&batch, &expressions,
batch_size)?;
- // Chunked output can retain shared buffers in every batch and
- // exceed the input estimate. Borrow only already-reserved spill
- // workspace; any remainder still uses the original sort consumer.
- let total_sorted_size: usize = sorted_batches
- .iter()
- .map(get_record_batch_memory_size)
- .sum();
- let mut workspace =
-
merge_pool.borrow(total_sorted_size.saturating_sub(reservation.size()));
+ // The chunks share the input buffers, so count each buffer once
+ let mut counter = RecordBatchMemoryCounter::new();
Review Comment:
`ReservationStream` releases a shared buffer when the *first* chunk that
references it is emitted. The chunks still in the stream keep that buffer alive
until the last one is emitted. So when the consumer doesn't reserve what it
receives (for example, the sort is the last operator or feeds a projection),
the view data goes unaccounted while the stream drains.
Here the 4096 `Utf8View` rows from the new test are sorted into 4 chunks,
and each chunk is dropped as it arrives:
| after chunk | `main` | this PR | still held by the stream |
|---|---|---|---|
| 1 | 1,572,864 | 49,152 | 557,056 |
| 2 | 1,048,576 | 32,768 | 540,672 |
| 3 | 524,288 | 16,384 | 524,288 |
Counting the chunks in reverse assigns each shared buffer to the last chunk
that references it. The total stays the same, and the reservation matches the
"still held" column:
```rust
let mut counter = RecordBatchMemoryCounter::new();
let mut sizes: Vec<usize> = sorted_batches
.iter()
.rev()
.map(|batch| counter.count_batch(batch))
.collect();
sizes.reverse();
reservation
.try_resize(counter.memory_usage())
.map_err(Self::err_with_oom_context)?;
let batches = sorted_batches
.into_iter()
.zip(sizes)
.map(move |(batch, size)| {
reservation.shrink(size);
Ok(batch)
});
```
With that, `ReservationStream` has no production user left, so it could go
back to `main`'s version or be removed. I ran this version against the
sorts/stream/spill unit tests, the `memory_limit` tests and the sort/spill fuzz
tests, and all of them pass, including your new test.
One trade-off: when the consumer is the in-memory merge, `BatchBuilder` also
charges each batch it holds in full. A shared buffer is then counted by both
until the stream's last chunk. That's still less than `main`, which counted it
once per remaining chunk, but it could change the pass rates in your 2M-row
repro.
--
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]