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::sync::Arc;
18use std::time::Instant;
19
20use async_stream::try_stream;
21use datafusion::execution::memory_pool::{MemoryConsumer, MemoryPool};
22use datafusion::physical_expr::{LexOrdering, PhysicalSortExpr};
23use datafusion::physical_plan::expressions::Column;
24use datafusion::physical_plan::metrics::{BaselineMetrics, ExecutionPlanMetricsSet};
25use datafusion::physical_plan::sorts::streaming_merge::StreamingMergeBuilder;
26use datafusion::physical_plan::stream::RecordBatchStreamAdapter;
27use datafusion_common::DataFusionError;
28use datatypes::arrow::array::{Array, BinaryArray, BinaryBuilder};
29use datatypes::arrow::compute::SortOptions;
30use datatypes::arrow::datatypes::{DataType, Field, Schema, SchemaRef};
31use datatypes::arrow::record_batch::RecordBatch;
32use datatypes::prelude::ConcreteDataType;
33use futures::stream::BoxStream;
34use futures::{StreamExt, TryStreamExt};
35use mito_codec::row_converter::{PrimaryKeyFilter, SparsePrimaryKeyCodec};
36use snafu::{OptionExt, ResultExt, ensure};
37use store_api::codec::PrimaryKeyEncoding;
38use store_api::region_engine::PartitionRange;
39use store_api::storage::consts::{PRIMARY_KEY_COLUMN_NAME, ReservedColumnId};
40use tokio::sync::Semaphore;
41
42use crate::error::{
43    InvalidRequestSnafu, JoinSnafu, MergeCandidateSeriesSnafu, NewRecordBatchSnafu, Result,
44    UnexpectedSnafu,
45};
46use crate::read::BoxedRecordBatchStream;
47use crate::read::pruner::{PartitionPruner, Pruner};
48use crate::read::range::RowGroupIndex;
49use crate::read::range_cache::{
50    build_candidate_range_cache_key, cache_flat_range_stream, cached_flat_range_stream,
51};
52use crate::read::scan_region::StreamContext;
53use crate::read::scan_util::{PartitionMetrics, new_filter_metrics, scan_flat_mem_ranges};
54use crate::sst::parquet::DEFAULT_READ_BATCH_SIZE;
55use crate::sst::parquet::format::PrimaryKeyArray;
56use crate::sst::parquet::prefilter::{
57    CachedPrimaryKeyFilter, build_primary_key_filter, prefilter_flat_batch_by_primary_key,
58};
59use crate::sst::parquet::reader::ReaderMetrics;
60use crate::sst::parquet::row_group::ParquetFetchMetrics;
61
62const CANDIDATE_SERIES_BATCH_SIZE: usize = 500;
63
64/// Identifies one series in a physical metric region.
65#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)]
66pub(crate) struct MetricSeriesId {
67    pub(crate) table_id: u32,
68    pub(crate) tsid: u64,
69}
70
71pub(crate) type MetricSeriesIdStream = BoxStream<'static, Result<Vec<MetricSeriesId>>>;
72
73/// Builds candidate metric series from the ranges assigned to a [`SeriesScan`](super::series_scan::SeriesScan).
74#[allow(dead_code)]
75pub(crate) struct SeriesCandidateScanner {
76    stream_ctx: Arc<StreamContext>,
77    partitions: Vec<Vec<PartitionRange>>,
78    partition_pruner: Arc<PartitionPruner>,
79    range_semaphore: Arc<Semaphore>,
80    memory_pool: Arc<dyn MemoryPool>,
81    metrics_set: ExecutionPlanMetricsSet,
82    part_metrics: PartitionMetrics,
83}
84
85#[allow(dead_code)]
86impl SeriesCandidateScanner {
87    /// Creates a candidate-series scanner for native memtable and SST ranges.
88    ///
89    /// Callers must fall back to the legacy series-scan path when the scan contains
90    /// extension ranges. Candidate-series discovery does not support other range types.
91    ///
92    /// `pruner` must be built with `PrunerOptions::enable_predicate_prefilter` set to
93    /// `false`: both scan phases share its file range builders, and the two-stage path
94    /// applies simple filters on the precise-filter path instead.
95    pub(crate) fn try_new(
96        stream_ctx: Arc<StreamContext>,
97        partitions: Vec<Vec<PartitionRange>>,
98        pruner: Arc<Pruner>,
99        range_semaphore: Arc<Semaphore>,
100        memory_pool: Arc<dyn MemoryPool>,
101        metrics_set: ExecutionPlanMetricsSet,
102        part_metrics: PartitionMetrics,
103    ) -> Result<Self> {
104        validate_metric_metadata(&stream_ctx)?;
105        #[cfg(feature = "enterprise")]
106        ensure!(
107            stream_ctx.input.extension_ranges().is_empty(),
108            InvalidRequestSnafu {
109                region_id: stream_ctx.input.region_metadata().region_id,
110                reason: "candidate-series scan does not support extension ranges; use the legacy series-scan path",
111            }
112        );
113        ensure!(
114            !pruner.predicate_prefilter_enabled(),
115            UnexpectedSnafu {
116                reason: format!(
117                    "candidate-series scan for region {} requires a pruner without predicate prefiltering",
118                    stream_ctx.input.region_metadata().region_id
119                ),
120            }
121        );
122        let all_ranges = partitions.iter().flatten().copied().collect::<Vec<_>>();
123        pruner.add_partition_ranges(&all_ranges);
124        let partition_pruner = Arc::new(PartitionPruner::new(pruner, &all_ranges));
125        Ok(Self {
126            stream_ctx,
127            partitions,
128            partition_pruner,
129            range_semaphore,
130            memory_pool,
131            metrics_set,
132            part_metrics,
133        })
134    }
135
136    /// Builds a globally sorted stream of candidate metric-series IDs.
137    pub(crate) async fn build_stream(&self) -> Result<MetricSeriesIdStream> {
138        let all_ranges = self
139            .partitions
140            .iter()
141            .flatten()
142            .copied()
143            .collect::<Vec<_>>();
144        let range_builder = SeriesCandidateRangeBuilder {
145            stream_ctx: self.stream_ctx.clone(),
146            partition_pruner: self.partition_pruner.clone(),
147            range_semaphore: self.range_semaphore.clone(),
148            memory_pool: self.memory_pool.clone(),
149            metrics_set: self.metrics_set.clone(),
150            part_metrics: self.part_metrics.clone(),
151        };
152        let mut tasks = Vec::with_capacity(all_ranges.len());
153        for (range_idx, part_range) in all_ranges.into_iter().enumerate() {
154            let range_builder = range_builder.clone();
155            tasks.push(common_runtime::spawn_query(async move {
156                let _permit = range_builder
157                    .range_semaphore
158                    .clone()
159                    .acquire_owned()
160                    .await
161                    .map_err(|error| {
162                        UnexpectedSnafu {
163                            reason: format!("failed to acquire candidate range permit: {error}"),
164                        }
165                        .build()
166                    })?;
167                range_builder
168                    .build_range_stream(part_range, range_idx)
169                    .await
170            }));
171        }
172
173        let mut range_streams = Vec::with_capacity(tasks.len());
174        for task in tasks {
175            range_streams.push(task.await.context(JoinSnafu)??);
176        }
177
178        // Keep scanner-level merge metrics in the same synthetic partition as
179        // SeriesDistributor. Output partitions occupy 0..self.partitions.len().
180        let merged = merge_primary_key_streams(
181            range_streams,
182            self.memory_pool.clone(),
183            &self.metrics_set,
184            self.partitions.len(),
185            "SeriesCandidateScanner::final_merge",
186        )?;
187        decode_metric_series(merged, self.stream_ctx.input.region_metadata().clone())
188    }
189
190    /// Returns the partition pruner shared with the data phase.
191    pub(crate) fn partition_pruner(&self) -> Arc<PartitionPruner> {
192        self.partition_pruner.clone()
193    }
194}
195
196#[derive(Clone)]
197struct SeriesCandidateRangeBuilder {
198    stream_ctx: Arc<StreamContext>,
199    partition_pruner: Arc<PartitionPruner>,
200    range_semaphore: Arc<Semaphore>,
201    memory_pool: Arc<dyn MemoryPool>,
202    metrics_set: ExecutionPlanMetricsSet,
203    part_metrics: PartitionMetrics,
204}
205
206impl SeriesCandidateRangeBuilder {
207    async fn build_range_stream(
208        &self,
209        part_range: PartitionRange,
210        merge_partition: usize,
211    ) -> Result<BoxedRecordBatchStream> {
212        let cache_key = build_candidate_range_cache_key(&self.stream_ctx, &part_range);
213        if let Some(key) = cache_key.as_ref() {
214            if let Some(value) = self.stream_ctx.input.cache_strategy.get_range_result(key) {
215                self.part_metrics.inc_range_cache_hit();
216                return Ok(cached_flat_range_stream(value));
217            }
218            self.part_metrics.inc_range_cache_miss();
219        }
220
221        let range_meta = &self.stream_ctx.ranges[part_range.identifier];
222        let mut sources = Vec::with_capacity(range_meta.row_group_indices.len());
223        for index in &range_meta.row_group_indices {
224            let source = self.build_source(*index, range_meta.time_range).await?;
225            if let Some(source) = source {
226                sources.push(source);
227            }
228        }
229
230        let sources = self.stream_ctx.input.create_parallel_flat_sources(
231            sources,
232            self.range_semaphore.clone(),
233            2,
234        )?;
235        let stream = merge_primary_key_streams(
236            sources,
237            self.memory_pool.clone(),
238            &self.metrics_set,
239            merge_partition,
240            "SeriesCandidateScanner::range_merge",
241        )?;
242
243        Ok(match cache_key {
244            Some(key) => cache_flat_range_stream(
245                stream,
246                self.stream_ctx.input.cache_strategy.clone(),
247                key,
248                self.part_metrics.clone(),
249            ),
250            None => stream,
251        })
252    }
253
254    async fn build_source(
255        &self,
256        index: RowGroupIndex,
257        time_range: crate::sst::file::FileTimeRange,
258    ) -> Result<Option<BoxedRecordBatchStream>> {
259        let metadata = self.stream_ctx.input.region_metadata().clone();
260        if self.stream_ctx.is_mem_range_index(index) {
261            let raw = scan_flat_mem_ranges(
262                self.stream_ctx.clone(),
263                self.part_metrics.clone(),
264                index,
265                time_range,
266            );
267            let filter = build_primary_key_filter(
268                &metadata,
269                None,
270                self.stream_ctx.input.predicate_group().predicate(),
271            );
272            return Ok(Some(candidate_primary_key_stream(Box::pin(raw), filter)));
273        }
274
275        if self.stream_ctx.is_file_range_index(index) {
276            let file = self.stream_ctx.input.file_from_index(index);
277            let predicate = self.stream_ctx.input.predicate_for_file(file);
278            if self
279                .partition_pruner
280                .try_skip_manifest_pruned_file_range(index, &self.part_metrics)
281            {
282                return Ok(None);
283            }
284            let mut reader_metrics = ReaderMetrics {
285                filter_metrics: new_filter_metrics(self.part_metrics.explain_verbose()),
286                ..Default::default()
287            };
288            let ranges = self
289                .partition_pruner
290                .build_file_ranges(index, &self.part_metrics, &mut reader_metrics)
291                .await?;
292            self.part_metrics.inc_num_file_ranges(ranges.len());
293            self.part_metrics
294                .merge_reader_metrics(&reader_metrics, None);
295
296            // Build a fresh encoded-PK filter from this SST's metadata so schema
297            // compatibility is evaluated independently for each file.
298            let filter = ranges.first().and_then(|range| {
299                build_primary_key_filter(
300                    range.region_metadata(),
301                    Some(metadata.as_ref()),
302                    predicate.as_ref(),
303                )
304            });
305            let part_metrics = self.part_metrics.clone();
306            let raw = Box::pin(try_stream! {
307                let fetch_metrics = part_metrics
308                    .explain_verbose()
309                    .then(|| Arc::new(ParquetFetchMetrics::default()));
310                let mut reader_metrics = ReaderMetrics {
311                    fetch_metrics: fetch_metrics.clone(),
312                    ..Default::default()
313                };
314                for range in ranges {
315                    let build_start = Instant::now();
316                    let Some(mut reader) = range
317                        .primary_key_reader(fetch_metrics.as_deref())
318                        .await?
319                    else {
320                        continue;
321                    };
322                    reader_metrics.build_cost += build_start.elapsed();
323
324                    let scan_start = Instant::now();
325                    while let Some(batch) = reader.try_next().await? {
326                        reader_metrics.num_record_batches += 1;
327                        reader_metrics.num_batches += 1;
328                        reader_metrics.num_rows += batch.num_rows();
329                        yield batch;
330                    }
331                    reader_metrics.scan_cost += scan_start.elapsed();
332                }
333                reader_metrics.observe_rows("candidate_series");
334                part_metrics.merge_reader_metrics(&reader_metrics, None);
335            });
336            return Ok(Some(candidate_primary_key_stream(raw, filter)));
337        }
338
339        UnexpectedSnafu {
340            reason: format!(
341                "candidate-series scan received unsupported range index {}",
342                index.index
343            ),
344        }
345        .fail()
346    }
347}
348
349pub(crate) fn validate_metric_metadata(stream_ctx: &StreamContext) -> Result<()> {
350    let metadata = stream_ctx.input.region_metadata();
351    let valid_prefix = metadata
352        .primary_key
353        .starts_with(&[ReservedColumnId::table_id(), ReservedColumnId::tsid()]);
354    let valid_types = metadata
355        .column_by_id(ReservedColumnId::table_id())
356        .zip(metadata.column_by_id(ReservedColumnId::tsid()))
357        .is_some_and(|(table_id, tsid)| {
358            table_id.column_schema.data_type == ConcreteDataType::uint32_datatype()
359                && tsid.column_schema.data_type == ConcreteDataType::uint64_datatype()
360        });
361    ensure!(
362        metadata.primary_key_encoding == PrimaryKeyEncoding::Sparse && valid_prefix && valid_types,
363        InvalidRequestSnafu {
364            region_id: metadata.region_id,
365            reason: "candidate-series scan requires sparse (__table_id, __tsid) primary keys",
366        }
367    );
368    Ok(())
369}
370
371fn primary_key_schema() -> SchemaRef {
372    Arc::new(Schema::new(vec![Field::new(
373        PRIMARY_KEY_COLUMN_NAME,
374        DataType::Binary,
375        false,
376    )]))
377}
378
379/// Filters a source by encoded-primary-key predicates and emits one binary row per local key.
380fn candidate_primary_key_stream(
381    mut input: BoxedRecordBatchStream,
382    mut filter: Option<CachedPrimaryKeyFilter>,
383) -> BoxedRecordBatchStream {
384    Box::pin(try_stream! {
385        let mut last_primary_key = Vec::new();
386        let mut has_last = false;
387        while let Some(batch) = input.try_next().await? {
388            if let Some(batch) = normalize_candidate_batch(
389                batch,
390                filter.as_mut(),
391                &mut last_primary_key,
392                &mut has_last,
393            )? {
394                yield batch;
395            }
396        }
397    })
398}
399
400fn normalize_candidate_batch(
401    mut batch: RecordBatch,
402    filter: Option<&mut CachedPrimaryKeyFilter>,
403    last_primary_key: &mut Vec<u8>,
404    has_last: &mut bool,
405) -> Result<Option<RecordBatch>> {
406    let pk_idx = batch
407        .schema()
408        .column_with_name(PRIMARY_KEY_COLUMN_NAME)
409        .map(|(idx, _)| idx)
410        .context(UnexpectedSnafu {
411            reason: "candidate source does not contain __primary_key",
412        })?;
413    if let Some(filter) = filter {
414        let Some(filtered) = prefilter_flat_batch_by_primary_key(
415            batch,
416            pk_idx,
417            filter as &mut dyn PrimaryKeyFilter,
418        )?
419        else {
420            return Ok(None);
421        };
422        batch = filtered;
423    }
424
425    let pk_column = batch.column(pk_idx);
426    let mut builder = BinaryBuilder::new();
427    if let Some(array) = pk_column.as_any().downcast_ref::<PrimaryKeyArray>() {
428        let values = array
429            .values()
430            .as_any()
431            .downcast_ref::<BinaryArray>()
432            .context(UnexpectedSnafu {
433                reason: "dictionary primary-key values are not binary",
434            })?;
435        for key in array.keys().values() {
436            append_unique_primary_key(
437                values.value(*key as usize),
438                &mut builder,
439                last_primary_key,
440                has_last,
441            );
442        }
443    } else if let Some(array) = pk_column.as_any().downcast_ref::<BinaryArray>() {
444        for value in array.iter().flatten() {
445            append_unique_primary_key(value, &mut builder, last_primary_key, has_last);
446        }
447    } else {
448        return UnexpectedSnafu {
449            reason: format!(
450                "primary-key column is neither binary nor dictionary, got {:?}",
451                pk_column.data_type()
452            ),
453        }
454        .fail();
455    }
456
457    let array = builder.finish();
458    if array.is_empty() {
459        return Ok(None);
460    }
461    let batch = RecordBatch::try_new(primary_key_schema(), vec![Arc::new(array)])
462        .context(NewRecordBatchSnafu)?;
463    Ok(Some(batch))
464}
465
466fn append_unique_primary_key(
467    value: &[u8],
468    builder: &mut BinaryBuilder,
469    last_primary_key: &mut Vec<u8>,
470    has_last: &mut bool,
471) {
472    if !*has_last || last_primary_key != value {
473        builder.append_value(value);
474        last_primary_key.clear();
475        last_primary_key.extend_from_slice(value);
476        *has_last = true;
477    }
478}
479
480fn merge_primary_key_streams(
481    sources: Vec<BoxedRecordBatchStream>,
482    memory_pool: Arc<dyn MemoryPool>,
483    metrics_set: &ExecutionPlanMetricsSet,
484    partition: usize,
485    consumer_name: &'static str,
486) -> Result<BoxedRecordBatchStream> {
487    if sources.is_empty() {
488        return Ok(Box::pin(futures::stream::empty()));
489    }
490    if sources.len() == 1 {
491        return Ok(sources.into_iter().next().unwrap());
492    }
493
494    let schema = primary_key_schema();
495    let df_sources = sources
496        .into_iter()
497        .map(|source| {
498            let stream = source.map_err(|error| DataFusionError::External(Box::new(error)));
499            Box::pin(RecordBatchStreamAdapter::new(schema.clone(), stream)) as _
500        })
501        .collect();
502    let ordering = LexOrdering::new([PhysicalSortExpr {
503        expr: Arc::new(Column::new(PRIMARY_KEY_COLUMN_NAME, 0)),
504        options: SortOptions {
505            descending: false,
506            nulls_first: false,
507        },
508    }])
509    // Safe to unwrap because `LexOrdering::new` returns `None` only for empty
510    // input, and this array always contains one sort expression.
511    .unwrap();
512    let reservation = MemoryConsumer::new(consumer_name).register(&memory_pool);
513    let mut merged = StreamingMergeBuilder::new()
514        .with_streams(df_sources)
515        .with_schema(schema)
516        .with_expressions(&ordering)
517        .with_metrics(BaselineMetrics::new(metrics_set, partition))
518        .with_batch_size(DEFAULT_READ_BATCH_SIZE)
519        .with_reservation(reservation)
520        .build()
521        .context(MergeCandidateSeriesSnafu)?;
522
523    Ok(Box::pin(try_stream! {
524        while let Some(batch) = merged.next().await {
525            yield batch.context(MergeCandidateSeriesSnafu)?;
526        }
527    }))
528}
529
530fn decode_metric_series(
531    mut input: BoxedRecordBatchStream,
532    metadata: store_api::metadata::RegionMetadataRef,
533) -> Result<MetricSeriesIdStream> {
534    let codec = SparsePrimaryKeyCodec::new(&metadata);
535    Ok(Box::pin(try_stream! {
536        let mut last_series = None;
537        let mut output = Vec::with_capacity(CANDIDATE_SERIES_BATCH_SIZE);
538        while let Some(batch) = input.try_next().await? {
539            let array = batch
540                .column(0)
541                .as_any()
542                .downcast_ref::<BinaryArray>()
543                .context(UnexpectedSnafu {
544                    reason: "merged candidate primary key is not binary",
545                })?;
546            for primary_key in array.iter().flatten() {
547                let (table_id, tsid) = codec
548                    .decode_ids(primary_key)
549                    .context(crate::error::DecodeSnafu)?;
550                let series = MetricSeriesId { table_id, tsid };
551                if last_series == Some(series) {
552                    continue;
553                }
554                last_series = Some(series);
555                output.push(series);
556                if output.len() == CANDIDATE_SERIES_BATCH_SIZE {
557                    yield std::mem::replace(
558                        &mut output,
559                        Vec::with_capacity(CANDIDATE_SERIES_BATCH_SIZE),
560                    );
561                }
562            }
563        }
564        if !output.is_empty() {
565            yield output;
566        }
567    }))
568}
569
570#[cfg(test)]
571mod tests {
572    use std::time::Instant;
573
574    use datafusion::execution::memory_pool::UnboundedMemoryPool;
575    use datafusion_expr::{col, lit};
576    use datatypes::arrow::array::{ArrayRef, DictionaryArray, UInt32Array};
577    use datatypes::arrow::datatypes::UInt32Type;
578    use futures::TryStreamExt;
579    use store_api::codec::PrimaryKeyEncoding;
580    use table::predicate::Predicate;
581
582    use super::*;
583    use crate::read::flat_projection::FlatProjectionMapper;
584    use crate::read::scan_region::ScanInput;
585    use crate::read::scan_util::PartitionMetrics;
586    use crate::test_util::scheduler_util::SchedulerEnv;
587    use crate::test_util::sst_util::sst_region_metadata_with_encoding;
588
589    #[tokio::test]
590    async fn candidate_scanner_rejects_predicate_prefilter_pruner() {
591        let env = SchedulerEnv::new().await;
592        let metadata = Arc::new(sst_region_metadata_with_encoding(
593            PrimaryKeyEncoding::Sparse,
594        ));
595        let mapper =
596            FlatProjectionMapper::new(&metadata, 0..metadata.column_metadatas.len()).unwrap();
597        let stream_ctx = Arc::new(StreamContext::seq_scan_ctx(ScanInput::new(
598            env.access_layer.clone(),
599            mapper,
600        )));
601        let pruner = Arc::new(Pruner::new(stream_ctx.clone(), 1));
602        let metrics_set = ExecutionPlanMetricsSet::new();
603        let part_metrics = PartitionMetrics::new(
604            metadata.region_id,
605            0,
606            "candidate-test",
607            Instant::now(),
608            false,
609            &metrics_set,
610        );
611
612        let error = SeriesCandidateScanner::try_new(
613            stream_ctx,
614            Vec::new(),
615            pruner,
616            Arc::new(Semaphore::new(1)),
617            Arc::new(UnboundedMemoryPool::default()),
618            metrics_set,
619            part_metrics,
620        )
621        .err()
622        .unwrap();
623
624        assert!(matches!(error, crate::error::Error::Unexpected { .. }));
625        assert!(
626            error
627                .to_string()
628                .contains("requires a pruner without predicate prefiltering")
629        );
630    }
631
632    fn binary_batch(values: &[&[u8]]) -> RecordBatch {
633        RecordBatch::try_new(
634            primary_key_schema(),
635            vec![Arc::new(BinaryArray::from_iter_values(
636                values.iter().copied(),
637            ))],
638        )
639        .unwrap()
640    }
641
642    fn dictionary_batch(values: &[&[u8]], keys: &[u32]) -> RecordBatch {
643        let dict_values: ArrayRef = Arc::new(BinaryArray::from_iter_values(values.iter().copied()));
644        let dict =
645            DictionaryArray::<UInt32Type>::try_new(UInt32Array::from(keys.to_vec()), dict_values)
646                .unwrap();
647        let schema = Arc::new(Schema::new(vec![Field::new_dictionary(
648            PRIMARY_KEY_COLUMN_NAME,
649            DataType::UInt32,
650            DataType::Binary,
651            false,
652        )]));
653        RecordBatch::try_new(schema, vec![Arc::new(dict)]).unwrap()
654    }
655
656    #[tokio::test]
657    async fn candidate_stream_normalizes_and_deduplicates_primary_keys() {
658        let input = Box::pin(futures::stream::iter(vec![
659            Ok(dictionary_batch(&[b"a", b"b"], &[0, 0, 1])),
660            Ok(binary_batch(&[b"b", b"c", b"c"])),
661        ]));
662        let batches = candidate_primary_key_stream(input, None)
663            .try_collect::<Vec<_>>()
664            .await
665            .unwrap();
666        let actual = batches
667            .iter()
668            .flat_map(|batch| {
669                batch
670                    .column(0)
671                    .as_any()
672                    .downcast_ref::<BinaryArray>()
673                    .unwrap()
674                    .iter()
675                    .flatten()
676            })
677            .collect::<Vec<_>>();
678        assert_eq!(actual, vec![b"a".as_slice(), b"b", b"c"]);
679    }
680
681    #[tokio::test]
682    async fn candidate_stream_filters_primary_keys_before_merge() {
683        let metadata = Arc::new(sst_region_metadata_with_encoding(
684            PrimaryKeyEncoding::Sparse,
685        ));
686        let codec = SparsePrimaryKeyCodec::new(&metadata);
687        let mut table_1 = Vec::new();
688        let mut table_2 = Vec::new();
689        codec.encode_internal(1, 10, &mut table_1).unwrap();
690        codec.encode_internal(2, 20, &mut table_2).unwrap();
691
692        let predicate = Predicate::new(vec![
693            col(store_api::metric_engine_consts::DATA_SCHEMA_TABLE_ID_COLUMN_NAME).eq(lit(1_u32)),
694        ]);
695        let filter = build_primary_key_filter(&metadata, None, Some(&predicate));
696        let input = Box::pin(futures::stream::iter(vec![Ok(dictionary_batch(
697            &[table_1.as_slice(), table_2.as_slice()],
698            &[0, 1],
699        ))]));
700        let batches = candidate_primary_key_stream(input, filter)
701            .try_collect::<Vec<_>>()
702            .await
703            .unwrap();
704
705        assert_eq!(batches.len(), 1);
706        let array = batches[0]
707            .column(0)
708            .as_any()
709            .downcast_ref::<BinaryArray>()
710            .unwrap();
711        assert_eq!(array.len(), 1);
712        assert_eq!(array.value(0), table_1);
713    }
714
715    #[tokio::test]
716    async fn decode_metric_series_yields_groups_of_500() {
717        let metadata = Arc::new(sst_region_metadata_with_encoding(
718            PrimaryKeyEncoding::Sparse,
719        ));
720        let codec = SparsePrimaryKeyCodec::new(&metadata);
721        let primary_keys = (0..501_u64)
722            .map(|tsid| {
723                let mut primary_key = Vec::new();
724                codec.encode_internal(1, tsid, &mut primary_key).unwrap();
725                primary_key
726            })
727            .collect::<Vec<_>>();
728        let batch = binary_batch(&primary_keys.iter().map(Vec::as_slice).collect::<Vec<_>>());
729        let source = Box::pin(futures::stream::iter(vec![Ok(batch)]));
730        let metrics = ExecutionPlanMetricsSet::new();
731        let pool = Arc::new(UnboundedMemoryPool::default());
732        let merged =
733            merge_primary_key_streams(vec![source], pool, &metrics, 0, "candidate-test").unwrap();
734        let groups = decode_metric_series(merged, metadata)
735            .unwrap()
736            .try_collect::<Vec<_>>()
737            .await
738            .unwrap();
739
740        assert_eq!(groups.iter().map(Vec::len).collect::<Vec<_>>(), [500, 1]);
741        assert_eq!(
742            groups[0][0],
743            MetricSeriesId {
744                table_id: 1,
745                tsid: 0
746            }
747        );
748        assert_eq!(
749            groups[1][0],
750            MetricSeriesId {
751                table_id: 1,
752                tsid: 500
753            }
754        );
755    }
756
757    #[tokio::test]
758    async fn merge_primary_keys_globally_sorts_and_deduplicates_series() {
759        let metadata = Arc::new(sst_region_metadata_with_encoding(
760            PrimaryKeyEncoding::Sparse,
761        ));
762        let codec = SparsePrimaryKeyCodec::new(&metadata);
763        let encode = |tsid| {
764            let mut primary_key = Vec::new();
765            codec.encode_internal(1, tsid, &mut primary_key).unwrap();
766            primary_key
767        };
768        let keys_1 = [encode(1), encode(3)];
769        let mut alternate_key_for_series_1 = encode(1);
770        alternate_key_for_series_1.push(0);
771        let keys_2 = [alternate_key_for_series_1, encode(2)];
772        let sources = vec![
773            Box::pin(futures::stream::iter(vec![Ok(binary_batch(
774                &keys_1.iter().map(Vec::as_slice).collect::<Vec<_>>(),
775            ))])) as BoxedRecordBatchStream,
776            Box::pin(futures::stream::iter(vec![Ok(binary_batch(
777                &keys_2.iter().map(Vec::as_slice).collect::<Vec<_>>(),
778            ))])),
779        ];
780
781        let merged = merge_primary_key_streams(
782            sources,
783            Arc::new(UnboundedMemoryPool::default()),
784            &ExecutionPlanMetricsSet::new(),
785            0,
786            "candidate-merge-test",
787        )
788        .unwrap();
789        let groups = decode_metric_series(merged, metadata)
790            .unwrap()
791            .try_collect::<Vec<_>>()
792            .await
793            .unwrap();
794
795        assert_eq!(
796            groups,
797            vec![vec![
798                MetricSeriesId {
799                    table_id: 1,
800                    tsid: 1,
801                },
802                MetricSeriesId {
803                    table_id: 1,
804                    tsid: 2,
805                },
806                MetricSeriesId {
807                    table_id: 1,
808                    tsid: 3,
809                },
810            ]]
811        );
812    }
813}