Skip to content

ExternalSorter memory accounting: fixes missing from branch-55 and remaining gaps on main #25804

Description

@andygrove

Describe the bug

I audited how SortExec / ExternalSorter, and the merge and spill code it uses (StreamingMergeBuilder, MultiLevelMergeBuilder, BatchBuilder, the cursor streams), interact with the MemoryPool. I checked both branch-55 (55.1.0) and main (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 to Additional 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 through BatchBuilder::push_batch inside ExternalSorter::sort_and_spill_in_mem_batches, which is finding 1 below.

Several fixes for this landed on main after branch-55 was cut, and none of them are on branch-55. Some gaps remain on main with no tracking issue.

Findings

# Finding 55.1.0 main Tracking
1 sort_and_spill_in_mem_batches frees the sort_spill_reservation_bytes headroom to the pool and then merges with a new_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. Present Spill path fixed by #24740. The final in-memory merge, and the concat path during a spill, still release the workspace.
2 MultiLevelMergeBuilder frees the headroom it received from sort() at the end of each intermediate pass (StreamAttachedReservation) and on the read-ahead fallback in get_sorted_spill_files_to_merge, so only the first pass is protected against contention. Present (repro below) Fixed by #24740
3 consume_and_spill_append frees the reservation before it writes the sorted batches, so they are unaccounted while the spill file is written. Present Fixed as part of #24923
4 ReusableRows keeps each input stream's Rows after the cursor, and with it the cursor's reservation, is dropped. The buffer stays allocated until the merge finishes. Present (two buffers per stream) One buffer per stream since #23802, still unreserved #25372, #23760
5 reserve_memory_for_batch_and_maybe_spill charges each zero-copy slice of a parent batch (for example AggregateExec EmitTo::All output) the parent's full buffer capacity. Present (repro below) Present; not changed by #25800 #22526 reports it from the aggregate side; no sort-side fix
6 sort_batch_stream sums get_record_batch_memory_size over the sorted chunks, re-counting dictionary values and view buffers that take shares between chunks. A 1.3 MB dictionary batch asks for 62 MB. Present Present #25791, #25800 (also fixes the dictionary case)
7 FieldCursorStream::convert_batch reserves the sort-key column that BatchBuilder::push_batch has already charged as part of the batch, so single-column merges count the key twice. Present Present Mentioned in #25271
8 Spill-merge admission charges 2 × largest batch × read-ahead per run. That leaves out the cursor rows and the prefetched batch, while the merge itself runs against an unbounded pool. Present Present #23760, partly #25565
9 The final spill merge grows its reservation until the pool refuses (fan-in is unlimited by default), leaving nothing for downstream consumers in the same pool. Present Present for sort; #25383 added replay headroom for aggregate #19216, #20715
10 FairSpillPool::try_grow checks a spillable request against the fair share using only that reservation's size, not the consumer's total or the pool total, while ExternalSorter holds several sibling reservations (split, take, new_empty). Present Present #25172 (closed)

To Reproduce

The tests below go in a module at the end of datafusion/physical-plan/src/sorts/sort.rs, since ExternalSorter is private. Both pass on branch-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. On branch-55 the second merge pass fails with Failed 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 MiB GreedyMemoryPool and spills 21 times. It does the same on main.
Test code
/// 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(())
    }
}

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:

Related: #22758, #25758

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

Labels

bugSomething isn't working

Type

No type

Projects

No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions