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 <killme2008@gmail.com>

* perf(index): speed up inverted index building with hashed buffers and sync pushes

Signed-off-by: Dennis Zhuang <killme2008@gmail.com>

* chore(index): use BuildHasher::hash_one for element hashes

Signed-off-by: Dennis Zhuang <killme2008@gmail.com>

* chore(index): require callers to act on spill requests

Signed-off-by: Dennis Zhuang <killme2008@gmail.com>

* fix(index): insert bloom hashes one by one to keep segment set capacity bounded

Signed-off-by: Dennis Zhuang <killme2008@gmail.com>

* test(index): add index build and bloom search benchmarks

Signed-off-by: Dennis Zhuang <killme2008@gmail.com>

* test(index): keep applier setup out of the bloom search benchmark

Signed-off-by: Dennis Zhuang <killme2008@gmail.com>

* fix(index): skip empty spills and keep the old inverted sort memory estimate

Signed-off-by: Dennis Zhuang <killme2008@gmail.com>

* fix(index): stream fulltext token hashes and test spill dispatch

Signed-off-by: Dennis Zhuang <killme2008@gmail.com>

* docs(index): note non-ASCII tokens still allocate in analyze_text_hashes

Signed-off-by: Dennis Zhuang <killme2008@gmail.com>

* docs(index): narrow the allocation note to case-insensitive non-ASCII tokens

Signed-off-by: Dennis Zhuang <killme2008@gmail.com>

---------

Signed-off-by: Dennis Zhuang <killme2008@gmail.com>
This commit is contained in:
dennis zhuang
2026-09-28 08:22:42 +00:00
committed by GitHub
parent bc5ad4b416
commit ffdd6d09a6
15 changed files with 879 additions and 228 deletions
+5
View File
@@ -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
+424
View File
@@ -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<R>(
&mut self,
_key: &str,
raw_data: R,
_options: PutOptions,
_properties: HashMap<String, String>,
) -> puffin::error::Result<u64>
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<String, String>,
) -> puffin::error::Result<u64> {
unreachable!("bloom fulltext index only writes blobs")
}
fn set_footer_lz4_compressed(&mut self, _lz4_compressed: bool) {}
async fn finish(self) -> puffin::error::Result<u64> {
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<Writer, index::error::Error> {
unreachable!("no memory limit is set")
}
async fn read_all(&self, _: &str) -> Result<Vec<(String, Reader)>, index::error::Error> {
Ok(vec![])
}
}
fn uuid(rng: &mut ChaCha8Rng) -> String {
let h = format!("{:032x}", rng.random::<u128>());
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<String> {
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<String> {
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<String> {
let mut rng = ChaCha8Rng::seed_from_u64(3);
let traces = (0..rows / 10)
.map(|_| format!("{:032x}", rng.random::<u128>()))
.collect::<Vec<_>>();
(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::<Vec<_>>();
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::<Vec<_>>();
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);
+68
View File
@@ -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<fastbloom::DefaultHasher> =
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<PrehashedHasher>;
/// 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::<Vec<_>>();
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());
}
}
}
+9 -8
View File
@@ -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<BloomFilter>)> {
) -> Result<(Vec<(u64, usize)>, Vec<PrehashedBloomFilter>)> {
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<BloomFilter>,
bloom_filters: Vec<PrehashedBloomFilter>,
predicates: &[InListPredicate],
) -> Vec<Range<usize>> {
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::<Vec<_>>())
.collect::<Vec<_>>();
// 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 {
+80 -84
View File
@@ -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<Bytes>,
/// 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<u64, PrehashedBuildHasher>,
/// 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<Item = Bytes>,
) -> Result<()> {
if nrows == 0 {
return Ok(());
}
if nrows == 1 {
return self.push_row_elems(elems).await;
}
let elems = elems.into_iter().collect::<Vec<_>>();
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::<Vec<_>>();
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<Item = Bytes>) -> 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<Item = u64>) -> 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<Item = u64>) {
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::<u64>();
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());
@@ -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<Item = Bytes>,
elem_hashes: impl IntoIterator<Item = u64>,
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<S: std::hash::BuildHasher>(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.
+4 -3
View File
@@ -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<Vec<BloomFilter>> {
) -> Result<Vec<PrehashedBloomFilter>> {
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);
}
@@ -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)?;
+47
View File
@@ -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<u8>,
) -> impl Iterator<Item = u64> + 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<Vec<Bytes>> {
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<dyn Tokenizer>, 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::<Vec<_>>();
let hashes = analyzer
.analyze_text_hashes(text, &mut Vec::new())
.collect::<Vec<_>>();
assert_eq!(expected, hashes);
}
}
#[test]
fn test_english_tokenizer() {
let tokenizer = EnglishTokenizer;
+14 -20
View File
@@ -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<BytesRef<'_>>,
) -> 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<BytesRef<'_>>) -> 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<BytesRef<'_>>,
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<BytesRef<'_>>, 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
+8 -8
View File
@@ -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<BytesRef<'_>>) -> 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<BytesRef<'_>>, 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<BytesRef<'_>>, 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<SortOutput>;
@@ -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::<u16>();
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<Bytes, (Bitmap, usize)>,
values_buffer: HashMap<Bytes, Posting, ahash::RandomState>,
/// 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<BytesRef<'_>>, 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<BytesRef<'_>>, 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<SortOutput> {
@@ -110,9 +137,7 @@ impl Sorter for ExternalSorter {
let mut tree_nodes: VecDeque<SortedStream> = 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>,
) -> 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<Bytes, Posting, ahash::RandomState>,
) -> impl Iterator<Item = (Bytes, Bitmap)> {
let mut values = values.into_iter().collect::<Vec<_>>();
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
@@ -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<IndexName>,
}
#[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<BytesRef<'_>>,
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<Mutex<Vec<String>>>,
}
#[async_trait]
impl Sorter for SpillRecordingSorter {
fn push_n(&mut self, _value: Option<BytesRef<'_>>, _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<SortOutput> {
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<BytesRef<'_>>) -> Result<()> {
let segment_index = self.total_row_count / self.segment_row_count;
self.total_row_count += 1;
fn push_n(&mut self, value: Option<BytesRef<'_>>, 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<BytesRef<'_>>, n: usize) -> Result<()> {
for _ in 0..n {
self.push(value).await?;
}
async fn spill(&mut self) -> Result<()> {
Ok(())
}
@@ -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)?;
}
}
}
}