1use async_stream::try_stream;
16use common_recordbatch::filter::SimpleFilterEvaluator;
17use common_time::range::TimestampRange;
18use common_time::timestamp::TimeUnit;
19use datafusion_expr::{Expr, col, lit};
20use datatypes::arrow::array::{ArrayRef, UInt32Array, UInt64Array};
21use datatypes::arrow::buffer::BooleanBuffer;
22use datatypes::arrow::datatypes::{DataType, SchemaRef};
23use datatypes::arrow::record_batch::RecordBatch;
24use datatypes::value::timestamp_to_scalar_value;
25use futures::TryStreamExt;
26use object_store::ObjectStore;
27use snafu::{OptionExt, ResultExt, ensure};
28use store_api::metadata::RegionMetadataRef;
29use table::predicate::Predicate;
30
31use crate::error::{InvalidRecordBatchSnafu, RecordBatchSnafu, Result, UnexpectedSnafu};
32use crate::series_index::{
33 MAX_TS_COLUMN, METRIC_SERIES_ID_BATCH_SIZE, MIN_TS_COLUMN, MetricSeriesId,
34 MetricSeriesIdStream, ROW_COUNT_COLUMN, SeriesIndexFileHandle, TABLE_ID_COLUMN, TSID_COLUMN,
35 series_index_path, series_index_schema,
36};
37use crate::sst::parquet::index_reader::ParquetIndexReader;
38use crate::sst::parquet::prefilter::simple_tag_filters;
39
40pub struct SeriesIndexSearcher {
42 file_handle: SeriesIndexFileHandle,
44 reader: Option<ParquetIndexReader>,
47 pruning_predicate: Predicate,
49 filters: Vec<SimpleFilterEvaluator>,
50}
51
52impl SeriesIndexSearcher {
53 pub(crate) async fn try_new(
58 metadata: RegionMetadataRef,
59 object_store: ObjectStore,
60 file_handle: SeriesIndexFileHandle,
61 predicate: Option<&Predicate>,
62 time_range: Option<TimestampRange>,
63 ) -> Result<Self> {
64 series_index_schema(&metadata)?;
66 if time_range.as_ref().is_some_and(TimestampRange::is_empty) {
67 return Ok(Self {
68 file_handle,
69 reader: None,
70 pruning_predicate: Predicate::new(Vec::new()),
71 filters: Vec::new(),
72 });
73 }
74
75 let file_id = file_handle.file_id();
76 let path = series_index_path(file_id.region_id(), file_id.file_id());
77 let reader = ParquetIndexReader::open(object_store, &path).await?;
78 let unit = validate_index_schema(reader.schema())?;
79
80 let mut filters = simple_tag_filters(&metadata, None, predicate);
81 for expr in time_range_exprs(unit, time_range.as_ref()) {
82 let filter = SimpleFilterEvaluator::try_new(&expr).context(UnexpectedSnafu {
83 reason: "failed to build an internal series-index time filter",
84 })?;
85 filters.push((expr, filter));
86 }
87 let (pruning_predicate, filters) = filters_for_schema(reader.schema(), &filters);
90
91 Ok(Self {
92 file_handle,
93 reader: Some(reader),
94 pruning_predicate,
95 filters,
96 })
97 }
98
99 pub fn search(&self) -> Result<MetricSeriesIdStream> {
102 let Some(reader) = self.reader.as_ref() else {
103 return Ok(Box::pin(futures::stream::empty()));
104 };
105 let mut projection_columns = Vec::with_capacity(self.filters.len() + 2);
106 projection_columns.extend([TABLE_ID_COLUMN, TSID_COLUMN]);
107 projection_columns.extend(self.filters.iter().map(SimpleFilterEvaluator::column_name));
108 let mut batches = reader.read(&self.pruning_predicate, &projection_columns)?;
109 let filters = self.filters.clone();
110 let file_handle = self.file_handle.clone();
111
112 Ok(Box::pin(try_stream! {
113 let mut last_series = None;
114 let mut output = Vec::with_capacity(METRIC_SERIES_ID_BATCH_SIZE);
115 while let Some(batch) = batches.try_next().await? {
116 let mut mask = BooleanBuffer::new_set(batch.num_rows());
117 for filter in &filters {
118 let column = column(&batch, filter.column_name())?;
119 let evaluated = filter.evaluate_array(column).context(RecordBatchSnafu)?;
120 mask = &mask & &evaluated;
121 }
122
123 let table_ids = column(&batch, TABLE_ID_COLUMN)?
124 .as_any()
125 .downcast_ref::<UInt32Array>()
126 .context(InvalidRecordBatchSnafu {
127 reason: "series index __table_id is not UInt32",
128 })?;
129 let tsids = column(&batch, TSID_COLUMN)?
130 .as_any()
131 .downcast_ref::<UInt64Array>()
132 .context(InvalidRecordBatchSnafu {
133 reason: "series index __tsid is not UInt64",
134 })?;
135
136 for (row, matched) in mask.iter().enumerate() {
137 if !matched {
138 continue;
139 }
140 let series = MetricSeriesId {
141 table_id: table_ids.value(row),
142 tsid: tsids.value(row),
143 };
144 if last_series == Some(series) {
145 continue;
146 }
147 last_series = Some(series);
148 output.push(series);
149 if output.len() == METRIC_SERIES_ID_BATCH_SIZE {
150 yield std::mem::replace(
151 &mut output,
152 Vec::with_capacity(METRIC_SERIES_ID_BATCH_SIZE),
153 );
154 }
155 }
156 }
157 if !output.is_empty() {
158 yield output;
159 }
160 drop(file_handle);
161 }))
162 }
163}
164
165fn filters_for_schema(
166 schema: &SchemaRef,
167 filters: &[(Expr, SimpleFilterEvaluator)],
168) -> (Predicate, Vec<SimpleFilterEvaluator>) {
169 let (exprs, filters): (Vec<_>, Vec<_>) = filters
170 .iter()
171 .filter(|(_, filter)| schema.field_with_name(filter.column_name()).is_ok())
172 .cloned()
173 .unzip();
174 (Predicate::new(exprs), filters)
175}
176
177fn time_range_exprs(unit: TimeUnit, time_range: Option<&TimestampRange>) -> Vec<Expr> {
182 let Some(time_range) = time_range else {
183 return Vec::new();
184 };
185 let mut exprs = Vec::with_capacity(2);
186 if let Some(start) = time_range
189 .start()
190 .and_then(|start| start.convert_to_ceil(unit))
191 {
192 exprs.push(
193 col(MAX_TS_COLUMN).gt_eq(lit(timestamp_to_scalar_value(unit, Some(start.value())))),
194 );
195 }
196 if let Some(end) = time_range.end().and_then(|end| end.convert_to_ceil(unit)) {
199 exprs.push(col(MIN_TS_COLUMN).lt(lit(timestamp_to_scalar_value(unit, Some(end.value())))));
200 }
201 exprs
202}
203
204fn validate_index_schema(schema: &SchemaRef) -> Result<TimeUnit> {
209 for (name, data_type) in [
210 (ROW_COUNT_COLUMN, DataType::UInt64),
211 (TABLE_ID_COLUMN, DataType::UInt32),
212 (TSID_COLUMN, DataType::UInt64),
213 ] {
214 let field = schema
215 .field_with_name(name)
216 .ok()
217 .with_context(|| InvalidRecordBatchSnafu {
218 reason: format!("series index is missing internal column {name}"),
219 })?;
220 ensure!(
221 field.data_type() == &data_type && !field.is_nullable(),
222 InvalidRecordBatchSnafu {
223 reason: format!(
224 "series index internal column {name} must be non-nullable {data_type:?}, got {:?}",
225 field.data_type()
226 ),
227 }
228 );
229 }
230 let unit = |name: &str| {
231 let field = schema
232 .field_with_name(name)
233 .ok()
234 .context(InvalidRecordBatchSnafu {
235 reason: format!("series index is missing internal column {name}"),
236 })?;
237 ensure!(
238 !field.is_nullable(),
239 InvalidRecordBatchSnafu {
240 reason: format!("series index column {name} must be non-nullable"),
241 }
242 );
243 match field.data_type() {
244 DataType::Timestamp(unit, _) => Ok(unit.into()),
245 data_type => InvalidRecordBatchSnafu {
246 reason: format!(
247 "series index column {name} must be a Timestamp, got {data_type:?}"
248 ),
249 }
250 .fail(),
251 }
252 };
253 let min_unit = unit(MIN_TS_COLUMN)?;
254 let max_unit = unit(MAX_TS_COLUMN)?;
255 ensure!(
256 min_unit == max_unit,
257 InvalidRecordBatchSnafu {
258 reason: format!(
259 "series index columns {MIN_TS_COLUMN} and {MAX_TS_COLUMN} have different time units"
260 ),
261 }
262 );
263 Ok(min_unit)
264}
265
266fn column<'a>(batch: &'a RecordBatch, name: &str) -> Result<&'a ArrayRef> {
267 let index = batch
268 .schema()
269 .index_of(name)
270 .ok()
271 .with_context(|| InvalidRecordBatchSnafu {
272 reason: format!("series index batch is missing column {name}"),
273 })?;
274 Ok(batch.column(index))
275}
276
277#[cfg(test)]
278mod tests {
279 use std::sync::Arc;
280
281 use api::v1::SemanticType;
282 use datafusion_expr::{col, lit};
283 use datatypes::arrow::array::{
284 BinaryArray, TimestampMicrosecondArray, TimestampMillisecondArray,
285 TimestampNanosecondArray, TimestampSecondArray, UInt8Array,
286 };
287 use datatypes::arrow::datatypes::{Field, Schema, TimeUnit as ArrowTimeUnit};
288 use datatypes::arrow::record_batch::RecordBatch;
289 use datatypes::prelude::ConcreteDataType;
290 use datatypes::schema::ColumnSchema;
291 use futures::TryStreamExt;
292 use object_store::services::Memory;
293 use store_api::codec::PrimaryKeyEncoding;
294 use store_api::metadata::{ColumnMetadata, RegionMetadataBuilder};
295 use store_api::storage::FileId;
296 use tokio::sync::mpsc::UnboundedReceiver;
297
298 use super::*;
299 use crate::series_index::purger::PurgeRequest;
300 use crate::series_index::{
301 SeriesIndexEntry, SeriesIndexWriter, SeriesIndexWriterOptions, series_index_channel,
302 };
303 use crate::test_util::sst_util::{new_sparse_primary_key, sst_region_metadata_with_encoding};
304
305 fn object_store() -> ObjectStore {
306 ObjectStore::new(Memory::default()).unwrap()
307 }
308
309 fn index_handle(
310 metadata: &RegionMetadataRef,
311 store: &ObjectStore,
312 ) -> (SeriesIndexFileHandle, UnboundedReceiver<PurgeRequest>) {
313 let (purger, receiver) = series_index_channel(store.clone());
314 let entry = SeriesIndexEntry {
315 file_size: 0,
316 index_uuid: FileId::random(),
317 bucket_start: common_time::Timestamp::new_second(0),
318 bucket_end: common_time::Timestamp::new_second(60),
319 source_file_ids: Vec::new(),
320 min_file_sequence: 0,
321 max_file_sequence: 0,
322 compaction_window_secs: 60,
323 window_sequences: Default::default(),
324 };
325 (
326 SeriesIndexFileHandle::new(metadata.region_id, entry, purger),
327 receiver,
328 )
329 }
330
331 fn flat_batch_with_time_unit(
332 primary_keys: &[Vec<u8>],
333 timestamps: &[i64],
334 unit: ArrowTimeUnit,
335 ) -> RecordBatch {
336 let ts_column = match unit {
337 ArrowTimeUnit::Second => {
338 Arc::new(TimestampSecondArray::from(timestamps.to_vec())) as ArrayRef
339 }
340 ArrowTimeUnit::Millisecond => {
341 Arc::new(TimestampMillisecondArray::from(timestamps.to_vec())) as ArrayRef
342 }
343 ArrowTimeUnit::Microsecond => {
344 Arc::new(TimestampMicrosecondArray::from(timestamps.to_vec())) as ArrayRef
345 }
346 ArrowTimeUnit::Nanosecond => {
347 Arc::new(TimestampNanosecondArray::from(timestamps.to_vec())) as ArrayRef
348 }
349 };
350 let schema = Arc::new(Schema::new(vec![
351 Field::new("ts", DataType::Timestamp(unit, None), false),
352 Field::new("__primary_key", DataType::Binary, false),
353 Field::new("__sequence", DataType::UInt64, false),
354 Field::new("__op_type", DataType::UInt8, false),
355 ]));
356 RecordBatch::try_new(
357 schema,
358 vec![
359 ts_column,
360 Arc::new(BinaryArray::from_iter_values(
361 primary_keys.iter().map(Vec::as_slice),
362 )),
363 Arc::new(UInt64Array::from(vec![1; timestamps.len()])),
364 Arc::new(UInt8Array::from(vec![0; timestamps.len()])),
365 ],
366 )
367 .unwrap()
368 }
369
370 async fn write_index(
371 metadata: RegionMetadataRef,
372 object_store: ObjectStore,
373 rows: &[(u32, u64, &str, &str, i64)],
374 row_group_size: usize,
375 ) -> (SeriesIndexFileHandle, UnboundedReceiver<PurgeRequest>) {
376 write_index_with_time_unit(
377 metadata,
378 object_store,
379 rows,
380 row_group_size,
381 ArrowTimeUnit::Millisecond,
382 )
383 .await
384 }
385
386 async fn write_index_with_time_unit(
387 metadata: RegionMetadataRef,
388 object_store: ObjectStore,
389 rows: &[(u32, u64, &str, &str, i64)],
390 row_group_size: usize,
391 unit: ArrowTimeUnit,
392 ) -> (SeriesIndexFileHandle, UnboundedReceiver<PurgeRequest>) {
393 let (file_handle, receiver) = index_handle(&metadata, &object_store);
394 let file_id = file_handle.file_id();
395 let path = series_index_path(file_id.region_id(), file_id.file_id());
396 let primary_keys = rows
397 .iter()
398 .map(|(table_id, tsid, tag_0, tag_1, _)| {
399 new_sparse_primary_key(&[*tag_0, *tag_1], &metadata, *table_id, *tsid)
400 })
401 .collect::<Vec<_>>();
402 let timestamps = rows.iter().map(|row| row.4).collect::<Vec<_>>();
403 let mut writer = SeriesIndexWriter::try_new(
404 metadata,
405 object_store,
406 &path,
407 SeriesIndexWriterOptions { row_group_size },
408 None,
409 )
410 .await
411 .unwrap();
412 writer
413 .write(&flat_batch_with_time_unit(&primary_keys, ×tamps, unit))
414 .await
415 .unwrap();
416 writer.finish().await.unwrap();
417 (file_handle, receiver)
418 }
419
420 async fn collect_ids(stream: MetricSeriesIdStream) -> Vec<MetricSeriesId> {
421 stream
422 .try_collect::<Vec<_>>()
423 .await
424 .unwrap()
425 .into_iter()
426 .flatten()
427 .collect()
428 }
429
430 #[tokio::test]
431 async fn search_applies_candidate_tag_filters_and_time_overlap() {
432 let metadata = Arc::new(sst_region_metadata_with_encoding(
433 PrimaryKeyEncoding::Sparse,
434 ));
435 let object_store = object_store();
436 let (index, _receiver) = write_index(
437 metadata.clone(),
438 object_store.clone(),
439 &[
440 (1, 10, "a", "x", 10),
441 (1, 20, "b", "x", 20),
442 (1, 30, "a", "y", 30),
443 ],
444 2,
445 )
446 .await;
447
448 let predicate = Predicate::new(vec![
451 col("tag_0").eq(lit("a")),
452 col("field_0").gt(lit(0_u64)),
453 ]);
454 let time_range = TimestampRange::new(
455 common_time::Timestamp::new_millisecond(20),
456 common_time::Timestamp::new_millisecond(31),
457 )
458 .unwrap();
459 let searcher = SeriesIndexSearcher::try_new(
460 metadata.clone(),
461 object_store.clone(),
462 index.clone(),
463 Some(&predicate),
464 Some(time_range),
465 )
466 .await
467 .unwrap();
468 let ids = collect_ids(searcher.search().unwrap()).await;
469 assert_eq!(
470 ids,
471 vec![MetricSeriesId {
472 table_id: 1,
473 tsid: 30
474 }]
475 );
476
477 let time_range = TimestampRange::new(
480 common_time::Timestamp::new_microsecond(20_001),
481 common_time::Timestamp::new_microsecond(30_001),
482 )
483 .unwrap();
484 let searcher = SeriesIndexSearcher::try_new(
485 metadata.clone(),
486 object_store.clone(),
487 index.clone(),
488 None,
489 Some(time_range),
490 )
491 .await
492 .unwrap();
493 let ids = collect_ids(searcher.search().unwrap()).await;
494 assert_eq!(
495 ids,
496 vec![MetricSeriesId {
497 table_id: 1,
498 tsid: 30
499 }]
500 );
501
502 let time_range = TimestampRange::new(
505 common_time::Timestamp::new_millisecond(20),
506 common_time::Timestamp::new_millisecond(30),
507 )
508 .unwrap();
509 let searcher =
510 SeriesIndexSearcher::try_new(metadata, object_store, index, None, Some(time_range))
511 .await
512 .unwrap();
513 let ids = collect_ids(searcher.search().unwrap()).await;
514 assert_eq!(
515 ids,
516 vec![MetricSeriesId {
517 table_id: 1,
518 tsid: 20
519 }]
520 );
521 }
522
523 #[tokio::test]
524 async fn search_skips_filters_for_columns_missing_from_older_index() {
525 let old_metadata = Arc::new(sst_region_metadata_with_encoding(
526 PrimaryKeyEncoding::Sparse,
527 ));
528 let object_store = object_store();
529 let (index, _receiver) = write_index(
530 old_metadata.clone(),
531 object_store.clone(),
532 &[
533 (1, 10, "a", "x", 10),
534 (1, 20, "b", "x", 20),
535 (1, 30, "a", "y", 30),
536 ],
537 2,
538 )
539 .await;
540
541 let mut builder = RegionMetadataBuilder::from_existing(old_metadata.as_ref().clone());
542 builder.push_column_metadata(ColumnMetadata {
543 column_schema: ColumnSchema::new("tag_2", ConcreteDataType::string_datatype(), true),
544 semantic_type: SemanticType::Tag,
545 column_id: 4,
546 });
547 let mut primary_key = old_metadata.primary_key.clone();
548 primary_key.push(4);
549 builder.primary_key(primary_key);
550 let current_metadata = Arc::new(builder.build().unwrap());
551
552 let predicate =
553 Predicate::new(vec![col("tag_0").eq(lit("a")), col("tag_2").eq(lit("new"))]);
554 let searcher = SeriesIndexSearcher::try_new(
555 current_metadata,
556 object_store,
557 index,
558 Some(&predicate),
559 None,
560 )
561 .await
562 .unwrap();
563 let ids = collect_ids(searcher.search().unwrap()).await;
564 assert_eq!(
565 ids,
566 vec![
567 MetricSeriesId {
568 table_id: 1,
569 tsid: 10,
570 },
571 MetricSeriesId {
572 table_id: 1,
573 tsid: 30,
574 },
575 ]
576 );
577 }
578
579 #[tokio::test]
580 async fn search_uses_file_time_unit_after_time_index_widen() {
581 let metadata = Arc::new(sst_region_metadata_with_encoding(
583 PrimaryKeyEncoding::Sparse,
584 ));
585 let object_store = object_store();
586 let (index, _receiver) = write_index(
587 metadata.clone(),
588 object_store.clone(),
589 &[
590 (1, 10, "a", "x", 10),
591 (1, 20, "b", "x", 20),
592 (1, 30, "c", "x", 30),
593 ],
594 2,
595 )
596 .await;
597
598 let mut widened = (*metadata).clone();
600 for column in &mut widened.column_metadatas {
601 if column.column_schema.name == "ts" {
602 column.column_schema.data_type = ConcreteDataType::timestamp_microsecond_datatype();
603 }
604 }
605 let widened = Arc::new(widened);
606
607 let time_range = TimestampRange::new(
612 common_time::Timestamp::new_microsecond(10_001),
613 common_time::Timestamp::new_microsecond(25_000),
614 )
615 .unwrap();
616 let searcher = SeriesIndexSearcher::try_new(
617 widened,
618 object_store.clone(),
619 index.clone(),
620 None,
621 Some(time_range),
622 )
623 .await
624 .unwrap();
625 let ids = collect_ids(searcher.search().unwrap()).await;
626 assert_eq!(
627 ids,
628 vec![MetricSeriesId {
629 table_id: 1,
630 tsid: 20
631 }]
632 );
633
634 let time_range = TimestampRange::new(
637 common_time::Timestamp::new_millisecond(20),
638 common_time::Timestamp::new_millisecond(30),
639 )
640 .unwrap();
641 let searcher =
642 SeriesIndexSearcher::try_new(metadata, object_store, index, None, Some(time_range))
643 .await
644 .unwrap();
645 let ids = collect_ids(searcher.search().unwrap()).await;
646 assert_eq!(
647 ids,
648 vec![MetricSeriesId {
649 table_id: 1,
650 tsid: 20
651 }]
652 );
653 }
654
655 #[tokio::test]
656 async fn search_reads_each_file_in_its_recorded_unit() {
657 let metadata = Arc::new(sst_region_metadata_with_encoding(
660 PrimaryKeyEncoding::Sparse,
661 ));
662 let object_store = object_store();
663 let (index, _receiver) = write_index(
664 metadata.clone(),
665 object_store.clone(),
666 &[(1, 10, "a", "x", 10), (1, 20, "b", "x", 20)],
667 2,
668 )
669 .await;
670
671 let mut widened = (*metadata).clone();
674 for column in &mut widened.column_metadatas {
675 if column.column_schema.name == "ts" {
676 column.column_schema.data_type = ConcreteDataType::timestamp_microsecond_datatype();
677 }
678 }
679 let widened = Arc::new(widened);
680 let (new_index, _new_receiver) = write_index_with_time_unit(
681 widened.clone(),
682 object_store.clone(),
683 &[(1, 30, "c", "x", 15_000), (1, 40, "d", "x", 15)],
684 2,
685 ArrowTimeUnit::Microsecond,
686 )
687 .await;
688
689 let time_range = TimestampRange::new(
696 common_time::Timestamp::new_microsecond(10_500),
697 common_time::Timestamp::new_microsecond(25_000),
698 )
699 .unwrap();
700 let searcher = SeriesIndexSearcher::try_new(
701 widened.clone(),
702 object_store.clone(),
703 index,
704 None,
705 Some(time_range),
706 )
707 .await
708 .unwrap();
709 let ids = collect_ids(searcher.search().unwrap()).await;
710 assert_eq!(
711 ids,
712 vec![MetricSeriesId {
713 table_id: 1,
714 tsid: 20
715 }]
716 );
717 let searcher =
718 SeriesIndexSearcher::try_new(widened, object_store, new_index, None, Some(time_range))
719 .await
720 .unwrap();
721 let ids = collect_ids(searcher.search().unwrap()).await;
722 assert_eq!(
723 ids,
724 vec![MetricSeriesId {
725 table_id: 1,
726 tsid: 30
727 }]
728 );
729 }
730
731 fn ts_field(name: &str, unit: Option<ArrowTimeUnit>) -> Field {
732 let data_type = match unit {
733 Some(unit) => DataType::Timestamp(unit, None),
734 None => DataType::Int64,
735 };
736 Field::new(name, data_type, false)
737 }
738
739 fn index_file_schema(
740 min_unit: Option<ArrowTimeUnit>,
741 max_unit: Option<ArrowTimeUnit>,
742 ) -> SchemaRef {
743 Arc::new(Schema::new(vec![
744 ts_field(MIN_TS_COLUMN, min_unit),
745 ts_field(MAX_TS_COLUMN, max_unit),
746 Field::new(ROW_COUNT_COLUMN, DataType::UInt64, false),
747 Field::new(TABLE_ID_COLUMN, DataType::UInt32, false),
748 Field::new(TSID_COLUMN, DataType::UInt64, false),
749 ]))
750 }
751
752 #[test]
753 fn validate_index_schema_rejects_unusable_columns() {
754 let err = validate_index_schema(&index_file_schema(None, None))
756 .unwrap_err()
757 .to_string();
758 assert!(err.contains("must be a Timestamp, got Int64"), "{err}");
759
760 let err = validate_index_schema(&index_file_schema(
762 Some(ArrowTimeUnit::Millisecond),
763 Some(ArrowTimeUnit::Microsecond),
764 ))
765 .unwrap_err()
766 .to_string();
767 assert!(err.contains("have different time units"), "{err}");
768
769 assert_eq!(
770 TimeUnit::Nanosecond,
771 validate_index_schema(&index_file_schema(
772 Some(ArrowTimeUnit::Nanosecond),
773 Some(ArrowTimeUnit::Nanosecond),
774 ))
775 .unwrap()
776 );
777 }
778
779 #[tokio::test]
780 async fn search_streams_fixed_size_batches() {
781 let metadata = Arc::new(sst_region_metadata_with_encoding(
782 PrimaryKeyEncoding::Sparse,
783 ));
784 let object_store = object_store();
785 let rows = (0..501_u64)
786 .map(|tsid| (1, tsid, "a", "x", tsid as i64))
787 .collect::<Vec<_>>();
788 let (index, _receiver) =
789 write_index(metadata.clone(), object_store.clone(), &rows, 100).await;
790
791 let searcher = SeriesIndexSearcher::try_new(metadata, object_store, index, None, None)
792 .await
793 .unwrap();
794 let batches = searcher
795 .search()
796 .unwrap()
797 .try_collect::<Vec<_>>()
798 .await
799 .unwrap();
800 assert_eq!(batches.iter().map(Vec::len).collect::<Vec<_>>(), [500, 1]);
801 assert_eq!(
802 batches[0][0],
803 MetricSeriesId {
804 table_id: 1,
805 tsid: 0
806 }
807 );
808 assert_eq!(
809 batches[1][0],
810 MetricSeriesId {
811 table_id: 1,
812 tsid: 500
813 }
814 );
815 }
816
817 #[tokio::test]
818 async fn search_streams_pin_deleted_index_until_released() {
819 let metadata = Arc::new(sst_region_metadata_with_encoding(
820 PrimaryKeyEncoding::Sparse,
821 ));
822 let store = object_store();
823 let rows = (0..501_u64)
824 .map(|tsid| (1, tsid, "a", "x", tsid as i64))
825 .collect::<Vec<_>>();
826 let (index, mut receiver) = write_index(metadata.clone(), store.clone(), &rows, 100).await;
827 let file_id = index.file_id();
828 let searcher = SeriesIndexSearcher::try_new(metadata, store, index.clone(), None, None)
829 .await
830 .unwrap();
831 index.mark_deleted();
832 drop(index);
833 assert!(receiver.try_recv().is_err());
834
835 let mut stream = searcher.search().unwrap();
836 let cancelled = searcher.search().unwrap();
837 drop(searcher);
838 assert!(receiver.try_recv().is_err());
840 assert_eq!(stream.try_next().await.unwrap().unwrap().len(), 500);
841 assert!(receiver.try_recv().is_err());
842 assert_eq!(
843 collect_ids(stream).await,
844 vec![MetricSeriesId {
845 table_id: 1,
846 tsid: 500
847 }]
848 );
849 assert!(receiver.try_recv().is_err());
850
851 drop(cancelled);
853 assert_eq!(receiver.try_recv().unwrap().file_id, file_id);
854 }
855
856 #[tokio::test]
857 async fn search_prunes_row_groups_and_empty_ranges() {
858 let metadata = Arc::new(sst_region_metadata_with_encoding(
859 PrimaryKeyEncoding::Sparse,
860 ));
861 let object_store = object_store();
862 let (index, _receiver) = write_index(
863 metadata.clone(),
864 object_store.clone(),
865 &[
866 (1, 0, "a", "x", 0),
867 (1, 1, "b", "x", 1),
868 (1, 2, "m", "x", 2),
869 (1, 3, "m", "x", 3),
870 (1, 4, "y", "x", 4),
871 (1, 5, "z", "x", 5),
872 ],
873 2,
874 )
875 .await;
876
877 let file_id = index.file_id();
878 let path = series_index_path(file_id.region_id(), file_id.file_id());
879 let reader = ParquetIndexReader::open(object_store.clone(), &path)
880 .await
881 .unwrap();
882 let unit = validate_index_schema(reader.schema()).unwrap();
883
884 let predicate = Predicate::new(vec![col("tag_0").eq(lit("m"))]);
886 let mut filters = simple_tag_filters(&metadata, None, Some(&predicate));
887 let (pruning_predicate, _) = filters_for_schema(reader.schema(), &filters);
888 assert_eq!(reader.row_groups_to_read(&pruning_predicate), vec![1]);
889
890 let time_range = TimestampRange::new(
893 common_time::Timestamp::new_millisecond(2),
894 common_time::Timestamp::new_millisecond(4),
895 )
896 .unwrap();
897 for expr in time_range_exprs(unit, Some(&time_range)) {
898 let filter = SimpleFilterEvaluator::try_new(&expr)
899 .context(UnexpectedSnafu {
900 reason: "failed to build an internal series-index time filter",
901 })
902 .unwrap();
903 filters.push((expr, filter));
904 }
905 let (pruning_predicate, _) = filters_for_schema(reader.schema(), &filters);
906 assert_eq!(reader.row_groups_to_read(&pruning_predicate), vec![1]);
907
908 let (missing_index, _receiver) = index_handle(&metadata, &object_store);
910 let empty = SeriesIndexSearcher::try_new(
911 metadata,
912 object_store,
913 missing_index,
914 None,
915 Some(TimestampRange::empty()),
916 )
917 .await
918 .unwrap();
919 assert!(
920 empty
921 .search()
922 .unwrap()
923 .try_collect::<Vec<_>>()
924 .await
925 .unwrap()
926 .is_empty()
927 );
928 }
929}