diff --git a/src/index/src/bloom_filter/creator.rs b/src/index/src/bloom_filter/creator.rs index 8bbd4cf1c13..df23ca2352a 100644 --- a/src/index/src/bloom_filter/creator.rs +++ b/src/index/src/bloom_filter/creator.rs @@ -159,13 +159,17 @@ impl BloomFilterCreator { nrows -= rows_to_push; self.accumulated_row_count += rows_to_push; - if let Some(elem) = elem - && !self.cur_seg_distinct_elems.contains(elem) - { - self.cur_seg_distinct_elems.insert(elem.to_vec()); - self.cur_seg_distinct_elems_mem_usage += elem.len(); - self.global_memory_usage - .fetch_add(elem.len(), Ordering::Relaxed); + if let Some(elem) = elem { + let old_len = self.cur_seg_distinct_elems.len(); + // A borrowed entry lookup avoids hashing new values twice and only + // allocates when the value is absent from the current segment. + self.cur_seg_distinct_elems + .get_or_insert_with(elem, <[u8]>::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 diff --git a/src/index/src/lib.rs b/src/index/src/lib.rs index 7969ece8913..c3ae35e9072 100644 --- a/src/index/src/lib.rs +++ b/src/index/src/lib.rs @@ -13,6 +13,7 @@ // limitations under the License. #![feature(iter_partition_in_place)] +#![feature(hash_set_entry)] pub mod bitmap; pub mod bloom_filter; diff --git a/src/mito-codec/src/index.rs b/src/mito-codec/src/index.rs index e0f2c819afe..53ff5cd3c04 100644 --- a/src/mito-codec/src/index.rs +++ b/src/mito-codec/src/index.rs @@ -36,6 +36,28 @@ use crate::row_converter::{PrimaryKeyCodec, SortField, build_primary_key_codec_w pub struct IndexValueCodec; impl IndexValueCodec { + /// Returns the index bytes of a value, or `None` for NULL. + /// Strings borrow their original UTF-8 bytes. For other non-null values, this + /// clears and reuses `buffer` rather than appending as [`Self::encode_nonnull_value`] does. + pub fn encode_value<'a>( + value: ValueRef<'a>, + field: &SortField, + buffer: &'a mut Vec, + ) -> Result> { + if value.is_null() { + return Ok(None); + } + if field.encode_data_type().is_string() { + return Ok(value + .try_into_string() + .context(FieldTypeMismatchSnafu)? + .map(str::as_bytes)); + } + buffer.clear(); + Self::encode_nonnull_value(value, field, buffer)?; + Ok(Some(buffer)) + } + /// Extracts one sparse PK column in index format, without constructing a Value. /// Numeric reserved fields borrow the PK; strings are unchunked into the reusable buffer. /// Missing and null labels return None, whereas an empty string returns Some(&[]). @@ -170,6 +192,40 @@ mod tests { use crate::error::Error; use crate::row_converter::{DensePrimaryKeyCodec, PrimaryKeyCodecExt, SortField}; + #[test] + fn borrowed_values_preserve_index_encoding_with_reused_buffer() { + let mut buffer = vec![0xff; 64]; + for value in [ + Value::from("中文\0abcdefgh"), + Value::from(""), + Value::Int64(-42), + Value::UInt64(u64::MAX), + Value::Boolean(true), + Value::Binary(vec![0, 1, 255].into()), + Value::Null, + Value::from("after-null"), + ] { + let field = SortField::new(value.data_type()); + let mut expected = Vec::new(); + if !value.is_null() { + IndexValueCodec::encode_nonnull_value(value.as_value_ref(), &field, &mut expected) + .unwrap(); + } + let encoded = + IndexValueCodec::encode_value(value.as_value_ref(), &field, &mut buffer).unwrap(); + assert_eq!(encoded, (!value.is_null()).then_some(expected.as_slice())); + } + for (value, field) in [ + (ValueRef::UInt64(1), ConcreteDataType::string_datatype()), + (ValueRef::String("x"), ConcreteDataType::uint64_datatype()), + ] { + assert!(matches!( + IndexValueCodec::encode_value(value, &SortField::new(field), &mut buffer), + Err(Error::FieldTypeMismatch { .. }) + )); + } + } + #[test] fn test_encode_value_basic() { let value = ValueRef::from("hello"); diff --git a/src/mito2/Cargo.toml b/src/mito2/Cargo.toml index ccc3fd6f176..66d1bea5f57 100644 --- a/src/mito2/Cargo.toml +++ b/src/mito2/Cargo.toml @@ -150,5 +150,10 @@ name = "bench_index_update" harness = false required-features = ["testing"] +[[bench]] +name = "bench_dense_index_update" +harness = false +required-features = ["testing"] + [package.metadata.cargo-udeps.ignore] normal = ["aquamarine"] diff --git a/src/mito2/benches/bench_dense_index_update.rs b/src/mito2/benches/bench_dense_index_update.rs new file mode 100644 index 00000000000..f81d30eb4a6 --- /dev/null +++ b/src/mito2/benches/bench_dense_index_update.rs @@ -0,0 +1,342 @@ +// 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. + +//! Real Dense index updates, excluding batch generation, creator setup and cleanup. +//! `CARGO_PROFILE_BENCH_DEBUG=0 cargo bench -p mito2 --features testing +//! --bench bench_dense_index_update -- --save-baseline before` +//! Repeat with `--baseline before` after changing the implementation. +//! Set `DENSE_INDEX_ALLOCATIONS=1` to print allocation counts/bytes for one update +//! instead of timing. The current-thread runtime keeps counting scoped to the update. + +use std::alloc::{GlobalAlloc, Layout, System}; +use std::cell::Cell; +use std::collections::HashSet; +use std::num::NonZeroUsize; +use std::sync::Arc; +use std::time::{Duration, Instant}; + +use api::v1::SemanticType; +use common_test_util::temp_dir::create_temp_dir; +use criterion::{Criterion, Throughput, criterion_group, criterion_main}; +use datatypes::arrow::array::{ + ArrayRef, BinaryDictionaryBuilder, StringDictionaryBuilder, TimestampMillisecondArray, + UInt8Array, UInt64Array, +}; +use datatypes::arrow::datatypes::UInt32Type; +use datatypes::arrow::record_batch::RecordBatch; +use datatypes::data_type::ConcreteDataType; +use datatypes::schema::{ColumnSchema, SkippingIndexOptions}; +use datatypes::value::Value; +use mito_codec::row_converter::{DensePrimaryKeyCodec, PrimaryKeyCodecExt}; +use mito2::sst::index::intermediate::IntermediateManager; +use mito2::sst::{FlatSchemaOptions, to_flat_sst_arrow_schema}; +use mito2::test_util::bench_util::{BloomFilterIndexer, InvertedIndexer}; +use store_api::codec::PrimaryKeyEncoding; +use store_api::metadata::{ColumnMetadata, RegionMetadataBuilder, RegionMetadataRef}; +use store_api::storage::{ColumnId, FileId, RegionId}; + +thread_local! { + static ALLOCATIONS: Cell> = const { Cell::new(None) }; +} + +struct CountingAllocator; + +fn record_allocation(bytes: usize) { + ALLOCATIONS.with(|counter| { + if let Some((count, total)) = counter.get() { + counter.set(Some((count + 1, total + bytes))); + } + }); +} + +// SAFETY: All allocations are delegated unchanged to System, including layout and ownership. +unsafe impl GlobalAlloc for CountingAllocator { + unsafe fn alloc(&self, layout: Layout) -> *mut u8 { + record_allocation(layout.size()); + unsafe { System.alloc(layout) } + } + + unsafe fn alloc_zeroed(&self, layout: Layout) -> *mut u8 { + record_allocation(layout.size()); + unsafe { System.alloc_zeroed(layout) } + } + + unsafe fn realloc(&self, ptr: *mut u8, layout: Layout, new_size: usize) -> *mut u8 { + record_allocation(new_size); + unsafe { System.realloc(ptr, layout, new_size) } + } + + unsafe fn dealloc(&self, ptr: *mut u8, layout: Layout) { + unsafe { System.dealloc(ptr, layout) } + } +} + +#[global_allocator] +static ALLOCATOR: CountingAllocator = CountingAllocator; + +// Allocation mode is a benchmark report, not a server log. +#[allow(clippy::print_stdout)] +fn report_allocations(name: &str) { + let (count, bytes) = ALLOCATIONS.replace(None).unwrap(); + println!("{name},allocations={count},bytes={bytes}"); +} + +const ROWS: usize = 4096; +const SEGMENT_ROWS: usize = 1024; + +struct Shape { + name: String, + tags: u32, + indexed_tags: u32, + inverted: bool, + bloom: bool, + rows_per_key: usize, + cardinality: usize, + numeric: bool, + field: bool, +} + +fn input(shape: &Shape) -> (RegionMetadataRef, HashSet, RecordBatch) { + let mut builder = RegionMetadataBuilder::new(RegionId::new(1, 1)); + let mut inverted_columns = HashSet::new(); + for id in 0..shape.tags + u32::from(shape.field) { + let field = id == shape.tags; + let indexed = if shape.field { + field + } else { + id < shape.indexed_tags + }; + let data_type = if shape.numeric || field { + ConcreteDataType::uint64_datatype() + } else { + ConcreteDataType::string_datatype() + }; + let mut schema = ColumnSchema::new(format!("col_{id}"), data_type, true); + if indexed && shape.inverted { + inverted_columns.insert(id); + } + if indexed && shape.bloom { + schema = schema + .with_skipping_options(SkippingIndexOptions { + granularity: SEGMENT_ROWS as _, + ..Default::default() + }) + .unwrap(); + } + builder.push_column_metadata(ColumnMetadata { + column_schema: schema, + semantic_type: if field { + SemanticType::Field + } else { + SemanticType::Tag + }, + column_id: id, + }); + } + builder + .push_column_metadata(ColumnMetadata { + column_schema: ColumnSchema::new( + "ts", + ConcreteDataType::timestamp_millisecond_datatype(), + false, + ), + semantic_type: SemanticType::Timestamp, + column_id: shape.tags + 1, + }) + .primary_key((0..shape.tags).collect()) + .primary_key_encoding(PrimaryKeyEncoding::Dense); + let metadata = Arc::new(builder.build().unwrap()); + // Last tag is unique per series, so low-cardinality indexed tags still have many PKs. + let values: Vec> = (0..shape.tags) + .map(|id| { + (0..ROWS / shape.rows_per_key) + .map(|series| { + let value = if id == shape.tags - 1 { + series + } else { + series % shape.cardinality + }; + if shape.numeric { + Value::UInt64(value as u64) + } else { + Value::from(format!( + "tag-{id:03}-value-{value:010}-abcdefghijklmnopqrstuvwxyz" + )) + } + }) + .collect() + }) + .collect(); + let mut columns: Vec = values + .iter() + .map(|values| { + if shape.numeric { + Arc::new(UInt64Array::from_iter_values( + (0..ROWS).map(|row| values[row / shape.rows_per_key].as_u64().unwrap()), + )) as ArrayRef + } else { + let mut dict = StringDictionaryBuilder::::new(); + for row in 0..ROWS { + dict.append(values[row / shape.rows_per_key].as_string().unwrap()) + .unwrap(); + } + Arc::new(dict.finish()) as ArrayRef + } + }) + .collect(); + if shape.field { + columns.push(Arc::new(UInt64Array::from_iter_values( + (0..ROWS).map(|row| (row % shape.cardinality) as u64), + ))); + } + let codec = DensePrimaryKeyCodec::new(&metadata); + let mut keys = BinaryDictionaryBuilder::::new(); + for series in 0..ROWS / shape.rows_per_key { + let key = codec + .encode(values.iter().map(|column| column[series].as_value_ref())) + .unwrap(); + for _ in 0..shape.rows_per_key { + keys.append(&key).unwrap(); + } + } + columns.extend([ + Arc::new(TimestampMillisecondArray::from_iter_values(0..ROWS as i64)) as ArrayRef, + Arc::new(keys.finish()), + Arc::new(UInt64Array::from(vec![1; ROWS])), + Arc::new(UInt8Array::from(vec![1; ROWS])), + ]); + let schema = to_flat_sst_arrow_schema( + &metadata, + &FlatSchemaOptions::from_encoding(PrimaryKeyEncoding::Dense), + ); + ( + metadata, + inverted_columns, + RecordBatch::try_new(schema, columns).unwrap(), + ) +} + +fn bench_dense_index_update(c: &mut Criterion) { + let runtime = tokio::runtime::Builder::new_current_thread() + .enable_all() + .build() + .unwrap(); + let dir = create_temp_dir("bench_dense_index_update"); + let intermediate = runtime + .block_on(IntermediateManager::init_fs(dir.path().to_str().unwrap())) + .unwrap(); + let allocations = std::env::var_os("DENSE_INDEX_ALLOCATIONS").is_some(); + let mut group = c.benchmark_group("dense_index_update"); + group.sample_size(30); + group.warm_up_time(Duration::from_secs(1)); + group.measurement_time(Duration::from_secs(3)); + group.throughput(Throughput::Elements(ROWS as u64)); + let mut shapes = Vec::new(); + for rows_per_key in [1, 8, 64] { + for (mode, inverted, bloom) in [ + ("inverted", true, false), + ("bloom", false, true), + ("both", true, true), + ] { + shapes.push(Shape { + name: format!("{mode}_string_one_10tags_{rows_per_key}rpk"), + tags: 10, + indexed_tags: 1, + inverted, + bloom, + rows_per_key, + cardinality: ROWS, + numeric: false, + field: false, + }); + } + } + for (name, tags, indexed_tags, rows_per_key, cardinality, numeric, field) in [ + ("both_string_all_40tags_1rpk", 40, 40, 1, ROWS, false, false), + ("both_string_all_40tags_8rpk", 40, 40, 8, ROWS, false, false), + ("both_string_low_cardinality", 10, 1, 1, 8, false, false), + ("both_numeric_1rpk", 10, 1, 1, ROWS, true, false), + ("both_numeric_64rpk", 10, 1, 64, ROWS, true, false), + ("both_numeric_varying_field", 10, 0, 64, 8, true, true), + ] { + shapes.push(Shape { + name: name.into(), + tags, + indexed_tags, + inverted: true, + bloom: true, + rows_per_key, + cardinality, + numeric, + field, + }); + } + for shape in shapes { + let (metadata, inverted_columns, batch) = input(&shape); + let run = |iterations| { + runtime.block_on(async { + let mut elapsed = Duration::ZERO; + for _ in 0..iterations { + let file_id = FileId::random(); + let mut inverted = shape.inverted.then(|| { + InvertedIndexer::new( + file_id, + &metadata, + intermediate.clone(), + None, + NonZeroUsize::new(SEGMENT_ROWS).unwrap(), + inverted_columns.clone(), + ) + }); + let mut bloom = if shape.bloom { + BloomFilterIndexer::new(file_id, &metadata, intermediate.clone(), None) + .unwrap() + } else { + None + }; + if allocations { + ALLOCATIONS.set(Some((0, 0))); + } + let start = Instant::now(); + if let Some(indexer) = &mut inverted { + indexer.update_flat(&batch).await.unwrap(); + } + if let Some(indexer) = &mut bloom { + indexer.update_flat(&batch).await.unwrap(); + } + elapsed += start.elapsed(); + if allocations { + report_allocations(&shape.name); + } + if let Some(indexer) = &mut inverted { + indexer.abort().await.unwrap(); + } + if let Some(indexer) = &mut bloom { + indexer.abort().await.unwrap(); + } + } + elapsed + }) + }; + if allocations { + run(1); + } else { + group.bench_function(&shape.name, |b| b.iter_custom(&run)); + } + } + group.finish(); +} + +criterion_group!(benches, bench_dense_index_update); +criterion_main!(benches); diff --git a/src/mito2/src/sst/index.rs b/src/mito2/src/sst/index.rs index c8fcd6fbc7e..67e8c4d8e1e 100644 --- a/src/mito2/src/sst/index.rs +++ b/src/mito2/src/sst/index.rs @@ -13,6 +13,9 @@ // limitations under the License. pub(crate) mod bloom_filter; +mod column; +#[cfg(test)] +mod column_test; pub(crate) mod fulltext_index; mod indexer; pub mod intermediate; diff --git a/src/mito2/src/sst/index/bloom_filter/creator.rs b/src/mito2/src/sst/index/bloom_filter/creator.rs index cfdd751bf7f..2776ddacc0d 100644 --- a/src/mito2/src/sst/index/bloom_filter/creator.rs +++ b/src/mito2/src/sst/index/bloom_filter/creator.rs @@ -41,6 +41,7 @@ use crate::error::{ use crate::read::Batch; use crate::sst::index::TYPE_BLOOM_FILTER_INDEX; use crate::sst::index::bloom_filter::INDEX_BLOB_TYPE; +use crate::sst::index::column::column_index_rows; use crate::sst::index::intermediate::{ IntermediateLocation, IntermediateManager, TempFileProvider, }; @@ -63,6 +64,7 @@ pub struct BloomFilterIndexer { codec: IndexValuesCodec, /// Scratch storage for extracting indexed tags from sparse primary keys. pk_offsets: SparseOffsetsCache, + /// Reusable buffer for encoding materialized values and extracting sparse labels. value_buf: Vec, /// Whether the indexing process has been aborted. @@ -313,19 +315,16 @@ impl BloomFilterIndexer { .context(crate::error::ConvertVectorSnafu)?; let sort_field = SortField::new(vector.data_type()); - for i in 0..n { - let value = vector.get_ref(i); - let elems = (!value.is_null()) - .then(|| { - let mut buf = vec![]; - IndexValueCodec::encode_nonnull_value(value, &sort_field, &mut buf) - .context(EncodeSnafu)?; - Ok(buf) - }) - .transpose()?; + for (row, count) in column_index_rows(batch, column_meta.semantic_type) { + let elem = IndexValueCodec::encode_value( + vector.get_ref(row), + &sort_field, + &mut self.value_buf, + ) + .context(EncodeSnafu)?; creator - .push_row_elems(elems) + .push_n_row_elem(count, elem) .await .context(PushBloomFilterValueSnafu)?; } diff --git a/src/mito2/src/sst/index/column.rs b/src/mito2/src/sst/index/column.rs new file mode 100644 index 00000000000..e2f9aaf60d1 --- /dev/null +++ b/src/mito2/src/sst/index/column.rs @@ -0,0 +1,51 @@ +// 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. + +use api::v1::SemanticType; +use datatypes::arrow::array::Array; +use datatypes::arrow::record_batch::RecordBatch; +use store_api::storage::consts::PRIMARY_KEY_COLUMN_NAME; + +use crate::sst::parquet::format::PrimaryKeyArray; + +/// Yields (first row, row count) for indexing a materialized column. +/// Tags are constant within consecutive equal PK dictionary keys. Fields and +/// timestamps must still be visited row by row. Distinct dictionary entries may +/// contain equal PK bytes; keeping those runs separate is correct and avoids decoding. +/// Inputs without a non-null PK dictionary fall back to visiting individual rows. +pub(crate) fn column_index_rows( + batch: &RecordBatch, + semantic_type: SemanticType, +) -> impl Iterator + '_ { + let keys = (semantic_type == SemanticType::Tag) + .then(|| batch.column_by_name(PRIMARY_KEY_COLUMN_NAME)) + .flatten() + .and_then(|array| array.as_any().downcast_ref::()) + .filter(|array| array.null_count() == 0 && array.values().null_count() == 0) + .map(|array| array.keys().values()); + let mut row = 0; + std::iter::from_fn(move || { + if row == batch.num_rows() { + return None; + } + let start = row; + row += 1; + if let Some(keys) = keys { + while row < keys.len() && keys[row] == keys[start] { + row += 1; + } + } + Some((start, row - start)) + }) +} diff --git a/src/mito2/src/sst/index/column_test.rs b/src/mito2/src/sst/index/column_test.rs new file mode 100644 index 00000000000..146bed120c0 --- /dev/null +++ b/src/mito2/src/sst/index/column_test.rs @@ -0,0 +1,340 @@ +// 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. + +use std::collections::BTreeMap; +use std::sync::Arc; + +use api::v1::SemanticType; +use datatypes::arrow::array::{ + Array, ArrayRef, BinaryDictionaryBuilder, DictionaryArray, StringArray, + TimestampMillisecondArray, UInt8Array, UInt32Array, UInt64Array, +}; +use datatypes::arrow::datatypes::UInt32Type; +use datatypes::arrow::record_batch::RecordBatch; +use datatypes::data_type::ConcreteDataType; +use datatypes::schema::{ColumnSchema, SkippingIndexOptions}; +use datatypes::value::ValueRef; +use index::bitmap::{Bitmap, BitmapType}; +use index::bloom_filter::reader::{BloomFilterReader, BloomFilterReaderImpl}; +use index::inverted_index::format::reader::{InvertedIndexBlobReader, InvertedIndexReader}; +use mito_codec::index::IndexValueCodec; +use mito_codec::row_converter::{DensePrimaryKeyCodec, PrimaryKeyCodecExt, SortField}; +use object_store::ObjectStore; +use object_store::services::Memory; +use prost::Message; +use puffin::puffin_manager::{PuffinManager, PuffinReader}; +use store_api::codec::PrimaryKeyEncoding; +use store_api::metadata::{ColumnMetadata, RegionMetadataBuilder, RegionMetadataRef}; +use store_api::storage::{FileId, RegionId}; + +use super::{IndexBuildType, IndexerBuilder, IndexerBuilderImpl}; +use crate::read::BatchBuilder; +use crate::region::options::IndexOptions; +use crate::sst::file::{RegionFileId, RegionIndexId}; +use crate::sst::index::bloom_filter::creator::tests::TestPathProvider; +use crate::sst::index::column::column_index_rows; +use crate::sst::index::intermediate::IntermediateManager; +use crate::sst::index::puffin_manager::PuffinManagerFactory; +use crate::sst::{FlatSchemaOptions, to_flat_sst_arrow_schema}; + +fn input() -> (RegionMetadataRef, RecordBatch, Vec>) { + let mut builder = RegionMetadataBuilder::new(RegionId::new(1, 1)); + for (id, name, data_type, semantic_type) in [ + ( + 0, + "tag_str", + ConcreteDataType::string_datatype(), + SemanticType::Tag, + ), + ( + 1, + "tag_num", + ConcreteDataType::uint64_datatype(), + SemanticType::Tag, + ), + ( + 2, + "field", + ConcreteDataType::uint64_datatype(), + SemanticType::Field, + ), + ( + 3, + "ts", + ConcreteDataType::timestamp_millisecond_datatype(), + SemanticType::Timestamp, + ), + ] { + let mut schema = ColumnSchema::new(name, data_type, id != 3); + if id != 3 { + schema.set_inverted_index(true); + schema + .set_skipping_options(&SkippingIndexOptions { + granularity: 3, + ..Default::default() + }) + .unwrap(); + } + builder.push_column_metadata(ColumnMetadata { + column_schema: schema, + semantic_type, + column_id: id, + }); + } + builder.primary_key(vec![0, 1]); + let metadata = Arc::new(builder.build().unwrap()); + // Null dictionary keys AND valid keys referencing null dictionary values. + // Repeated PKs cross batches/segments and recur non-consecutively. Fields vary within PKs. + let tag_keys = vec![ + Some(2), + Some(4), + Some(2), + Some(4), + Some(2), + None, + Some(0), + None, + Some(1), + Some(1), + Some(1), + Some(1), + Some(1), + Some(1), + Some(3), + Some(2), + ]; + let labels = [ + None, + Some(""), + Some("中文\0abcdefgh"), + Some("a-longer-string-value"), + Some("中文\0abcdefgh"), + ]; + let mut encoded = Vec::new(); + let mut pk = BinaryDictionaryBuilder::::new(); + let codec = DensePrimaryKeyCodec::new(&metadata); + let numbers: Vec> = tag_keys + .iter() + .map(|key| { + if key.is_none() || *key == Some(0) { + None + } else { + Some(42) + } + }) + .collect(); + for (key, number) in tag_keys.iter().zip(&numbers) { + let label = key.and_then(|key| labels[key as usize]); + let bytes = codec + .encode( + [ + label.map_or(ValueRef::Null, ValueRef::String), + number.map_or(ValueRef::Null, ValueRef::UInt64), + ] + .into_iter(), + ) + .unwrap(); + pk.append(&bytes).unwrap(); + encoded.push(bytes); + } + let n = tag_keys.len(); + let tag = DictionaryArray::::new( + UInt32Array::from(tag_keys), + Arc::new(StringArray::from(labels.to_vec())), + ); + let columns: Vec = vec![ + Arc::new(tag), + Arc::new(UInt64Array::from(numbers)), + Arc::new(UInt64Array::from_iter( + (0..n).map(|i| (i % 4 != 0).then_some(i as u64)), + )), + Arc::new(TimestampMillisecondArray::from_iter_values(0..n as i64)), + Arc::new(pk.finish()), + Arc::new(UInt64Array::from(vec![1; n])), + Arc::new(UInt8Array::from(vec![1; n])), + ]; + let schema = to_flat_sst_arrow_schema( + &metadata, + &FlatSchemaOptions::from_encoding(PrimaryKeyEncoding::Dense), + ); + ( + metadata, + RecordBatch::try_new(schema, columns).unwrap(), + encoded, + ) +} + +#[tokio::test] +async fn materialized_runs_match_legacy_index_bytes_and_bloom_lookups() { + let (metadata, batch, encoded) = input(); + let (dir, factory) = PuffinManagerFactory::new_for_test_async("materialized_index_bytes").await; + let puffin = factory.build( + ObjectStore::new(Memory::default()).unwrap(), + TestPathProvider, + ); + let mut options = IndexOptions::default(); + options.inverted_index.segment_row_count = 3; + let builder = IndexerBuilderImpl { + build_type: IndexBuildType::Flush, + metadata: metadata.clone(), + puffin_manager: puffin.clone(), + write_cache_enabled: false, + intermediate_manager: IntermediateManager::init_fs(dir.path().to_str().unwrap()) + .await + .unwrap(), + index_options: options, + inverted_index_config: Default::default(), + fulltext_index_config: Default::default(), + bloom_filter_index_config: Default::default(), + #[cfg(feature = "vector_index")] + vector_index_config: Default::default(), + }; + let mut expected = None; + // The legacy Batch path is an independent whole-PK decode/owned-encoding oracle. + // One-row slices disable run merging; a tag-only projection also exercises the no-PK fallback. + for mode in ["legacy", "runs", "single_rows", "no_pk"] { + let file = RegionFileId::new(metadata.region_id, FileId::random()); + let mut indexer = builder.build(file, 0, None).await; + match mode { + "legacy" => { + for (row, key) in encoded.iter().enumerate() { + let mut old = BatchBuilder::new(key.clone()); + old.push_field_array(2, batch.column(2).slice(row, 1)) + .unwrap(); + old.timestamps_array(batch.column(3).slice(row, 1)).unwrap(); + old.sequences_array(batch.column(5).slice(row, 1)).unwrap(); + old.op_types_array(batch.column(6).slice(row, 1)).unwrap(); + indexer.update(&mut old.build().unwrap()).await; + } + } + "single_rows" => { + for row in 0..batch.num_rows() { + indexer.update_flat(&batch.slice(row, 1)).await; + } + } + _ => { + let batch = if mode == "no_pk" { + batch.project(&[0, 1, 2, 3]).unwrap() + } else { + batch.clone() + }; + for (offset, count) in [(0, 0), (0, 2), (2, 7), (9, 7)] { + indexer.update_flat(&batch.slice(offset, count)).await; + } + } + } + let output = indexer.finish().await; + assert_eq!(output.inverted_index.row_count, batch.num_rows(), "{mode}"); + assert_eq!(output.bloom_filter.row_count, batch.num_rows(), "{mode}"); + let reader = puffin.reader(&RegionIndexId::new(file, 0)).await.unwrap(); + let blob = reader.blob("greptime-inverted-index-v1").await.unwrap(); + let inverted = InvertedIndexBlobReader::new(blob.reader().await.unwrap()); + let metas = inverted.metadata(None).await.unwrap(); + assert_eq!(metas.metas.len(), 3); + assert_eq!(metas.total_row_count, batch.num_rows() as u64); + assert_eq!(metas.segment_row_count, 3); + // Equality on the empty string must exclude NULL-only segments. + let string_meta = &metas.metas["0"]; + let fst = inverted + .fst( + string_meta.base_offset + string_meta.relative_fst_offset as u64, + string_meta.fst_size, + None, + ) + .await + .unwrap(); + let [offset, size] = bytemuck::cast::(fst.get(b"").unwrap()); + let bitmap = inverted + .bitmap( + string_meta.base_offset + offset as u64, + size, + BitmapType::Roaring, + None, + ) + .await + .unwrap(); + assert_eq!( + bitmap, + Bitmap::from_lsb0_bytes(&[0b0001_1100], BitmapType::Roaring) + ); + let mut index_bytes = BTreeMap::new(); + // Blob column order is unspecified; compare each column's complete FST and bitmaps. + for (name, meta) in &metas.metas { + index_bytes.insert( + format!("inverted/{name}"), + inverted + .range_read(meta.base_offset, meta.inverted_index_size as u32, None) + .await + .unwrap(), + ); + } + for id in 0..3 { + let blob = reader + .blob(&format!("greptime-bloom-filter-v1-{id}")) + .await + .unwrap(); + let bloom = BloomFilterReaderImpl::new(blob.reader().await.unwrap()); + let meta = bloom.metadata(None).await.unwrap(); + assert_eq!(meta.row_count, batch.num_rows() as u64); + assert_eq!(meta.segment_count, 6); + let mut bytes = bloom + .range_read(0, meta.bloom_filter_size as u32, None) + .await + .unwrap() + .to_vec(); + bytes.extend_from_slice(&meta.encode_to_vec()); + index_bytes.insert(format!("bloom/{id}"), bytes); + let vector = + datatypes::vectors::Helper::try_into_vector(batch.column(id).clone()).unwrap(); + let field = SortField::new(vector.data_type()); + for row in 0..batch.num_rows() { + let value = vector.get_ref(row); + if !value.is_null() { + let mut bytes = Vec::new(); + IndexValueCodec::encode_nonnull_value(value, &field, &mut bytes).unwrap(); + let loc = &meta.bloom_filter_locs[meta.segment_loc_indices[row / 3] as usize]; + assert!( + bloom + .bloom_filter(loc, None) + .await + .unwrap() + .contains(&bytes), + "{mode}: {id}/{row}" + ); + } + } + } + if let Some(expected) = &expected { + assert_eq!(&index_bytes, expected, "{mode}"); + } else { + expected = Some(index_bytes); + } + } +} + +#[test] +fn tag_runs_preserve_sliced_row_positions_and_field_rows() { + let (_, batch, _) = input(); + let sliced = batch.slice(2, 12); + assert_eq!( + column_index_rows(&sliced, SemanticType::Tag).collect::>(), + vec![(0, 3), (3, 3), (6, 6)] + ); + for semantic_type in [SemanticType::Field, SemanticType::Timestamp] { + assert_eq!( + column_index_rows(&sliced, semantic_type).collect::>(), + (0..12).map(|i| (i, 1)).collect::>() + ); + } +} diff --git a/src/mito2/src/sst/index/inverted_index/creator.rs b/src/mito2/src/sst/index/inverted_index/creator.rs index 0243016fc7a..9ea5258ee3b 100644 --- a/src/mito2/src/sst/index/inverted_index/creator.rs +++ b/src/mito2/src/sst/index/inverted_index/creator.rs @@ -44,6 +44,7 @@ use crate::error::{ }; use crate::read::Batch; use crate::sst::index::TYPE_INVERTED_INDEX; +use crate::sst::index::column::column_index_rows; use crate::sst::index::intermediate::{ IntermediateLocation, IntermediateManager, TempFileProvider, }; @@ -196,27 +197,17 @@ impl InvertedIndexer { .context(crate::error::ConvertVectorSnafu)?; let sort_field = SortField::new(vector.data_type()); - for row in 0..batch.num_rows() { - self.value_buf.clear(); - let value_ref = vector.get_ref(row); - - if value_ref.is_null() { - self.index_creator - .push_with_name(target_key, None) - .await - .context(PushIndexValueSnafu)?; - } else { - IndexValueCodec::encode_nonnull_value( - value_ref, - &sort_field, - &mut self.value_buf, - ) - .context(EncodeSnafu)?; - self.index_creator - .push_with_name(target_key, Some(&self.value_buf)) - .await - .context(PushIndexValueSnafu)?; - } + for (row, count) in column_index_rows(batch, column_meta.semantic_type) { + let elem = IndexValueCodec::encode_value( + vector.get_ref(row), + &sort_field, + &mut self.value_buf, + ) + .context(EncodeSnafu)?; + self.index_creator + .push_with_name_n(target_key, elem, count) + .await + .context(PushIndexValueSnafu)?; } } else if is_sparse && column_meta.semantic_type == SemanticType::Tag { if self.codec.pk_col_info(*col_id).is_some() {