2010YOUY01 commented on code in PR #25820:
URL: https://github.com/apache/datafusion/pull/25820#discussion_r4140264962
##########
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:
I see, thanks for the explanation.
--
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]