Repository navigation
Support releasing batches in RecordBatchMemoryCounter #26140
Description
Activity
- addedenhancementNew feature or requestNew feature or requesthelp wantedExtra attention is neededExtra attention is needed
on Oct 8, 2026 Hi @jayzhan211, I would like to work on this issue. Is anyone already working on, or planning, the release API?
My proposed scope is limited to RecordBatchMemoryCounter, regression tests, and the record_batch_memory benchmark, with no operator or caller changes. I would preserve the existing count_* results and inline buffer-tracking fast path, share traversal between count and uncount, and cover shared buffers, Array-object overhead, nested/view/dictionary arrays, inline-to-map promotion, and balanced randomized count/uncount sequences.
Happy to coordinate if there is already work in progress.
I'm not sure, you can go ahead
@jayzhan211 I'd like to take this. I have a working implementation with all the acceptance criteria tests passing — PR: #26147
- added a commit that references this issue
on Oct 9, 2026 A concrete data point for TopK, one of the operators this issue lists as not able to use the counter yet.
TopK charges each retained input batch its
get_record_batch_memory_size(RecordBatchStore::insert), so N zero-copy slices of one allocation are charged N × that allocation. You can hit this without any slicing inside DataFusion:arrow_ipc::reader::StreamDecoderdecodes zero-copy, so every record batch of a multi-batch IPC stream points into the one stream buffer.Repro on 53.1.0: a 128 MiB pool, and 2 000 rows passed as 2 000 one-row slices of one batch:
use std::sync::Arc; use datafusion::arrow::array::{Int64Array, RecordBatch, StringArray}; use datafusion::arrow::datatypes::{DataType, Field, Schema}; use datafusion::datasource::MemTable; use datafusion::execution::runtime_env::RuntimeEnvBuilder; use datafusion::prelude::{SessionConfig, SessionContext}; #[tokio::main(flavor = "current_thread")] async fn main() -> datafusion::error::Result<()> { // 2 000 rows, passed as 2 000 one-row zero-copy slices of one batch let n = 2_000; let schema = Arc::new(Schema::new(vec![ Field::new("k", DataType::Int64, false), Field::new("v", DataType::Utf8, false), ])); let batch = RecordBatch::try_new( Arc::clone(&schema), vec![ Arc::new(Int64Array::from_iter_values(0..n as i64)), Arc::new(StringArray::from_iter_values((0..n).map(|i| format!("{i:0>100}")))), ], )?; let slices: Vec<RecordBatch> = (0..n).map(|i| batch.slice(i, 1)).collect(); let runtime = RuntimeEnvBuilder::new() .with_memory_limit(128 * 1024 * 1024, 1.0) .build_arc()?; let ctx = SessionContext::new_with_config_rt(SessionConfig::new(), runtime); ctx.register_table("t", Arc::new(MemTable::try_new(schema, vec![slices])?))?; // k = 100 000 > n: every slice keeps a row, so TopK retains all of them ctx.sql("SELECT k, v FROM t ORDER BY k DESC LIMIT 100000") .await? .collect() .await?; Ok(()) }
ResourcesExhausted("Additional allocation failed for TopK[0] with top memory consumers (across reservations) as: TopK[0]#0(can spill: false) consumed 127.8 MB, peak 127.8 MB. Error: Failed to allocate additional 279.5 KB for TopK[0] with 127.8 MB already allocated for this reservation - 243.8 KB remain available for the total pool")Each one-row slice is charged the parent's full ~280 KB of buffer capacity. I haven't run this on
main, but the path looks unchanged there:register_batchkeeps the input batch as-is.RecordBatchStore::insertaddsget_record_batch_memory_sizeper entry.maybe_compactonly runs once unused rows reach20 * batch_size + k, so with a largekit never compacts.
We ran into this on stored IPC streams of hundreds of one-row batches. An
ORDER BYover 577 rows with a largeLIMITwas charged ~0.75 MB per batch and failed at 128 MiB. On 53.1 the hash-join build side failed the same way: a join of two ~1 000-row tables was charged ~0.6 MB per one-row batch. #22862 already fixes that onmain. Our workaround is toconcat_batcheseach decoded IPC stream before handing it to DataFusion.
Is your feature request related to a problem or challenge?
Operators that keep
RecordBatches in memory across polls (sort, window, joins, ...) report that memory to the memory pool in inconsistent ways. Two examples on currentmain:Over-counting. A sort over aggregate output charges every slice the size of its whole parent buffer:
The query needs ~346 MB without a limit.
EXPLAIN ANALYZEshows the finalAggregateExecreportingoutput_bytes=1488.0 MBagainst45.8 MBfor theSortExecabove it, for the same rows (Hash aggregation produces batches reporting huge memory size #22526).Not counted. Window operators hold no reservation at all:
RecordBatchMemoryCounter(#22862) already solves the "count each shared buffer once" part, but it can only add: there is no way to stop counting a batch. So it only fits operators that build once and keep everything until the end (hash join build side,AsofJoin). Operators that retain batches incrementally and drop them later (sort, window, sort-merge join, TopK) can't use it, and keep their own estimates instead.Describe the solution you'd like
Let
RecordBatchMemoryCounterrelease batches, as Samyak2 suggested in the #22862 review: track a reference count per buffer instead of a set of seen buffers.uncount_batch(&mut self, batch: &RecordBatch) -> usize(plusuncount_array, anduncount_batch_with_array_overheadas the inverse ofcount_batch_with_array_overhead) decrements the counts and returns the bytes released: the sizes of buffers whose count reaches 0.Example: two zero-copy slices of one 4 MB batch.
memory_usage()count_batch(slice1)count_batch(slice2)uncount_batch(slice1)uncount_batch(slice2)Implementation notes:
count_*methods andget_record_batch_memory_sizemust return exactly what they return today, and no caller changes are needed.ArrayDatafallback) should be shared by counting and uncounting rather than duplicated.record_batch_memorybenchmark should show no regression for counting.Acceptance criteria:
record_batch_memorybenchmark before/after in the PR description, plus a case for count + uncount.Describe alternatives you've considered
Arrow's
claim()API (#22898) tracks buffers inside Arrow itself, but it doesn't fit per-operator budgets: reservations cannot fail, the last claimer is charged, and enabling thepoolfeature adds a mutex to every buffer in the process.Additional context
Part of #22758. This is the first step towards a per-operator container that owns the batches an operator retains and keeps its
MemoryReservationequal to the unique buffers it holds; window operators andExternalSorterwould be the first users.