diff --git a/rust/lance-encoding/benches/decoder.rs b/rust/lance-encoding/benches/decoder.rs index cc0404e1bb3..f4aecf69522 100644 --- a/rust/lance-encoding/benches/decoder.rs +++ b/rust/lance-encoding/benches/decoder.rs @@ -2,7 +2,7 @@ // SPDX-FileCopyrightText: Copyright The Lance Authors use std::{collections::HashMap, sync::Arc}; -use arrow_array::{RecordBatch, UInt32Array}; +use arrow_array::{RecordBatch, StringArray, UInt32Array}; use arrow_schema::{DataType, Field, Schema, TimeUnit}; use arrow_select::take::take; use criterion::{Criterion, criterion_group, criterion_main}; @@ -10,6 +10,7 @@ use futures::StreamExt; use lance_core::cache::LanceCache; use lance_datagen::ArrayGeneratorExt; use lance_encoding::{ + BufferScheduler, EncodingsIo, decoder::{ DecodeBatchScheduler, DecoderConfig, DecoderPlugins, FilterExpression, create_decode_stream, }, @@ -563,6 +564,140 @@ fn bench_decode_compressed_parallel(c: &mut Criterion) { } } +async fn decode_take( + encoded: &lance_encoding::encoder::EncodedBatch, + indices: &[u64], + cache: Arc, +) -> RecordBatch { + let io_scheduler = Arc::new(BufferScheduler::new(encoded.data.clone())) as Arc; + let filter = FilterExpression::no_filter(); + let mut decode_scheduler = DecodeBatchScheduler::try_new( + encoded.schema.as_ref(), + &encoded.top_level_columns, + &encoded.page_table, + &vec![], + encoded.num_rows, + Arc::::default(), + io_scheduler.clone(), + cache, + &filter, + &DecoderConfig::default(), + ) + .await + .unwrap(); + let (tx, rx) = unbounded_channel(); + decode_scheduler.schedule_take(indices, &filter, tx, io_scheduler); + let mut stream = create_decode_stream( + &encoded.schema, + indices.len() as u64, + indices.len() as u32, + true, + false, + true, + rx, + None, + ) + .unwrap(); + stream.next().await.unwrap().task.await.unwrap() +} + +fn bench_variable_offsets_decode(c: &mut Criterion) { + const NUM_ROWS: usize = 262_144; + const NUM_TAKES: usize = 512; + + let metadata = HashMap::from([ + ( + "lance-encoding:structural-encoding".to_string(), + "miniblock".to_string(), + ), + ( + "lance-encoding:dict-divisor".to_string(), + "100000".to_string(), + ), + ("lance-encoding:compression".to_string(), "none".to_string()), + ]); + let schema = Arc::new(Schema::new(vec![ + Field::new("value", DataType::Utf8, false).with_metadata(metadata), + ])); + let lance_schema = Arc::new(lance_core::datatypes::Schema::try_from(schema.as_ref()).unwrap()); + let corpora = [ + ( + "range", + Arc::new(StringArray::from_iter_values( + (0..NUM_ROWS).map(|index| format!("row_{index:012}")), + )) as Arc, + ), + ( + "delta", + Arc::new(StringArray::from_iter_values( + (0..NUM_ROWS).map(|index| "x".repeat(4 + index % 64)), + )) as Arc, + ), + ]; + let mut take_indices = (0..NUM_TAKES) + .map(|index| (index as u64).wrapping_mul(104_729).wrapping_add(8_191) % NUM_ROWS as u64) + .collect::>(); + take_indices.sort_unstable(); + take_indices.dedup(); + + let rt = tokio::runtime::Runtime::new().unwrap(); + let mut group = c.benchmark_group("variable_offsets_decode"); + + for (workload, array) in corpora { + let batch = RecordBatch::try_new(schema.clone(), vec![array]).unwrap(); + for version in [LanceFileVersion::V2_2, LanceFileVersion::V2_3] { + let strategy = default_encoding_strategy(version); + let options = EncodingOptions { + version, + ..Default::default() + }; + let encoded = rt + .block_on(encode_batch( + &batch, + lance_schema.clone(), + strategy.as_ref(), + &options, + )) + .unwrap(); + let encoded_bytes = encoded.data.len(); + let benchmark_suffix = format!("{workload}/{version}/{encoded_bytes}B"); + + group.throughput(criterion::Throughput::Elements(NUM_ROWS as u64)); + group.bench_function(format!("scan/{benchmark_suffix}"), |bencher| { + bencher.iter(|| { + let decoded = rt + .block_on(lance_encoding::decoder::decode_batch( + &encoded, + &FilterExpression::no_filter(), + Arc::::default(), + false, + version, + Some(Arc::new(LanceCache::no_cache())), + )) + .unwrap(); + assert_eq!(decoded.num_rows(), NUM_ROWS); + }) + }); + + for cache_mode in ["cold", "warm"] { + let cache = if cache_mode == "cold" { + Arc::new(LanceCache::no_cache()) + } else { + Arc::new(LanceCache::with_capacity(64 * 1024 * 1024)) + }; + group.throughput(criterion::Throughput::Elements(take_indices.len() as u64)); + group.bench_function(format!("take/{benchmark_suffix}/{cache_mode}"), |bencher| { + bencher.iter(|| { + let decoded = + rt.block_on(decode_take(&encoded, &take_indices, cache.clone())); + assert_eq!(decoded.num_rows(), take_indices.len()); + }) + }); + } + } + } +} + #[cfg(target_os = "linux")] criterion_group!( name=benches; @@ -570,7 +705,7 @@ criterion_group!( .with_profiler(lance_testing::pprof::PProfProfiler::new(100, lance_testing::pprof::Output::Flamegraph(None))); targets = bench_decode, bench_decode_fsl, bench_decode_str_with_dict_encoding, bench_decode_packed_struct, bench_decode_str_with_fixed_size_binary_encoding, bench_decode_compressed, - bench_decode_compressed_parallel); + bench_decode_compressed_parallel, bench_variable_offsets_decode); // Non-linux version does not support pprof. #[cfg(not(target_os = "linux"))] @@ -578,5 +713,6 @@ criterion_group!( name=benches; config = Criterion::default().significance_level(0.1).sample_size(10); targets = bench_decode, bench_decode_fsl, bench_decode_str_with_dict_encoding, bench_decode_packed_struct, - bench_decode_compressed, bench_decode_compressed_parallel); + bench_decode_compressed, bench_decode_compressed_parallel, + bench_variable_offsets_decode); criterion_main!(benches); diff --git a/rust/lance-encoding/benches/encoder.rs b/rust/lance-encoding/benches/encoder.rs index 08eb89d32fe..e766a557339 100644 --- a/rust/lance-encoding/benches/encoder.rs +++ b/rust/lance-encoding/benches/encoder.rs @@ -1,12 +1,12 @@ // SPDX-License-Identifier: Apache-2.0 // SPDX-FileCopyrightText: Copyright The Lance Authors -use std::{collections::HashMap, sync::Arc}; +use std::{collections::HashMap, hint::black_box, sync::Arc}; -use arrow_array::{ArrayRef, BooleanArray, ListArray, RecordBatch}; +use arrow_array::{ArrayRef, BooleanArray, ListArray, RecordBatch, StringArray}; use arrow_buffer::{OffsetBuffer, ScalarBuffer}; use arrow_schema::{DataType, Field, Schema}; -use criterion::{Criterion, criterion_group, criterion_main}; +use criterion::{Criterion, Throughput, criterion_group, criterion_main}; use lance_encoding::{ encoder::{EncodingOptions, default_encoding_strategy, encode_batch}, version::LanceFileVersion, @@ -162,17 +162,90 @@ fn bench_encode_structural_pages(c: &mut Criterion) { }); } +fn bench_variable_offsets(c: &mut Criterion) { + const NUM_ROWS: usize = 262_144; + let metadata = HashMap::from([ + ( + "lance-encoding:structural-encoding".to_string(), + "miniblock".to_string(), + ), + ( + "lance-encoding:dict-divisor".to_string(), + "100000".to_string(), + ), + ("lance-encoding:compression".to_string(), "none".to_string()), + ]); + let schema = Arc::new(Schema::new(vec![ + Field::new("value", DataType::Utf8, false).with_metadata(metadata), + ])); + let lance_schema = Arc::new(lance_core::datatypes::Schema::try_from(schema.as_ref()).unwrap()); + let runtime = tokio::runtime::Runtime::new().unwrap(); + let mut group = c.benchmark_group("variable_offsets"); + group.throughput(Throughput::Elements(NUM_ROWS as u64)); + let corpora: [(&str, ArrayRef); 2] = [ + ( + "range", + Arc::new(StringArray::from_iter_values( + (0..NUM_ROWS).map(|index| format!("row_{index:012}")), + )), + ), + ( + "delta", + Arc::new(StringArray::from_iter_values( + (0..NUM_ROWS).map(|index| "x".repeat(4 + index % 64)), + )), + ), + ]; + for (workload, array) in corpora { + let batch = RecordBatch::try_new(schema.clone(), vec![array]).unwrap(); + for version in [LanceFileVersion::V2_2, LanceFileVersion::V2_3] { + let strategy = default_encoding_strategy(version); + let options = EncodingOptions { + version, + ..Default::default() + }; + let encoded_bytes = runtime + .block_on(encode_batch( + &batch, + lance_schema.clone(), + strategy.as_ref(), + &options, + )) + .unwrap() + .data + .len(); + group.bench_function( + format!("{workload}/{version}/{encoded_bytes}B"), + |bencher| { + bencher.iter(|| { + black_box( + runtime + .block_on(encode_batch( + &batch, + lance_schema.clone(), + strategy.as_ref(), + &options, + )) + .unwrap(), + ) + }) + }, + ); + } + } +} + #[cfg(target_os = "linux")] criterion_group!( name=benches; config = Criterion::default().significance_level(0.1).sample_size(10) .with_profiler(lance_testing::pprof::PProfProfiler::new(100, lance_testing::pprof::Output::Flamegraph(None))); - targets = bench_encode_compressed, bench_encode_structural_pages); + targets = bench_encode_compressed, bench_encode_structural_pages, bench_variable_offsets); #[cfg(not(target_os = "linux"))] criterion_group!( name=benches; config = Criterion::default().significance_level(0.1).sample_size(10); - targets = bench_encode_compressed, bench_encode_structural_pages); + targets = bench_encode_compressed, bench_encode_structural_pages, bench_variable_offsets); criterion_main!(benches);