1use 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
62pub(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 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 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 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 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 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
371fn 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 .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}