mod common;
use std::{
sync::{Arc, OnceLock},
time::Duration,
};
use common::{LOW_CARDINALITY_COUNT, TOTAL_ROWS};
use criterion::{BenchmarkId, Criterion, black_box, criterion_group, criterion_main};
use datafusion_common::ScalarValue;
use lance_core::cache::LanceCache;
use lance_index::metrics::NoOpMetricsCollector;
use lance_index::pbold;
use lance_index::scalar::lance_format::LanceIndexStore;
use lance_index::scalar::registry::ScalarIndexPlugin;
use lance_index::scalar::{SargableQuery, ScalarIndex, bitmap::BitmapIndexPlugin};
use lance_io::object_store::ObjectStore;
use object_store::path::Path;
#[cfg(target_os = "linux")]
use pprof::criterion::{Output, PProfProfiler};
static RUNTIME: OnceLock<tokio::runtime::Runtime> = OnceLock::new();
static CACHE: OnceLock<Arc<LanceCache>> = OnceLock::new();
static INT_UNIQUE_INDEX_NO_CACHE: OnceLock<Arc<dyn ScalarIndex>> = OnceLock::new();
static INT_UNIQUE_INDEX_CACHED: OnceLock<Arc<dyn ScalarIndex>> = OnceLock::new();
static INT_LOW_CARD_INDEX_NO_CACHE: OnceLock<Arc<dyn ScalarIndex>> = OnceLock::new();
static INT_LOW_CARD_INDEX_CACHED: OnceLock<Arc<dyn ScalarIndex>> = OnceLock::new();
static STRING_UNIQUE_INDEX_NO_CACHE: OnceLock<Arc<dyn ScalarIndex>> = OnceLock::new();
static STRING_UNIQUE_INDEX_CACHED: OnceLock<Arc<dyn ScalarIndex>> = OnceLock::new();
static STRING_LOW_CARD_INDEX_NO_CACHE: OnceLock<Arc<dyn ScalarIndex>> = OnceLock::new();
static STRING_LOW_CARD_INDEX_CACHED: OnceLock<Arc<dyn ScalarIndex>> = OnceLock::new();
fn get_runtime() -> &'static tokio::runtime::Runtime {
RUNTIME.get_or_init(|| tokio::runtime::Builder::new_multi_thread().build().unwrap())
}
fn get_cache(use_cache: bool, key_prefix: &str) -> Arc<LanceCache> {
if use_cache {
Arc::new(
CACHE
.get_or_init(|| Arc::new(LanceCache::with_capacity(1024 * 1024 * 1024)))
.with_key_prefix(key_prefix),
)
} else {
Arc::new(LanceCache::no_cache())
}
}
async fn create_int_unique_index(
store: Arc<LanceIndexStore>,
use_cache: bool,
) -> Arc<dyn ScalarIndex> {
let stream = common::generate_int_unique_stream();
BitmapIndexPlugin::train_bitmap_index(stream, store.as_ref())
.await
.unwrap();
let details = prost_types::Any::from_msg(&pbold::BitmapIndexDetails::default()).unwrap();
(BitmapIndexPlugin
.load_index(store, &details, None, &get_cache(use_cache, "int_unique"))
.await
.unwrap()) as _
}
async fn create_int_low_card_index(
store: Arc<LanceIndexStore>,
use_cache: bool,
) -> Arc<dyn ScalarIndex> {
let stream = common::generate_int_low_cardinality_stream();
BitmapIndexPlugin::train_bitmap_index(stream, store.as_ref())
.await
.unwrap();
let details = prost_types::Any::from_msg(&pbold::BitmapIndexDetails::default()).unwrap();
(BitmapIndexPlugin
.load_index(store, &details, None, &get_cache(use_cache, "int_low_card"))
.await
.unwrap()) as _
}
async fn create_string_unique_index(
store: Arc<LanceIndexStore>,
use_cache: bool,
) -> Arc<dyn ScalarIndex> {
let stream = common::generate_string_unique_stream();
BitmapIndexPlugin::train_bitmap_index(stream, store.as_ref())
.await
.unwrap();
let details = prost_types::Any::from_msg(&pbold::BitmapIndexDetails::default()).unwrap();
(BitmapIndexPlugin
.load_index(
store,
&details,
None,
&get_cache(use_cache, "string_unique"),
)
.await
.unwrap()) as _
}
async fn create_string_low_card_index(
store: Arc<LanceIndexStore>,
use_cache: bool,
) -> Arc<dyn ScalarIndex> {
let stream = common::generate_string_low_cardinality_stream();
BitmapIndexPlugin::train_bitmap_index(stream, store.as_ref())
.await
.unwrap();
let details = prost_types::Any::from_msg(&pbold::BitmapIndexDetails::default()).unwrap();
(BitmapIndexPlugin
.load_index(
store,
&details,
None,
&get_cache(use_cache, "string_low_card"),
)
.await
.unwrap()) as _
}
fn setup_int_unique_index(rt: &tokio::runtime::Runtime, use_cache: bool) -> Arc<dyn ScalarIndex> {
let static_ref = if use_cache {
&INT_UNIQUE_INDEX_CACHED
} else {
&INT_UNIQUE_INDEX_NO_CACHE
};
static_ref
.get_or_init(|| {
rt.block_on(async {
let tempdir = tempfile::tempdir().unwrap();
let store = Arc::new(LanceIndexStore::new(
Arc::new(ObjectStore::local()),
Path::from_filesystem_path(tempdir.path()).unwrap(),
get_cache(use_cache, "int_unique"),
));
let index = create_int_unique_index(store, use_cache).await;
let _ = tempdir.keep();
index
})
})
.clone()
}
fn setup_int_low_card_index(rt: &tokio::runtime::Runtime, use_cache: bool) -> Arc<dyn ScalarIndex> {
let static_ref = if use_cache {
&INT_LOW_CARD_INDEX_CACHED
} else {
&INT_LOW_CARD_INDEX_NO_CACHE
};
static_ref
.get_or_init(|| {
rt.block_on(async {
let tempdir = tempfile::tempdir().unwrap();
let store = Arc::new(LanceIndexStore::new(
Arc::new(ObjectStore::local()),
Path::from_filesystem_path(tempdir.path()).unwrap(),
get_cache(use_cache, "int_low_card"),
));
let index = create_int_low_card_index(store, use_cache).await;
let _ = tempdir.keep();
index
})
})
.clone()
}
fn setup_string_unique_index(
rt: &tokio::runtime::Runtime,
use_cache: bool,
) -> Arc<dyn ScalarIndex> {
let static_ref = if use_cache {
&STRING_UNIQUE_INDEX_CACHED
} else {
&STRING_UNIQUE_INDEX_NO_CACHE
};
static_ref
.get_or_init(|| {
rt.block_on(async {
let tempdir = tempfile::tempdir().unwrap();
let store = Arc::new(LanceIndexStore::new(
Arc::new(ObjectStore::local()),
Path::from_filesystem_path(tempdir.path()).unwrap(),
get_cache(use_cache, "string_unique"),
));
let index = create_string_unique_index(store, use_cache).await;
let _ = tempdir.keep();
index
})
})
.clone()
}
fn setup_string_low_card_index(
rt: &tokio::runtime::Runtime,
use_cache: bool,
) -> Arc<dyn ScalarIndex> {
let static_ref = if use_cache {
&STRING_LOW_CARD_INDEX_CACHED
} else {
&STRING_LOW_CARD_INDEX_NO_CACHE
};
static_ref
.get_or_init(|| {
rt.block_on(async {
let tempdir = tempfile::tempdir().unwrap();
let store = Arc::new(LanceIndexStore::new(
Arc::new(ObjectStore::local()),
Path::from_filesystem_path(tempdir.path()).unwrap(),
get_cache(use_cache, "string_low_card"),
));
let index = create_string_low_card_index(store, use_cache).await;
let _ = tempdir.keep();
index
})
})
.clone()
}
fn bench_equality(c: &mut Criterion) {
let rt = get_runtime();
let int_unique_value = (TOTAL_ROWS / 2) as i64;
let string_unique_value = format!("string_{:010}", TOTAL_ROWS / 2);
let int_low_card_value = (LOW_CARDINALITY_COUNT / 2) as i64;
let string_low_card_value = format!("value_{:03}", LOW_CARDINALITY_COUNT / 2);
let mut group = c.benchmark_group("bitmap_equality");
group
.sample_size(10)
.measurement_time(Duration::from_secs(10));
for use_cache in [false, true] {
let cache_label = if use_cache { "cached" } else { "no_cache" };
group.bench_function(BenchmarkId::new("int_unique", cache_label), |b| {
let index = setup_int_unique_index(rt, use_cache);
b.to_async(rt).iter(|| {
let index = index.clone();
let value = int_unique_value;
async move {
let query = SargableQuery::Equals(ScalarValue::Int64(Some(value)));
black_box(index.search(&query, &NoOpMetricsCollector).await.unwrap());
}
})
});
group.bench_function(BenchmarkId::new("int_low_card", cache_label), |b| {
let index = setup_int_low_card_index(rt, use_cache);
b.to_async(rt).iter(|| {
let index = index.clone();
let value = int_low_card_value;
async move {
let query = SargableQuery::Equals(ScalarValue::Int64(Some(value)));
black_box(index.search(&query, &NoOpMetricsCollector).await.unwrap());
}
})
});
group.bench_function(BenchmarkId::new("string_unique", cache_label), |b| {
let index = setup_string_unique_index(rt, use_cache);
let value = string_unique_value.clone();
b.to_async(rt).iter(|| {
let index = index.clone();
let value = value.clone();
async move {
let query = SargableQuery::Equals(ScalarValue::Utf8(Some(value)));
black_box(index.search(&query, &NoOpMetricsCollector).await.unwrap());
}
})
});
group.bench_function(BenchmarkId::new("string_low_card", cache_label), |b| {
let index = setup_string_low_card_index(rt, use_cache);
let value = string_low_card_value.clone();
b.to_async(rt).iter(|| {
let index = index.clone();
let value = value.clone();
async move {
let query = SargableQuery::Equals(ScalarValue::Utf8(Some(value)));
black_box(index.search(&query, &NoOpMetricsCollector).await.unwrap());
}
})
});
}
group.finish();
}
fn bench_in(c: &mut Criterion) {
let rt = get_runtime();
let value_counts = [1, 3, 5];
for &num_values in &value_counts {
let mut group = c.benchmark_group(format!("bitmap_in_{}", num_values));
group
.sample_size(10)
.measurement_time(Duration::from_secs(10));
let mid_int = (TOTAL_ROWS / 2) as i64;
let mid_string = TOTAL_ROWS / 2;
let mid_low_card = LOW_CARDINALITY_COUNT / 2;
let int_values: Vec<ScalarValue> = (0..num_values)
.map(|i| ScalarValue::Int64(Some(mid_int + i as i64 - num_values as i64 / 2)))
.collect();
let int_low_card_values: Vec<ScalarValue> = (0..num_values)
.map(|i| ScalarValue::Int64(Some((mid_low_card + i - num_values / 2) as i64)))
.collect();
let string_values: Vec<ScalarValue> = (0..num_values)
.map(|i| {
ScalarValue::Utf8(Some(format!(
"string_{:010}",
(mid_string as i64 + i as i64 - num_values as i64 / 2) as u64
)))
})
.collect();
let string_low_card_values: Vec<ScalarValue> = (0..num_values)
.map(|i| {
ScalarValue::Utf8(Some(format!(
"value_{:03}",
(mid_low_card as i32 + i as i32 - num_values as i32 / 2) as usize
)))
})
.collect();
for use_cache in [false, true] {
let cache_label = if use_cache { "cached" } else { "no_cache" };
group.bench_function(BenchmarkId::new("int_unique", cache_label), |b| {
let index = setup_int_unique_index(rt, use_cache);
let values = int_values.clone();
b.to_async(rt).iter(|| {
let index = index.clone();
let values = values.clone();
async move {
let query = SargableQuery::IsIn(values);
black_box(index.search(&query, &NoOpMetricsCollector).await.unwrap());
}
})
});
group.bench_function(BenchmarkId::new("int_low_card", cache_label), |b| {
let index = setup_int_low_card_index(rt, use_cache);
let values = int_low_card_values.clone();
b.to_async(rt).iter(|| {
let index = index.clone();
let values = values.clone();
async move {
let query = SargableQuery::IsIn(values);
black_box(index.search(&query, &NoOpMetricsCollector).await.unwrap());
}
})
});
group.bench_function(BenchmarkId::new("string_unique", cache_label), |b| {
let index = setup_string_unique_index(rt, use_cache);
let values = string_values.clone();
b.to_async(rt).iter(|| {
let index = index.clone();
let values = values.clone();
async move {
let query = SargableQuery::IsIn(values);
black_box(index.search(&query, &NoOpMetricsCollector).await.unwrap());
}
})
});
group.bench_function(BenchmarkId::new("string_low_card", cache_label), |b| {
let index = setup_string_low_card_index(rt, use_cache);
let values = string_low_card_values.clone();
b.to_async(rt).iter(|| {
let index = index.clone();
let values = values.clone();
async move {
let query = SargableQuery::IsIn(values);
black_box(index.search(&query, &NoOpMetricsCollector).await.unwrap());
}
})
});
}
group.finish();
}
}
fn bench_bitmap(c: &mut Criterion) {
bench_equality(c);
bench_in(c);
}
#[cfg(target_os = "linux")]
criterion_group!(
name=benches;
config = Criterion::default()
.measurement_time(Duration::from_secs(10))
.sample_size(10)
.with_profiler(PProfProfiler::new(100, Output::Flamegraph(None)));
targets = bench_bitmap);
#[cfg(not(target_os = "linux"))]
criterion_group!(
name=benches;
config = Criterion::default()
.measurement_time(Duration::from_secs(10))
.sample_size(10);
targets = bench_bitmap);
criterion_main!(benches);