1use std::collections::HashMap;
16use std::sync::Arc;
17use std::sync::atomic::AtomicUsize;
18
19use api::v1::SemanticType;
20use common_telemetry::{debug, warn};
21use datatypes::arrow::record_batch::RecordBatch;
22use datatypes::schema::SkippingIndexType;
23use datatypes::vectors::Helper;
24use index::bloom_filter::creator::BloomFilterCreator;
25use index::target::IndexTarget;
26use mito_codec::index::{IndexValueCodec, IndexValuesCodec};
27use mito_codec::row_converter::sparse::SparsePrimaryKeyView;
28use mito_codec::row_converter::{SortField, SparseOffsetsCache};
29use puffin::puffin_manager::{PuffinWriter, PutOptions};
30use smallvec::SmallVec;
31use snafu::{ResultExt, ensure};
32use store_api::codec::PrimaryKeyEncoding;
33use store_api::metadata::RegionMetadataRef;
34use store_api::storage::{ColumnId, FileId};
35use tokio_util::compat::{TokioAsyncReadCompatExt, TokioAsyncWriteCompatExt};
36
37use crate::error::{
38 BiErrorsSnafu, BloomFilterFinishSnafu, DecodeSnafu, EncodeSnafu, IndexOptionsSnafu,
39 OperateAbortedIndexSnafu, PuffinAddBlobSnafu, PushBloomFilterValueSnafu, Result,
40};
41use crate::read::Batch;
42use crate::sst::index::TYPE_BLOOM_FILTER_INDEX;
43use crate::sst::index::bloom_filter::INDEX_BLOB_TYPE;
44use crate::sst::index::column::column_index_rows;
45use crate::sst::index::intermediate::{
46 IntermediateLocation, IntermediateManager, TempFileProvider,
47};
48use crate::sst::index::primary_key::PrimaryKeyRuns;
49use crate::sst::index::puffin_manager::SstPuffinWriter;
50use crate::sst::index::statistics::{ByteCount, RowCount, Statistics};
51
52const PIPE_BUFFER_SIZE_FOR_SENDING_BLOB: usize = 8192;
54
55pub struct BloomFilterIndexer {
57 creators: HashMap<ColumnId, BloomFilterCreator>,
59
60 temp_file_provider: Arc<TempFileProvider>,
62
63 codec: IndexValuesCodec,
65 pk_offsets: SparseOffsetsCache,
67 value_buf: Vec<u8>,
69
70 aborted: bool,
72
73 stats: Statistics,
75
76 global_memory_usage: Arc<AtomicUsize>,
78
79 metadata: RegionMetadataRef,
81}
82
83impl BloomFilterIndexer {
84 pub fn new(
86 sst_file_id: FileId,
87 metadata: &RegionMetadataRef,
88 intermediate_manager: IntermediateManager,
89 memory_usage_threshold: Option<usize>,
90 ) -> Result<Option<Self>> {
91 let mut creators = HashMap::new();
92
93 let temp_file_provider = Arc::new(TempFileProvider::new(
94 IntermediateLocation::new(&metadata.region_id, &sst_file_id),
95 intermediate_manager,
96 ));
97 let global_memory_usage = Arc::new(AtomicUsize::new(0));
98
99 for column in &metadata.column_metadatas {
100 let options =
101 column
102 .column_schema
103 .skipping_index_options()
104 .context(IndexOptionsSnafu {
105 column_name: &column.column_schema.name,
106 })?;
107
108 let options = match options {
109 Some(options) if options.index_type == SkippingIndexType::BloomFilter => options,
110 _ => continue,
111 };
112
113 let creator = BloomFilterCreator::new(
114 options.granularity as _,
115 options.false_positive_rate(),
116 temp_file_provider.clone(),
117 global_memory_usage.clone(),
118 memory_usage_threshold,
119 );
120 creators.insert(column.column_id, creator);
121 }
122
123 if creators.is_empty() {
124 return Ok(None);
125 }
126
127 let codec = IndexValuesCodec::from_tag_columns(
128 metadata.primary_key_encoding,
129 metadata.primary_key_columns(),
130 );
131 let indexer = Self {
132 creators,
133 temp_file_provider,
134 codec,
135 pk_offsets: SparseOffsetsCache::new(),
136 value_buf: Vec::new(),
137 aborted: false,
138 stats: Statistics::new(TYPE_BLOOM_FILTER_INDEX),
139 global_memory_usage,
140 metadata: metadata.clone(),
141 };
142 Ok(Some(indexer))
143 }
144
145 pub async fn update(&mut self, batch: &mut Batch) -> Result<()> {
150 ensure!(!self.aborted, OperateAbortedIndexSnafu);
151
152 if self.creators.is_empty() {
153 return Ok(());
154 }
155
156 if let Err(update_err) = self.do_update(batch).await {
157 if let Err(err) = self.do_cleanup().await {
159 if cfg!(any(test, feature = "test")) {
160 panic!("Failed to clean up index creator, err: {err:?}",);
161 } else {
162 warn!(err; "Failed to clean up index creator");
163 }
164 }
165 return Err(update_err);
166 }
167
168 Ok(())
169 }
170
171 pub async fn update_flat(&mut self, batch: &RecordBatch) -> Result<()> {
173 ensure!(!self.aborted, OperateAbortedIndexSnafu);
174
175 if self.creators.is_empty() || batch.num_rows() == 0 {
176 return Ok(());
177 }
178
179 if let Err(update_err) = self.do_update_flat(batch).await {
180 if let Err(err) = self.do_cleanup().await {
182 if cfg!(any(test, feature = "test")) {
183 panic!("Failed to clean up index creator, err: {err:?}",);
184 } else {
185 warn!(err; "Failed to clean up index creator");
186 }
187 }
188 return Err(update_err);
189 }
190
191 Ok(())
192 }
193
194 pub(crate) async fn finish(
199 &mut self,
200 puffin_writer: &mut SstPuffinWriter,
201 ) -> Result<(RowCount, ByteCount)> {
202 ensure!(!self.aborted, OperateAbortedIndexSnafu);
203
204 if self.stats.row_count() == 0 {
205 return Ok((0, 0));
207 }
208
209 let finish_res = self.do_finish(puffin_writer).await;
210 if let Err(err) = self.do_cleanup().await {
212 if cfg!(any(test, feature = "test")) {
213 panic!("Failed to clean up index creator, err: {err:?}",);
214 } else {
215 warn!(err; "Failed to clean up index creator");
216 }
217 }
218
219 finish_res.map(|_| (self.stats.row_count(), self.stats.byte_count()))
220 }
221
222 pub async fn abort(&mut self) -> Result<()> {
226 if self.aborted {
227 return Ok(());
228 }
229 self.aborted = true;
230
231 self.do_cleanup().await
232 }
233
234 async fn do_update(&mut self, batch: &mut Batch) -> Result<()> {
235 let mut guard = self.stats.record_update();
236
237 let n = batch.num_rows();
238 guard.inc_row_count(n);
239
240 for (col_id, creator) in &mut self.creators {
241 match self.codec.pk_col_info(*col_id) {
242 Some(col_info) => {
244 let pk_idx = col_info.idx;
245 let field = &col_info.field;
246 let elems = batch
247 .pk_col_value(self.codec.decoder(), pk_idx, *col_id)?
248 .filter(|v| !v.is_null())
249 .map(|v| {
250 let mut buf = vec![];
251 IndexValueCodec::encode_nonnull_value(
252 v.as_value_ref(),
253 field,
254 &mut buf,
255 )
256 .context(EncodeSnafu)?;
257 Ok(buf)
258 })
259 .transpose()?;
260 creator
261 .push_n_row_elems(n, elems)
262 .await
263 .context(PushBloomFilterValueSnafu)?;
264 }
265 None => {
267 let Some(values) = batch.field_col_value(*col_id) else {
268 debug!(
269 "Column {} not found in the batch during building bloom filter index",
270 col_id
271 );
272 continue;
273 };
274 let sort_field = SortField::new(values.data.data_type());
275 for i in 0..n {
276 let value = values.data.get_ref(i);
277 let elems = (!value.is_null())
278 .then(|| {
279 let mut buf = vec![];
280 IndexValueCodec::encode_nonnull_value(value, &sort_field, &mut buf)
281 .context(EncodeSnafu)?;
282 Ok(buf)
283 })
284 .transpose()?;
285
286 creator
287 .push_row_elems(elems)
288 .await
289 .context(PushBloomFilterValueSnafu)?;
290 }
291 }
292 }
293 }
294
295 Ok(())
296 }
297
298 async fn do_update_flat(&mut self, batch: &RecordBatch) -> Result<()> {
299 let mut guard = self.stats.record_update();
300
301 let n = batch.num_rows();
302 guard.inc_row_count(n);
303
304 let is_sparse = self.metadata.primary_key_encoding == PrimaryKeyEncoding::Sparse;
305 let mut sparse_columns: SmallVec<[(ColumnId, &mut BloomFilterCreator); 8]> =
306 SmallVec::new();
307
308 for (col_id, creator) in &mut self.creators {
309 let column_meta = self.metadata.column_by_id(*col_id).unwrap();
311 let column_name = &column_meta.column_schema.name;
312 if let Some(column_array) = batch.column_by_name(column_name) {
313 let vector = Helper::try_into_vector(column_array.clone())
315 .context(crate::error::ConvertVectorSnafu)?;
316 let sort_field = SortField::new(vector.data_type());
317
318 for (row, count) in column_index_rows(batch, column_meta.semantic_type) {
319 let elem = IndexValueCodec::encode_value(
320 vector.get_ref(row),
321 &sort_field,
322 &mut self.value_buf,
323 )
324 .context(EncodeSnafu)?;
325
326 creator
327 .push_n_row_elem(count, elem)
328 .await
329 .context(PushBloomFilterValueSnafu)?;
330 }
331 } else if is_sparse && column_meta.semantic_type == SemanticType::Tag {
332 if self.codec.pk_col_info(*col_id).is_some() {
333 sparse_columns.push((*col_id, creator));
334 }
335 } else {
336 debug!(
337 "Column {} not found in the batch during building bloom filter index",
338 column_name
339 );
340 }
341 }
342
343 if !sparse_columns.is_empty() {
344 for (pk, count) in PrimaryKeyRuns::try_new(batch)? {
345 let mut view =
346 SparsePrimaryKeyView::new(pk, &mut self.pk_offsets).context(DecodeSnafu)?;
347 for (col_id, creator) in &mut sparse_columns {
348 let value = IndexValueCodec::encode_sparse_value(
349 &mut view,
350 *col_id,
351 &mut self.value_buf,
352 )
353 .context(DecodeSnafu)?;
354 creator
355 .push_n_row_elem(count, value)
356 .await
357 .context(PushBloomFilterValueSnafu)?;
358 }
359 }
360 }
361
362 Ok(())
363 }
364
365 async fn do_finish(&mut self, puffin_writer: &mut SstPuffinWriter) -> Result<()> {
367 let mut guard = self.stats.record_finish();
368
369 for (id, creator) in &mut self.creators {
370 let written_bytes = Self::do_finish_single_creator(id, creator, puffin_writer).await?;
371 guard.inc_byte_count(written_bytes);
372 }
373
374 Ok(())
375 }
376
377 async fn do_cleanup(&mut self) -> Result<()> {
378 let mut _guard = self.stats.record_cleanup();
379
380 self.creators.clear();
381 self.temp_file_provider.cleanup().await
382 }
383
384 async fn do_finish_single_creator(
405 col_id: &ColumnId,
406 creator: &mut BloomFilterCreator,
407 puffin_writer: &mut SstPuffinWriter,
408 ) -> Result<ByteCount> {
409 let (tx, rx) = tokio::io::duplex(PIPE_BUFFER_SIZE_FOR_SENDING_BLOB);
410
411 let target_key = IndexTarget::ColumnId(*col_id);
412 let blob_name = format!("{INDEX_BLOB_TYPE}-{target_key}");
413 let (index_finish, puffin_add_blob) = futures::join!(
414 creator.finish(tx.compat_write()),
415 puffin_writer.put_blob(
416 &blob_name,
417 rx.compat(),
418 PutOptions::default(),
419 Default::default(),
420 )
421 );
422
423 match (
424 puffin_add_blob.context(PuffinAddBlobSnafu),
425 index_finish.context(BloomFilterFinishSnafu),
426 ) {
427 (Err(e1), Err(e2)) => BiErrorsSnafu {
428 first: Box::new(e1),
429 second: Box::new(e2),
430 }
431 .fail()?,
432
433 (Ok(_), e @ Err(_)) => e?,
434 (e @ Err(_), Ok(_)) => e.map(|_| ())?,
435 (Ok(written_bytes), Ok(_)) => {
436 return Ok(written_bytes);
437 }
438 }
439
440 Ok(0)
441 }
442
443 pub fn memory_usage(&self) -> usize {
445 self.global_memory_usage
446 .load(std::sync::atomic::Ordering::Relaxed)
447 }
448
449 pub fn column_ids(&self) -> impl Iterator<Item = ColumnId> + use<'_> {
451 self.creators.keys().copied()
452 }
453}
454
455#[cfg(test)]
456pub(crate) mod tests {
457
458 use api::v1::SemanticType;
459 use datatypes::data_type::ConcreteDataType;
460 use datatypes::schema::{ColumnSchema, SkippingIndexOptions};
461 use datatypes::value::ValueRef;
462 use datatypes::vectors::{UInt8Vector, UInt64Vector};
463 use index::bloom_filter::reader::{BloomFilterReader, BloomFilterReaderImpl};
464 use mito_codec::row_converter::{DensePrimaryKeyCodec, PrimaryKeyCodecExt};
465 use object_store::ObjectStore;
466 use object_store::services::Memory;
467 use puffin::puffin_manager::{PuffinManager, PuffinReader};
468 use store_api::metadata::{ColumnMetadata, RegionMetadataBuilder};
469 use store_api::storage::RegionId;
470
471 use super::*;
472 use crate::access_layer::FilePathProvider;
473 use crate::read::BatchColumn;
474 use crate::sst::file::{RegionFileId, RegionIndexId};
475 use crate::sst::index::puffin_manager::PuffinManagerFactory;
476
477 pub fn mock_object_store() -> ObjectStore {
478 ObjectStore::new(Memory::default()).unwrap()
479 }
480
481 pub async fn new_intm_mgr(path: impl AsRef<str>) -> IntermediateManager {
482 IntermediateManager::init_fs(path).await.unwrap()
483 }
484
485 pub struct TestPathProvider;
486
487 impl FilePathProvider for TestPathProvider {
488 fn build_index_file_path(&self, file_id: RegionFileId) -> String {
489 file_id.file_id().to_string()
490 }
491
492 fn build_index_file_path_with_version(&self, index_id: RegionIndexId) -> String {
493 index_id.file_id.file_id().to_string()
494 }
495
496 fn build_sst_file_path(&self, file_id: RegionFileId) -> String {
497 file_id.file_id().to_string()
498 }
499 }
500
501 pub fn mock_region_metadata() -> RegionMetadataRef {
518 let mut builder = RegionMetadataBuilder::new(RegionId::new(1, 2));
519 builder
520 .push_column_metadata(ColumnMetadata {
521 column_schema: ColumnSchema::new(
522 "tag_str",
523 ConcreteDataType::string_datatype(),
524 false,
525 )
526 .with_skipping_options(SkippingIndexOptions::new_unchecked(
527 2,
528 0.01,
529 SkippingIndexType::BloomFilter,
530 ))
531 .unwrap(),
532 semantic_type: SemanticType::Tag,
533 column_id: 1,
534 })
535 .push_column_metadata(ColumnMetadata {
536 column_schema: ColumnSchema::new(
537 "ts",
538 ConcreteDataType::timestamp_millisecond_datatype(),
539 false,
540 ),
541 semantic_type: SemanticType::Timestamp,
542 column_id: 2,
543 })
544 .push_column_metadata(ColumnMetadata {
545 column_schema: ColumnSchema::new(
546 "field_u64",
547 ConcreteDataType::uint64_datatype(),
548 false,
549 )
550 .with_skipping_options(SkippingIndexOptions::new_unchecked(
551 4,
552 0.01,
553 SkippingIndexType::BloomFilter,
554 ))
555 .unwrap(),
556 semantic_type: SemanticType::Field,
557 column_id: 3,
558 })
559 .primary_key(vec![1]);
560
561 Arc::new(builder.build().unwrap())
562 }
563
564 pub fn new_batch(str_tag: impl AsRef<str>, u64_field: impl IntoIterator<Item = u64>) -> Batch {
565 let fields = vec![(0, SortField::new(ConcreteDataType::string_datatype()))];
566 let codec = DensePrimaryKeyCodec::with_fields(fields);
567 let row: [ValueRef; 1] = [str_tag.as_ref().into()];
568 let primary_key = codec.encode(row.into_iter()).unwrap();
569
570 let u64_field = BatchColumn {
571 column_id: 3,
572 data: Arc::new(UInt64Vector::from_iter_values(u64_field)),
573 };
574 let num_rows = u64_field.data.len();
575
576 Batch::new(
577 primary_key,
578 Arc::new(UInt64Vector::from_iter_values(std::iter::repeat_n(
579 0, num_rows,
580 ))),
581 Arc::new(UInt64Vector::from_iter_values(std::iter::repeat_n(
582 0, num_rows,
583 ))),
584 Arc::new(UInt8Vector::from_iter_values(std::iter::repeat_n(
585 1, num_rows,
586 ))),
587 vec![u64_field],
588 )
589 .unwrap()
590 }
591
592 #[tokio::test]
593 async fn test_bloom_filter_indexer() {
594 let prefix = "test_bloom_filter_indexer_";
595 let tempdir = common_test_util::temp_dir::create_temp_dir(prefix);
596 let object_store = mock_object_store();
597 let intm_mgr = new_intm_mgr(tempdir.path().to_string_lossy()).await;
598 let region_metadata = mock_region_metadata();
599 let memory_usage_threshold = Some(1024);
600
601 let file_id = FileId::random();
602 let mut indexer =
603 BloomFilterIndexer::new(file_id, ®ion_metadata, intm_mgr, memory_usage_threshold)
604 .unwrap()
605 .unwrap();
606
607 let mut batch = new_batch("tag1", 0..10);
609 indexer.update(&mut batch).await.unwrap();
610
611 let mut batch = new_batch("tag2", 10..20);
612 indexer.update(&mut batch).await.unwrap();
613
614 let (_d, factory) = PuffinManagerFactory::new_for_test_async(prefix).await;
615 let puffin_manager = factory.build(object_store, TestPathProvider);
616
617 let file_id = RegionFileId::new(region_metadata.region_id, file_id);
618 let file_id = RegionIndexId::new(file_id, 0);
619 let mut puffin_writer = puffin_manager.writer(&file_id).await.unwrap();
620 let (row_count, byte_count) = indexer.finish(&mut puffin_writer).await.unwrap();
621 assert_eq!(row_count, 20);
622 assert!(byte_count > 0);
623 puffin_writer.finish().await.unwrap();
624
625 let puffin_reader = puffin_manager.reader(&file_id).await.unwrap();
626
627 {
629 let blob_guard = puffin_reader
630 .blob("greptime-bloom-filter-v1-1")
631 .await
632 .unwrap();
633 let reader = blob_guard.reader().await.unwrap();
634 let bloom_filter = BloomFilterReaderImpl::new(reader);
635 let metadata = bloom_filter.metadata(None).await.unwrap();
636
637 assert_eq!(metadata.segment_count, 10);
638 for i in 0..5 {
639 let loc = &metadata.bloom_filter_locs[metadata.segment_loc_indices[i] as usize];
640 let bf = bloom_filter.bloom_filter(loc, None).await.unwrap();
641 assert!(bf.contains(b"tag1"));
642 }
643 for i in 5..10 {
644 let loc = &metadata.bloom_filter_locs[metadata.segment_loc_indices[i] as usize];
645 let bf = bloom_filter.bloom_filter(loc, None).await.unwrap();
646 assert!(bf.contains(b"tag2"));
647 }
648 }
649
650 {
652 let sort_field = SortField::new(ConcreteDataType::uint64_datatype());
653
654 let blob_guard = puffin_reader
655 .blob("greptime-bloom-filter-v1-3")
656 .await
657 .unwrap();
658 let reader = blob_guard.reader().await.unwrap();
659 let bloom_filter = BloomFilterReaderImpl::new(reader);
660 let metadata = bloom_filter.metadata(None).await.unwrap();
661
662 assert_eq!(metadata.segment_count, 5);
663 for i in 0u64..20 {
664 let idx = i as usize / 4;
665 let loc = &metadata.bloom_filter_locs[metadata.segment_loc_indices[idx] as usize];
666 let bf = bloom_filter.bloom_filter(loc, None).await.unwrap();
667 let mut buf = vec![];
668 IndexValueCodec::encode_nonnull_value(ValueRef::UInt64(i), &sort_field, &mut buf)
669 .unwrap();
670
671 assert!(bf.contains(&buf));
672 }
673 }
674 }
675}