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::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#[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#[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 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 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 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 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 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
379fn 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 .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}