tohuya6 commented on code in PR #25800:
URL: https://github.com/apache/datafusion/pull/25800#discussion_r4116998859
##########
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:
Thanks @andygrove! I switched to reverse counting and removed
`ReservationStream` rather than restoring main's version, since nothing uses it
now and clippy flags it as dead code.
The new test also checks the shared buffers stay reserved after the first
chunk, and I listed the other `get_sliced_size` callers in the description.
--
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]