1use std::collections::HashSet;
16use std::num::NonZeroUsize;
17use std::sync::Arc;
18use std::sync::atomic::AtomicUsize;
19
20use api::v1::SemanticType;
21use common_telemetry::{debug, warn};
22use datatypes::arrow::record_batch::RecordBatch;
23use datatypes::vectors::Helper;
24use index::inverted_index::create::InvertedIndexCreator;
25use index::inverted_index::create::sort::external_sort::ExternalSorter;
26use index::inverted_index::create::sort_create::SortIndexCreator;
27use index::inverted_index::format::writer::InvertedIndexBlobWriter;
28use index::target::IndexTarget;
29use mito_codec::index::{IndexValueCodec, IndexValuesCodec};
30use mito_codec::row_converter::sparse::SparsePrimaryKeyView;
31use mito_codec::row_converter::{SortField, SparseOffsetsCache};
32use puffin::puffin_manager::{PuffinWriter, PutOptions};
33use smallvec::SmallVec;
34use snafu::{ResultExt, ensure};
35use store_api::codec::PrimaryKeyEncoding;
36use store_api::metadata::RegionMetadataRef;
37use store_api::storage::{ColumnId, FileId};
38use tokio::io::duplex;
39use tokio_util::compat::{TokioAsyncReadCompatExt, TokioAsyncWriteCompatExt};
40
41use crate::error::{
42 BiErrorsSnafu, DecodeSnafu, EncodeSnafu, IndexFinishSnafu, OperateAbortedIndexSnafu,
43 PuffinAddBlobSnafu, PushIndexValueSnafu, Result,
44};
45use crate::read::Batch;
46use crate::sst::index::TYPE_INVERTED_INDEX;
47use crate::sst::index::column::column_index_rows;
48use crate::sst::index::intermediate::{
49 IntermediateLocation, IntermediateManager, TempFileProvider,
50};
51use crate::sst::index::inverted_index::INDEX_BLOB_TYPE;
52use crate::sst::index::primary_key::PrimaryKeyRuns;
53use crate::sst::index::puffin_manager::SstPuffinWriter;
54use crate::sst::index::statistics::{ByteCount, RowCount, Statistics};
55
56const MIN_MEMORY_USAGE_THRESHOLD_PER_COLUMN: usize = 1024 * 1024; const PIPE_BUFFER_SIZE_FOR_SENDING_BLOB: usize = 8192;
61
62pub struct InvertedIndexer {
64 index_creator: Box<dyn InvertedIndexCreator>,
66 temp_file_provider: Arc<TempFileProvider>,
68
69 codec: IndexValuesCodec,
71 value_buf: Vec<u8>,
73 pk_offsets: SparseOffsetsCache,
75
76 stats: Statistics,
78 aborted: bool,
80
81 memory_usage: Arc<AtomicUsize>,
83
84 indexed_column_ids: Vec<(ColumnId, String)>,
86
87 metadata: RegionMetadataRef,
89}
90
91impl InvertedIndexer {
92 pub fn new(
95 sst_file_id: FileId,
96 metadata: &RegionMetadataRef,
97 intermediate_manager: IntermediateManager,
98 memory_usage_threshold: Option<usize>,
99 segment_row_count: NonZeroUsize,
100 indexed_column_ids: HashSet<ColumnId>,
101 ) -> Self {
102 let temp_file_provider = Arc::new(TempFileProvider::new(
103 IntermediateLocation::new(&metadata.region_id, &sst_file_id),
104 intermediate_manager,
105 ));
106
107 let memory_usage = Arc::new(AtomicUsize::new(0));
108
109 let sorter = ExternalSorter::factory(
110 temp_file_provider.clone() as _,
111 Some(MIN_MEMORY_USAGE_THRESHOLD_PER_COLUMN),
112 memory_usage.clone(),
113 memory_usage_threshold,
114 );
115 let index_creator = Box::new(SortIndexCreator::new(sorter, segment_row_count));
116
117 let codec = IndexValuesCodec::from_tag_columns(
118 metadata.primary_key_encoding,
119 metadata.primary_key_columns(),
120 );
121 let indexed_column_ids = indexed_column_ids
122 .into_iter()
123 .map(|col_id| {
124 let target_key = format!("{}", IndexTarget::ColumnId(col_id));
125 (col_id, target_key)
126 })
127 .collect();
128 Self {
129 codec,
130 index_creator,
131 temp_file_provider,
132 value_buf: vec![],
133 pk_offsets: SparseOffsetsCache::new(),
134 stats: Statistics::new(TYPE_INVERTED_INDEX),
135 aborted: false,
136 memory_usage,
137 indexed_column_ids,
138 metadata: metadata.clone(),
139 }
140 }
141
142 pub async fn update(&mut self, batch: &mut Batch) -> Result<()> {
145 ensure!(!self.aborted, OperateAbortedIndexSnafu);
146
147 if batch.is_empty() {
148 return Ok(());
149 }
150
151 if let Err(update_err) = self.do_update(batch).await {
152 if let Err(err) = self.do_cleanup().await {
154 if cfg!(any(test, feature = "test")) {
155 panic!("Failed to clean up index creator, err: {err}",);
156 } else {
157 warn!(err; "Failed to clean up index creator");
158 }
159 }
160 return Err(update_err);
161 }
162
163 Ok(())
164 }
165
166 pub async fn update_flat(&mut self, batch: &RecordBatch) -> Result<()> {
168 ensure!(!self.aborted, OperateAbortedIndexSnafu);
169
170 if batch.num_rows() == 0 {
171 return Ok(());
172 }
173
174 self.do_update_flat(batch).await
175 }
176
177 async fn do_update_flat(&mut self, batch: &RecordBatch) -> Result<()> {
178 let mut guard = self.stats.record_update();
179
180 guard.inc_row_count(batch.num_rows());
181
182 let is_sparse = self.metadata.primary_key_encoding == PrimaryKeyEncoding::Sparse;
183 let mut sparse_columns: SmallVec<[(ColumnId, &str); 8]> = SmallVec::new();
184
185 for (col_id, target_key) in &self.indexed_column_ids {
186 let Some(column_meta) = self.metadata.column_by_id(*col_id) else {
187 debug!(
188 "Column {} not found in the metadata during building inverted index",
189 col_id
190 );
191 continue;
192 };
193 let column_name = &column_meta.column_schema.name;
194 if let Some(column_array) = batch.column_by_name(column_name) {
195 let vector = Helper::try_into_vector(column_array.clone())
197 .context(crate::error::ConvertVectorSnafu)?;
198 let sort_field = SortField::new(vector.data_type());
199
200 for (row, count) in column_index_rows(batch, column_meta.semantic_type) {
201 let elem = IndexValueCodec::encode_value(
202 vector.get_ref(row),
203 &sort_field,
204 &mut self.value_buf,
205 )
206 .context(EncodeSnafu)?;
207 if self.index_creator.push_with_name_n(target_key, elem, count) {
208 self.index_creator
209 .spill()
210 .await
211 .context(PushIndexValueSnafu)?;
212 }
213 }
214 } else if is_sparse && column_meta.semantic_type == SemanticType::Tag {
215 if self.codec.pk_col_info(*col_id).is_some() {
216 sparse_columns.push((*col_id, target_key));
217 }
218 } else {
219 debug!(
220 "Column {} not found in the batch during building inverted index",
221 col_id
222 );
223 }
224 }
225
226 if !sparse_columns.is_empty() {
227 for (pk, count) in PrimaryKeyRuns::try_new(batch)? {
228 let mut view =
229 SparsePrimaryKeyView::new(pk, &mut self.pk_offsets).context(DecodeSnafu)?;
230 for &(col_id, target_key) in &sparse_columns {
232 let value = IndexValueCodec::encode_sparse_value(
233 &mut view,
234 col_id,
235 &mut self.value_buf,
236 )
237 .context(DecodeSnafu)?;
238 if self
239 .index_creator
240 .push_with_name_n(target_key, value, count)
241 {
242 self.index_creator
243 .spill()
244 .await
245 .context(PushIndexValueSnafu)?;
246 }
247 }
248 }
249 }
250
251 Ok(())
252 }
253
254 pub(crate) async fn finish(
257 &mut self,
258 puffin_writer: &mut SstPuffinWriter,
259 ) -> Result<(RowCount, ByteCount)> {
260 ensure!(!self.aborted, OperateAbortedIndexSnafu);
261
262 if self.stats.row_count() == 0 {
263 return Ok((0, 0));
265 }
266
267 let finish_res = self.do_finish(puffin_writer).await;
268 if let Err(err) = self.do_cleanup().await {
270 if cfg!(any(test, feature = "test")) {
271 panic!("Failed to clean up index creator, err: {err}",);
272 } else {
273 warn!(err; "Failed to clean up index creator");
274 }
275 }
276
277 finish_res.map(|_| (self.stats.row_count(), self.stats.byte_count()))
278 }
279
280 pub async fn abort(&mut self) -> Result<()> {
282 if self.aborted {
283 return Ok(());
284 }
285 self.aborted = true;
286
287 self.do_cleanup().await
288 }
289
290 async fn do_update(&mut self, batch: &mut Batch) -> Result<()> {
291 let mut guard = self.stats.record_update();
292
293 let n = batch.num_rows();
294 guard.inc_row_count(n);
295
296 for (col_id, target_key) in &self.indexed_column_ids {
297 match self.codec.pk_col_info(*col_id) {
298 Some(col_info) => {
300 let pk_idx = col_info.idx;
301 let field = &col_info.field;
302 let value = batch
303 .pk_col_value(self.codec.decoder(), pk_idx, *col_id)?
304 .filter(|v| !v.is_null())
305 .map(|v| {
306 self.value_buf.clear();
307 IndexValueCodec::encode_nonnull_value(
308 v.as_value_ref(),
309 field,
310 &mut self.value_buf,
311 )
312 .context(EncodeSnafu)?;
313 Ok(self.value_buf.as_slice())
314 })
315 .transpose()?;
316
317 if self.index_creator.push_with_name_n(target_key, value, n) {
318 self.index_creator
319 .spill()
320 .await
321 .context(PushIndexValueSnafu)?;
322 }
323 }
324 None => {
326 let Some(values) = batch.field_col_value(*col_id) else {
327 debug!(
328 "Column {} not found in the batch during building inverted index",
329 col_id
330 );
331 continue;
332 };
333 let sort_field = SortField::new(values.data.data_type());
334 for i in 0..n {
335 self.value_buf.clear();
336 let value = values.data.get_ref(i);
337 if value.is_null() {
338 if self.index_creator.push_with_name(target_key, None) {
339 self.index_creator
340 .spill()
341 .await
342 .context(PushIndexValueSnafu)?;
343 }
344 } else {
345 IndexValueCodec::encode_nonnull_value(
346 value,
347 &sort_field,
348 &mut self.value_buf,
349 )
350 .context(EncodeSnafu)?;
351 if self
352 .index_creator
353 .push_with_name(target_key, Some(&self.value_buf))
354 {
355 self.index_creator
356 .spill()
357 .await
358 .context(PushIndexValueSnafu)?;
359 }
360 }
361 }
362 }
363 }
364 }
365
366 Ok(())
367 }
368
369 async fn do_finish(&mut self, puffin_writer: &mut SstPuffinWriter) -> Result<()> {
388 let mut guard = self.stats.record_finish();
389
390 let (tx, rx) = duplex(PIPE_BUFFER_SIZE_FOR_SENDING_BLOB);
391 let mut index_writer = InvertedIndexBlobWriter::new(tx.compat_write());
392
393 let (index_finish, puffin_add_blob) = futures::join!(
394 self.index_creator
396 .finish(&mut index_writer, index::bitmap::BitmapType::Roaring),
397 puffin_writer.put_blob(
398 INDEX_BLOB_TYPE,
399 rx.compat(),
400 PutOptions::default(),
401 Default::default(),
402 )
403 );
404
405 match (
406 puffin_add_blob.context(PuffinAddBlobSnafu),
407 index_finish.context(IndexFinishSnafu),
408 ) {
409 (Err(e1), Err(e2)) => BiErrorsSnafu {
410 first: Box::new(e1),
411 second: Box::new(e2),
412 }
413 .fail()?,
414
415 (Ok(_), e @ Err(_)) => e?,
416 (e @ Err(_), Ok(_)) => e.map(|_| ())?,
417 (Ok(written_bytes), Ok(_)) => {
418 guard.inc_byte_count(written_bytes);
419 }
420 }
421
422 Ok(())
423 }
424
425 async fn do_cleanup(&mut self) -> Result<()> {
426 let _guard = self.stats.record_cleanup();
427
428 self.temp_file_provider.cleanup().await
429 }
430
431 pub fn column_ids(&self) -> impl Iterator<Item = ColumnId> + '_ {
432 self.indexed_column_ids.iter().map(|(col_id, _)| *col_id)
433 }
434
435 pub fn memory_usage(&self) -> usize {
436 self.memory_usage.load(std::sync::atomic::Ordering::Relaxed)
437 }
438}
439
440#[cfg(test)]
441mod tests {
442 use std::collections::BTreeSet;
443
444 use api::v1::SemanticType;
445 use datafusion_expr::{Expr as DfExpr, Operator, binary_expr, col, lit};
446 use datatypes::data_type::ConcreteDataType;
447 use datatypes::schema::ColumnSchema;
448 use datatypes::value::ValueRef;
449 use datatypes::vectors::{UInt8Vector, UInt64Vector};
450 use futures::future::BoxFuture;
451 use mito_codec::row_converter::{DensePrimaryKeyCodec, PrimaryKeyCodecExt};
452 use object_store::ObjectStore;
453 use object_store::services::Memory;
454 use puffin::puffin_manager::PuffinManager;
455 use puffin::puffin_manager::cache::PuffinMetadataCache;
456 use store_api::metadata::{ColumnMetadata, RegionMetadataBuilder};
457 use store_api::region_request::PathType;
458 use store_api::storage::RegionId;
459
460 use super::*;
461 use crate::access_layer::RegionFilePathFactory;
462 use crate::cache::index::inverted_index::InvertedIndexCache;
463 use crate::metrics::CACHE_BYTES;
464 use crate::read::BatchColumn;
465 use crate::sst::file::{RegionFileId, RegionIndexId};
466 use crate::sst::index::inverted_index::applier::builder::InvertedIndexApplierBuilder;
467 use crate::sst::index::puffin_manager::PuffinManagerFactory;
468
469 fn mock_object_store() -> ObjectStore {
470 ObjectStore::new(Memory::default()).unwrap()
471 }
472
473 async fn new_intm_mgr(path: impl AsRef<str>) -> IntermediateManager {
474 IntermediateManager::init_fs(path).await.unwrap()
475 }
476
477 fn mock_region_metadata() -> RegionMetadataRef {
478 let mut builder = RegionMetadataBuilder::new(RegionId::new(1, 2));
479 builder
480 .push_column_metadata(ColumnMetadata {
481 column_schema: ColumnSchema::new(
482 "tag_str",
483 ConcreteDataType::string_datatype(),
484 false,
485 ),
486 semantic_type: SemanticType::Tag,
487 column_id: 1,
488 })
489 .push_column_metadata(ColumnMetadata {
490 column_schema: ColumnSchema::new(
491 "tag_i32",
492 ConcreteDataType::int32_datatype(),
493 false,
494 ),
495 semantic_type: SemanticType::Tag,
496 column_id: 2,
497 })
498 .push_column_metadata(ColumnMetadata {
499 column_schema: ColumnSchema::new(
500 "ts",
501 ConcreteDataType::timestamp_millisecond_datatype(),
502 false,
503 ),
504 semantic_type: SemanticType::Timestamp,
505 column_id: 3,
506 })
507 .push_column_metadata(ColumnMetadata {
508 column_schema: ColumnSchema::new(
509 "field_u64",
510 ConcreteDataType::uint64_datatype(),
511 false,
512 ),
513 semantic_type: SemanticType::Field,
514 column_id: 4,
515 })
516 .primary_key(vec![1, 2]);
517
518 Arc::new(builder.build().unwrap())
519 }
520
521 fn new_batch(
522 str_tag: impl AsRef<str>,
523 i32_tag: impl Into<i32>,
524 u64_field: impl IntoIterator<Item = u64>,
525 ) -> Batch {
526 let fields = vec![
527 (0, SortField::new(ConcreteDataType::string_datatype())),
528 (1, SortField::new(ConcreteDataType::int32_datatype())),
529 ];
530 let codec = DensePrimaryKeyCodec::with_fields(fields);
531 let row: [ValueRef; 2] = [str_tag.as_ref().into(), i32_tag.into().into()];
532 let primary_key = codec.encode(row.into_iter()).unwrap();
533
534 let u64_field = BatchColumn {
535 column_id: 4,
536 data: Arc::new(UInt64Vector::from_iter_values(u64_field)),
537 };
538 let num_rows = u64_field.data.len();
539
540 Batch::new(
541 primary_key,
542 Arc::new(UInt64Vector::from_iter_values(std::iter::repeat_n(
543 0, num_rows,
544 ))),
545 Arc::new(UInt64Vector::from_iter_values(std::iter::repeat_n(
546 0, num_rows,
547 ))),
548 Arc::new(UInt8Vector::from_iter_values(std::iter::repeat_n(
549 1, num_rows,
550 ))),
551 vec![u64_field],
552 )
553 .unwrap()
554 }
555
556 async fn build_applier_factory(
557 prefix: &str,
558 rows: BTreeSet<(&'static str, i32, [u64; 2])>,
559 ) -> impl Fn(DfExpr) -> BoxFuture<'static, Vec<usize>> {
560 let (d, factory) = PuffinManagerFactory::new_for_test_async(prefix).await;
561 let table_dir = "table0".to_string();
562 let sst_file_id = FileId::random();
563 let object_store = mock_object_store();
564 let region_metadata = mock_region_metadata();
565 let intm_mgr = new_intm_mgr(d.path().to_string_lossy()).await;
566 let memory_threshold = None;
567 let segment_row_count = 2;
568 let indexed_column_ids = HashSet::from_iter([1, 2, 4]);
569
570 let mut creator = InvertedIndexer::new(
571 sst_file_id,
572 ®ion_metadata,
573 intm_mgr,
574 memory_threshold,
575 NonZeroUsize::new(segment_row_count).unwrap(),
576 indexed_column_ids.clone(),
577 );
578
579 for (str_tag, i32_tag, u64_field) in &rows {
580 let mut batch = new_batch(str_tag, *i32_tag, u64_field.iter().copied());
581 creator.update(&mut batch).await.unwrap();
582 }
583
584 let puffin_manager = factory.build(
585 object_store.clone(),
586 RegionFilePathFactory::new(table_dir.clone(), PathType::Bare),
587 );
588
589 let sst_file_id = RegionFileId::new(region_metadata.region_id, sst_file_id);
590 let index_id = RegionIndexId::new(sst_file_id, 0);
591 let mut writer = puffin_manager.writer(&index_id).await.unwrap();
592 let (row_count, _) = creator.finish(&mut writer).await.unwrap();
593 assert_eq!(row_count, rows.len() * segment_row_count);
594 writer.finish().await.unwrap();
595
596 move |expr| {
597 let _d = &d;
598 let cache = Arc::new(InvertedIndexCache::new(10, 10, 100));
599 let puffin_metadata_cache = Arc::new(PuffinMetadataCache::new(10, &CACHE_BYTES));
600 let applier = InvertedIndexApplierBuilder::new(
601 table_dir.clone(),
602 PathType::Bare,
603 object_store.clone(),
604 ®ion_metadata,
605 indexed_column_ids.clone(),
606 factory.clone(),
607 )
608 .with_inverted_index_cache(Some(cache))
609 .with_puffin_metadata_cache(Some(puffin_metadata_cache))
610 .build(&[expr])
611 .unwrap()
612 .unwrap();
613 let sst_metadata = Arc::new(region_metadata.clone());
614 let plan = applier.plan_for_sst(&sst_metadata).unwrap().unwrap();
615 Box::pin(async move {
616 applier
617 .apply(index_id, None, &plan.index_applier, None)
618 .await
619 .unwrap()
620 .matched_segment_ids
621 .iter_ones()
622 .collect()
623 })
624 }
625 }
626
627 #[tokio::test]
628 async fn test_create_and_query_get_key() {
629 let rows = BTreeSet::from_iter([
630 ("aaa", 1, [1, 2]),
631 ("aaa", 2, [2, 3]),
632 ("aaa", 3, [3, 4]),
633 ("aab", 1, [4, 5]),
634 ("aab", 2, [5, 6]),
635 ("aab", 3, [6, 7]),
636 ("abc", 1, [7, 8]),
637 ("abc", 2, [8, 9]),
638 ("abc", 3, [9, 10]),
639 ]);
640
641 let applier_factory = build_applier_factory("test_create_and_query_get_key_", rows).await;
642
643 let expr = col("tag_str").eq(lit("aaa"));
644 let res = applier_factory(expr).await;
645 assert_eq!(res, vec![0, 1, 2]);
646
647 let expr = col("tag_i32").eq(lit(2));
648 let res = applier_factory(expr).await;
649 assert_eq!(res, vec![1, 4, 7]);
650
651 let expr = col("tag_str").eq(lit("aaa")).and(col("tag_i32").eq(lit(2)));
652 let res = applier_factory(expr).await;
653 assert_eq!(res, vec![1]);
654
655 let expr = col("tag_str")
656 .eq(lit("aaa"))
657 .or(col("tag_str").eq(lit("abc")));
658 let res = applier_factory(expr).await;
659 assert_eq!(res, vec![0, 1, 2, 6, 7, 8]);
660
661 let expr = col("tag_str").in_list(vec![lit("aaa"), lit("abc")], false);
662 let res = applier_factory(expr).await;
663 assert_eq!(res, vec![0, 1, 2, 6, 7, 8]);
664
665 let expr = col("field_u64").eq(lit(2u64));
666 let res = applier_factory(expr).await;
667 assert_eq!(res, vec![0, 1]);
668 }
669
670 #[tokio::test]
671 async fn test_create_and_query_range() {
672 let rows = BTreeSet::from_iter([
673 ("aaa", 1, [1, 2]),
674 ("aaa", 2, [2, 3]),
675 ("aaa", 3, [3, 4]),
676 ("aab", 1, [4, 5]),
677 ("aab", 2, [5, 6]),
678 ("aab", 3, [6, 7]),
679 ("abc", 1, [7, 8]),
680 ("abc", 2, [8, 9]),
681 ("abc", 3, [9, 10]),
682 ]);
683
684 let applier_factory = build_applier_factory("test_create_and_query_range_", rows).await;
685
686 let expr = col("tag_str").between(lit("aaa"), lit("aab"));
687 let res = applier_factory(expr).await;
688 assert_eq!(res, vec![0, 1, 2, 3, 4, 5]);
689
690 let expr = col("tag_i32").between(lit(2), lit(3));
691 let res = applier_factory(expr).await;
692 assert_eq!(res, vec![1, 2, 4, 5, 7, 8]);
693
694 let expr = col("tag_str").between(lit("aaa"), lit("aaa"));
695 let res = applier_factory(expr).await;
696 assert_eq!(res, vec![0, 1, 2]);
697
698 let expr = col("tag_i32").between(lit(2), lit(2));
699 let res = applier_factory(expr).await;
700 assert_eq!(res, vec![1, 4, 7]);
701
702 let expr = col("field_u64").between(lit(2u64), lit(5u64));
703 let res = applier_factory(expr).await;
704 assert_eq!(res, vec![0, 1, 2, 3, 4]);
705 }
706
707 #[tokio::test]
708 async fn test_create_and_query_comparison() {
709 let rows = BTreeSet::from_iter([
710 ("aaa", 1, [1, 2]),
711 ("aaa", 2, [2, 3]),
712 ("aaa", 3, [3, 4]),
713 ("aab", 1, [4, 5]),
714 ("aab", 2, [5, 6]),
715 ("aab", 3, [6, 7]),
716 ("abc", 1, [7, 8]),
717 ("abc", 2, [8, 9]),
718 ("abc", 3, [9, 10]),
719 ]);
720
721 let applier_factory =
722 build_applier_factory("test_create_and_query_comparison_", rows).await;
723
724 let expr = col("tag_str").lt(lit("aab"));
725 let res = applier_factory(expr).await;
726 assert_eq!(res, vec![0, 1, 2]);
727
728 let expr = col("tag_i32").lt(lit(2));
729 let res = applier_factory(expr).await;
730 assert_eq!(res, vec![0, 3, 6]);
731
732 let expr = col("field_u64").lt(lit(2u64));
733 let res = applier_factory(expr).await;
734 assert_eq!(res, vec![0]);
735
736 let expr = col("tag_str").gt(lit("aab"));
737 let res = applier_factory(expr).await;
738 assert_eq!(res, vec![6, 7, 8]);
739
740 let expr = col("tag_i32").gt(lit(2));
741 let res = applier_factory(expr).await;
742 assert_eq!(res, vec![2, 5, 8]);
743
744 let expr = col("field_u64").gt(lit(8u64));
745 let res = applier_factory(expr).await;
746 assert_eq!(res, vec![7, 8]);
747
748 let expr = col("tag_str").lt_eq(lit("aab"));
749 let res = applier_factory(expr).await;
750 assert_eq!(res, vec![0, 1, 2, 3, 4, 5]);
751
752 let expr = col("tag_i32").lt_eq(lit(2));
753 let res = applier_factory(expr).await;
754 assert_eq!(res, vec![0, 1, 3, 4, 6, 7]);
755
756 let expr = col("field_u64").lt_eq(lit(2u64));
757 let res = applier_factory(expr).await;
758 assert_eq!(res, vec![0, 1]);
759
760 let expr = col("tag_str").gt_eq(lit("aab"));
761 let res = applier_factory(expr).await;
762 assert_eq!(res, vec![3, 4, 5, 6, 7, 8]);
763
764 let expr = col("tag_i32").gt_eq(lit(2));
765 let res = applier_factory(expr).await;
766 assert_eq!(res, vec![1, 2, 4, 5, 7, 8]);
767
768 let expr = col("field_u64").gt_eq(lit(8u64));
769 let res = applier_factory(expr).await;
770 assert_eq!(res, vec![6, 7, 8]);
771
772 let expr = col("tag_str")
773 .gt(lit("aaa"))
774 .and(col("tag_str").lt(lit("abc")));
775 let res = applier_factory(expr).await;
776 assert_eq!(res, vec![3, 4, 5]);
777
778 let expr = col("tag_i32").gt(lit(1)).and(col("tag_i32").lt(lit(3)));
779 let res = applier_factory(expr).await;
780 assert_eq!(res, vec![1, 4, 7]);
781
782 let expr = col("field_u64")
783 .gt(lit(2u64))
784 .and(col("field_u64").lt(lit(9u64)));
785 let res = applier_factory(expr).await;
786 assert_eq!(res, vec![1, 2, 3, 4, 5, 6, 7]);
787 }
788
789 #[tokio::test]
790 async fn test_create_and_query_regex() {
791 let rows = BTreeSet::from_iter([
792 ("aaa", 1, [1, 2]),
793 ("aaa", 2, [2, 3]),
794 ("aaa", 3, [3, 4]),
795 ("aab", 1, [4, 5]),
796 ("aab", 2, [5, 6]),
797 ("aab", 3, [6, 7]),
798 ("abc", 1, [7, 8]),
799 ("abc", 2, [8, 9]),
800 ("abc", 3, [9, 10]),
801 ]);
802
803 let applier_factory = build_applier_factory("test_create_and_query_regex_", rows).await;
804
805 let expr = binary_expr(col("tag_str"), Operator::RegexMatch, lit(".*"));
806 let res = applier_factory(expr).await;
807 assert_eq!(res, vec![0, 1, 2, 3, 4, 5, 6, 7, 8]);
808
809 let expr = binary_expr(col("tag_str"), Operator::RegexMatch, lit("a.*c"));
810 let res = applier_factory(expr).await;
811 assert_eq!(res, vec![6, 7, 8]);
812
813 let expr = binary_expr(col("tag_str"), Operator::RegexMatch, lit("a.*b$"));
814 let res = applier_factory(expr).await;
815 assert_eq!(res, vec![3, 4, 5]);
816
817 let expr = binary_expr(col("tag_str"), Operator::RegexMatch, lit("\\w"));
818 let res = applier_factory(expr).await;
819 assert_eq!(res, vec![0, 1, 2, 3, 4, 5, 6, 7, 8]);
820
821 let expr = binary_expr(col("tag_str"), Operator::RegexMatch, lit("\\d"));
822 let res = applier_factory(expr).await;
823 assert!(res.is_empty());
824
825 let expr = binary_expr(col("tag_str"), Operator::RegexMatch, lit("^aaa$"));
826 let res = applier_factory(expr).await;
827 assert_eq!(res, vec![0, 1, 2]);
828 }
829}