Repository navigation
Conversation
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
jayzhan211
left a comment
There was a problem hiding this comment.
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( |
There was a problem hiding this comment.
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], |
There was a problem hiding this comment.
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() { |
There was a problem hiding this comment.
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 Report❌ Patch coverage is
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. 🚀 New features to boost your workflow:
|
- 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
Which issue does this PR close?
Closes #26140.
Rationale for this change
RecordBatchMemoryCountercan 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.uncount_batch(),uncount_array(),uncount_batch_with_array_overhead()-- each returns the bytes released.visit_array_buffers(array, op)withBufferOp::Count | Uncountdirection, 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).Are these changes tested?
Yes -- comprehensive test suite covering all acceptance criteria:
test_uncount_batch_round_trip,test_uncount_array_round_trip,test_uncount_batch_with_array_overhead_round_triptest_uncount_two_slices_example_from_issuetest_uncount_view_array_slices_sharing_data_buffers,test_uncount_binary_view_array_shared_data_bufferstest_uncount_dictionaries_sharing_valuestest_uncount_with_overflow_promotionAre there any user-facing changes?
No. Purely additive API -- existing
count_*methods andget_record_batch_memory_sizereturn identical results. No caller changes needed.Benchmark
record_batch_memorybenchmark shows no regression for counting: