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
115 changes: 69 additions & 46 deletions datafusion/pruning/benches/string_in_list_pruning.rs
Original file line number Diff line number Diff line change
Expand Up @@ -15,18 +15,20 @@
// specific language governing permissions and limitations
// under the License.

//! Compare compact string IN-list pruning with per-value min/max expansion.
//! Compare string IN-list pruning with per-value min/max expansion.
//!
//! Both cases raise `max_in_list_size` to the domain size. On a baseline
//! without compact pruning, `in_list` measures the ordinary raised-cap path.
//! The explicit `expanded_or` is a balanced tree of equalities, which produces
//! the same per-value statistics checks without making the baseline depend on
//! a deeply nested expression. Half of the statistics intervals hit a domain
//! member and half fall in a sparse gap. Bloom filters are not involved.
//! Both cases raise `max_in_list_size` to the domain size. `in_list` follows the
//! production representation for each domain size, while the explicit
//! `expanded_or` remains a stable balanced tree of equalities. This makes the
//! small-domain cases suitable for measuring compact representation threshold
//! changes without changing the comparison baseline. Half of the statistics
//! intervals hit a domain member and half fall in a sparse gap. Bloom filters
//! are not involved.
//!
//! Run with `cargo bench -p datafusion-pruning --bench string_in_list_pruning`.
//! The construction benchmarks reuse their input physical expressions; the
//! evaluation benchmarks reuse their already-built pruning predicates.
//! Construction is independent of the statistics batch size, so it varies only
//! the IN-list domain size. Evaluation varies both the domain size and number
//! of containers, reusing already-built pruning predicates and statistics.

use std::collections::HashSet;
use std::hint::black_box;
Expand All @@ -41,8 +43,9 @@ use datafusion_physical_expr::PhysicalExprRef;
use datafusion_physical_expr::expressions::{BinaryExpr, col, in_list, lit};
use datafusion_pruning::{PruningPredicate, PruningPredicateBuilder, PruningStatistics};

const DOMAIN_SIZES: [usize; 4] = [20, 21, 256, 1024];
const CONTAINERS: usize = 4096;
const DOMAIN_SIZES: [usize; 9] = [1, 2, 4, 8, 16, 20, 21, 256, 1024];
const CONTAINER_COUNTS: [usize; 3] = [16, 256, 4096];
const BASELINE_CONTAINER_COUNT: usize = 4096;

fn value(index: usize) -> String {
format!("key{index:08}")
Expand Down Expand Up @@ -80,20 +83,20 @@ struct IntervalStatistics {
}

impl IntervalStatistics {
fn new(domain_size: usize) -> Self {
let min = StringViewArray::from_iter_values((0..CONTAINERS).map(|index| {
fn new(domain_size: usize, container_count: usize) -> Self {
let min = StringViewArray::from_iter_values((0..container_count).map(|index| {
let start = (index / 2 % domain_size) * 10;
value(start + if index % 2 == 0 { 0 } else { 3 })
}));
let max = StringViewArray::from_iter_values((0..CONTAINERS).map(|index| {
let max = StringViewArray::from_iter_values((0..container_count).map(|index| {
let start = (index / 2 % domain_size) * 10;
value(start + if index % 2 == 0 { 0 } else { 7 })
}));
Self {
min: Arc::new(min),
max: Arc::new(max),
null_counts: Arc::new(UInt64Array::from(vec![0; CONTAINERS])),
row_counts: Arc::new(UInt64Array::from(vec![128; CONTAINERS])),
null_counts: Arc::new(UInt64Array::from(vec![0; container_count])),
row_counts: Arc::new(UInt64Array::from(vec![128; container_count])),
}
}
}
Expand All @@ -108,7 +111,7 @@ impl PruningStatistics for IntervalStatistics {
}

fn num_containers(&self) -> usize {
CONTAINERS
self.min.len()
}

fn null_counts(&self, column: &Column) -> Option<ArrayRef> {
Expand All @@ -135,7 +138,6 @@ struct BenchmarkCase {
expanded_or: PhysicalExprRef,
in_list_predicate: PruningPredicate,
expanded_or_predicate: PruningPredicate,
statistics: IntervalStatistics,
}

impl BenchmarkCase {
Expand Down Expand Up @@ -168,28 +170,40 @@ impl BenchmarkCase {
.to_string()
.contains("IN_SET_INTERSECTS")
);
let statistics = IntervalStatistics::new(size);

// Check that both benchmark paths do the same useful work, rather
// than comparing compact pruning with an always-true fallback.
let expected = (0..CONTAINERS)
.map(|index| index % 2 == 0)
.collect::<Vec<_>>();
assert_eq!(in_list_predicate.prune(&statistics).unwrap(), expected);
assert_eq!(expanded_or_predicate.prune(&statistics).unwrap(), expected);

Self {
size,
schema,
in_list,
expanded_or,
in_list_predicate,
expanded_or_predicate,
statistics,
}
}
}

fn expected_results(container_count: usize) -> Vec<bool> {
(0..container_count).map(|index| index % 2 == 0).collect()
}

fn assert_equivalent_results(case: &BenchmarkCase, statistics: &IntervalStatistics) {
let expected = expected_results(statistics.num_containers());
assert_eq!(case.in_list_predicate.prune(statistics).unwrap(), expected);
assert_eq!(
case.expanded_or_predicate.prune(statistics).unwrap(),
expected
);
}

fn evaluation_group_name(container_count: usize) -> String {
// Keep the original benchmark IDs for the pre-existing 4,096-container
// cases so Criterion baselines remain directly comparable.
if container_count == BASELINE_CONTAINER_COUNT {
"string_in_list_pruning/evaluate".to_string()
} else {
format!("string_in_list_pruning/evaluate/{container_count}_containers")
}
}

fn criterion_benchmark(criterion: &mut Criterion) {
let cases = DOMAIN_SIZES.map(BenchmarkCase::new);
let mut construction = criterion.benchmark_group("string_in_list_pruning/construct");
Expand All @@ -216,25 +230,34 @@ fn criterion_benchmark(criterion: &mut Criterion) {
}
construction.finish();

let mut evaluation = criterion.benchmark_group("string_in_list_pruning/evaluate");
evaluation.throughput(Throughput::Elements(CONTAINERS as u64));
for case in &cases {
for (name, predicate) in [
("in_list", &case.in_list_predicate),
("expanded_or", &case.expanded_or_predicate),
] {
evaluation.bench_with_input(
BenchmarkId::new(name, case.size),
predicate,
|bencher, predicate| {
bencher.iter(|| {
black_box(predicate.prune(black_box(&case.statistics)).unwrap())
});
},
);
for container_count in CONTAINER_COUNTS {
let mut evaluation =
criterion.benchmark_group(evaluation_group_name(container_count));
evaluation.throughput(Throughput::Elements(container_count as u64));
for case in &cases {
let statistics = IntervalStatistics::new(case.size, container_count);

// Check every matrix cell outside the timed loop so both paths do
// the same useful work rather than measuring an always-true fallback.
assert_equivalent_results(case, &statistics);

for (name, predicate) in [
("in_list", &case.in_list_predicate),
("expanded_or", &case.expanded_or_predicate),
] {
evaluation.bench_with_input(
BenchmarkId::new(name, case.size),
predicate,
|bencher, predicate| {
bencher.iter(|| {
black_box(predicate.prune(black_box(&statistics)).unwrap())
});
},
);
}
}
evaluation.finish();
}
evaluation.finish();
}

criterion_group!(benches, criterion_benchmark);
Expand Down