sunchao commented on PR #11279:
URL: https://github.com/apache/arrow-rs/pull/11279#issuecomment-5963439900
<details>
<summary>Standalone allocation reproducer (separate from Criterion
timing)</summary>
This probe uses the exact same 42-case fixture as the benchmark companion.
It counts requested sizes through Rust's System allocator, with all input
arrays
and indices already live before each operation. It validates every result
value
and null. It measures:
- `peak_extra_live_bytes`: maximum live requested allocation bytes during the
operation, minus the pre-operation baseline. Includes scratch and the
result.
- `output_extra_live_bytes`: live increment after return while the result
lives.
- `output_data_capacity`: capacity of unique output data allocations,
excluding
views, null bitmap, and other result metadata. This is not the live-byte
metric.
- `allocation_calls` and `requested_bytes`: successful
allocations/reallocations
and their requested sizes; a reallocation counts its new total requested
size.
- `retained_source_data_capacity`: input data capacity still referenced by
output.
- `after_drop_extra_live_bytes`: live increment after dropping the result.
- `input_fingerprint`: deterministic FNV-1a over selected indices, validity,
lengths and bytes, excluding pointer values, for workload parity.
Inputs are excluded from incremental allocation counters, but the output is
included. These numbers are not RSS, allocator arena sizes, or allocator
metadata.
Reallocations account for the requested live-size delta, not a moving
allocator's
internal temporary old-plus-new footprint. The counters are global; this
runner
executes one operation at a time without background worker threads.
Instrumented
timings are deliberately not used. This manifest uses
`default-features=false`;
Criterion uses the repository benchmark's defaults plus `test_utils`, so
these
figures describe the selected operations, not Criterion's whole process.
Use the exact tested commits (the base checkout adds only the benchmark
fixture to
kernel base `705e6405b56cc01efb1c267203d240d0d85804ab`):
```sh
git clone https://github.com/apache/arrow-rs.git arrow-rs-bench
cd arrow-rs-bench
git fetch origin pull/11346/head:bench-base pull/11279/head:compact-head
git worktree add ../arrow-base 363eb96df844e7b7ecca8289d88cfc2c12e3ceb5
git worktree add ../arrow-head d8726603c966e1bad0f781de443ca067e553d0dc
cd ..
mkdir -p allocation-probe/base allocation-probe/head
rustup toolchain install 1.99.0 --profile minimal
```
Both checkouts have the identical fixture, SHA256
`b27220df51dbebe4d8c824db19e39227e8ff3a76dbc04d709ad386140f0223d3`.
The source below intentionally reads that shared fixture from `arrow-base`;
crate imports still resolve to the base or head manifest's dependencies.
`allocation-probe/base/Cargo.toml`:
```toml
[package]
name = "byte-view-allocation-base"
version = "0.1.0"
edition = "2024"
[workspace]
[features]
compact = []
[[bin]]
name = "byte-view-allocation-base"
path = "../measure.rs"
[dependencies]
arrow = { path = "../../arrow-base/arrow", default-features = false }
arrow-buffer = { path = "../../arrow-base/arrow-buffer" }
arrow-select = { path = "../../arrow-base/arrow-select" }
rand = "=0.10.3"
```
`allocation-probe/head/Cargo.toml`:
```toml
[package]
name = "byte-view-allocation-head"
version = "0.1.0"
edition = "2024"
[workspace]
[features]
compact = []
[[bin]]
name = "byte-view-allocation-head"
path = "../measure.rs"
[dependencies]
arrow = { path = "../../arrow-head/arrow", default-features = false }
arrow-buffer = { path = "../../arrow-head/arrow-buffer" }
arrow-select = { path = "../../arrow-head/arrow-select" }
rand = "=0.10.3"
```
`allocation-probe/measure.rs`:
```rust
use std::alloc::{GlobalAlloc, Layout, System};
use std::collections::{HashMap, HashSet};
use std::sync::Arc;
use std::sync::atomic::{AtomicUsize, Ordering::Relaxed};
use arrow::array::{Array, ArrayRef, StringViewArray};
use arrow_select::interleave::interleave;
#[cfg(feature = "compact")]
use arrow_select::interleave::Interleaver;
#[path = "../arrow-base/arrow/benches/interleave_byte_view_cases/mod.rs"]
mod cases;
struct TrackingAllocator;
static LIVE: AtomicUsize = AtomicUsize::new(0);
static PEAK: AtomicUsize = AtomicUsize::new(0);
static ALLOCATIONS: AtomicUsize = AtomicUsize::new(0);
static REQUESTED: AtomicUsize = AtomicUsize::new(0);
fn allocated(size: usize) {
ALLOCATIONS.fetch_add(1, Relaxed);
REQUESTED.fetch_add(size, Relaxed);
let live = LIVE.fetch_add(size, Relaxed) + size;
PEAK.fetch_max(live, Relaxed);
}
unsafe impl GlobalAlloc for TrackingAllocator {
unsafe fn alloc(&self, layout: Layout) -> *mut u8 {
let ptr = unsafe { System.alloc(layout) };
if !ptr.is_null() {
allocated(layout.size());
}
ptr
}
unsafe fn alloc_zeroed(&self, layout: Layout) -> *mut u8 {
let ptr = unsafe { System.alloc_zeroed(layout) };
if !ptr.is_null() {
allocated(layout.size());
}
ptr
}
unsafe fn dealloc(&self, ptr: *mut u8, layout: Layout) {
LIVE.fetch_sub(layout.size(), Relaxed);
unsafe { System.dealloc(ptr, layout) };
}
unsafe fn realloc(&self, ptr: *mut u8, layout: Layout, size: usize) ->
*mut u8 {
let new = unsafe { System.realloc(ptr, layout, size) };
if !new.is_null() {
LIVE.fetch_sub(layout.size(), Relaxed);
allocated(size);
}
new
}
}
#[global_allocator]
static ALLOCATOR: TrackingAllocator = TrackingAllocator;
fn hash_bytes(hash: &mut u64, bytes: &[u8]) {
for byte in bytes {
*hash ^= u64::from(*byte);
*hash = hash.wrapping_mul(0x100000001b3);
}
}
fn main() {
println!("case,operation,input_fingerprint,source_buffers,unique_source_data_capacity,selected_payload_bytes,unique_selected_ranges,output_buffers,output_data_len,output_data_capacity,retained_source_data_capacity,allocation_calls,requested_bytes,peak_extra_live_bytes,output_extra_live_bytes,after_drop_extra_live_bytes");
cases::for_each_case(|case| {
let values: Vec<&dyn Array> = case.arrays.iter().map(|a| a as &dyn
Array).collect();
let mut sources = HashMap::new();
let mut source_buffers = 0;
for array in &case.arrays {
for buffer in array.data_buffers().iter() {
source_buffers += 1;
sources.insert(buffer.data_ptr().as_ptr() as usize,
buffer.capacity());
}
}
let source_capacity: usize = sources.values().sum();
let mut fingerprint = 0xcbf29ce484222325;
let mut selected_bytes = 0;
let mut unique_ranges = HashSet::new();
for &(source, row) in &case.indices {
hash_bytes(&mut fingerprint, &(source as u64).to_le_bytes());
hash_bytes(&mut fingerprint, &(row as u64).to_le_bytes());
let array = &case.arrays[source];
let valid = array.is_valid(row);
hash_bytes(&mut fingerprint, &[u8::from(valid)]);
if valid {
let value = array.value(row);
hash_bytes(&mut fingerprint, &(value.len() as
u64).to_le_bytes());
hash_bytes(&mut fingerprint, value.as_bytes());
if value.len() > 12 {
selected_bytes += value.len();
unique_ranges.insert((value.as_ptr() as usize,
value.len()));
}
}
}
let mut operations = vec!["interleave", "interleave_gc"];
#[cfg(feature = "compact")]
operations.extend(["compact", "compact_shared"]);
for operation in operations.drain(..) {
let before = LIVE.load(Relaxed);
let calls_before = ALLOCATIONS.load(Relaxed);
let requested_before = REQUESTED.load(Relaxed);
PEAK.store(before, Relaxed);
let output: ArrayRef = match operation {
"interleave" => interleave(&values, &case.indices).unwrap(),
"interleave_gc" => {
let selected = interleave(&values,
&case.indices).unwrap();
let selected =
selected.as_any().downcast_ref::<StringViewArray>().unwrap();
Arc::new(selected.gc())
}
#[cfg(feature = "compact")]
"compact" => Interleaver::new()
.with_compact_byte_views(true)
.interleave(&values, &case.indices)
.unwrap(),
#[cfg(feature = "compact")]
"compact_shared" => Interleaver::new()
.with_compact_byte_views(true)
.with_preserve_byte_view_sharing(true)
.interleave(&values, &case.indices)
.unwrap(),
_ => unreachable!(),
};
let output_live = LIVE.load(Relaxed) - before;
let peak = PEAK.load(Relaxed) - before;
let calls = ALLOCATIONS.load(Relaxed) - calls_before;
let requested = REQUESTED.load(Relaxed) - requested_before;
let array =
output.as_any().downcast_ref::<StringViewArray>().unwrap();
assert_eq!(array.len(), case.indices.len());
for (i, &(source, row)) in case.indices.iter().enumerate() {
assert_eq!(array.is_null(i),
case.arrays[source].is_null(row));
if array.is_valid(i) {
assert_eq!(array.value(i),
case.arrays[source].value(row));
}
}
let output_buffers = array.data_buffers().len();
let output_len: usize = array.data_buffers().iter().map(|b|
b.len()).sum();
let mut output_allocations = HashMap::new();
for buffer in array.data_buffers().iter() {
output_allocations.insert(buffer.data_ptr().as_ptr() as
usize, buffer.capacity());
}
let output_capacity: usize = output_allocations.values().sum();
let retained: usize = output_allocations.keys().filter_map(|p|
sources.get(p)).sum();
drop(output_allocations);
drop(output);
let after_drop = LIVE.load(Relaxed) as isize - before as isize;
println!("{},{operation},{fingerprint:016x},{source_buffers},{source_capacity},{selected_bytes},{},{output_buffers},{output_len},{output_capacity},{retained},{calls},{requested},{peak},{output_live},{after_drop}",
case.name, unique_ranges.len());
}
});
}
```
Seed the harness lockfiles from the tested checkouts, then run the two
independent
release executables. Cargo adds the harness package and removes unused
workspace
packages; the measured harness's dependency versions all occur in these
workspace
lockfiles.
```sh
cp arrow-base/Cargo.lock allocation-probe/base/Cargo.lock
cp arrow-head/Cargo.lock allocation-probe/head/Cargo.lock
cargo +1.99.0 run --release --manifest-path allocation-probe/base/Cargo.toml
> allocations-base.csv
cargo +1.99.0 run --release --manifest-path allocation-probe/head/Cargo.toml
--features compact > allocations-head.csv
```
The base CSV has 84 rows plus its header; head has 168 rows plus its header.
Match `case,input_fingerprint` across files. Compare base `interleave_gc`
with head `compact` or `compact_shared`, and keep plain `interleave`
as the ownership control. For every compact/gc row the measured retained
source
capacity and after-drop increment are zero. Dependency lockfiles should be
retained with any reproduced run; future dependency releases can otherwise
change allocator/code-generation details.
</details>
--
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]