andygrove opened a new issue, #25804:
URL: https://github.com/apache/datafusion/issues/25804
### 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` (f029b944c).
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 #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`.
<details><summary>Test code</summary>
```rust
/// 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(())
}
}
```
</details>
### 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:
- [ ] Backport #24740 to `branch-55` (findings 1 and 2)
- [ ] Backport the reservation fix from #24923 to `branch-55` (finding 3).
#24923 itself adds an async spill-writing API, so only that change is
backported.
- [ ] Review #25372 (finding 4) and #25800 (finding 6)
- [ ] Sort-side fix for finding 5: count shared buffers once across the
buffered batches, as #22862 did for the hash join build side
- [ ] Bound the sort's final spill-merge reservation (finding 9), as #25383
does for aggregate replay
- [ ] Per-consumer accounting in `FairSpillPool` (finding 10)
Related: #22758, #25758
--
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]