diff --git a/Cargo.toml b/Cargo.toml index 40e99d6da..ead328b92 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -193,3 +193,12 @@ harness = false [[bench]] name = "str_search_and_get" harness = false + +[[bench]] +name = "merge_segments" +harness = false + +[[bench]] +name = "regex_all_terms" +harness = false + diff --git a/benches/merge_segments.rs b/benches/merge_segments.rs new file mode 100644 index 000000000..7e911f7ec --- /dev/null +++ b/benches/merge_segments.rs @@ -0,0 +1,224 @@ +// Benchmarks segment merging +// +// Notes: +// - Input segments are kept intact (no deletes / no IndexWriter merge). +// - Output is written to a `NullDirectory` that discards all files except +// fieldnorms (needed for merging). + +use std::collections::HashMap; +use std::io::{self, Write}; +use std::path::{Path, PathBuf}; +use std::sync::{Arc, RwLock}; + +use binggan::{black_box, BenchRunner}; +use rand::prelude::*; +use rand::rngs::StdRng; +use rand::SeedableRng; +use tantivy::directory::error::{DeleteError, OpenReadError, OpenWriteError}; +use tantivy::directory::{ + AntiCallToken, Directory, FileHandle, OwnedBytes, TerminatingWrite, WatchCallback, WatchHandle, + WritePtr, +}; +use tantivy::indexer::{merge_filtered_segments, NoMergePolicy}; +use tantivy::schema::{Schema, TEXT}; +use tantivy::{doc, HasLen, Index, IndexSettings, Segment}; + +#[derive(Clone, Default, Debug)] +struct NullDirectory { + blobs: Arc>>, +} + +struct NullWriter; + +impl Write for NullWriter { + fn write(&mut self, buf: &[u8]) -> io::Result { + Ok(buf.len()) + } + + fn flush(&mut self) -> io::Result<()> { + Ok(()) + } +} + +impl TerminatingWrite for NullWriter { + fn terminate_ref(&mut self, _token: AntiCallToken) -> io::Result<()> { + Ok(()) + } +} + +struct InMemoryWriter { + path: PathBuf, + buffer: Vec, + blobs: Arc>>, +} + +impl Write for InMemoryWriter { + fn write(&mut self, buf: &[u8]) -> io::Result { + self.buffer.extend_from_slice(buf); + Ok(buf.len()) + } + + fn flush(&mut self) -> io::Result<()> { + Ok(()) + } +} + +impl TerminatingWrite for InMemoryWriter { + fn terminate_ref(&mut self, _token: AntiCallToken) -> io::Result<()> { + let bytes = OwnedBytes::new(std::mem::take(&mut self.buffer)); + self.blobs.write().unwrap().insert(self.path.clone(), bytes); + Ok(()) + } +} + +#[derive(Debug, Default)] +struct NullFileHandle; +impl HasLen for NullFileHandle { + fn len(&self) -> usize { + 0 + } +} +impl FileHandle for NullFileHandle { + fn read_bytes(&self, _range: std::ops::Range) -> io::Result { + unimplemented!() + } +} + +impl Directory for NullDirectory { + fn get_file_handle(&self, path: &Path) -> Result, OpenReadError> { + if let Some(bytes) = self.blobs.read().unwrap().get(path) { + return Ok(Arc::new(bytes.clone())); + } + Ok(Arc::new(NullFileHandle)) + } + + fn delete(&self, _path: &Path) -> Result<(), DeleteError> { + Ok(()) + } + + fn exists(&self, _path: &Path) -> Result { + Ok(true) + } + + fn open_write(&self, path: &Path) -> Result { + let path_buf = path.to_path_buf(); + if path.to_string_lossy().ends_with(".fieldnorm") { + let writer = InMemoryWriter { + path: path_buf, + buffer: Vec::new(), + blobs: Arc::clone(&self.blobs), + }; + Ok(io::BufWriter::new(Box::new(writer))) + } else { + Ok(io::BufWriter::new(Box::new(NullWriter))) + } + } + + fn atomic_read(&self, path: &Path) -> Result, OpenReadError> { + if let Some(bytes) = self.blobs.read().unwrap().get(path) { + return Ok(bytes.as_slice().to_vec()); + } + Err(OpenReadError::FileDoesNotExist(path.to_path_buf())) + } + + fn atomic_write(&self, _path: &Path, _data: &[u8]) -> io::Result<()> { + Ok(()) + } + + fn sync_directory(&self) -> io::Result<()> { + Ok(()) + } + + fn watch(&self, _watch_callback: WatchCallback) -> tantivy::Result { + Ok(WatchHandle::empty()) + } +} + +struct MergeScenario { + #[allow(dead_code)] + index: Index, + segments: Vec, + settings: IndexSettings, + label: String, +} + +fn build_index( + num_segments: usize, + docs_per_segment: usize, + tokens_per_doc: usize, + vocab_size: usize, +) -> MergeScenario { + let mut schema_builder = Schema::builder(); + let body = schema_builder.add_text_field("body", TEXT); + let schema = schema_builder.build(); + let index = Index::create_in_ram(schema.clone()); + + assert!(vocab_size > 0); + let total_tokens = num_segments * docs_per_segment * tokens_per_doc; + let use_unique_terms = vocab_size >= total_tokens; + let mut rng = StdRng::from_seed([7u8; 32]); + let mut next_token_id: u64 = 0; + + { + let mut writer = index.writer_with_num_threads(1, 256_000_000).unwrap(); + writer.set_merge_policy(Box::new(NoMergePolicy)); + for _ in 0..num_segments { + for _ in 0..docs_per_segment { + let mut tokens = Vec::with_capacity(tokens_per_doc); + for _ in 0..tokens_per_doc { + let token_id = if use_unique_terms { + let id = next_token_id; + next_token_id += 1; + id + } else { + rng.random_range(0..vocab_size as u64) + }; + tokens.push(format!("term_{token_id}")); + } + writer.add_document(doc!(body => tokens.join(" "))).unwrap(); + } + writer.commit().unwrap(); + } + } + + let segments = index.searchable_segments().unwrap(); + let settings = index.settings().clone(); + let label = format!( + "segments={}, docs/seg={}, tokens/doc={}, vocab={}", + num_segments, docs_per_segment, tokens_per_doc, vocab_size + ); + + MergeScenario { + index, + segments, + settings, + label, + } +} + +fn main() { + let scenarios = vec![ + build_index(8, 50_000, 12, 8), + build_index(16, 50_000, 12, 8), + build_index(16, 100_000, 12, 8), + build_index(8, 50_000, 8, 8 * 50_000 * 8), + ]; + + let mut runner = BenchRunner::new(); + for scenario in scenarios { + let mut group = runner.new_group(); + group.set_name(format!("merge_segments inv_index — {}", scenario.label)); + let segments = scenario.segments.clone(); + let settings = scenario.settings.clone(); + group.register("merge", move |_| { + let output_dir = NullDirectory::default(); + let filter_doc_ids = vec![None; segments.len()]; + let merged_index = + merge_filtered_segments(&segments, settings.clone(), filter_doc_ids, output_dir) + .unwrap(); + black_box(merged_index); + }); + + group.run(); + } +} diff --git a/benches/regex_all_terms.rs b/benches/regex_all_terms.rs new file mode 100644 index 000000000..c97c7b0e5 --- /dev/null +++ b/benches/regex_all_terms.rs @@ -0,0 +1,113 @@ +// Benchmarks regex query that matches all terms in a synthetic index. +// +// Corpus model: +// - N unique terms: t000000, t000001, ... +// - M docs +// - K tokens per doc: doc i gets terms derived from (i, token_index) +// +// Query: +// - Regex "t.*" to match all terms +// +// Run with: +// - cargo bench --bench regex_all_terms +// + +use std::fmt::Write; + +use binggan::{black_box, BenchRunner}; +use tantivy::collector::Count; +use tantivy::query::RegexQuery; +use tantivy::schema::{Schema, TEXT}; +use tantivy::{doc, Index, ReloadPolicy}; + +const HEAP_SIZE_BYTES: usize = 200_000_000; + +#[derive(Clone, Copy)] +struct BenchConfig { + num_terms: usize, + num_docs: usize, + tokens_per_doc: usize, +} + +fn main() { + let configs = default_configs(); + + let mut runner = BenchRunner::new(); + for config in configs { + let (index, text_field) = build_index(config, HEAP_SIZE_BYTES); + let reader = index + .reader_builder() + .reload_policy(ReloadPolicy::Manual) + .try_into() + .expect("reader"); + let searcher = reader.searcher(); + let query = RegexQuery::from_pattern("t.*", text_field).expect("regex query"); + + let mut group = runner.new_group(); + group.set_name(format!( + "regex_all_terms_t{}_d{}_k{}", + config.num_terms, config.num_docs, config.tokens_per_doc + )); + group.register("regex_count", move |_| { + let count = searcher.search(&query, &Count).expect("search"); + black_box(count); + }); + group.run(); + } +} + +fn default_configs() -> Vec { + vec![ + BenchConfig { + num_terms: 10_000, + num_docs: 100_000, + tokens_per_doc: 1, + }, + BenchConfig { + num_terms: 10_000, + num_docs: 100_000, + tokens_per_doc: 8, + }, + BenchConfig { + num_terms: 100_000, + num_docs: 100_000, + tokens_per_doc: 1, + }, + BenchConfig { + num_terms: 100_000, + num_docs: 100_000, + tokens_per_doc: 8, + }, + ] +} + +fn build_index(config: BenchConfig, heap_size_bytes: usize) -> (Index, tantivy::schema::Field) { + let mut schema_builder = Schema::builder(); + let text_field = schema_builder.add_text_field("text", TEXT); + let schema = schema_builder.build(); + let index = Index::create_in_ram(schema); + + let term_width = config.num_terms.to_string().len(); + { + let mut writer = index + .writer_with_num_threads(1, heap_size_bytes) + .expect("writer"); + let mut buffer = String::new(); + for doc_id in 0..config.num_docs { + buffer.clear(); + for token_idx in 0..config.tokens_per_doc { + if token_idx > 0 { + buffer.push(' '); + } + let term_id = (doc_id * config.tokens_per_doc + token_idx) % config.num_terms; + write!(&mut buffer, "t{term_id:0term_width$}").expect("write token"); + } + writer + .add_document(doc!(text_field => buffer.as_str())) + .expect("add_document"); + } + writer.commit().expect("commit"); + } + + (index, text_field) +}