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