/// Reproductions for the ExternalSorter memory-pool audit. Each test asserts
/// the CURRENT behavior, so it passes today and documents the defect.
#[cfg(test)]
mod memory_audit_repro {
use super::*;
use crate::metrics::ExecutionPlanMetricsSet;
use arrow::array::{Int32Array, Int64Array};
use arrow::datatypes::{DataType, Field, Schema};
use datafusion_execution::memory_pool::{
GreedyMemoryPool, MemoryLimit, MemoryPool,
};
use datafusion_execution::runtime_env::RuntimeEnvBuilder;
use datafusion_physical_expr::expressions::Column;
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
fn new_sorter(
schema: &SchemaRef,
pool: &Arc<dyn MemoryPool>,
batch_size: usize,
sort_spill_reservation_bytes: usize,
sort_in_place_threshold_bytes: usize,
) -> Result<ExternalSorter> {
let runtime = RuntimeEnvBuilder::new()
.with_memory_pool(Arc::clone(pool))
.build_arc()?;
ExternalSorter::new(
0,
Arc::clone(schema),
[PhysicalSortExpr::new_default(Arc::new(Column::new("x", 0)))].into(),
batch_size,
sort_spill_reservation_bytes,
sort_in_place_threshold_bytes,
SpillCompression::Uncompressed,
&ExecutionPlanMetricsSet::new(),
runtime,
)
}
/// Zero-copy slices of one parent batch (what `AggregateExec` emits for
/// `EmitTo::All`) are each charged the parent's full buffer capacity.
#[tokio::test]
async fn sliced_input_is_charged_full_parent_per_slice() -> Result<()> {
let schema = Arc::new(Schema::new(vec![Field::new("x", DataType::Int64, false)]));
let n = 64 * 1024;
let parent = RecordBatch::try_new(
Arc::clone(&schema),
vec![Arc::new(Int64Array::from_iter_values((0..n as i64).rev()))],
)?;
let parent_bytes = get_record_batch_memory_size(&parent); // 512 KiB
// 4x the real data: enough to sort everything in memory.
let pool: Arc<dyn MemoryPool> = Arc::new(GreedyMemoryPool::new(4 * parent_bytes));
let mut sorter = new_sorter(&schema, &pool, 1024, 0, 0)?;
for i in 0..64 {
sorter.insert_batch(parent.slice(i * 1024, 1024)).await?;
}
// Each 8 KiB slice was reserved as ~520 KiB, so the sorter spilled
// every 3 slices even though all the data is 512 KiB.
println!(
"parent={parent_bytes} pool={} spills={}",
4 * parent_bytes,
sorter.spill_count()
);
assert!(sorter.spill_count() >= 15);
Ok(())
}
/// Models an aggressive concurrent consumer: once armed, every byte a
/// reservation releases is immediately taken by someone else.
#[derive(Debug)]
struct StealingPool {
inner: GreedyMemoryPool,
armed: AtomicBool,
stolen: AtomicUsize,
}
impl fmt::Display for StealingPool {
fn fmt(&self, f: &mut Formatter<'_>) -> fmt::Result {
write!(f, "stealing({})", self.inner)
}
}
impl MemoryPool for StealingPool {
fn name(&self) -> &str {
"stealing"
}
fn grow(&self, r: &MemoryReservation, n: usize) {
self.inner.grow(r, n)
}
fn shrink(&self, r: &MemoryReservation, n: usize) {
if self.armed.load(Ordering::Relaxed) {
self.stolen.fetch_add(n, Ordering::Relaxed);
} else {
self.inner.shrink(r, n)
}
}
fn try_grow(&self, r: &MemoryReservation, n: usize) -> Result<()> {
self.inner.try_grow(r, n)
}
fn reserved(&self) -> usize {
self.inner.reserved()
}
fn memory_limit(&self) -> MemoryLimit {
self.inner.memory_limit()
}
}
/// The `sort_spill_reservation_bytes` headroom transferred to the merge by
/// `sort()` (#20642) is released to the pool when the first merge pass
/// finishes, so later passes must re-acquire it from the pool.
#[tokio::test]
async fn merge_headroom_is_lost_after_first_pass() -> Result<()> {
// Two spilled runs of 128-row Int32 batches need 2 * 2 KiB per pass.
let headroom = 4 * 1024;
let pool_size = headroom + 40 * 1024;
let stealing = Arc::new(StealingPool {
inner: GreedyMemoryPool::new(pool_size),
armed: AtomicBool::new(false),
stolen: AtomicUsize::new(0),
});
let pool: Arc<dyn MemoryPool> = Arc::clone(&stealing) as _;
let schema = Arc::new(Schema::new(vec![Field::new("x", DataType::Int32, false)]));
let mut sorter = new_sorter(&schema, &pool, 128, headroom, usize::MAX)?;
for i in 0..200 {
let values: Vec<i32> = ((i * 100)..((i + 1) * 100)).rev().collect();
let batch = RecordBatch::try_new(
Arc::clone(&schema),
vec![Arc::new(Int32Array::from(values))],
)?;
sorter.insert_batch(batch).await?;
}
let spill_files = sorter.spill_count();
assert!(spill_files >= 3, "need a multi-pass merge, got {spill_files}");
let merge_stream = sorter.sort().await?;
drop(sorter);
// Same contention as `test_sort_merge_reservation_transferred_not_freed`,
// except the contender keeps taking whatever is released.
let contender = MemoryConsumer::new("CompetingPartition").register(&pool);
contender.try_grow(pool_size - pool.reserved())?;
stealing.armed.store(true, Ordering::Relaxed);
let err = merge_stream
.try_collect::<Vec<_>>()
.await
.expect_err("second merge pass should fail to seat two runs");
println!(
"spill_files={spill_files} stolen={} err={err}",
stealing.stolen.load(Ordering::Relaxed)
);
assert!(err.to_string().contains("ExternalSorterMerge[0]"));
assert_eq!(stealing.stolen.load(Ordering::Relaxed), headroom);
Ok(())
}
}
Describe the bug
I audited how
SortExec/ExternalSorter, and the merge and spill code it uses (StreamingMergeBuilder,MultiLevelMergeBuilder,BatchBuilder, the cursor streams), interact with theMemoryPool. I checked bothbranch-55(55.1.0) andmain(f029b94).DataFusion Comet pins 55.1.0 and runs every Spark sort through
ExternalSorter, including both inputs of every sort-merge join. Comet has lost tasks toAdditional allocation failed for ExternalSorterMerge[N]in TPC-H and TPC-DS runs (apache/datafusion-comet#2452, apache/datafusion-comet#5666, apache/datafusion-comet#6128). The backtrace in apache/datafusion-comet#2452 goes throughBatchBuilder::push_batchinsideExternalSorter::sort_and_spill_in_mem_batches, which is finding 1 below.Several fixes for this landed on
mainafterbranch-55was cut, and none of them are onbranch-55. Some gaps remain onmainwith no tracking issue.Findings
mainsort_and_spill_in_mem_batchesfrees thesort_spill_reservation_bytesheadroom to the pool and then merges with anew_empty()reservation, so the merge it depends on has to win that memory back. Re-reserving the headroom after the spill can also fail, with nothing left to spill.sort()does the same before the final in-memory merge.MultiLevelMergeBuilderfrees the headroom it received fromsort()at the end of each intermediate pass (StreamAttachedReservation) and on the read-ahead fallback inget_sorted_spill_files_to_merge, so only the first pass is protected against contention.consume_and_spill_appendfrees the reservation before it writes the sorted batches, so they are unaccounted while the spill file is written.ReusableRowskeeps each input stream'sRowsafter the cursor, and with it the cursor's reservation, is dropped. The buffer stays allocated until the merge finishes.reserve_memory_for_batch_and_maybe_spillcharges each zero-copy slice of a parent batch (for exampleAggregateExecEmitTo::Alloutput) the parent's full buffer capacity.sort_batch_streamsumsget_record_batch_memory_sizeover the sorted chunks, re-counting dictionary values and view buffers thattakeshares between chunks. A 1.3 MB dictionary batch asks for 62 MB.FieldCursorStream::convert_batchreserves the sort-key column thatBatchBuilder::push_batchhas already charged as part of the batch, so single-column merges count the key twice.FairSpillPool::try_growchecks a spillable request against the fair share using only that reservation's size, not the consumer's total or the pool total, whileExternalSorterholds several sibling reservations (split,take,new_empty).To Reproduce
The tests below go in a module at the end of
datafusion/physical-plan/src/sorts/sort.rs, sinceExternalSorteris private. Both pass onbranch-55, which means both bugs are present:merge_headroom_is_lost_after_first_pass(finding 2) uses a pool that hands every released byte to another consumer, standing in for concurrent partitions or Spark tasks. Onbranch-55the second merge pass fails withFailed to allocate additional 1024.0 B for ExternalSorterMerge[0] ... 0.0 B remain available. With fix: preserve external sort workspace across spilling #24740 the merge completes.sliced_input_is_charged_full_parent_per_slice(finding 5) sorts 64 slices of one 512 KiB batch in a 2 MiBGreedyMemoryPooland spills 21 times. It does the same onmain.Test code
Expected behavior
Reservations should match what the sort actually holds. Headroom reserved for a merge should stay with that merge, bytes should stay reserved until their batches are written or dropped, and buffers shared between batches should be counted once.
Additional context
Plan:
branch-55(findings 1 and 2): [branch-55] fix: preserve external sort workspace across spilling (#24740) #25805branch-55(finding 3): [branch-55] fix: keep sort spill reservation until batches are written (#24923) #25806. feat: add an async API for spill file writing #24923 itself adds an async spill-writing API, so only that change is backported.FairSpillPool(finding 10)Related: #22758, #25758