From ffdd6d09a67f7c71bc109229c41f91889c1df2a0 Mon Sep 17 00:00:00 2001 From: dennis zhuang Date: Mon, 28 Sep 2026 08:22:42 +0000 Subject: [PATCH] perf(index): cut allocations when building bloom and inverted indexes (#9359) * perf(index): build bloom filters from element hashes without per-token allocation Signed-off-by: Dennis Zhuang * perf(index): speed up inverted index building with hashed buffers and sync pushes Signed-off-by: Dennis Zhuang * chore(index): use BuildHasher::hash_one for element hashes Signed-off-by: Dennis Zhuang * chore(index): require callers to act on spill requests Signed-off-by: Dennis Zhuang * fix(index): insert bloom hashes one by one to keep segment set capacity bounded Signed-off-by: Dennis Zhuang * test(index): add index build and bloom search benchmarks Signed-off-by: Dennis Zhuang * test(index): keep applier setup out of the bloom search benchmark Signed-off-by: Dennis Zhuang * fix(index): skip empty spills and keep the old inverted sort memory estimate Signed-off-by: Dennis Zhuang * fix(index): stream fulltext token hashes and test spill dispatch Signed-off-by: Dennis Zhuang * docs(index): note non-ASCII tokens still allocate in analyze_text_hashes Signed-off-by: Dennis Zhuang * docs(index): narrow the allocation note to case-insensitive non-ASCII tokens Signed-off-by: Dennis Zhuang --------- Signed-off-by: Dennis Zhuang --- Cargo.lock | 1 + src/index/Cargo.toml | 5 + src/index/benches/index_build_bench.rs | 424 ++++++++++++++++++ src/index/src/bloom_filter.rs | 68 +++ src/index/src/bloom_filter/applier.rs | 17 +- src/index/src/bloom_filter/creator.rs | 164 ++++--- .../bloom_filter/creator/finalize_segment.rs | 21 +- src/index/src/bloom_filter/reader.rs | 7 +- .../src/fulltext_index/create/bloom_filter.rs | 5 +- src/index/src/fulltext_index/tokenizer.rs | 47 ++ src/index/src/inverted_index/create.rs | 34 +- src/index/src/inverted_index/create/sort.rs | 16 +- .../create/sort/external_sort.rs | 131 ++++-- .../src/inverted_index/create/sort_create.rs | 115 +++-- .../src/sst/index/inverted_index/creator.rs | 52 ++- 15 files changed, 879 insertions(+), 228 deletions(-) create mode 100644 src/index/benches/index_build_bench.rs diff --git a/Cargo.lock b/Cargo.lock index 597ad270305..bb7f69d1551 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -6813,6 +6813,7 @@ dependencies = [ name = "index" version = "1.3.0-alpha.1" dependencies = [ + "ahash 0.8.12", "async-trait", "asynchronous-codec", "bytemuck", diff --git a/src/index/Cargo.toml b/src/index/Cargo.toml index bb1225c12c8..078f9e42c56 100644 --- a/src/index/Cargo.toml +++ b/src/index/Cargo.toml @@ -8,6 +8,7 @@ license.workspace = true workspace = true [dependencies] +ahash.workspace = true async-trait.workspace = true asynchronous-codec = "0.7.0" bytemuck.workspace = true @@ -57,3 +58,7 @@ harness = false [[bench]] name = "bytes_to_u64_vec" harness = false + +[[bench]] +name = "index_build_bench" +harness = false diff --git a/src/index/benches/index_build_bench.rs b/src/index/benches/index_build_bench.rs new file mode 100644 index 00000000000..4c6c3fac020 --- /dev/null +++ b/src/index/benches/index_build_bench.rs @@ -0,0 +1,424 @@ +// Copyright 2023 Greptime Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +//! Index build and bloom search benchmarks over synthetic observability data. + +use std::collections::{BTreeSet, HashMap}; +use std::hint::black_box; +use std::num::NonZeroUsize; +use std::path::PathBuf; +use std::sync::Arc; +use std::sync::atomic::AtomicUsize; + +use async_trait::async_trait; +use criterion::{BatchSize, Criterion, Throughput, criterion_group, criterion_main}; +use futures::{AsyncRead, AsyncReadExt}; +use index::bitmap::BitmapType; +use index::bloom_filter::applier::{BloomFilterApplier, InListPredicate}; +use index::bloom_filter::creator::BloomFilterCreator; +use index::bloom_filter::reader::BloomFilterReaderImpl; +use index::external_provider::{ExternalTempFileProvider, Reader, Writer}; +use index::fulltext_index::Config; +use index::fulltext_index::create::{BloomFilterFulltextIndexCreator, FulltextIndexCreator}; +use index::fulltext_index::tokenizer::{Analyzer, EnglishTokenizer}; +use index::inverted_index::create::InvertedIndexCreator; +use index::inverted_index::create::sort::external_sort::ExternalSorter; +use index::inverted_index::create::sort_create::SortIndexCreator; +use index::inverted_index::format::writer::InvertedIndexBlobWriter; +use puffin::puffin_manager::{PuffinWriter, PutOptions}; +use rand::seq::IndexedRandom; +use rand::{Rng, SeedableRng}; +use rand_chacha::ChaCha8Rng; + +const ROWS: usize = 100_000; +const SEGMENT_ROWS: usize = 10240; +const ROW_GROUP_ROWS: usize = 102400; + +/// Drains blobs so the bloom fulltext creator can finish without a puffin file. +struct DrainPuffinWriter; + +#[async_trait] +impl PuffinWriter for DrainPuffinWriter { + async fn put_blob( + &mut self, + _key: &str, + raw_data: R, + _options: PutOptions, + _properties: HashMap, + ) -> puffin::error::Result + where + R: AsyncRead + Send, + { + let mut buf = Vec::new(); + Box::pin(raw_data).read_to_end(&mut buf).await.unwrap(); + Ok(buf.len() as u64) + } + + async fn put_dir( + &mut self, + _key: &str, + _dir: PathBuf, + _options: PutOptions, + _properties: HashMap, + ) -> puffin::error::Result { + unreachable!("bloom fulltext index only writes blobs") + } + + fn set_footer_lz4_compressed(&mut self, _lz4_compressed: bool) {} + + async fn finish(self) -> puffin::error::Result { + Ok(0) + } +} + +/// The benchmarks set no memory limit, so nothing is spilled. +struct NoSpill; + +#[async_trait] +impl ExternalTempFileProvider for NoSpill { + async fn create(&self, _: &str, _: &str) -> Result { + unreachable!("no memory limit is set") + } + + async fn read_all(&self, _: &str) -> Result, index::error::Error> { + Ok(vec![]) + } +} + +fn uuid(rng: &mut ChaCha8Rng) -> String { + let h = format!("{:032x}", rng.random::()); + format!( + "{}-{}-{}-{}-{}", + &h[0..8], + &h[8..12], + &h[12..16], + &h[16..20], + &h[20..32] + ) +} + +/// Access, auth, timeout and error log lines with request ids, IPs and numbers: most +/// tokens of a segment are distinct. +fn service_logs(rows: usize) -> Vec { + let mut rng = ChaCha8Rng::seed_from_u64(1); + let paths = ["users", "orders", "items", "carts", "payments", "sessions"]; + (0..rows) + .map(|_| { + let ip = format!( + "10.{}.{}.{}", + rng.random_range(0..16), + rng.random_range(0..256), + rng.random_range(1..255) + ); + match rng.random_range(0..10) { + 0..=5 => format!( + "INFO GET /api/v1/{}/{}/detail?page={} 200 {}ms request_id={} client={}", + paths.choose(&mut rng).unwrap(), + rng.random_range(0..200000), + rng.random_range(0..50), + rng.random_range(1..2000), + uuid(&mut rng), + ip + ), + 6..=7 => format!( + "WARN connection to {}:{} timed out after {}ms retry={}", + ip, + rng.random_range(1024..65535), + rng.random_range(100..30000), + rng.random_range(0..5) + ), + _ => format!( + "ERROR failed to process message {} from queue orders-{}: deadline exceeded", + uuid(&mut rng), + rng.random_range(0..64) + ), + } + }) + .collect() +} + +/// Java stack traces: long lines where a few frame names repeat many times. +fn stack_trace_logs(rows: usize) -> Vec { + let mut rng = ChaCha8Rng::seed_from_u64(2); + let frames = [ + "at org.apache.kafka.clients.consumer.KafkaConsumer.poll(KafkaConsumer.java:1250)", + "at com.example.orders.OrderService.process(OrderService.java:88)", + "at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1136)", + "at java.base/java.lang.Thread.run(Thread.java:833)", + ]; + (0..rows) + .map(|_| { + let mut line = format!( + "ERROR request {} failed: java.lang.IllegalStateException: retry exhausted", + uuid(&mut rng) + ); + for _ in 0..40 { + line.push(' '); + line.push_str(frames.choose(&mut rng).unwrap()); + } + line + }) + .collect() +} + +/// Trace ids with temporal locality: each trace spreads over about a hundred rows. +fn trace_ids(rows: usize) -> Vec { + let mut rng = ChaCha8Rng::seed_from_u64(3); + let traces = (0..rows / 10) + .map(|_| format!("{:032x}", rng.random::())) + .collect::>(); + (0..rows) + .map(|i| { + let j = (i / 10) as i64 + rng.random_range(-50i64..=50); + traces[j.clamp(0, traces.len() as i64 - 1) as usize].clone() + }) + .collect() +} + +/// Prometheus-like series sorted by labels; each series has `samples` rows. +fn series(num_series: usize) -> Vec<[String; 6]> { + let mut rng = ChaCha8Rng::seed_from_u64(4); + let mut series = (0..num_series) + .map(|s| { + let pod = s / 2; + let ns = pod % 400; + [ + format!("region-{}", ns % 4), + format!("az-{}", ns % 12), + format!("cluster-{}", ns % 40), + format!("namespace-{ns}"), + format!("pod-{:05x}-{pod}", rng.random_range(0..0xfffff)), + format!("container-{}", s % 2), + ] + }) + .collect::>(); + series.sort(); + series +} + +fn bloom_creator() -> BloomFilterCreator { + BloomFilterCreator::new( + SEGMENT_ROWS, + 0.01, + Arc::new(NoSpill), + Arc::new(AtomicUsize::new(0)), + None, + ) +} + +fn bench_bloom_fulltext_build(c: &mut Criterion) { + let rt = tokio::runtime::Runtime::new().unwrap(); + let mut group = c.benchmark_group("bloom_fulltext_build"); + group.sample_size(10); + + for (name, lines) in [ + ("service_logs", service_logs(ROWS)), + ("stack_traces", stack_trace_logs(ROWS / 10)), + ] { + group.throughput(Throughput::Elements(lines.len() as u64)); + group.bench_function(name, |b| { + b.iter_batched( + || { + BloomFilterFulltextIndexCreator::new( + Config::default(), + SEGMENT_ROWS, + 0.01, + Arc::new(NoSpill), + Arc::new(AtomicUsize::new(0)), + None, + ) + }, + |mut creator| { + let lines = &lines; + rt.block_on(async move { + for line in lines { + creator.push_text(line).await.unwrap(); + } + let written = creator + .finish(&mut DrainPuffinWriter, "blob", PutOptions::default()) + .await + .unwrap(); + black_box(written) + }) + }, + BatchSize::LargeInput, + ) + }); + } + group.finish(); +} + +fn bench_bloom_skipping_build(c: &mut Criterion) { + let rt = tokio::runtime::Runtime::new().unwrap(); + let mut group = c.benchmark_group("bloom_skipping_build"); + group.sample_size(10); + + let trace_ids = trace_ids(ROWS * 10); + group.throughput(Throughput::Elements(trace_ids.len() as u64)); + group.bench_function("trace_id", |b| { + b.iter_batched( + bloom_creator, + |mut creator| { + let trace_ids = &trace_ids; + rt.block_on(async move { + for id in trace_ids { + creator + .push_n_row_elem(1, Some(id.as_bytes())) + .await + .unwrap(); + } + let mut out = Vec::new(); + creator.finish(&mut out).await.unwrap(); + black_box(out) + }) + }, + BatchSize::LargeInput, + ) + }); + + // A sorted tag pushed as one run per series, as the flat SST path does. + let series = series(ROWS); + let samples = 10; + group.throughput(Throughput::Elements((series.len() * samples) as u64)); + group.bench_function("sorted_tag_runs", |b| { + b.iter_batched( + bloom_creator, + |mut creator| { + let series = &series; + rt.block_on(async move { + for s in series { + creator + .push_n_row_elem(samples, Some(s[4].as_bytes())) + .await + .unwrap(); + } + let mut out = Vec::new(); + creator.finish(&mut out).await.unwrap(); + black_box(out) + }) + }, + BatchSize::LargeInput, + ) + }); + group.finish(); +} + +fn bench_inverted_build(c: &mut Criterion) { + let rt = tokio::runtime::Runtime::new().unwrap(); + let mut group = c.benchmark_group("inverted_build"); + group.sample_size(10); + + let series = series(ROWS); + let samples = 10; + let names = ["region", "az", "cluster", "namespace", "pod", "container"]; + group.throughput(Throughput::Elements((series.len() * samples) as u64)); + group.bench_function("six_tags", |b| { + b.iter_batched( + || { + let factory = ExternalSorter::factory( + Arc::new(NoSpill), + None, + Arc::new(AtomicUsize::new(0)), + None, + ); + SortIndexCreator::new(factory, NonZeroUsize::new(1024).unwrap()) + }, + |mut creator| { + let series = &series; + rt.block_on(async move { + for s in series { + for (name, value) in names.iter().zip(s) { + let spill = + creator.push_with_name_n(name, Some(value.as_bytes()), samples); + assert!(!spill, "no memory limit is set"); + } + } + let mut out = Vec::new(); + let mut writer = InvertedIndexBlobWriter::new(&mut out); + creator + .finish(&mut writer, BitmapType::Roaring) + .await + .unwrap(); + black_box(out) + }) + }, + BatchSize::LargeInput, + ) + }); + group.finish(); +} + +#[allow(clippy::single_range_in_vec_init)] +fn bench_bloom_search(c: &mut Criterion) { + let rt = tokio::runtime::Runtime::new().unwrap(); + let mut group = c.benchmark_group("bloom_search"); + + let lines = service_logs(ROWS * 10); + let analyzer = Analyzer::new(Box::new(EnglishTokenizer), false); + let blob = rt.block_on(async { + let mut creator = bloom_creator(); + for line in &lines { + creator + .push_row_elems(analyzer.analyze_text(line).unwrap()) + .await + .unwrap(); + } + let mut out = Vec::new(); + creator.finish(&mut out).await.unwrap(); + out + }); + + let num_rows = lines.len(); + // Each row group tests every segment's filter against the three probes. + let predicates = ["timed", "deadline", "nonexistent"] + .iter() + .map(|term| InListPredicate { + list: BTreeSet::from([term.as_bytes().to_vec()]), + }) + .collect::>(); + group.bench_function("three_terms_per_row_group", |b| { + b.iter_batched( + || { + let reader = BloomFilterReaderImpl::new(blob.clone()); + rt.block_on(BloomFilterApplier::new(Box::new(reader))) + .unwrap() + }, + |mut applier| { + rt.block_on(async { + let mut matched = 0; + for start in (0..num_rows).step_by(ROW_GROUP_ROWS) { + let end = (start + ROW_GROUP_ROWS).min(num_rows); + matched += applier + .search(&predicates, &[start..end], None) + .await + .unwrap() + .len(); + } + black_box(matched) + }) + }, + BatchSize::SmallInput, + ) + }); + group.finish(); +} + +criterion_group!( + benches, + bench_bloom_fulltext_build, + bench_bloom_skipping_build, + bench_inverted_build, + bench_bloom_search +); +criterion_main!(benches); diff --git a/src/index/src/bloom_filter.rs b/src/index/src/bloom_filter.rs index eb818a0f5a0..156be09e69e 100644 --- a/src/index/src/bloom_filter.rs +++ b/src/index/src/bloom_filter.rs @@ -17,5 +17,73 @@ pub mod creator; pub mod error; pub mod reader; +use std::hash::{BuildHasher, BuildHasherDefault, Hasher}; +use std::sync::LazyLock; + /// The seed used for the Bloom filter. pub const SEED: u128 = 42; + +static ELEMENT_HASHER: LazyLock = + LazyLock::new(|| fastbloom::DefaultHasher::seeded(&SEED.to_be_bytes())); + +/// Returns the hash fastbloom derives for `elem` with the persisted [`SEED`]. +/// +/// A filter built from these hashes with [`PrehashedBuildHasher`] has exactly the same +/// bits as one built by inserting the elements themselves, so files stay readable by +/// both paths. +pub fn element_hash(elem: &[u8]) -> u64 { + ELEMENT_HASHER.hash_one(elem) +} + +/// Hasher that returns an already computed [`element_hash`] unchanged. +#[derive(Default)] +pub struct PrehashedHasher(u64); + +impl Hasher for PrehashedHasher { + fn finish(&self) -> u64 { + self.0 + } + + fn write(&mut self, _bytes: &[u8]) { + unreachable!("PrehashedHasher only accepts u64 hashes") + } + + fn write_u64(&mut self, hash: u64) { + self.0 = hash; + } +} + +pub type PrehashedBuildHasher = BuildHasherDefault; + +/// A persisted bloom filter probed by [`element_hash`] values. +pub type PrehashedBloomFilter = fastbloom::BloomFilter<512, PrehashedBuildHasher>; + +#[cfg(test)] +mod tests { + use fastbloom::BloomFilter; + + use super::*; + + #[test] + fn test_prehashed_filter_matches_seeded_filter() { + // Persisted filters are read back with `.seed(&SEED)`, so building them from + // precomputed hashes must set exactly the same bits. + for count in [0usize, 1, 7, 100, 5000] { + let elems = (0..count) + .map(|i| format!("elem-{i}").into_bytes()) + .chain([Vec::new()]) + .collect::>(); + let mut seeded = BloomFilter::with_false_pos(0.01) + .seed(&SEED) + .expected_items(elems.len()); + let mut prehashed = BloomFilter::with_false_pos(0.01) + .hasher(PrehashedBuildHasher::default()) + .expected_items(elems.len()); + for elem in &elems { + seeded.insert(elem); + prehashed.insert(&element_hash(elem)); + } + assert_eq!(seeded.as_slice(), prehashed.as_slice()); + } + } +} diff --git a/src/index/src/bloom_filter/applier.rs b/src/index/src/bloom_filter/applier.rs index fe1aafec2d6..96a1c6ff17f 100644 --- a/src/index/src/bloom_filter/applier.rs +++ b/src/index/src/bloom_filter/applier.rs @@ -16,13 +16,13 @@ use std::collections::BTreeSet; use std::ops::Range; use std::sync::Arc; -use fastbloom::BloomFilter; use greptime_proto::v1::index::BloomFilterMeta; use itertools::Itertools; use crate::Bytes; use crate::bloom_filter::error::Result; use crate::bloom_filter::reader::{BloomFilterReadMetrics, BloomFilterReader}; +use crate::bloom_filter::{PrehashedBloomFilter, element_hash}; /// Filter bytes one batch of [`BloomFilterApplier::search_groups`] reads. A single row /// group larger than this is still read as one batch, so this is not a memory limit; a @@ -190,7 +190,7 @@ impl BloomFilterApplier { &mut self, segments: &[usize], metrics: Option<&mut BloomFilterReadMetrics>, - ) -> Result<(Vec<(u64, usize)>, Vec)> { + ) -> Result<(Vec<(u64, usize)>, Vec)> { let segment_locations = segments .iter() .map(|&seg| (self.meta.segment_loc_indices[seg], seg)) @@ -215,11 +215,15 @@ impl BloomFilterApplier { fn find_matching_rows( &self, segment_locations: Vec<(u64, usize)>, - bloom_filters: Vec, + bloom_filters: Vec, predicates: &[InListPredicate], ) -> Vec> { let rows_per_segment = self.meta.rows_per_segment as usize; let mut matching_row_ranges = Vec::with_capacity(bloom_filters.len()); + let predicate_hashes = predicates + .iter() + .map(|p| p.list.iter().map(|v| element_hash(v)).collect::>()) + .collect::>(); // Group segments by their location index (since they have the same bloom filter) and check if they match all predicates for ((_loc_index, group), bloom_filter) in segment_locations @@ -229,12 +233,9 @@ impl BloomFilterApplier { .zip(bloom_filters.iter()) { // Check if this bloom filter matches each predicate (AND semantics) - let matches_all_predicates = predicates.iter().all(|predicate| { + let matches_all_predicates = predicate_hashes.iter().all(|hashes| { // For each predicate, at least one probe must match (OR semantics) - predicate - .list - .iter() - .any(|probe| bloom_filter.contains(probe)) + hashes.iter().any(|hash| bloom_filter.contains(hash)) }); if !matches_all_predicates { diff --git a/src/index/src/bloom_filter/creator.rs b/src/index/src/bloom_filter/creator.rs index 07f4f977caf..83c0ae0c697 100644 --- a/src/index/src/bloom_filter/creator.rs +++ b/src/index/src/bloom_filter/creator.rs @@ -26,8 +26,8 @@ use prost::Message; use snafu::ResultExt; use crate::Bytes; -use crate::bloom_filter::SEED; use crate::bloom_filter::error::{IoSnafu, Result}; +use crate::bloom_filter::{PrehashedBuildHasher, element_hash}; use crate::external_provider::ExternalTempFileProvider; /// `BloomFilterCreator` is responsible for creating and managing bloom filters @@ -52,8 +52,11 @@ pub struct BloomFilterCreator { /// Row count that added to the bloom filter so far. accumulated_row_count: usize, - /// A set of distinct elements in the current segment. - cur_seg_distinct_elems: HashSet, + /// Distinct element hashes (see [`element_hash`]) in the current segment. + /// + /// Elements with equal hashes set the same bits, so deduplicating by hash loses + /// nothing and avoids copying the values. + cur_seg_distinct_elems: HashSet, /// The memory usage of the current segment's distinct elements. cur_seg_distinct_elems_mem_usage: usize, @@ -106,98 +109,36 @@ impl BloomFilterCreator { /// reaches `rows_per_segment`, it finalizes the current segment. pub async fn push_n_row_elems( &mut self, - mut nrows: usize, + nrows: usize, elems: impl IntoIterator, ) -> Result<()> { - if nrows == 0 { - return Ok(()); - } - if nrows == 1 { - return self.push_row_elems(elems).await; - } - - let elems = elems.into_iter().collect::>(); - while nrows > 0 { - let rows_to_seg_end = - self.rows_per_segment - (self.accumulated_row_count % self.rows_per_segment); - let rows_to_push = nrows.min(rows_to_seg_end); - nrows -= rows_to_push; - - self.accumulated_row_count += rows_to_push; - - let mut mem_diff = 0; - for elem in &elems { - let len = elem.len(); - let is_new = self.cur_seg_distinct_elems.insert(elem.clone()); - if is_new { - mem_diff += len; - } - } - self.cur_seg_distinct_elems_mem_usage += mem_diff; - self.global_memory_usage - .fetch_add(mem_diff, Ordering::Relaxed); - - if self - .accumulated_row_count - .is_multiple_of(self.rows_per_segment) - { - self.finalize_segment().await?; - self.finalized_row_count = self.accumulated_row_count; - } - } - - Ok(()) + let hashes = elems + .into_iter() + .map(|e| element_hash(&e)) + .collect::>(); + self.push_n_row_hashes(nrows, &hashes).await } - /// Adds `nrows` copies of a single borrowed value (or null), copying it only when - /// it is new to a segment. Row counts advance for nulls as well. - pub async fn push_n_row_elem(&mut self, mut nrows: usize, elem: Option<&[u8]>) -> Result<()> { - while nrows > 0 { - let rows_to_seg_end = - self.rows_per_segment - (self.accumulated_row_count % self.rows_per_segment); - let rows_to_push = nrows.min(rows_to_seg_end); - nrows -= rows_to_push; - self.accumulated_row_count += rows_to_push; - - if let Some(elem) = elem { - let old_len = self.cur_seg_distinct_elems.len(); - // Only allocate when the value is absent from the current segment. - if !self.cur_seg_distinct_elems.contains(elem) { - self.cur_seg_distinct_elems.insert(elem.to_vec()); - } - if self.cur_seg_distinct_elems.len() != old_len { - self.cur_seg_distinct_elems_mem_usage += elem.len(); - self.global_memory_usage - .fetch_add(elem.len(), Ordering::Relaxed); - } - } - if self - .accumulated_row_count - .is_multiple_of(self.rows_per_segment) - { - self.finalize_segment().await?; - self.finalized_row_count = self.accumulated_row_count; - } + /// Adds `nrows` copies of a single borrowed value (or null). Row counts advance for + /// nulls as well. + pub async fn push_n_row_elem(&mut self, nrows: usize, elem: Option<&[u8]>) -> Result<()> { + match elem { + Some(elem) => self.push_n_row_hashes(nrows, &[element_hash(elem)]).await, + None => self.push_n_row_hashes(nrows, &[]).await, } - Ok(()) } /// Adds a row of elements to the bloom filter. If the number of accumulated rows /// reaches `rows_per_segment`, it finalizes the current segment. pub async fn push_row_elems(&mut self, elems: impl IntoIterator) -> Result<()> { - self.accumulated_row_count += 1; + self.push_row_hashes(elems.into_iter().map(|e| element_hash(&e))) + .await + } - let mut mem_diff = 0; - for elem in elems.into_iter() { - let len = elem.len(); - let is_new = self.cur_seg_distinct_elems.insert(elem); - if is_new { - mem_diff += len; - } - } - self.cur_seg_distinct_elems_mem_usage += mem_diff; - self.global_memory_usage - .fetch_add(mem_diff, Ordering::Relaxed); + /// Adds a row of element hashes computed by [`element_hash`]. + pub async fn push_row_hashes(&mut self, hashes: impl IntoIterator) -> Result<()> { + self.accumulated_row_count += 1; + self.insert_hashes(hashes); if self .accumulated_row_count @@ -210,6 +151,43 @@ impl BloomFilterCreator { Ok(()) } + /// Adds `nrows` rows that all contain the element hashes in `hashes`. + pub async fn push_n_row_hashes(&mut self, mut nrows: usize, hashes: &[u64]) -> Result<()> { + while nrows > 0 { + let rows_to_seg_end = + self.rows_per_segment - (self.accumulated_row_count % self.rows_per_segment); + let rows_to_push = nrows.min(rows_to_seg_end); + nrows -= rows_to_push; + self.accumulated_row_count += rows_to_push; + self.insert_hashes(hashes.iter().copied()); + + if self + .accumulated_row_count + .is_multiple_of(self.rows_per_segment) + { + self.finalize_segment().await?; + self.finalized_row_count = self.accumulated_row_count; + } + } + + Ok(()) + } + + fn insert_hashes(&mut self, hashes: impl IntoIterator) { + let old_len = self.cur_seg_distinct_elems.len(); + // Not `extend`: it reserves for the iterator's length, which counts duplicate + // tokens, and the capacity survives `drain` at segment boundaries. + for hash in hashes { + self.cur_seg_distinct_elems.insert(hash); + } + let mem_diff = (self.cur_seg_distinct_elems.len() - old_len) * size_of::(); + if mem_diff > 0 { + self.cur_seg_distinct_elems_mem_usage += mem_diff; + self.global_memory_usage + .fetch_add(mem_diff, Ordering::Relaxed); + } + } + /// Finalizes any remaining segments and writes the bloom filters and metadata to the provided writer. pub async fn finish(&mut self, mut writer: impl AsyncWrite + Unpin) -> Result<()> { if self.accumulated_row_count > self.finalized_row_count { @@ -286,6 +264,7 @@ mod tests { use futures::io::Cursor; use super::*; + use crate::bloom_filter::SEED; use crate::external_provider::MockExternalTempFileProvider; /// Converts a slice of bytes to a vector of `u64`. @@ -296,6 +275,23 @@ mod tests { .collect() } + #[tokio::test] + async fn test_duplicate_hashes_do_not_grow_segment_set() { + let mut creator = BloomFilterCreator::new( + 4, + 0.01, + Arc::new(MockExternalTempFileProvider::new()), + Arc::new(AtomicUsize::new(0)), + None, + ); + creator + .push_row_hashes(std::iter::repeat_n(7, 1_000_000)) + .await + .unwrap(); + assert_eq!(creator.cur_seg_distinct_elems.len(), 1); + assert!(creator.cur_seg_distinct_elems.capacity() < 16); + } + #[tokio::test] async fn test_bloom_filter_creator() { let mut writer = Cursor::new(Vec::new()); diff --git a/src/index/src/bloom_filter/creator/finalize_segment.rs b/src/index/src/bloom_filter/creator/finalize_segment.rs index dd7d03864e4..4cbce66edce 100644 --- a/src/index/src/bloom_filter/creator/finalize_segment.rs +++ b/src/index/src/bloom_filter/creator/finalize_segment.rs @@ -22,8 +22,7 @@ use futures::stream::StreamExt; use futures::{AsyncWriteExt, Stream, stream}; use snafu::ResultExt; -use crate::Bytes; -use crate::bloom_filter::creator::SEED; +use crate::bloom_filter::PrehashedBuildHasher; use crate::bloom_filter::creator::intermediate_codec::IntermediateBloomFilterCodecV1; use crate::bloom_filter::error::{IntermediateSnafu, IoSnafu, Result}; use crate::external_provider::ExternalTempFileProvider; @@ -98,14 +97,14 @@ impl FinalizedBloomFilterStorage { /// If the memory usage exceeds the threshold, flushes the in-memory Bloom filters to disk. pub async fn add( &mut self, - elems: impl IntoIterator, + elem_hashes: impl IntoIterator, element_count: usize, ) -> Result<()> { let mut bf = BloomFilter::with_false_pos(self.false_positive_rate) - .seed(&SEED) + .hasher(PrehashedBuildHasher::default()) .expected_items(element_count); - for elem in elems.into_iter() { - bf.insert(&elem); + for hash in elem_hashes.into_iter() { + bf.insert(&hash); } let fbf = FinalizedBloomFilterSegment::from(bf, element_count); @@ -231,7 +230,7 @@ pub struct FinalizedBloomFilterSegment { } impl FinalizedBloomFilterSegment { - fn from(bf: BloomFilter, elem_count: usize) -> Self { + fn from(bf: BloomFilter<512, S>, elem_count: usize) -> Self { let bf_slice = bf.as_slice(); let mut bloom_filter_bytes = Vec::with_capacity(std::mem::size_of_val(bf_slice)); for &x in bf_slice { @@ -256,6 +255,7 @@ mod tests { use super::*; use crate::bloom_filter::creator::tests::u64_vec_from_bytes; + use crate::bloom_filter::{SEED, element_hash}; use crate::external_provider::MockExternalTempFileProvider; #[tokio::test] @@ -300,11 +300,12 @@ mod tests { let dup_batch = 200; for i in 0..(batch - dup_batch) { - let elems = (elem_count * i..elem_count * (i + 1)).map(|x| x.to_string().into_bytes()); + let elems = (elem_count * i..elem_count * (i + 1)) + .map(|x| element_hash(x.to_string().as_bytes())); storage.add(elems, elem_count).await.unwrap(); } for _ in 0..dup_batch { - storage.add(Some(vec![]), 1).await.unwrap(); + storage.add(Some(element_hash(&[])), 1).await.unwrap(); } // Flush happens. @@ -356,7 +357,7 @@ mod tests { let batch = 1000; for _ in 0..batch { - storage.add(Some(vec![]), 1).await.unwrap(); + storage.add(Some(element_hash(&[])), 1).await.unwrap(); } // Drain the storage. diff --git a/src/index/src/bloom_filter/reader.rs b/src/index/src/bloom_filter/reader.rs index 04eae84f005..15cf957b7b6 100644 --- a/src/index/src/bloom_filter/reader.rs +++ b/src/index/src/bloom_filter/reader.rs @@ -25,10 +25,10 @@ use greptime_proto::v1::index::{BloomFilterLoc, BloomFilterMeta}; use prost::Message; use snafu::{ResultExt, ensure}; -use crate::bloom_filter::SEED; use crate::bloom_filter::error::{ DecodeProtoSnafu, FileSizeTooSmallSnafu, IoSnafu, Result, UnexpectedMetaSizeSnafu, }; +use crate::bloom_filter::{PrehashedBloomFilter, PrehashedBuildHasher, SEED}; /// Minimum size of the bloom filter, which is the size of the length of the bloom filter. const BLOOM_META_LEN_SIZE: u64 = 4; @@ -181,11 +181,12 @@ pub trait BloomFilterReader: Sync { Ok(bm) } + /// Reads multiple bloom filters; probe them with [`crate::bloom_filter::element_hash`]. async fn bloom_filter_vec( &self, locs: &[BloomFilterLoc], metrics: Option<&mut BloomFilterReadMetrics>, - ) -> Result> { + ) -> Result> { let ranges = locs .iter() .map(|l| l.offset..l.offset + l.size) @@ -196,7 +197,7 @@ pub trait BloomFilterReader: Sync { for (bs, loc) in bss.into_iter().zip(locs.iter()) { let vec = bytes_to_u64_vec(&bs); let bm = BloomFilter::from_vec(vec) - .seed(&SEED) + .hasher(PrehashedBuildHasher::default()) .expected_items(loc.element_count as _); result.push(bm); } diff --git a/src/index/src/fulltext_index/create/bloom_filter.rs b/src/index/src/fulltext_index/create/bloom_filter.rs index 28648f9b540..a2db31418eb 100644 --- a/src/index/src/fulltext_index/create/bloom_filter.rs +++ b/src/index/src/fulltext_index/create/bloom_filter.rs @@ -74,11 +74,12 @@ impl BloomFilterFulltextIndexCreator { #[async_trait] impl FulltextIndexCreator for BloomFilterFulltextIndexCreator { async fn push_text(&mut self, text: &str) -> Result<()> { - let tokens = self.analyzer.analyze_text(text)?; + let mut token_buf = Vec::new(); + let hashes = self.analyzer.analyze_text_hashes(text, &mut token_buf); self.inner .as_mut() .context(AbortedSnafu)? - .push_row_elems(tokens) + .push_row_hashes(hashes) .await .map_err(BoxedError::new) .context(ExternalSnafu)?; diff --git a/src/index/src/fulltext_index/tokenizer.rs b/src/index/src/fulltext_index/tokenizer.rs index 3afc826e6f1..e6752a40420 100644 --- a/src/index/src/fulltext_index/tokenizer.rs +++ b/src/index/src/fulltext_index/tokenizer.rs @@ -13,6 +13,7 @@ // limitations under the License. use crate::Bytes; +use crate::bloom_filter::element_hash; use crate::fulltext_index::error::Result; lazy_static::lazy_static! { @@ -140,6 +141,30 @@ impl Analyzer { } } + /// Returns the bloom filter hash of each token in the given text. + /// + /// Equivalent to hashing every token returned by [`Analyzer::analyze_text`] with + /// [`element_hash`]. Only case-insensitive non-ASCII tokens allocate, in `to_lowercase`; + /// case-insensitive ASCII tokens are lowercased in `buf`. + pub fn analyze_text_hashes<'a>( + &self, + text: &'a str, + buf: &'a mut Vec, + ) -> impl Iterator + use<'a> { + let case_sensitive = self.case_sensitive; + self.tokenizer.tokenize(text).into_iter().map(move |token| { + if case_sensitive { + element_hash(token.as_bytes()) + } else if token.is_ascii() { + buf.clear(); + buf.extend(token.bytes().map(|b| b.to_ascii_lowercase())); + element_hash(buf) + } else { + element_hash(token.to_lowercase().as_bytes()) + } + }) + } + /// Analyzes the given text into a list of tokens. pub fn analyze_text(&self, text: &str) -> Result> { let res = self @@ -162,6 +187,28 @@ impl Analyzer { mod tests { use super::*; + #[test] + fn test_analyze_text_hashes_matches_analyze_text() { + let text = "Hello, WORLD ship_Ship 清洁表面 ÄÖÜ straße İstanbul x"; + for (tokenizer, case_sensitive) in [ + (Box::new(EnglishTokenizer) as Box, false), + (Box::new(EnglishTokenizer), true), + (Box::new(ChineseTokenizer), false), + ] { + let analyzer = Analyzer::new(tokenizer, case_sensitive); + let expected = analyzer + .analyze_text(text) + .unwrap() + .iter() + .map(|t| element_hash(t)) + .collect::>(); + let hashes = analyzer + .analyze_text_hashes(text, &mut Vec::new()) + .collect::>(); + assert_eq!(expected, hashes); + } + } + #[test] fn test_english_tokenizer() { let tokenizer = EnglishTokenizer; diff --git a/src/index/src/inverted_index/create.rs b/src/index/src/inverted_index/create.rs index dade9705ac8..74da3612232 100644 --- a/src/index/src/inverted_index/create.rs +++ b/src/index/src/inverted_index/create.rs @@ -27,30 +27,24 @@ use crate::inverted_index::format::writer::InvertedIndexWriter; pub trait InvertedIndexCreator: Send { /// Adds a value to the named index. A `None` value represents an absence of data (null) /// - /// - `index_name`: Identifier for the index being built - /// - `value`: The data to be indexed, or `None` for a null entry - /// /// It should be equivalent to calling `push_with_name_n` with `n = 1` - async fn push_with_name( - &mut self, - index_name: &str, - value: Option>, - ) -> Result<()> { - self.push_with_name_n(index_name, value, 1).await + #[must_use = "a true result requires calling `spill` before pushing more"] + fn push_with_name(&mut self, index_name: &str, value: Option>) -> bool { + self.push_with_name_n(index_name, value, 1) } - /// Adds `n` identical values to the named index. `None` values represent absence of data (null) + /// Buffers `n` identical values for the named index. `None` values represent absence of + /// data (null). /// - /// - `index_name`: Identifier for the index being built - /// - `value`: The data to be indexed, or `None` for a null entry - /// - /// It should be equivalent to calling `push_with_name` `n` times - async fn push_with_name_n( - &mut self, - index_name: &str, - value: Option>, - n: usize, - ) -> Result<()>; + /// Returns true when buffered data exceeds the memory limit; the caller must then call + /// [`InvertedIndexCreator::spill`] before pushing more. Pushing is synchronous so the + /// per-row path does not allocate a future. + #[must_use = "a true result requires calling `spill` before pushing more"] + fn push_with_name_n(&mut self, index_name: &str, value: Option>, n: usize) + -> bool; + + /// Moves the buffers that asked for it to external storage. + async fn spill(&mut self) -> Result<()>; /// Finalizes the index creation process, ensuring all data is properly indexed and stored /// in the provided writer diff --git a/src/index/src/inverted_index/create/sort.rs b/src/index/src/inverted_index/create/sort.rs index bdc1ec21f62..fbf81adbaa4 100644 --- a/src/index/src/inverted_index/create/sort.rs +++ b/src/index/src/inverted_index/create/sort.rs @@ -42,15 +42,15 @@ pub struct SortOutput { /// Handles data sorting, supporting incremental input and retrieval of sorted output #[async_trait] pub trait Sorter: Send { - /// Inputs a non-null or null value into the sorter. - /// Should be equivalent to calling `push_n` with n = 1 - async fn push(&mut self, value: Option>) -> Result<()> { - self.push_n(value, 1).await - } + /// Buffers `n` identical non-null or null values in memory. + /// + /// Returns true when the buffer should be spilled with [`Sorter::spill`] before more + /// values are pushed. Kept synchronous so the per-row path does not allocate a future. + #[must_use = "a true result requires calling `spill` before pushing more"] + fn push_n(&mut self, value: Option>, n: usize) -> bool; - /// Pushing n identical non-null or null values into the sorter. - /// Should be equivalent to calling `push` n times - async fn push_n(&mut self, value: Option>, n: usize) -> Result<()>; + /// Moves the in-memory buffer to external storage. + async fn spill(&mut self) -> Result<()>; /// Completes the sorting process and returns the sorted data async fn output(&mut self) -> Result; diff --git a/src/index/src/inverted_index/create/sort/external_sort.rs b/src/index/src/inverted_index/create/sort/external_sort.rs index 12c49ca89de..cef07df9f8f 100644 --- a/src/index/src/inverted_index/create/sort/external_sort.rs +++ b/src/index/src/inverted_index/create/sort/external_sort.rs @@ -12,7 +12,7 @@ // See the License for the specific language governing permissions and // limitations under the License. -use std::collections::{BTreeMap, VecDeque}; +use std::collections::{HashMap, VecDeque}; use std::mem; use std::num::NonZeroUsize; use std::ops::RangeInclusive; @@ -22,6 +22,7 @@ use std::sync::atomic::{AtomicUsize, Ordering}; use async_trait::async_trait; use common_telemetry::{debug, error}; use futures::stream; +use roaring::RoaringBitmap; use snafu::ResultExt; use crate::bitmap::Bitmap; @@ -35,6 +36,28 @@ use crate::inverted_index::create::sort_create::SorterFactory; use crate::inverted_index::error::{IntermediateSnafu, Result}; use crate::{Bytes, BytesRef}; +/// Segments of one value. Rows arrive in order, so segment ids only grow. +struct Posting { + segments: RoaringBitmap, + last_segment: u32, +} + +/// Estimated cost of one segment id in a roaring array container. +const SEGMENT_SIZE: usize = size_of::(); + +impl Posting { + /// Adds `start..=end` and returns the estimated memory growth. + fn push(&mut self, start: u32, end: u32) -> usize { + if end <= self.last_segment { + return 0; + } + let start = start.max(self.last_segment + 1); + self.segments.insert_range(start..=end); + self.last_segment = end; + (end - start + 1) as usize * SEGMENT_SIZE + } +} + /// `ExternalSorter` manages the sorting of data using both in-memory structures and external files. /// It dumps data to external files when the in-memory buffer crosses a certain memory threshold. pub struct ExternalSorter { @@ -48,7 +71,7 @@ pub struct ExternalSorter { segment_null_bitmap: Bitmap, /// In-memory buffer to hold values and their corresponding bitmaps until memory threshold is exceeded - values_buffer: BTreeMap, + values_buffer: HashMap, /// Count of all rows ingested so far total_row_count: usize, @@ -78,11 +101,11 @@ pub struct ExternalSorter { #[async_trait] impl Sorter for ExternalSorter { - /// Pushes n identical values into the sorter, adding them to the in-memory buffer and dumping - /// the buffer to an external file if necessary - async fn push_n(&mut self, value: Option>, n: usize) -> Result<()> { + /// Pushes n identical values into the in-memory buffer; returns whether it should be + /// spilled. + fn push_n(&mut self, value: Option>, n: usize) -> bool { if n == 0 { - return Ok(()); + return false; } let segment_index_range = self.segment_index_range(n); @@ -90,13 +113,17 @@ impl Sorter for ExternalSorter { if let Some(value) = value { let memory_diff = self.push_not_null(value, segment_index_range); - self.may_dump_buffer(memory_diff).await + self.account_memory(memory_diff) } else { self.segment_null_bitmap.insert_range(segment_index_range); - Ok(()) + false } } + async fn spill(&mut self) -> Result<()> { + self.dump_buffer().await + } + /// Finalizes the sorting operation, merging data from both in-memory buffer and external files /// into a sorted stream async fn output(&mut self) -> Result { @@ -110,9 +137,7 @@ impl Sorter for ExternalSorter { let mut tree_nodes: VecDeque = VecDeque::with_capacity(readers.len() + 1); tree_nodes.push_back(Box::new(stream::iter( - mem::take(&mut self.values_buffer) - .into_iter() - .map(|(value, (bitmap, _))| Ok((value, bitmap))), + Self::sorted(mem::take(&mut self.values_buffer)).map(Ok), ))); for (_, reader) in readers { tree_nodes.push_back(IntermediateReader::new(reader).into_stream().await?); @@ -149,7 +174,7 @@ impl ExternalSorter { temp_file_provider, segment_null_bitmap: Bitmap::new_bitvec(), // bitvec is more efficient for many null values - values_buffer: BTreeMap::new(), + values_buffer: HashMap::default(), total_row_count: 0, segment_row_count, @@ -188,51 +213,65 @@ impl ExternalSorter { value: BytesRef<'_>, segment_index_range: RangeInclusive, ) -> usize { + let (start, end) = ( + *segment_index_range.start() as u32, + *segment_index_range.end() as u32, + ); match self.values_buffer.get_mut(value) { - Some((bitmap, mem_usage)) => { - bitmap.insert_range(segment_index_range); - let new_usage = bitmap.memory_usage() + value.len(); - let diff = new_usage - *mem_usage; - *mem_usage = new_usage; - - diff - } + Some(posting) => posting.push(start, end), None => { - let mut bitmap = Bitmap::new_roaring(); - bitmap.insert_range(segment_index_range); - - let mem_usage = bitmap.memory_usage() + value.len(); - self.values_buffer - .insert(value.to_vec(), (bitmap, mem_usage)); - - mem_usage + let mut segments = RoaringBitmap::new(); + segments.insert_range(start..=end); + let posting = Posting { + segments, + last_segment: end, + }; + self.values_buffer.insert(value.to_vec(), posting); + value.len() + (end - start + 1) as usize * SEGMENT_SIZE } } } - /// Checks if the in-memory buffer exceeds the threshold and offloads it to external storage if necessary - async fn may_dump_buffer(&mut self, memory_diff: usize) -> Result<()> { + /// Drains `values` sorted by value. + fn sorted( + values: HashMap, + ) -> impl Iterator { + let mut values = values.into_iter().collect::>(); + values.sort_unstable_by(|a, b| a.0.cmp(&b.0)); + values + .into_iter() + .map(|(value, posting)| (value, Bitmap::Roaring(posting.segments))) + } + + /// Records `memory_diff` and returns whether the buffer exceeds the thresholds and + /// should be offloaded to external storage. + fn account_memory(&mut self, memory_diff: usize) -> bool { self.current_memory_usage += memory_diff; let memory_usage = self.current_memory_usage; self.global_memory_usage .fetch_add(memory_diff, Ordering::Relaxed); - if self.global_memory_usage_sort_limit.is_none() { - return Ok(()); + let Some(limit) = self.global_memory_usage_sort_limit else { + return false; + }; + if self.global_memory_usage.load(Ordering::Relaxed) < limit { + return false; } - - if self.global_memory_usage.load(Ordering::Relaxed) - < self.global_memory_usage_sort_limit.unwrap() - { - return Ok(()); - } - if let Some(current_threshold) = self.current_memory_usage_threshold && memory_usage < current_threshold { + return false; + } + true + } + + async fn dump_buffer(&mut self) -> Result<()> { + // A second spill request before any new value would write an empty file with the + // same id and replace the first one. + if self.values_buffer.is_empty() { return Ok(()); } - + let memory_usage = self.current_memory_usage; let file_id = &format!("{:012}", self.total_row_count); let index_name = &self.index_name; let writer = self @@ -247,7 +286,7 @@ impl ExternalSorter { self.current_memory_usage = 0; let entries = values.len(); - IntermediateWriter::new(writer).write_all(values.into_iter().map(|(k, (b, _))| (k, b))).await.inspect(|_| + IntermediateWriter::new(writer).write_all(Self::sorted(values)).await.inspect(|_| debug!("Dumped {entries} entries ({memory_usage} bytes) to intermediate file {file_id} for index {index_name}") ).inspect_err(|e| error!(e; "Failed to dump {entries} entries to intermediate file {file_id} for index {index_name}") @@ -271,7 +310,7 @@ impl ExternalSorter { #[cfg(test)] mod tests { - use std::collections::HashMap; + use std::collections::{BTreeMap, HashMap}; use std::iter; use std::sync::Mutex; @@ -329,7 +368,9 @@ mod tests { dictionary_values_and_sorted_result(row_count, segment_row_count); for (value, n) in dic_values { - sorter.push_n(value.as_deref(), n).await.unwrap(); + if sorter.push_n(value.as_deref(), n) { + sorter.spill().await.unwrap(); + } } sorted_result @@ -338,7 +379,9 @@ mod tests { shuffle_values_and_sorted_result(row_count, segment_row_count); for value in mock_values { - sorter.push(value.as_deref()).await.unwrap(); + if sorter.push_n(value.as_deref(), 1) { + sorter.spill().await.unwrap(); + } } sorted_result diff --git a/src/index/src/inverted_index/create/sort_create.rs b/src/index/src/inverted_index/create/sort_create.rs index cb22179a698..9218dd76ab6 100644 --- a/src/index/src/inverted_index/create/sort_create.rs +++ b/src/index/src/inverted_index/create/sort_create.rs @@ -42,6 +42,9 @@ pub struct SortIndexCreator { /// Number of rows in each segment, used to produce sorters segment_row_count: NonZeroUsize, + + /// Indexes whose sorter asked to spill its buffer. + pending_spills: Vec, } #[async_trait] @@ -50,22 +53,33 @@ impl InvertedIndexCreator for SortIndexCreator { /// /// If the index does not exist, a new index is created even if `n` is 0. /// Caller may leverage this behavior to create indexes with no data. - async fn push_with_name_n( + fn push_with_name_n( &mut self, index_name: &str, value: Option>, n: usize, - ) -> Result<()> { - match self.sorters.get_mut(index_name) { - Some(sorter) => sorter.push_n(value, n).await, + ) -> bool { + let sorter = match self.sorters.get_mut(index_name) { + Some(sorter) => sorter, None => { - let index_name = index_name.to_string(); - let mut sorter = (self.sorter_factory)(index_name.clone(), self.segment_row_count); - sorter.push_n(value, n).await?; - self.sorters.insert(index_name, sorter); - Ok(()) + let sorter = (self.sorter_factory)(index_name.to_string(), self.segment_row_count); + self.sorters.entry(index_name.to_string()).or_insert(sorter) + } + }; + if sorter.push_n(value, n) { + self.pending_spills.push(index_name.to_string()); + return true; + } + false + } + + async fn spill(&mut self) -> Result<()> { + for index_name in std::mem::take(&mut self.pending_spills) { + if let Some(sorter) = self.sorters.get_mut(&index_name) { + sorter.spill().await?; } } + Ok(()) } /// Finalizes the sorting for all indexes and writes them using the inverted index writer @@ -110,6 +124,7 @@ impl SortIndexCreator { sorter_factory, sorters: HashMap::new(), segment_row_count, + pending_spills: Vec::new(), } } } @@ -117,6 +132,7 @@ impl SortIndexCreator { #[cfg(test)] mod tests { use std::collections::BTreeMap; + use std::sync::{Arc, Mutex}; use common_base::BitVec; use futures::{StreamExt, stream}; @@ -140,10 +156,7 @@ mod tests { for (index_name, values) in index_values { for value in values { - creator - .push_with_name(index_name, Some(value)) - .await - .unwrap(); + assert!(!creator.push_with_name(index_name, Some(value))); } } @@ -189,10 +202,7 @@ mod tests { for (index_name, values) in index_values { for value in values { - creator - .push_with_name(index_name, Some(value)) - .await - .unwrap(); + assert!(!creator.push_with_name(index_name, Some(value))); } } @@ -221,9 +231,9 @@ mod tests { let mut creator = SortIndexCreator::new(NaiveSorter::factory(), NonZeroUsize::new(1).unwrap()); - creator.push_with_name_n("a", None, 0).await.unwrap(); - creator.push_with_name_n("b", None, 0).await.unwrap(); - creator.push_with_name_n("c", None, 0).await.unwrap(); + assert!(!creator.push_with_name_n("a", None, 0)); + assert!(!creator.push_with_name_n("b", None, 0)); + assert!(!creator.push_with_name_n("c", None, 0)); let mut mock_writer = MockInvertedIndexWriter::new(); mock_writer @@ -250,6 +260,51 @@ mod tests { .unwrap(); } + #[tokio::test] + async fn test_sort_index_creator_spills_requested_sorters() { + let spilled = Arc::new(Mutex::new(Vec::new())); + let factory: SorterFactory = { + let spilled = spilled.clone(); + Box::new(move |index_name, _| { + Box::new(SpillRecordingSorter { + index_name, + spilled: spilled.clone(), + }) + }) + }; + let mut creator = SortIndexCreator::new(factory, NonZeroUsize::new(1).unwrap()); + + assert!(creator.push_with_name("a", Some(b"1"))); + assert!(!creator.push_with_name("b", Some(b"1"))); + creator.spill().await.unwrap(); + assert_eq!(*spilled.lock().unwrap(), vec!["a"]); + + creator.spill().await.unwrap(); + assert_eq!(*spilled.lock().unwrap(), vec!["a"]); + } + + /// Requests a spill on every push to index `a` and records the spills it receives. + struct SpillRecordingSorter { + index_name: String, + spilled: Arc>>, + } + + #[async_trait] + impl Sorter for SpillRecordingSorter { + fn push_n(&mut self, _value: Option>, _n: usize) -> bool { + self.index_name == "a" + } + + async fn spill(&mut self) -> Result<()> { + self.spilled.lock().unwrap().push(self.index_name.clone()); + Ok(()) + } + + async fn output(&mut self) -> Result { + unreachable!() + } + } + fn set_bit(bit_vec: &mut BitVec, index: usize) { if index >= bit_vec.len() { bit_vec.resize(index + 1, false); @@ -277,20 +332,18 @@ mod tests { #[async_trait] impl Sorter for NaiveSorter { - async fn push(&mut self, value: Option>) -> Result<()> { - let segment_index = self.total_row_count / self.segment_row_count; - self.total_row_count += 1; + fn push_n(&mut self, value: Option>, n: usize) -> bool { + for _ in 0..n { + let segment_index = self.total_row_count / self.segment_row_count; + self.total_row_count += 1; - let bitmap = self.values.entry(value.map(Into::into)).or_default(); - set_bit(bitmap, segment_index); - - Ok(()) + let bitmap = self.values.entry(value.map(Into::into)).or_default(); + set_bit(bitmap, segment_index); + } + false } - async fn push_n(&mut self, value: Option>, n: usize) -> Result<()> { - for _ in 0..n { - self.push(value).await?; - } + async fn spill(&mut self) -> Result<()> { Ok(()) } diff --git a/src/mito2/src/sst/index/inverted_index/creator.rs b/src/mito2/src/sst/index/inverted_index/creator.rs index 9ea5258ee3b..d1dee226515 100644 --- a/src/mito2/src/sst/index/inverted_index/creator.rs +++ b/src/mito2/src/sst/index/inverted_index/creator.rs @@ -204,10 +204,12 @@ impl InvertedIndexer { &mut self.value_buf, ) .context(EncodeSnafu)?; - self.index_creator - .push_with_name_n(target_key, elem, count) - .await - .context(PushIndexValueSnafu)?; + if self.index_creator.push_with_name_n(target_key, elem, count) { + self.index_creator + .spill() + .await + .context(PushIndexValueSnafu)?; + } } } else if is_sparse && column_meta.semantic_type == SemanticType::Tag { if self.codec.pk_col_info(*col_id).is_some() { @@ -233,10 +235,15 @@ impl InvertedIndexer { &mut self.value_buf, ) .context(DecodeSnafu)?; - self.index_creator + if self + .index_creator .push_with_name_n(target_key, value, count) - .await - .context(PushIndexValueSnafu)?; + { + self.index_creator + .spill() + .await + .context(PushIndexValueSnafu)?; + } } } } @@ -307,10 +314,12 @@ impl InvertedIndexer { }) .transpose()?; - self.index_creator - .push_with_name_n(target_key, value, n) - .await - .context(PushIndexValueSnafu)?; + if self.index_creator.push_with_name_n(target_key, value, n) { + self.index_creator + .spill() + .await + .context(PushIndexValueSnafu)?; + } } // fields None => { @@ -326,10 +335,12 @@ impl InvertedIndexer { self.value_buf.clear(); let value = values.data.get_ref(i); if value.is_null() { - self.index_creator - .push_with_name(target_key, None) - .await - .context(PushIndexValueSnafu)?; + if self.index_creator.push_with_name(target_key, None) { + self.index_creator + .spill() + .await + .context(PushIndexValueSnafu)?; + } } else { IndexValueCodec::encode_nonnull_value( value, @@ -337,10 +348,15 @@ impl InvertedIndexer { &mut self.value_buf, ) .context(EncodeSnafu)?; - self.index_creator + if self + .index_creator .push_with_name(target_key, Some(&self.value_buf)) - .await - .context(PushIndexValueSnafu)?; + { + self.index_creator + .spill() + .await + .context(PushIndexValueSnafu)?; + } } } }