Skip to main content

mito2/read/
series_candidate.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15//! Candidate metric-series discovery for the two-stage series scan.
16
17use 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
67/// Builds candidate metric series from the ranges assigned to a [`SeriesScan`](super::series_scan::SeriesScan).
68pub(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    /// Creates a candidate-series scanner for native memtable and SST ranges.
82    ///
83    /// Callers must fall back to the legacy series-scan path when the scan contains
84    /// extension ranges. Candidate-series discovery does not support other range types.
85    ///
86    /// `pruner` must be built with `PrunerOptions::enable_predicate_prefilter` set to
87    /// `false`: both scan phases share its file range builders, and the two-stage path
88    /// applies simple filters on the precise-filter path instead.
89    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    /// Builds a globally sorted stream of candidate metric-series IDs.
141    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        // Keep scanner-level merge metrics in the same synthetic partition as
201        // SeriesDistributor. Output partitions occupy 0..self.partitions.len().
202        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    /// Returns the partition pruner shared with the data phase.
213    pub(crate) fn partition_pruner(&self) -> Arc<PartitionPruner> {
214        self.partition_pruner.clone()
215    }
216}
217
218/// A scanner-wide replacement plan, independent of partition-range boundaries.
219#[derive(Default)]
220struct SeriesIndexCoverage {
221    indexes: Vec<SeriesIndexFileHandle>,
222    /// Indices into `ScanInput.files`, shared by every occurrence of an SST.
223    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
301/// Reads an index once for the entire scan, rather than once per partition range.
302fn 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        // A cache entry describes the complete original range. Never cache a
364        // partial range whose missing candidates are supplied by a global index.
365        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                // Leave the range reference for the data phase so its first
438                // read can cache the builder. Retaining builders only prevents
439                // eviction; it does not allow caching at zero references.
440                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            // Build a fresh encoded-PK filter from this SST's metadata so schema
463            // compatibility is evaluated independently for each file.
464            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
550/// Filters a source by encoded-primary-key predicates and emits one binary row per local key.
551fn 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    // Safe to unwrap because `LexOrdering::new` returns `None` only for empty
681    // input, and this array always contains one sort expression.
682    .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    /// Uses absent SST objects so any covered-SST read fails the test.
774    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        // Covered SSTs have no builder yet. Keep their references so the data
930        // phase can cache each builder on its first read and reuse it later.
931        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        // Index-backed ranges must not publish their empty residual streams as
953        // complete candidate results for a future scan without the index.
954        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        // Opening the missing selected index must also fail, rather than emit
986        // an incomplete candidate set or fall back to the absent SSTs.
987        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}