Skip to content

feat: support releasing batches in RecordBatchMemoryCounter - #26149

Open
txwyy123 wants to merge 2 commits into
apache:mainfrom
txwyy123:support-uncount-batch
Open

txwyy123 wants to merge 2 commits into
apache:mainfrom
txwyy123:support-uncount-batch

Conversation

@txwyy123

@txwyy123 txwyy123 commented Oct 9, 2026

Copy link
Copy Markdown

Which issue does this PR close?

Closes #26140.

Rationale for this change

RecordBatchMemoryCounter can only count (add) batches but never release them. Operators that retain batches incrementally and drop them later (sort, window, sort-merge join, TopK) can't use it -- they keep their own estimates instead, leading to over-counting (sort: 1488 MB reported for 346 MB actual) or no counting at all (window: 200 MB limit ignored, 1.27 GB RSS).

What changes are included in this PR?

  • BufferIdSet → BufferIdMap: internal storage changed from an insert-only set to a reference-counted map ((NonZero<usize>, u32) tuples inline, HashMap<NonZero<usize>, u32> overflow). Inline fast path for ≤16 buffers preserved.
  • New public API: uncount_batch(), uncount_array(), uncount_batch_with_array_overhead() -- each returns the bytes released.
  • Shared visitor: Extracted visit_array_buffers(array, op) with BufferOp::Count | Uncount direction, replacing duplicated per-type buffer walks. Same type dispatch (primitives, boolean, binary, utf8, views, lists, list-views, fixed-size, struct, union, dictionary, map, run-end-encoded, fallback).
  • 52 tests (22 pre-existing + 30 new) covering every acceptance criterion.

Are these changes tested?

Yes -- comprehensive test suite covering all acceptance criteria:

Criterion Tests
Count/uncount round trip test_uncount_batch_round_trip, test_uncount_array_round_trip, test_uncount_batch_with_array_overhead_round_trip
Two-slice example from issue test_uncount_two_slices_example_from_issue
View arrays sharing data buffers test_uncount_view_array_slices_sharing_data_buffers, test_uncount_binary_view_array_shared_data_buffers
Dictionaries sharing values test_uncount_dictionaries_sharing_values
Nested types struct, list (shared child), map (shared children), union, union slices, run-end-encoded, run-end slices, fixed-size-binary, fixed-size-list, deeply nested (List<Struct<Dict>>), list-view, large-list-view
>16 buffers + removals test_uncount_with_overflow_promotion
Randomized sequence vs reference model 200-op sequence across 10 array types (Int32, Int64, String, StringView, Float64, Dict, List, Boolean + slices)
Edge cases null bitmaps, empty arrays, empty batches, NullArray, boolean (bit-packed), double-uncount (no underflow), never-counted no-op, array-overhead shared across batches

Are there any user-facing changes?

No. Purely additive API -- existing count_* methods and get_record_batch_memory_size return identical results. No caller changes needed.

Benchmark

record_batch_memory benchmark shows no regression for counting:

column_count/1       22 ns
column_count/64      2.1 µs
shared_slices/4      1.8 µs
shared_slices/64     30.1 µs
count_uncount/4      3.7 µs  (new -- full count+uncount cycle, 32 slices × 4 columns)

Add uncount_batch, uncount_array, and uncount_batch_with_array_overhead
methods to RecordBatchMemoryCounter. The internal BufferIdSet is changed
to a reference-counted BufferIdMap so that a buffer's capacity is
subtracted from memory_usage only when its refcount drops to zero.

The per-type buffer walk is shared between counting and uncounting via
a visit_array_buffers visitor, avoiding code duplication.

Comprehensive test suite covers all acceptance criteria from apache#26140:
- Count/uncount round trip (batch, array, batch_with_array_overhead)
- Two-slice shared-buffer example from the issue
- View arrays sharing data buffers (StringView slices, BinaryView)
- Dictionaries sharing values across arrays
- All nested types: List, Map, Union, RunEndEncoded, FixedSizeList,
  FixedSizeBinary, deeply nested (List<Struct<Dict>>)
- Slice-sharing for Union and RunEndEncoded
- Null bitmaps, empty arrays, empty batches, NullArray
- Binary/LargeBinary/Utf8/LargeUtf8 offset+data buffer pairs
- String slices sharing offset+data buffers
- ListView and LargeListView
- Boolean arrays (bit-packed)
- >16 distinct buffers with inline-to-hash promotion + removals
- count_batch_with_array_overhead shared across batches
- Double-uncount idempotency (no underflow)
- Randomized 200-op sequence across 10 diverse array types
  (Int32, Int64, String, StringView, Float64, Dictionary, List,
  Boolean, plus slices) verified against a reference model
- record_batch_memory benchmark with count_uncount case, no regression

Closes apache#26140
@github-actions github-actions Bot added the common Related to common crate label Oct 9, 2026
@jayzhan211
jayzhan211 self-requested a review October 10, 2026 01:02

@jayzhan211 jayzhan211 left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks @txwyy123 , here are some suggestions


/// Inverse of [`count_unique_array_object_memory_size`]: removes array
/// identities from the tracked set and returns the released overhead.
fn uncount_unique_array_object_memory_size(

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

counted_arrays is still a set, so array overhead is released on the first uncount even when another counted batch holds the same ArrayRef (batch.clone(), dictionary values() shared across reader batches). Buffers use refcounts; overhead should too.

Repro (fails on this PR: memory_usage() is 12 after the first uncount, expected 108):

#[test]
fn overhead_released_while_still_held() {
    let col: ArrayRef = Arc::new(Int32Array::from(vec![1, 2, 3]));
    let schema = Arc::new(Schema::new(vec![Field::new("v", DataType::Int32, false)]));
    let b1 = RecordBatch::try_new(Arc::clone(&schema), vec![Arc::clone(&col)]).unwrap();
    let b2 = b1.clone();
    let mut counter = RecordBatchMemoryCounter::new();
    let full = counter.count_batch_with_array_overhead(&b1);
    counter.count_batch_with_array_overhead(&b2);
    assert_eq!(counter.uncount_batch_with_array_overhead(&b1), 0);
    assert_eq!(counter.memory_usage(), full);
    assert_eq!(counter.uncount_batch_with_array_overhead(&b2), full);
}

Fix: change counted_arrays to HashMap<usize, u32> (crate::HashMap) and only recurse on 0→1 / 1→0, which keeps count and uncount symmetric:

-    if !counted_arrays.insert(array_ptr) {
-        return 0;
-    }
+    let count = counted_arrays.entry(array_ptr).or_insert(0);
+    *count += 1;
+    if *count > 1 {
+        return 0;
+    }
-    if !counted_arrays.remove(&array_ptr) {
-        return 0;
-    }
+    let Entry::Occupied(mut entry) = counted_arrays.entry(array_ptr) else {
+        return 0;
+    };
+    *entry.get_mut() -= 1;
+    if *entry.get() > 0 {
+        return 0;
+    }
+    entry.remove();

(use hashbrown::hash_map::Entry;.) test_uncount_batch_with_array_overhead_shared_across_is because .slice()creates newArrayRef`s; please add the test above.

inline: [Option<NonZero<usize>>; INLINE_BUFFER_IDS],
struct BufferIdMap {
/// Inline storage: `(address, count)` pairs.
inline: [(NonZero<usize>, u32); INLINE_BUFFER_IDS],

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Counting regresses against main, which #26140 asks to avoid. record_batch_memory, interleaved, 2 runs each:

case main PR
column_count/4 24.9 ns 27.4 ns (+10%)
column_count/16 120.2 ns 129.7 ns (+8%)
array_layout/list 58.1 ns 65.0 ns (+12%)
shared_slices/4 693 ns 805 ns (+16%)
shared_slices/16 3.64 µs 4.06 µs (+12%)

The padded [(NonZero<usize>, u32); 16] doubles the inline array (128 → 256 B), and get_record_batch_memory_size builds a new counter for every batch. Separate arrays plus position() bring the 16-buffer cases back to main, but the 4-buffer cases stay about 10% slower, so more is needed:

struct BufferIdMap {
    ids: [Option<NonZero<usize>>; INLINE_BUFFER_IDS],
    counts: [u32; INLINE_BUFFER_IDS],
    len: usize,
    overflow: Option<hashbrown::HashMap<NonZero<usize>, u32>>,
}

Please also add main vs PR numbers to the description, which the issue lists as an acceptance item

// ---- uncount tests ----

#[test]
fn test_uncount_batch_round_trip() {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Could you trim the tests? They're ~900 of the added lines, which makes the PR hard to review, and most per-type round trips repeat coverage the file already has. test_array_memory_size_matches_array_data_layouts already loops over every layout (view, run-end, union, map, dictionary, ...). Adding a round trip to its helper covers all of them:

 fn assert_array_memory_size_matches(array: &dyn Array) {
     let mut counter = RecordBatchMemoryCounter::new();
-    counter.visit_array_buffers(array, BufferOp::Count);
-    assert_eq!(counter.memory_usage(), array_data_memory_size(array));
+    let counted = counter.count_array(array);
+    assert_eq!(counted, array_data_memory_size(array));
+    assert_eq!(counter.uncount_array(array), counted);
+    assert_eq!(counter.memory_usage(), 0);
 }

With that, the standalone round trips can go: test_uncount_{batch,array}_round_trip, _nested_struct, _list_view_array, _large_list_view_array, _null_array, _boolean_array, _binary_array, _large_utf8_array, _fixed_size_binary, _fixed_size_list, _array_with_nulls, _string_array_with_nulls, _union_array, _run_end_encoded_array, _deeply_nested, _empty_array, _empty_batch.

The shared-buffer tests (two slices, view / dictionary / list / map / union / run-end / string sharing) could become one table of (first, second) arrays that share buffers. Each row would check: count both → uncount first releases nothing → uncount second releases everything. test_uncount_binary_view_array_shared_data_buffers duplicates the StringView case. Please keep the randomized model test and the overflow-promotion, never-counted, double-uncount and overhead tests.

@codecov-commenter

Copy link
Copy Markdown

Codecov Report

❌ Patch coverage is 96.68790% with 26 lines in your changes missing coverage. Please review.
✅ Project coverage is 82.79%. Comparing base (a23b89f) to head (db5cfd0).
⚠️ Report is 19 commits behind head on main.

Files with missing lines Patch % Lines
datafusion/common/src/utils/memory.rs 96.68% 22 Missing and 4 partials ⚠️
Additional details and impacted files
@@            Coverage Diff             @@
##             main   #26149      +/-   ##
==========================================
+ Coverage   82.75%   82.79%   +0.03%     
==========================================
  Files        1147     1147              
  Lines      449881   451295    +1414     
  Branches   449881   451295    +1414     
==========================================
+ Hits       372306   373631    +1325     
- Misses      54905    54949      +44     
- Partials    22670    22715      +45     

☔ View full report in Codecov by Harness.
📢 Have feedback on the report? Share it here.

🚀 New features to boost your workflow:
  • ❄️ Test Analytics: Detect flaky tests, report on failures, and find test suite problems.
  • 📦 JS Bundle Analysis: Save yourself from yourself by tracking and limiting bundle sizes in JS merges.

- Run cargo fmt to fix line-width formatting in tests
- Use #[expect(dead_code)] instead of #[allow(dead_code)] per project clippy config
- Collapse nested if-let into single condition to satisfy collapsible_if lint
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

common Related to common crate

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Support releasing batches in RecordBatchMemoryCounter

3 participants