andygrove commented on code in PR #25820:
URL: https://github.com/apache/datafusion/pull/25820#discussion_r4135921594


##########
datafusion/physical-plan/src/sorts/sort.rs:
##########
@@ -946,6 +1030,58 @@ pub(crate) fn get_reserved_bytes_for_record_batch(batch: 
&RecordBatch) -> Result
     })
 }
 
+/// Estimate how much memory is needed to sort `batch` when it is buffered with
+/// the batches `counter` has already counted.
+///
+/// Like [`get_reserved_bytes_for_record_batch`], but only counts the buffers 
of
+/// `batch` that `counter` has not seen. Buffered batches can share buffers, 
for
+/// example zero-copy slices of one larger batch, and a shared buffer only 
needs
+/// to be reserved once. Without a `counter`, every buffer of `batch` is 
counted.
+fn get_reserved_bytes_for_next_record_batch(
+    batch: &RecordBatch,
+    counter: Option<&mut RecordBatchMemoryCounter>,
+) -> Result<usize> {
+    let Some(counter) = counter else {
+        return get_reserved_bytes_for_record_batch(batch);
+    };
+    let sliced_size = batch.get_sliced_size()?;
+    Ok(get_reserved_bytes_for_record_batch_size(
+        counter.count_batch(batch),
+        sliced_size,
+    ))
+}
+
+/// Returns whether the sorted copy of an array of `data_type` keeps
+/// referencing some of the array's buffers.
+///
+/// Sorting takes the selected rows into new buffers, except that `take` keeps
+/// the data buffers of view arrays and the values of dictionaries and list
+/// views. Every sorted run is charged for those buffers in full, so buffered
+/// batches that share them must be charged for them each, too.
+fn sorted_copy_keeps_buffers(data_type: &DataType) -> bool {

Review Comment:
   It's there for correctness rather than speed. For view, dictionary and list 
view columns, `take` keeps referencing the input's data buffers or values, so 
each sorted run, and the merge, is charged for those buffers in full. If the 
sorter counted them once when buffering, it would buffer more batches than the 
merge of their sorted runs can hold. With the check removed, 
`sort_spill_is_compacted_for_view_arrays` fails with `Failed to allocate 
additional 234.0 KB for ExternalSorterMerge[0]`. It can go once the sorted runs 
and the merge count shared buffers once too (#25800, #25791). I added the 
reason to its doc comment in a8d14f5.
   



##########
datafusion/physical-plan/src/sorts/sort.rs:
##########
@@ -3763,6 +3899,60 @@ mod tests {
         Ok(())
     }
 
+    /// Zero-copy slices of one batch share its buffers, so buffering them
+    /// reserves those buffers once. Sorting the slices as separate runs must
+    /// keep the shared buffers reserved while any run still holds them.
+    #[tokio::test]
+    async fn test_sliced_runs_reserve_shared_buffers_once() -> Result<()> {
+        let schema = Arc::new(Schema::new(vec![Field::new("x", 
DataType::Int64, false)]));
+        let parent = RecordBatch::try_new(
+            Arc::clone(&schema),
+            vec![Arc::new(Int64Array::from_iter_values((0..4096).rev()))],
+        )?;
+        let slices: Vec<_> = (0..4).map(|i| parent.slice(i * 1024, 
1024)).collect();
+
+        let pool: Arc<dyn MemoryPool> = Arc::new(GreedyMemoryPool::new(1024 * 
1024));
+        let runtime = RuntimeEnvBuilder::new()
+            .with_memory_pool(Arc::clone(&pool))
+            .build_arc()?;
+        let ordering: LexOrdering =
+            [PhysicalSortExpr::new_default(Arc::new(Column::new("x", 
0)))].into();
+        let mut sorter = ExternalSorter::new(
+            0,
+            Arc::clone(&schema),
+            ordering.clone(),
+            1024,
+            0,
+            0, // Sort each slice as its own run, rather than concatenating 
them.
+            SpillCompression::Uncompressed,
+            &ExecutionPlanMetricsSet::new(),
+            runtime,
+        )?;
+        for slice in &slices {
+            sorter.insert_batch(slice.clone()).await?;
+        }
+        let parent_bytes = get_record_batch_memory_size(&parent);
+        let sliced_bytes = slices
+            .iter()
+            .map(|slice| slice.get_sliced_size())
+            .sum::<Result<usize>>()?;
+        assert_eq!(pool.reserved(), parent_bytes + sliced_bytes);
+
+        let runs = std::mem::take(&mut sorter.in_mem_batches);
+        let mut streams = sorter.sort_run_streams(runs)?;
+        let last = streams.pop().unwrap();
+        for stream in streams {
+            stream.try_collect::<Vec<_>>().await?;
+        }
+        // The last run still holds the parent's buffers until it is sorted
+        assert_eq!(pool.reserved(), parent_bytes + 
slices[3].get_sliced_size()?);

Review Comment:
   Removed in a8d14f5. `sort_after_group_by_reserves_shared_output_once` 
already runs the query under a memory limit, and each of its two orderings 
fails on `main` by itself, so it covers both the coalesced runs and one run per 
slice. The unit test also checked that the shared buffers stay reserved until 
the last run holding them is sorted. A query can't observe that, because 
reserving too little doesn't make it fail, so that part is no longer tested.
   



##########
datafusion/physical-plan/src/sorts/sort.rs:
##########
@@ -865,7 +941,10 @@ impl ExternalSorter {
         &mut self,
         input: &RecordBatch,
     ) -> Result<()> {
-        let size = get_reserved_bytes_for_record_batch(input)?;
+        let size = get_reserved_bytes_for_next_record_batch(

Review Comment:
   Agreed, resizing to a running total would be easier to follow than growing 
by a delta. The counter already keeps a running total of the unique buffers, so 
most of the state is there. I'll try it in a follow-up.
   



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