Skip to content
Open
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
179 changes: 179 additions & 0 deletions datafusion/functions-aggregate/benches/first_last.rs
Original file line number Diff line number Diff line change
Expand Up @@ -452,6 +452,172 @@ fn create_list_of_struct_array(n: usize, list_len: usize, null_density: f32) ->
))
}

/// Drive one coalesce-peers variant: `separate xN` (pre-rewrite, one primitive
/// accumulator per column) vs `coalesced struct` (post-rewrite, one struct
/// accumulator). Each accumulator is primed with `worst_ord` (the always-losing
/// ordering value), then fed one `ord` per iteration from `iter_ords`. Varying
/// `iter_ords` selects which path dominates — see [`coalesce_comparison_bench`].
#[expect(clippy::too_many_arguments)]
fn run_coalesce_variant(
c: &mut Criterion,
name: &str,
variant: &str,
column_values: &[ArrayRef],
struct_values: &ArrayRef,
worst_ord: &ArrayRef,
iter_ords: &[ArrayRef],
group_indices: &[usize],
num_groups: usize,
) {
// Pre-rewrite: one accumulator per column.
c.bench_function(
&format!("{name} separate x{} {variant}", column_values.len()),
|b| {
b.iter_batched(
|| {
column_values
.iter()
.map(|values| {
let mut acc = prepare_typed_groups_accumulator(
true,
values.data_type().clone(),
);
acc.update_batch(
&[Arc::clone(values), Arc::clone(worst_ord)],
group_indices,
None,
num_groups,
)
.unwrap();
acc
})
.collect::<Vec<_>>()
},
|mut accumulators| {
for ord in iter_ords {
for (acc, values) in accumulators.iter_mut().zip(column_values) {
#[expect(clippy::unit_arg)]
black_box(
acc.update_batch(
&[Arc::clone(values), Arc::clone(ord)],
group_indices,
None,
num_groups,
)
.unwrap(),
);
}
}
},
BatchSize::SmallInput,
)
},
);

// Post-rewrite: a single struct-valued accumulator.
c.bench_function(&format!("{name} coalesced struct {variant}"), |b| {
b.iter_batched(
|| {
let mut acc = prepare_typed_groups_accumulator(
true,
struct_values.data_type().clone(),
);
acc.update_batch(
&[Arc::clone(struct_values), Arc::clone(worst_ord)],
group_indices,
None,
num_groups,
)
.unwrap();
acc
},
|mut accumulator| {
for ord in iter_ords {
#[expect(clippy::unit_arg)]
black_box(
accumulator
.update_batch(
&[Arc::clone(struct_values), Arc::clone(ord)],
group_indices,
None,
num_groups,
)
.unwrap(),
);
}
},
BatchSize::SmallInput,
)
});
}

/// Head-to-head for the coalesce-peers rewrite: N independent primitive
/// `first_value` accumulators (the pre-rewrite plan) vs one struct-valued
/// accumulator carrying the same N columns (the post-rewrite plan). Runs two
/// variants, since the two plans differ on two separable costs:
///
/// - `(winner stable)` reuses one `ord` every iteration. After the first
/// iteration the running winner already holds the best value in `ord`, so
/// this is dominated by compare-and-reject — the per-row ordering comparison
/// the rewrite collapses (N compares -> 1).
/// - `(winner changes)` feeds a distinct, strictly-decreasing `ord` per
/// iteration so every row becomes a new winner (smaller ord wins; `worst_ord`
/// = i64::MAX primes the loser). This forces the running value to be
/// replaced+copied on every row, exercising the update path the winner-stable
/// case skips — where the struct plan copies one wider row vs N narrow ones.
#[expect(clippy::needless_pass_by_value)]
fn coalesce_comparison_bench(
c: &mut Criterion,
name: &str,
column_values: Vec<ArrayRef>,
struct_values: ArrayRef,
ord: ArrayRef,
num_groups: usize,
) {
const ITERS: usize = 100;
let n = ord.len();
let group_indices: Vec<usize> = (0..n).map(|i| i % num_groups).collect();
let worst_ord: ArrayRef = Arc::new(Int64Array::from(vec![i64::MAX; n]));

// Winner-stable: the same `ord` every iteration.
let stable_ords: Vec<ArrayRef> = (0..ITERS).map(|_| Arc::clone(&ord)).collect();

// Winner-changing: strictly decreasing across (iteration, row), and below
// `worst_ord`, so every row wins and updates the running value. Built once,
// outside the timed loop.
let changing_ords: Vec<ArrayRef> = (0..ITERS)
.map(|k| {
let base = i64::MAX - 1 - (k as i64) * (n as i64);
Arc::new(Int64Array::from(
(0..n as i64).map(|i| base - i).collect::<Vec<i64>>(),
)) as ArrayRef
})
.collect();

run_coalesce_variant(
c,
name,
"(winner stable)",
&column_values,
&struct_values,
&worst_ord,
&stable_ords,
&group_indices,
num_groups,
);
run_coalesce_variant(
c,
name,
"(winner changes)",
&column_values,
&struct_values,
&worst_ord,
&changing_ords,
&group_indices,
num_groups,
);
}

fn first_last_nested_benchmark(c: &mut Criterion) {
const N: usize = 65536;
const NUM_GROUPS: usize = 1024;
Expand Down Expand Up @@ -510,6 +676,19 @@ fn first_last_nested_benchmark(c: &mut Criterion) {
);
}
}
// Coalesce-peers head-to-head on null-free columns.
let a = Arc::new(create_primitive_array::<Int64Type>(N, 0.0)) as ArrayRef;
let b = Arc::new(create_string_array_with_len::<i32>(N, 0.0, 16)) as ArrayRef;
let d = Arc::new(create_primitive_array::<Float64Type>(N, 0.0)) as ArrayRef;
let struct_values = create_struct_array(N, 0.0);
coalesce_comparison_bench(
c,
"first_value coalesce_peers(i64,utf8,f64)",
vec![a, b, d],
struct_values,
ord,
NUM_GROUPS,
);
}

fn first_last_benchmark(c: &mut Criterion) {
Expand Down