1use std::cmp::Reverse;
18use std::collections::HashSet;
19use std::sync::Arc;
20use std::time::Instant;
21
22use async_stream::try_stream;
23use datafusion::execution::memory_pool::{MemoryConsumer, MemoryPool};
24use datafusion::physical_expr::{LexOrdering, PhysicalSortExpr};
25use datafusion::physical_plan::expressions::Column;
26use datafusion::physical_plan::metrics::{BaselineMetrics, ExecutionPlanMetricsSet, MetricBuilder};
27use datafusion::physical_plan::sorts::streaming_merge::StreamingMergeBuilder;
28use datafusion::physical_plan::stream::RecordBatchStreamAdapter;
29use datafusion_common::DataFusionError;
30use datatypes::arrow::array::{Array, BinaryArray, BinaryBuilder};
31use datatypes::arrow::compute::SortOptions;
32use datatypes::arrow::datatypes::{DataType, Field, Schema, SchemaRef};
33use datatypes::arrow::record_batch::RecordBatch;
34use datatypes::prelude::ConcreteDataType;
35use futures::{StreamExt, TryStreamExt};
36use mito_codec::row_converter::{PrimaryKeyFilter, SparsePrimaryKeyCodec};
37use snafu::{OptionExt, ResultExt, ensure};
38use store_api::codec::PrimaryKeyEncoding;
39use store_api::region_engine::PartitionRange;
40use store_api::storage::consts::{PRIMARY_KEY_COLUMN_NAME, ReservedColumnId};
41use tokio::sync::Semaphore;
42
43use crate::error::{
44 InvalidRequestSnafu, JoinSnafu, MergeCandidateSeriesSnafu, NewRecordBatchSnafu, Result,
45 UnexpectedSnafu,
46};
47use crate::read::BoxedRecordBatchStream;
48use crate::read::pruner::{PartitionPruner, Pruner};
49use crate::read::range::RowGroupIndex;
50use crate::read::range_cache::{
51 build_candidate_range_cache_key, cache_flat_range_stream, cached_flat_range_stream,
52};
53use crate::read::scan_region::StreamContext;
54use crate::read::scan_util::{PartitionMetrics, new_filter_metrics, scan_flat_mem_ranges};
55use crate::series_index::{
56 METRIC_SERIES_ID_BATCH_SIZE, MetricSeriesId, MetricSeriesIdStream, SeriesIndexFileHandle,
57 SeriesIndexReadContext, SeriesIndexSearcher,
58};
59use crate::sst::parquet::DEFAULT_READ_BATCH_SIZE;
60use crate::sst::parquet::format::PrimaryKeyArray;
61use crate::sst::parquet::prefilter::{
62 CachedPrimaryKeyFilter, build_primary_key_filter, prefilter_flat_batch_by_primary_key,
63};
64use crate::sst::parquet::reader::ReaderMetrics;
65use crate::sst::parquet::row_group::ParquetFetchMetrics;
66
67pub(crate) struct SeriesCandidateScanner {
69 stream_ctx: Arc<StreamContext>,
70 partitions: Vec<Vec<PartitionRange>>,
71 partition_pruner: Arc<PartitionPruner>,
72 candidate_pruner: Arc<PartitionPruner>,
73 coverage: Arc<SeriesIndexCoverage>,
74 range_semaphore: Arc<Semaphore>,
75 memory_pool: Arc<dyn MemoryPool>,
76 metrics_set: ExecutionPlanMetricsSet,
77 part_metrics: PartitionMetrics,
78}
79
80impl SeriesCandidateScanner {
81 pub(crate) fn try_new(
90 stream_ctx: Arc<StreamContext>,
91 partitions: Vec<Vec<PartitionRange>>,
92 pruner: Arc<Pruner>,
93 range_semaphore: Arc<Semaphore>,
94 memory_pool: Arc<dyn MemoryPool>,
95 metrics_set: ExecutionPlanMetricsSet,
96 part_metrics: PartitionMetrics,
97 ) -> Result<Self> {
98 validate_metric_metadata(&stream_ctx)?;
99 #[cfg(feature = "enterprise")]
100 ensure!(
101 stream_ctx.input.extension_ranges().is_empty(),
102 InvalidRequestSnafu {
103 region_id: stream_ctx.input.region_metadata().region_id,
104 reason: "candidate-series scan does not support extension ranges; use the legacy series-scan path",
105 }
106 );
107 ensure!(
108 !pruner.predicate_prefilter_enabled(),
109 UnexpectedSnafu {
110 reason: format!(
111 "candidate-series scan for region {} requires a pruner without predicate prefiltering",
112 stream_ctx.input.region_metadata().region_id
113 ),
114 }
115 );
116 let all_ranges = partitions.iter().flatten().copied().collect::<Vec<_>>();
117 pruner.add_partition_ranges(&all_ranges);
118 let coverage = Arc::new(SeriesIndexCoverage::new(&stream_ctx, &all_ranges));
119 let partition_pruner = Arc::new(PartitionPruner::new(pruner.clone(), &all_ranges));
120 let candidate_pruner = if coverage.covered_files.is_empty() {
121 partition_pruner.clone()
122 } else {
123 Arc::new(
124 PartitionPruner::new(pruner, &all_ranges).excluding_files(&coverage.covered_files),
125 )
126 };
127 Ok(Self {
128 stream_ctx,
129 partitions,
130 partition_pruner,
131 candidate_pruner,
132 coverage,
133 range_semaphore,
134 memory_pool,
135 metrics_set,
136 part_metrics,
137 })
138 }
139
140 pub(crate) async fn build_stream(&self) -> Result<MetricSeriesIdStream> {
142 let all_ranges = self
143 .partitions
144 .iter()
145 .flatten()
146 .copied()
147 .collect::<Vec<_>>();
148 let range_builder = SeriesCandidateRangeBuilder {
149 stream_ctx: self.stream_ctx.clone(),
150 partition_pruner: self.candidate_pruner.clone(),
151 coverage: self.coverage.clone(),
152 range_semaphore: self.range_semaphore.clone(),
153 memory_pool: self.memory_pool.clone(),
154 metrics_set: self.metrics_set.clone(),
155 part_metrics: self.part_metrics.clone(),
156 };
157 let mut tasks = Vec::with_capacity(all_ranges.len());
158 for (range_idx, part_range) in all_ranges.into_iter().enumerate() {
159 let range_builder = range_builder.clone();
160 tasks.push(common_runtime::spawn_query(async move {
161 let _permit = range_builder
162 .range_semaphore
163 .clone()
164 .acquire_owned()
165 .await
166 .map_err(|error| {
167 UnexpectedSnafu {
168 reason: format!("failed to acquire candidate range permit: {error}"),
169 }
170 .build()
171 })?;
172 range_builder
173 .build_range_stream(part_range, range_idx)
174 .await
175 }));
176 }
177
178 let mut range_streams = Vec::with_capacity(tasks.len());
179 for task in tasks {
180 range_streams.push(task.await.context(JoinSnafu)??);
181 }
182
183 if let Some(context) = &self.stream_ctx.input.series_index {
184 MetricBuilder::new(&self.metrics_set)
185 .counter("candidate_index_files", self.partitions.len())
186 .add(self.coverage.indexes.len());
187 MetricBuilder::new(&self.metrics_set)
188 .counter("candidate_index_covered_ssts", self.partitions.len())
189 .add(self.coverage.covered_files.len());
190 for index in &self.coverage.indexes {
191 range_streams.push(index_primary_key_stream(
192 self.stream_ctx.clone(),
193 context.clone(),
194 index.clone(),
195 self.range_semaphore.clone(),
196 ));
197 }
198 }
199
200 let merged = merge_primary_key_streams(
203 range_streams,
204 self.memory_pool.clone(),
205 &self.metrics_set,
206 self.partitions.len(),
207 "SeriesCandidateScanner::final_merge",
208 )?;
209 decode_metric_series(merged, self.stream_ctx.input.region_metadata().clone())
210 }
211
212 pub(crate) fn partition_pruner(&self) -> Arc<PartitionPruner> {
214 self.partition_pruner.clone()
215 }
216}
217
218#[derive(Default)]
220struct SeriesIndexCoverage {
221 indexes: Vec<SeriesIndexFileHandle>,
222 covered_files: HashSet<usize>,
224}
225
226impl SeriesIndexCoverage {
227 fn new(stream_ctx: &StreamContext, ranges: &[PartitionRange]) -> Self {
228 let Some(context) = &stream_ctx.input.series_index else {
229 return Self::default();
230 };
231 let mut uncovered: HashSet<_> = ranges
232 .iter()
233 .flat_map(|range| {
234 stream_ctx.ranges[range.identifier]
235 .row_group_indices
236 .iter()
237 .filter(|index| stream_ctx.is_file_range_index(**index))
238 .map(|index| index.index - stream_ctx.input.num_memtables())
239 })
240 .collect();
241 let region_id = stream_ctx.input.region_metadata().region_id;
242 let mut candidates: Vec<_> = context
243 .version
244 .series_indexes
245 .values()
246 .map(|index| {
247 let files: HashSet<_> = uncovered
248 .iter()
249 .copied()
250 .filter(|file_index| {
251 index
252 .entry()
253 .covers_file(stream_ctx.input.files[*file_index].meta_ref(), region_id)
254 })
255 .collect();
256 (index, files)
257 })
258 .collect();
259 let mut coverage = Self::default();
260 while let Some((position, count)) = candidates
261 .iter()
262 .enumerate()
263 .map(|(position, (index, files))| {
264 (
265 position,
266 files.intersection(&uncovered).count(),
267 index.entry(),
268 )
269 })
270 .max_by_key(|(_, count, entry)| {
271 (
272 *count,
273 entry.max_file_sequence,
274 Reverse(entry.index_uuid.as_bytes()),
275 )
276 })
277 .map(|(position, count, _)| (position, count))
278 {
279 if count == 0 {
280 break;
281 }
282 let (index, files) = candidates.swap_remove(position);
283 for file in files {
284 if uncovered.remove(&file) {
285 coverage.covered_files.insert(file);
286 }
287 }
288 coverage.indexes.push(index.clone());
289 }
290 coverage
291 }
292
293 fn covers_source(&self, stream_ctx: &StreamContext, index: RowGroupIndex) -> bool {
294 stream_ctx.is_file_range_index(index)
295 && self
296 .covered_files
297 .contains(&(index.index - stream_ctx.input.num_memtables()))
298 }
299}
300
301fn index_primary_key_stream(
303 stream_ctx: Arc<StreamContext>,
304 context: SeriesIndexReadContext,
305 index: SeriesIndexFileHandle,
306 semaphore: Arc<Semaphore>,
307) -> BoxedRecordBatchStream {
308 Box::pin(try_stream! {
309 let metadata = stream_ctx.input.region_metadata();
310 let codec = SparsePrimaryKeyCodec::new(metadata);
311 let mut series = {
312 let _permit = semaphore.acquire().await.map_err(|error| UnexpectedSnafu {
313 reason: format!("failed to acquire candidate index permit: {error}"),
314 }.build())?;
315 SeriesIndexSearcher::try_new(
316 metadata.clone(),
317 context.store.clone(),
318 index,
319 stream_ctx.input.predicate_group().predicate(),
320 stream_ctx.input.time_range,
321 ).await?.search()?
322 };
323 loop {
324 let batch = {
325 let _permit = semaphore.acquire().await.map_err(|error| UnexpectedSnafu {
326 reason: format!("failed to acquire candidate index permit: {error}"),
327 }.build())?;
328 series.try_next().await?
329 };
330 let Some(batch) = batch else { break };
331 let mut builder = BinaryBuilder::new();
332 let mut key = Vec::new();
333 for series in batch {
334 key.clear();
335 codec.encode_internal(series.table_id, series.tsid, &mut key)
336 .context(crate::error::EncodeSnafu)?;
337 builder.append_value(&key);
338 }
339 yield RecordBatch::try_new(primary_key_schema(), vec![Arc::new(builder.finish())])
340 .context(NewRecordBatchSnafu)?;
341 }
342 })
343}
344
345#[derive(Clone)]
346struct SeriesCandidateRangeBuilder {
347 stream_ctx: Arc<StreamContext>,
348 coverage: Arc<SeriesIndexCoverage>,
349 partition_pruner: Arc<PartitionPruner>,
350 range_semaphore: Arc<Semaphore>,
351 memory_pool: Arc<dyn MemoryPool>,
352 metrics_set: ExecutionPlanMetricsSet,
353 part_metrics: PartitionMetrics,
354}
355
356impl SeriesCandidateRangeBuilder {
357 async fn build_range_stream(
358 &self,
359 part_range: PartitionRange,
360 merge_partition: usize,
361 ) -> Result<BoxedRecordBatchStream> {
362 let range_meta = &self.stream_ctx.ranges[part_range.identifier];
363 let replaced_sources = range_meta
366 .row_group_indices
367 .iter()
368 .any(|index| self.coverage.covers_source(&self.stream_ctx, *index));
369 let cache_key = if replaced_sources {
370 None
371 } else {
372 build_candidate_range_cache_key(&self.stream_ctx, &part_range)
373 };
374 if let Some(key) = cache_key.as_ref() {
375 if let Some(value) = self.stream_ctx.input.cache_strategy.get_range_result(key) {
376 self.part_metrics.inc_range_cache_hit();
377 return Ok(cached_flat_range_stream(value));
378 }
379 self.part_metrics.inc_range_cache_miss();
380 }
381
382 let mut sources = Vec::with_capacity(range_meta.row_group_indices.len());
383 for index in &range_meta.row_group_indices {
384 let source = self.build_source(*index, range_meta.time_range).await?;
385 if let Some(source) = source {
386 sources.push(source);
387 }
388 }
389
390 let sources = self.stream_ctx.input.create_parallel_flat_sources(
391 sources,
392 self.range_semaphore.clone(),
393 2,
394 )?;
395 let stream = merge_primary_key_streams(
396 sources,
397 self.memory_pool.clone(),
398 &self.metrics_set,
399 merge_partition,
400 "SeriesCandidateScanner::range_merge",
401 )?;
402
403 Ok(match cache_key {
404 Some(key) => cache_flat_range_stream(
405 stream,
406 self.stream_ctx.input.cache_strategy.clone(),
407 key,
408 self.part_metrics.clone(),
409 ),
410 None => stream,
411 })
412 }
413
414 async fn build_source(
415 &self,
416 index: RowGroupIndex,
417 time_range: crate::sst::file::FileTimeRange,
418 ) -> Result<Option<BoxedRecordBatchStream>> {
419 let metadata = self.stream_ctx.input.region_metadata().clone();
420 if self.stream_ctx.is_mem_range_index(index) {
421 let raw = scan_flat_mem_ranges(
422 self.stream_ctx.clone(),
423 self.part_metrics.clone(),
424 index,
425 time_range,
426 );
427 let filter = build_primary_key_filter(
428 &metadata,
429 None,
430 self.stream_ctx.input.predicate_group().predicate(),
431 );
432 return Ok(Some(candidate_primary_key_stream(Box::pin(raw), filter)));
433 }
434
435 if self.stream_ctx.is_file_range_index(index) {
436 if self.coverage.covers_source(&self.stream_ctx, index) {
437 return Ok(None);
441 }
442 let file = self.stream_ctx.input.file_from_index(index);
443 let predicate = self.stream_ctx.input.predicate_for_file(file);
444 if self
445 .partition_pruner
446 .try_skip_manifest_pruned_file_range(index, &self.part_metrics)
447 {
448 return Ok(None);
449 }
450 let mut reader_metrics = ReaderMetrics {
451 filter_metrics: new_filter_metrics(self.part_metrics.explain_verbose()),
452 ..Default::default()
453 };
454 let ranges = self
455 .partition_pruner
456 .build_file_ranges(index, &self.part_metrics, &mut reader_metrics)
457 .await?;
458 self.part_metrics.inc_num_file_ranges(ranges.len());
459 self.part_metrics
460 .merge_reader_metrics(&reader_metrics, None);
461
462 let filter = ranges.first().and_then(|range| {
465 build_primary_key_filter(
466 range.region_metadata(),
467 Some(metadata.as_ref()),
468 predicate.as_ref(),
469 )
470 });
471 let part_metrics = self.part_metrics.clone();
472 let raw = Box::pin(try_stream! {
473 let fetch_metrics = part_metrics
474 .explain_verbose()
475 .then(|| Arc::new(ParquetFetchMetrics::default()));
476 let mut reader_metrics = ReaderMetrics {
477 fetch_metrics: fetch_metrics.clone(),
478 ..Default::default()
479 };
480 for range in ranges {
481 let build_start = Instant::now();
482 let Some(mut reader) = range
483 .primary_key_reader(fetch_metrics.as_deref())
484 .await?
485 else {
486 continue;
487 };
488 reader_metrics.build_cost += build_start.elapsed();
489
490 let scan_start = Instant::now();
491 while let Some(batch) = reader.try_next().await? {
492 reader_metrics.num_record_batches += 1;
493 reader_metrics.num_batches += 1;
494 reader_metrics.num_rows += batch.num_rows();
495 yield batch;
496 }
497 reader_metrics.scan_cost += scan_start.elapsed();
498 }
499 reader_metrics.observe_rows("candidate_series");
500 part_metrics.merge_reader_metrics(&reader_metrics, None);
501 });
502 return Ok(Some(candidate_primary_key_stream(raw, filter)));
503 }
504
505 UnexpectedSnafu {
506 reason: format!(
507 "candidate-series scan received unsupported range index {}",
508 index.index
509 ),
510 }
511 .fail()
512 }
513}
514
515pub(crate) fn is_sparse_metric_metadata(metadata: &store_api::metadata::RegionMetadataRef) -> bool {
516 let valid_prefix = metadata
517 .primary_key
518 .starts_with(&[ReservedColumnId::table_id(), ReservedColumnId::tsid()]);
519 let valid_types = metadata
520 .column_by_id(ReservedColumnId::table_id())
521 .zip(metadata.column_by_id(ReservedColumnId::tsid()))
522 .is_some_and(|(table_id, tsid)| {
523 table_id.column_schema.data_type == ConcreteDataType::uint32_datatype()
524 && tsid.column_schema.data_type == ConcreteDataType::uint64_datatype()
525 });
526
527 metadata.primary_key_encoding == PrimaryKeyEncoding::Sparse && valid_prefix && valid_types
528}
529
530pub(crate) fn validate_metric_metadata(stream_ctx: &StreamContext) -> Result<()> {
531 let metadata = stream_ctx.input.region_metadata();
532 ensure!(
533 is_sparse_metric_metadata(metadata),
534 InvalidRequestSnafu {
535 region_id: metadata.region_id,
536 reason: "candidate-series scan requires sparse (__table_id, __tsid) primary keys",
537 }
538 );
539 Ok(())
540}
541
542fn primary_key_schema() -> SchemaRef {
543 Arc::new(Schema::new(vec![Field::new(
544 PRIMARY_KEY_COLUMN_NAME,
545 DataType::Binary,
546 false,
547 )]))
548}
549
550fn candidate_primary_key_stream(
552 mut input: BoxedRecordBatchStream,
553 mut filter: Option<CachedPrimaryKeyFilter>,
554) -> BoxedRecordBatchStream {
555 Box::pin(try_stream! {
556 let mut last_primary_key = Vec::new();
557 let mut has_last = false;
558 while let Some(batch) = input.try_next().await? {
559 if let Some(batch) = normalize_candidate_batch(
560 batch,
561 filter.as_mut(),
562 &mut last_primary_key,
563 &mut has_last,
564 )? {
565 yield batch;
566 }
567 }
568 })
569}
570
571fn normalize_candidate_batch(
572 mut batch: RecordBatch,
573 filter: Option<&mut CachedPrimaryKeyFilter>,
574 last_primary_key: &mut Vec<u8>,
575 has_last: &mut bool,
576) -> Result<Option<RecordBatch>> {
577 let pk_idx = batch
578 .schema()
579 .column_with_name(PRIMARY_KEY_COLUMN_NAME)
580 .map(|(idx, _)| idx)
581 .context(UnexpectedSnafu {
582 reason: "candidate source does not contain __primary_key",
583 })?;
584 if let Some(filter) = filter {
585 let Some(filtered) = prefilter_flat_batch_by_primary_key(
586 batch,
587 pk_idx,
588 filter as &mut dyn PrimaryKeyFilter,
589 )?
590 else {
591 return Ok(None);
592 };
593 batch = filtered;
594 }
595
596 let pk_column = batch.column(pk_idx);
597 let mut builder = BinaryBuilder::new();
598 if let Some(array) = pk_column.as_any().downcast_ref::<PrimaryKeyArray>() {
599 let values = array
600 .values()
601 .as_any()
602 .downcast_ref::<BinaryArray>()
603 .context(UnexpectedSnafu {
604 reason: "dictionary primary-key values are not binary",
605 })?;
606 for key in array.keys().values() {
607 append_unique_primary_key(
608 values.value(*key as usize),
609 &mut builder,
610 last_primary_key,
611 has_last,
612 );
613 }
614 } else if let Some(array) = pk_column.as_any().downcast_ref::<BinaryArray>() {
615 for value in array.iter().flatten() {
616 append_unique_primary_key(value, &mut builder, last_primary_key, has_last);
617 }
618 } else {
619 return UnexpectedSnafu {
620 reason: format!(
621 "primary-key column is neither binary nor dictionary, got {:?}",
622 pk_column.data_type()
623 ),
624 }
625 .fail();
626 }
627
628 let array = builder.finish();
629 if array.is_empty() {
630 return Ok(None);
631 }
632 let batch = RecordBatch::try_new(primary_key_schema(), vec![Arc::new(array)])
633 .context(NewRecordBatchSnafu)?;
634 Ok(Some(batch))
635}
636
637fn append_unique_primary_key(
638 value: &[u8],
639 builder: &mut BinaryBuilder,
640 last_primary_key: &mut Vec<u8>,
641 has_last: &mut bool,
642) {
643 if !*has_last || last_primary_key != value {
644 builder.append_value(value);
645 last_primary_key.clear();
646 last_primary_key.extend_from_slice(value);
647 *has_last = true;
648 }
649}
650
651fn merge_primary_key_streams(
652 sources: Vec<BoxedRecordBatchStream>,
653 memory_pool: Arc<dyn MemoryPool>,
654 metrics_set: &ExecutionPlanMetricsSet,
655 partition: usize,
656 consumer_name: &'static str,
657) -> Result<BoxedRecordBatchStream> {
658 if sources.is_empty() {
659 return Ok(Box::pin(futures::stream::empty()));
660 }
661 if sources.len() == 1 {
662 return Ok(sources.into_iter().next().unwrap());
663 }
664
665 let schema = primary_key_schema();
666 let df_sources = sources
667 .into_iter()
668 .map(|source| {
669 let stream = source.map_err(|error| DataFusionError::External(Box::new(error)));
670 Box::pin(RecordBatchStreamAdapter::new(schema.clone(), stream)) as _
671 })
672 .collect();
673 let ordering = LexOrdering::new([PhysicalSortExpr {
674 expr: Arc::new(Column::new(PRIMARY_KEY_COLUMN_NAME, 0)),
675 options: SortOptions {
676 descending: false,
677 nulls_first: false,
678 },
679 }])
680 .unwrap();
683 let reservation = MemoryConsumer::new(consumer_name).register(&memory_pool);
684 let mut merged = StreamingMergeBuilder::new()
685 .with_streams(df_sources)
686 .with_schema(schema)
687 .with_expressions(&ordering)
688 .with_metrics(BaselineMetrics::new(metrics_set, partition))
689 .with_batch_size(DEFAULT_READ_BATCH_SIZE)
690 .with_reservation(reservation)
691 .build()
692 .context(MergeCandidateSeriesSnafu)?;
693
694 Ok(Box::pin(try_stream! {
695 while let Some(batch) = merged.next().await {
696 yield batch.context(MergeCandidateSeriesSnafu)?;
697 }
698 }))
699}
700
701fn decode_metric_series(
702 mut input: BoxedRecordBatchStream,
703 metadata: store_api::metadata::RegionMetadataRef,
704) -> Result<MetricSeriesIdStream> {
705 let codec = SparsePrimaryKeyCodec::new(&metadata);
706 Ok(Box::pin(try_stream! {
707 let mut last_series = None;
708 let mut output = Vec::with_capacity(METRIC_SERIES_ID_BATCH_SIZE);
709 while let Some(batch) = input.try_next().await? {
710 let array = batch
711 .column(0)
712 .as_any()
713 .downcast_ref::<BinaryArray>()
714 .context(UnexpectedSnafu {
715 reason: "merged candidate primary key is not binary",
716 })?;
717 for primary_key in array.iter().flatten() {
718 let (table_id, tsid) = codec
719 .decode_ids(primary_key)
720 .context(crate::error::DecodeSnafu)?;
721 let series = MetricSeriesId { table_id, tsid };
722 if last_series == Some(series) {
723 continue;
724 }
725 last_series = Some(series);
726 output.push(series);
727 if output.len() == METRIC_SERIES_ID_BATCH_SIZE {
728 yield std::mem::replace(
729 &mut output,
730 Vec::with_capacity(METRIC_SERIES_ID_BATCH_SIZE),
731 );
732 }
733 }
734 }
735 if !output.is_empty() {
736 yield output;
737 }
738 }))
739}
740
741#[cfg(test)]
742mod tests {
743 use std::num::NonZeroU64;
744 use std::time::Instant;
745
746 use common_time::Timestamp;
747 use datafusion::execution::memory_pool::UnboundedMemoryPool;
748 use datafusion_expr::{col, lit};
749 use datatypes::arrow::array::{
750 ArrayRef, DictionaryArray, TimestampMillisecondArray, UInt8Array, UInt32Array, UInt64Array,
751 };
752 use datatypes::arrow::datatypes::UInt32Type;
753 use futures::TryStreamExt;
754 use store_api::codec::PrimaryKeyEncoding;
755 use store_api::storage::FileId;
756 use table::predicate::Predicate;
757
758 use super::*;
759 use crate::cache::{CacheManager, CacheStrategy};
760 use crate::read::flat_projection::FlatProjectionMapper;
761 use crate::read::pruner::PrunerOptions;
762 use crate::read::scan_region::ScanInput;
763 use crate::read::scan_util::PartitionMetrics;
764 use crate::series_index::{
765 SeriesIndexEntry, SeriesIndexVersion, SeriesIndexWriter, SeriesIndexWriterOptions,
766 series_index_channel, series_index_path,
767 };
768 use crate::sst::file::{FileHandle, FileMeta};
769 use crate::test_util::new_noop_file_purger;
770 use crate::test_util::scheduler_util::SchedulerEnv;
771 use crate::test_util::sst_util::sst_region_metadata_with_encoding;
772
773 async fn indexed_scanner() -> (SchedulerEnv, SeriesCandidateScanner, Arc<Pruner>) {
775 let env = SchedulerEnv::new().await;
776 let metadata = Arc::new(sst_region_metadata_with_encoding(
777 PrimaryKeyEncoding::Sparse,
778 ));
779 let store = env.access_layer.object_store().clone();
780 let entry = SeriesIndexEntry {
781 file_size: 0,
782 index_uuid: FileId::random(),
783 bucket_start: Timestamp::new_millisecond(0),
784 bucket_end: Timestamp::new_millisecond(20),
785 source_file_ids: Vec::new(),
786 min_file_sequence: 1,
787 max_file_sequence: 2,
788 compaction_window_secs: 1,
789 window_sequences: Default::default(),
790 };
791 let path = series_index_path(metadata.region_id, entry.index_uuid);
792 let codec = SparsePrimaryKeyCodec::new(&metadata);
793 let keys: Vec<_> = (0..1001)
794 .map(|tsid| {
795 let mut key = Vec::new();
796 codec.encode_internal(1, tsid, &mut key).unwrap();
797 key
798 })
799 .collect();
800 let batch = RecordBatch::try_from_iter(vec![
801 (
802 "ts",
803 Arc::new(TimestampMillisecondArray::from(vec![10; keys.len()])) as ArrayRef,
804 ),
805 (
806 "__primary_key",
807 Arc::new(BinaryArray::from_iter_values(&keys)),
808 ),
809 (
810 "__sequence",
811 Arc::new(UInt64Array::from_value(1, keys.len())),
812 ),
813 ("__op_type", Arc::new(UInt8Array::from_value(0, keys.len()))),
814 ])
815 .unwrap();
816 let mut writer = SeriesIndexWriter::try_new(
817 metadata.clone(),
818 store.clone(),
819 &path,
820 SeriesIndexWriterOptions {
821 row_group_size: 500,
822 },
823 None,
824 )
825 .await
826 .unwrap();
827 writer.write(&batch).await.unwrap();
828 writer.finish().await.unwrap();
829 let (purger, _receiver) = series_index_channel(store.clone());
830 let handle = SeriesIndexFileHandle::new(metadata.region_id, entry.clone(), purger);
831 let context = SeriesIndexReadContext {
832 store,
833 version: Arc::new(SeriesIndexVersion {
834 series_indexes: [(entry.index_uuid, handle)].into(),
835 ..Default::default()
836 }),
837 };
838 let files = (0..2)
839 .map(|i| {
840 FileHandle::new(
841 FileMeta {
842 region_id: metadata.region_id,
843 file_id: FileId::random(),
844 time_range: (
845 Timestamp::new_millisecond(i * 10),
846 Timestamp::new_millisecond(i * 10 + 9),
847 ),
848 sequence: NonZeroU64::new(i as u64 + 1),
849 num_row_groups: 1,
850 ..Default::default()
851 },
852 new_noop_file_purger(),
853 )
854 })
855 .collect();
856 let mapper =
857 FlatProjectionMapper::new(&metadata, 0..metadata.column_metadatas.len()).unwrap();
858 let input = ScanInput::builder(env.access_layer.clone(), mapper)
859 .with_predicate(
860 crate::read::scan_region::PredicateGroup::new(
861 &metadata,
862 &[col("__table_id").eq(lit(1_u32))],
863 )
864 .unwrap(),
865 )
866 .with_files(files)
867 .with_series_index(Some(context))
868 .with_cache(CacheStrategy::EnableAll(Arc::new(
869 CacheManager::builder()
870 .range_result_cache_size(1024 * 1024)
871 .build(),
872 )))
873 .build();
874 let stream_ctx = Arc::new(StreamContext::seq_scan_ctx(input));
875 let ranges = stream_ctx.partition_ranges();
876 assert_eq!(2, ranges.len());
877 let metrics_set = ExecutionPlanMetricsSet::new();
878 let part_metrics = PartitionMetrics::new(
879 metadata.region_id,
880 2,
881 "candidate-test",
882 Instant::now(),
883 false,
884 &metrics_set,
885 );
886 let pruner = Arc::new(Pruner::new_with_options(
887 stream_ctx.clone(),
888 1,
889 PrunerOptions {
890 retain_builders: true,
891 enable_predicate_prefilter: false,
892 },
893 ));
894 let scanner = SeriesCandidateScanner::try_new(
895 stream_ctx,
896 ranges.into_iter().map(|range| vec![range]).collect(),
897 pruner.clone(),
898 Arc::new(Semaphore::new(1)),
899 Arc::new(UnboundedMemoryPool::default()),
900 metrics_set,
901 part_metrics,
902 )
903 .unwrap();
904 (env, scanner, pruner)
905 }
906
907 #[tokio::test]
908 async fn index_is_shared_across_ranges_without_caching_partial_candidates() {
909 let (_env, scanner, pruner) = indexed_scanner().await;
910 let keys: Vec<_> = scanner
911 .partitions
912 .iter()
913 .flatten()
914 .map(|range| build_candidate_range_cache_key(&scanner.stream_ctx, range).unwrap())
915 .collect();
916 let groups = scanner
917 .build_stream()
918 .await
919 .unwrap()
920 .try_collect::<Vec<_>>()
921 .await
922 .unwrap();
923 assert_eq!(
924 (0..1001)
925 .map(|tsid| MetricSeriesId { table_id: 1, tsid })
926 .collect::<Vec<_>>(),
927 groups.into_iter().flatten().collect::<Vec<_>>()
928 );
929 for file_index in 0..2 {
932 assert_eq!(1, pruner.test_remaining_ranges(file_index));
933 }
934 assert_eq!(
935 1,
936 scanner
937 .metrics_set
938 .clone_inner()
939 .sum_by_name("candidate_index_files")
940 .unwrap()
941 .as_usize()
942 );
943 assert_eq!(
944 2,
945 scanner
946 .metrics_set
947 .clone_inner()
948 .sum_by_name("candidate_index_covered_ssts")
949 .unwrap()
950 .as_usize()
951 );
952 assert!(!format!("{:?}", scanner.part_metrics).contains("range_cache_miss"));
955 for key in keys {
956 assert!(
957 scanner
958 .stream_ctx
959 .input
960 .cache_strategy
961 .get_range_result(&key)
962 .is_none()
963 );
964 }
965 }
966
967 #[tokio::test]
968 async fn index_read_failure_after_output_is_propagated() {
969 let (_env, scanner, _pruner) = indexed_scanner().await;
970 let context = scanner.stream_ctx.input.series_index.clone().unwrap();
971 let index = scanner.coverage.indexes[0].clone();
972 let path = series_index_path(
973 scanner.stream_ctx.input.region_metadata().region_id,
974 index.entry().index_uuid,
975 );
976 let mut stream = index_primary_key_stream(
977 scanner.stream_ctx.clone(),
978 context.clone(),
979 index,
980 scanner.range_semaphore.clone(),
981 );
982 assert_eq!(500, stream.try_next().await.unwrap().unwrap().num_rows());
983 context.store.delete(&path).await.unwrap();
984 assert!(stream.try_collect::<Vec<_>>().await.is_err());
985 assert!(
988 scanner
989 .build_stream()
990 .await
991 .unwrap()
992 .try_collect::<Vec<_>>()
993 .await
994 .is_err()
995 );
996 }
997
998 #[tokio::test]
999 async fn candidate_scanner_rejects_predicate_prefilter_pruner() {
1000 let env = SchedulerEnv::new().await;
1001 let metadata = Arc::new(sst_region_metadata_with_encoding(
1002 PrimaryKeyEncoding::Sparse,
1003 ));
1004 let mapper =
1005 FlatProjectionMapper::new(&metadata, 0..metadata.column_metadatas.len()).unwrap();
1006 let stream_ctx = Arc::new(StreamContext::seq_scan_ctx(
1007 ScanInput::builder(env.access_layer.clone(), mapper).build(),
1008 ));
1009 let pruner = Arc::new(Pruner::new(stream_ctx.clone(), 1));
1010 let metrics_set = ExecutionPlanMetricsSet::new();
1011 let part_metrics = PartitionMetrics::new(
1012 metadata.region_id,
1013 0,
1014 "candidate-test",
1015 Instant::now(),
1016 false,
1017 &metrics_set,
1018 );
1019
1020 let error = SeriesCandidateScanner::try_new(
1021 stream_ctx,
1022 Vec::new(),
1023 pruner,
1024 Arc::new(Semaphore::new(1)),
1025 Arc::new(UnboundedMemoryPool::default()),
1026 metrics_set,
1027 part_metrics,
1028 )
1029 .err()
1030 .unwrap();
1031
1032 assert!(matches!(error, crate::error::Error::Unexpected { .. }));
1033 assert!(
1034 error
1035 .to_string()
1036 .contains("requires a pruner without predicate prefiltering")
1037 );
1038 }
1039
1040 fn binary_batch(values: &[&[u8]]) -> RecordBatch {
1041 RecordBatch::try_new(
1042 primary_key_schema(),
1043 vec![Arc::new(BinaryArray::from_iter_values(
1044 values.iter().copied(),
1045 ))],
1046 )
1047 .unwrap()
1048 }
1049
1050 fn dictionary_batch(values: &[&[u8]], keys: &[u32]) -> RecordBatch {
1051 let dict_values: ArrayRef = Arc::new(BinaryArray::from_iter_values(values.iter().copied()));
1052 let dict =
1053 DictionaryArray::<UInt32Type>::try_new(UInt32Array::from(keys.to_vec()), dict_values)
1054 .unwrap();
1055 let schema = Arc::new(Schema::new(vec![Field::new_dictionary(
1056 PRIMARY_KEY_COLUMN_NAME,
1057 DataType::UInt32,
1058 DataType::Binary,
1059 false,
1060 )]));
1061 RecordBatch::try_new(schema, vec![Arc::new(dict)]).unwrap()
1062 }
1063
1064 #[tokio::test]
1065 async fn candidate_stream_normalizes_and_deduplicates_primary_keys() {
1066 let input = Box::pin(futures::stream::iter(vec![
1067 Ok(dictionary_batch(&[b"a", b"b"], &[0, 0, 1])),
1068 Ok(binary_batch(&[b"b", b"c", b"c"])),
1069 ]));
1070 let batches = candidate_primary_key_stream(input, None)
1071 .try_collect::<Vec<_>>()
1072 .await
1073 .unwrap();
1074 let actual = batches
1075 .iter()
1076 .flat_map(|batch| {
1077 batch
1078 .column(0)
1079 .as_any()
1080 .downcast_ref::<BinaryArray>()
1081 .unwrap()
1082 .iter()
1083 .flatten()
1084 })
1085 .collect::<Vec<_>>();
1086 assert_eq!(actual, vec![b"a".as_slice(), b"b", b"c"]);
1087 }
1088
1089 #[tokio::test]
1090 async fn candidate_stream_filters_primary_keys_before_merge() {
1091 let metadata = Arc::new(sst_region_metadata_with_encoding(
1092 PrimaryKeyEncoding::Sparse,
1093 ));
1094 let codec = SparsePrimaryKeyCodec::new(&metadata);
1095 let mut table_1 = Vec::new();
1096 let mut table_2 = Vec::new();
1097 codec.encode_internal(1, 10, &mut table_1).unwrap();
1098 codec.encode_internal(2, 20, &mut table_2).unwrap();
1099
1100 let predicate = Predicate::new(vec![
1101 col(store_api::metric_engine_consts::DATA_SCHEMA_TABLE_ID_COLUMN_NAME).eq(lit(1_u32)),
1102 ]);
1103 let filter = build_primary_key_filter(&metadata, None, Some(&predicate));
1104 let input = Box::pin(futures::stream::iter(vec![Ok(dictionary_batch(
1105 &[table_1.as_slice(), table_2.as_slice()],
1106 &[0, 1],
1107 ))]));
1108 let batches = candidate_primary_key_stream(input, filter)
1109 .try_collect::<Vec<_>>()
1110 .await
1111 .unwrap();
1112
1113 assert_eq!(batches.len(), 1);
1114 let array = batches[0]
1115 .column(0)
1116 .as_any()
1117 .downcast_ref::<BinaryArray>()
1118 .unwrap();
1119 assert_eq!(array.len(), 1);
1120 assert_eq!(array.value(0), table_1);
1121 }
1122
1123 #[tokio::test]
1124 async fn decode_metric_series_yields_groups_of_500() {
1125 let metadata = Arc::new(sst_region_metadata_with_encoding(
1126 PrimaryKeyEncoding::Sparse,
1127 ));
1128 let codec = SparsePrimaryKeyCodec::new(&metadata);
1129 let primary_keys = (0..501_u64)
1130 .map(|tsid| {
1131 let mut primary_key = Vec::new();
1132 codec.encode_internal(1, tsid, &mut primary_key).unwrap();
1133 primary_key
1134 })
1135 .collect::<Vec<_>>();
1136 let batch = binary_batch(&primary_keys.iter().map(Vec::as_slice).collect::<Vec<_>>());
1137 let source = Box::pin(futures::stream::iter(vec![Ok(batch)]));
1138 let metrics = ExecutionPlanMetricsSet::new();
1139 let pool = Arc::new(UnboundedMemoryPool::default());
1140 let merged =
1141 merge_primary_key_streams(vec![source], pool, &metrics, 0, "candidate-test").unwrap();
1142 let groups = decode_metric_series(merged, metadata)
1143 .unwrap()
1144 .try_collect::<Vec<_>>()
1145 .await
1146 .unwrap();
1147
1148 assert_eq!(groups.iter().map(Vec::len).collect::<Vec<_>>(), [500, 1]);
1149 assert_eq!(
1150 groups[0][0],
1151 MetricSeriesId {
1152 table_id: 1,
1153 tsid: 0
1154 }
1155 );
1156 assert_eq!(
1157 groups[1][0],
1158 MetricSeriesId {
1159 table_id: 1,
1160 tsid: 500
1161 }
1162 );
1163 }
1164
1165 #[tokio::test]
1166 async fn merge_primary_keys_globally_sorts_and_deduplicates_series() {
1167 let metadata = Arc::new(sst_region_metadata_with_encoding(
1168 PrimaryKeyEncoding::Sparse,
1169 ));
1170 let codec = SparsePrimaryKeyCodec::new(&metadata);
1171 let encode = |tsid| {
1172 let mut primary_key = Vec::new();
1173 codec.encode_internal(1, tsid, &mut primary_key).unwrap();
1174 primary_key
1175 };
1176 let keys_1 = [encode(1), encode(3)];
1177 let mut alternate_key_for_series_1 = encode(1);
1178 alternate_key_for_series_1.push(0);
1179 let keys_2 = [alternate_key_for_series_1, encode(2)];
1180 let sources = vec![
1181 Box::pin(futures::stream::iter(vec![Ok(binary_batch(
1182 &keys_1.iter().map(Vec::as_slice).collect::<Vec<_>>(),
1183 ))])) as BoxedRecordBatchStream,
1184 Box::pin(futures::stream::iter(vec![Ok(binary_batch(
1185 &keys_2.iter().map(Vec::as_slice).collect::<Vec<_>>(),
1186 ))])),
1187 ];
1188
1189 let merged = merge_primary_key_streams(
1190 sources,
1191 Arc::new(UnboundedMemoryPool::default()),
1192 &ExecutionPlanMetricsSet::new(),
1193 0,
1194 "candidate-merge-test",
1195 )
1196 .unwrap();
1197 let groups = decode_metric_series(merged, metadata)
1198 .unwrap()
1199 .try_collect::<Vec<_>>()
1200 .await
1201 .unwrap();
1202
1203 assert_eq!(
1204 groups,
1205 vec![vec![
1206 MetricSeriesId {
1207 table_id: 1,
1208 tsid: 1,
1209 },
1210 MetricSeriesId {
1211 table_id: 1,
1212 tsid: 2,
1213 },
1214 MetricSeriesId {
1215 table_id: 1,
1216 tsid: 3,
1217 },
1218 ]]
1219 );
1220 }
1221}