perf(mito2): optimize materialized column index updates (#9182)

Signed-off-by: Lei, HUANG <ratuthomm@gmail.com>
This commit is contained in:
Lei, HUANG
2026-09-16 16:28:12 +00:00
committed by GitHub
parent 4ecec69bee
commit d6a8974613
10 changed files with 831 additions and 39 deletions
+11 -7
View File
@@ -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
+1
View File
@@ -13,6 +13,7 @@
// limitations under the License.
#![feature(iter_partition_in_place)]
#![feature(hash_set_entry)]
pub mod bitmap;
pub mod bloom_filter;
+56
View File
@@ -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<u8>,
) -> Result<Option<&'a [u8]>> {
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");
+5
View File
@@ -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"]
@@ -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<Option<(usize, usize)>> = 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<ColumnId>, 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<Vec<Value>> = (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<ArrayRef> = 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::<UInt32Type>::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::<UInt32Type>::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);
+3
View File
@@ -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;
+10 -11
View File
@@ -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<u8>,
/// 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)?;
}
+51
View File
@@ -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<Item = (usize, usize)> + '_ {
let keys = (semantic_type == SemanticType::Tag)
.then(|| batch.column_by_name(PRIMARY_KEY_COLUMN_NAME))
.flatten()
.and_then(|array| array.as_any().downcast_ref::<PrimaryKeyArray>())
.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))
})
}
+340
View File
@@ -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<Vec<u8>>) {
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::<UInt32Type>::new();
let codec = DensePrimaryKeyCodec::new(&metadata);
let numbers: Vec<Option<u64>> = 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::<UInt32Type>::new(
UInt32Array::from(tag_keys),
Arc::new(StringArray::from(labels.to_vec())),
);
let columns: Vec<ArrayRef> = 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::<u64, [u32; 2]>(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<_>>(),
vec![(0, 3), (3, 3), (6, 6)]
);
for semantic_type in [SemanticType::Field, SemanticType::Timestamp] {
assert_eq!(
column_index_rows(&sliced, semantic_type).collect::<Vec<_>>(),
(0..12).map(|i| (i, 1)).collect::<Vec<_>>()
);
}
}
@@ -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() {