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]

Reply via email to