1use std::collections::HashSet;
18use std::sync::Arc;
19use std::time::Instant;
20
21use async_stream::try_stream;
22use futures::TryStreamExt;
23use mito_codec::row_converter::{PrimaryKeyFilter, SparsePrimaryKeyCodec};
24use snafu::ResultExt;
25use store_api::region_engine::PartitionRange;
26use tokio::sync::Semaphore;
27
28#[cfg(feature = "enterprise")]
29use crate::error::InvalidRequestSnafu;
30use crate::error::{JoinSnafu, Result, UnexpectedSnafu};
31use crate::read::BoxedRecordBatchStream;
32use crate::read::pruner::PartitionPruner;
33use crate::read::range_cache::{
34 build_series_range_cache_key, cache_flat_range_stream, cached_flat_range_stream,
35};
36use crate::read::scan_region::StreamContext;
37use crate::read::scan_util::{
38 PartitionMetrics, SplitRecordBatchStream, compute_average_batch_size,
39 compute_parallel_channel_size, new_filter_metrics, scan_flat_mem_ranges,
40 should_split_flat_batches_for_merge,
41};
42use crate::read::seq_scan::SeqScan;
43use crate::read::series_candidate::{MetricSeriesId, validate_metric_metadata};
44use crate::sst::parquet::DEFAULT_READ_BATCH_SIZE;
45use crate::sst::parquet::flat_format::primary_key_column_index;
46use crate::sst::parquet::prefilter::prefilter_flat_batch_by_primary_key;
47use crate::sst::parquet::reader::ReaderMetrics;
48use crate::sst::parquet::row_group::ParquetFetchMetrics;
49
50const TSID_DOMAIN_END: u128 = 1u128 << u64::BITS;
51
52#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
54pub(crate) struct SeriesRange {
55 start: u128,
56 end: u128,
57}
58
59impl SeriesRange {
60 pub(crate) fn new(partition: usize, partitions: usize) -> Option<Self> {
61 if partitions == 0 || partition >= partitions {
62 return None;
63 }
64
65 let partitions = partitions as u128;
66 let partition = partition as u128;
67 let boundary = |partition: u128| {
68 let numerator = partition * TSID_DOMAIN_END;
69 numerator.div_ceil(partitions)
70 };
71 Some(Self {
72 start: boundary(partition),
73 end: boundary(partition + 1),
74 })
75 }
76
77 fn partition_for(tsid: u64, partitions: usize) -> usize {
78 ((tsid as u128 * partitions as u128) >> u64::BITS) as usize
79 }
80}
81
82#[derive(Debug)]
84pub(crate) struct AssignedSeriesBatch {
85 range: SeriesRange,
86 series: Vec<MetricSeriesId>,
87}
88
89#[allow(dead_code)]
90impl AssignedSeriesBatch {
91 fn new(range: SeriesRange, series: Vec<MetricSeriesId>) -> Self {
92 Self { range, series }
93 }
94
95 pub(crate) fn range(&self) -> SeriesRange {
96 self.range
97 }
98
99 pub(crate) fn series(&self) -> &[MetricSeriesId] {
100 &self.series
101 }
102}
103
104#[allow(dead_code)]
106pub(crate) struct SeriesBatchCollector {
107 assignments: Vec<Vec<MetricSeriesId>>,
108}
109
110#[allow(dead_code)]
111impl SeriesBatchCollector {
112 pub(crate) fn new(partitions: usize) -> Option<Self> {
113 (partitions > 0).then(|| Self {
114 assignments: (0..partitions).map(|_| Vec::new()).collect(),
115 })
116 }
117
118 pub(crate) fn push(&mut self, batch: Vec<MetricSeriesId>) {
119 let partitions = self.assignments.len();
120 for series in batch {
121 let partition = SeriesRange::partition_for(series.tsid, partitions);
122 self.assignments[partition].push(series);
123 }
124 }
125
126 pub(crate) fn finish(self) -> Vec<AssignedSeriesBatch> {
127 let partitions = self.assignments.len();
128 self.assignments
129 .into_iter()
130 .enumerate()
131 .map(|(partition, series)| {
132 AssignedSeriesBatch::new(SeriesRange::new(partition, partitions).unwrap(), series)
133 })
134 .collect()
135 }
136}
137
138#[derive(Clone, Debug)]
140struct MetricSeriesFilter {
141 range: SeriesRange,
142 series: Arc<HashSet<MetricSeriesId>>,
143}
144
145impl MetricSeriesFilter {
146 fn new(assigned: &AssignedSeriesBatch) -> Self {
147 let series = assigned.series.iter().copied().collect();
148 Self {
149 range: assigned.range,
150 series: Arc::new(series),
151 }
152 }
153
154 fn primary_key_filter(&self, codec: SparsePrimaryKeyCodec) -> Box<dyn PrimaryKeyFilter> {
155 Box::new(MetricSeriesPrimaryKeyFilter {
156 codec,
157 series: self.series.clone(),
158 last_primary_key: Vec::new(),
159 last_match: None,
160 })
161 }
162
163 fn overlaps_encoded_bounds(
164 &self,
165 codec: &SparsePrimaryKeyCodec,
166 encoded_min: &[u8],
167 encoded_max: &[u8],
168 ) -> Option<bool> {
169 let (min_table_id, min_tsid) = codec.decode_ids(encoded_min).ok()?;
170 let (max_table_id, max_tsid) = codec.decode_ids(encoded_max).ok()?;
171 let min = MetricSeriesId {
172 table_id: min_table_id,
173 tsid: min_tsid,
174 };
175 let max = MetricSeriesId {
176 table_id: max_table_id,
177 tsid: max_tsid,
178 };
179 if min > max {
180 return None;
181 }
182 if min_table_id != max_table_id {
183 return Some(true);
184 }
185
186 Some(u128::from(min_tsid) < self.range.end && u128::from(max_tsid) >= self.range.start)
187 }
188}
189
190struct MetricSeriesPrimaryKeyFilter {
191 codec: SparsePrimaryKeyCodec,
192 series: Arc<HashSet<MetricSeriesId>>,
193 last_primary_key: Vec<u8>,
194 last_match: Option<bool>,
195}
196
197impl PrimaryKeyFilter for MetricSeriesPrimaryKeyFilter {
198 fn matches(&mut self, primary_key: &[u8]) -> mito_codec::error::Result<bool> {
199 if let Some(last_match) = self.last_match
200 && self.last_primary_key == primary_key
201 {
202 return Ok(last_match);
203 }
204
205 let (table_id, tsid) = self.codec.decode_ids(primary_key)?;
206 let matched = self.series.contains(&MetricSeriesId { table_id, tsid });
207 self.last_primary_key.clear();
208 self.last_primary_key.extend_from_slice(primary_key);
209 self.last_match = Some(matched);
210 Ok(matched)
211 }
212}
213
214fn filter_flat_stream_by_series(
215 mut input: BoxedRecordBatchStream,
216 codec: SparsePrimaryKeyCodec,
217 filter: MetricSeriesFilter,
218) -> BoxedRecordBatchStream {
219 Box::pin(try_stream! {
220 let mut primary_key_filter = filter.primary_key_filter(codec);
221 while let Some(batch) = input.try_next().await? {
222 let pk_idx = primary_key_column_index(batch.num_columns());
223 if let Some(batch) = prefilter_flat_batch_by_primary_key(
224 batch,
225 pk_idx,
226 primary_key_filter.as_mut(),
227 )? {
228 yield batch;
229 }
230 }
231 })
232}
233
234#[allow(dead_code)]
236pub(crate) struct SeriesReader {
237 stream_ctx: Arc<StreamContext>,
238 partition_ranges: Vec<PartitionRange>,
239 range: SeriesRange,
240 filter: MetricSeriesFilter,
241 codec: SparsePrimaryKeyCodec,
242 partition_pruner: Arc<PartitionPruner>,
243 range_semaphore: Arc<Semaphore>,
244 part_metrics: PartitionMetrics,
245}
246
247#[allow(dead_code)]
248impl SeriesReader {
249 pub(crate) fn try_new(
259 stream_ctx: Arc<StreamContext>,
260 partition_ranges: Vec<PartitionRange>,
261 assigned_series: AssignedSeriesBatch,
262 partition_pruner: Arc<PartitionPruner>,
263 range_semaphore: Arc<Semaphore>,
264 part_metrics: PartitionMetrics,
265 ) -> Result<Self> {
266 validate_metric_metadata(&stream_ctx)?;
267 #[cfg(feature = "enterprise")]
268 snafu::ensure!(
269 stream_ctx.input.extension_ranges().is_empty(),
270 InvalidRequestSnafu {
271 region_id: stream_ctx.input.region_metadata().region_id,
272 reason: "series reader does not support extension ranges",
273 }
274 );
275
276 let range = assigned_series.range();
277 let filter = MetricSeriesFilter::new(&assigned_series);
278 let codec = SparsePrimaryKeyCodec::new(stream_ctx.input.region_metadata());
279 Ok(Self {
280 stream_ctx,
281 partition_ranges,
282 range,
283 filter,
284 codec,
285 partition_pruner,
286 range_semaphore,
287 part_metrics,
288 })
289 }
290
291 pub(crate) async fn build_stream(&self) -> Result<BoxedRecordBatchStream> {
292 if self.partition_ranges.is_empty() || self.filter.series.is_empty() {
293 return Ok(Box::pin(futures::stream::empty()));
294 }
295
296 let mut tasks = Vec::with_capacity(self.partition_ranges.len());
297 for part_range in self.partition_ranges.iter().copied() {
298 let stream_ctx = self.stream_ctx.clone();
299 let filter = self.filter.clone();
300 let codec = self.codec.clone();
301 let partition_pruner = self.partition_pruner.clone();
302 let range_semaphore = self.range_semaphore.clone();
303 let part_metrics = self.part_metrics.clone();
304 let range = self.range;
305 tasks.push(common_runtime::spawn_query(async move {
306 let _permit = range_semaphore.acquire().await.map_err(|error| {
307 UnexpectedSnafu {
308 reason: format!("failed to acquire series range permit: {error}"),
309 }
310 .build()
311 })?;
312 build_series_partition_range(
313 stream_ctx,
314 part_range,
315 range,
316 filter,
317 codec,
318 partition_pruner,
319 part_metrics,
320 )
321 .await
322 }));
323 }
324
325 let mut range_streams = Vec::with_capacity(tasks.len());
326 let mut estimated_batch_sizes = Vec::with_capacity(tasks.len());
327 for task in tasks {
328 let (stream, estimated_batch_size) = task.await.context(JoinSnafu)??;
329 range_streams.push(stream);
330 estimated_batch_sizes.push(estimated_batch_size);
331 }
332
333 let estimated_batch_size = compute_average_batch_size(estimated_batch_sizes);
336 SeqScan::build_flat_reader_from_sources(
337 &self.stream_ctx,
338 range_streams,
339 Some(self.range_semaphore.clone()),
340 Some(&self.part_metrics),
341 true,
342 compute_parallel_channel_size(estimated_batch_size),
343 )
344 .await
345 }
346}
347
348async fn build_series_partition_range(
349 stream_ctx: Arc<StreamContext>,
350 part_range: PartitionRange,
351 range: SeriesRange,
352 filter: MetricSeriesFilter,
353 codec: SparsePrimaryKeyCodec,
354 partition_pruner: Arc<PartitionPruner>,
355 part_metrics: PartitionMetrics,
356) -> Result<(BoxedRecordBatchStream, usize)> {
357 let cache_key = build_series_range_cache_key(&stream_ctx, &part_range, range);
358 if let Some(key) = cache_key.as_ref() {
359 if let Some(value) = stream_ctx.input.cache_strategy.get_range_result(key) {
360 part_metrics.inc_range_cache_hit();
361 return Ok((cached_flat_range_stream(value), DEFAULT_READ_BATCH_SIZE));
362 }
363 part_metrics.inc_range_cache_miss();
364 }
365
366 let range_meta = &stream_ctx.ranges[part_range.identifier];
367 let split_batch_size = should_split_flat_batches_for_merge(&stream_ctx, range_meta);
368 let mut sources = Vec::with_capacity(range_meta.row_group_indices.len());
369
370 for index in range_meta.row_group_indices.iter().copied() {
371 if stream_ctx.is_mem_range_index(index) {
372 let stream = Box::pin(scan_flat_mem_ranges(
373 stream_ctx.clone(),
374 part_metrics.clone(),
375 index,
376 range_meta.time_range,
377 ));
378 sources.push(filter_flat_stream_by_series(
379 stream,
380 codec.clone(),
381 filter.clone(),
382 ));
383 continue;
384 }
385
386 if stream_ctx.is_file_range_index(index) {
387 let file = stream_ctx.input.file_from_index(index);
388 if matches!(
389 file.primary_key_range()
390 .and_then(|(min, max)| { filter.overlaps_encoded_bounds(&codec, &min, &max) }),
391 Some(false)
392 ) {
393 continue;
394 }
395
396 if partition_pruner.try_skip_manifest_pruned_file_range(index, &part_metrics) {
397 continue;
398 }
399 let mut reader_metrics = ReaderMetrics {
400 filter_metrics: new_filter_metrics(part_metrics.explain_verbose()),
401 ..Default::default()
402 };
403 let file_ranges = partition_pruner
404 .build_file_ranges(index, &part_metrics, &mut reader_metrics)
405 .await?;
406 part_metrics.inc_num_file_ranges(file_ranges.len());
407 part_metrics.merge_reader_metrics(&reader_metrics, None);
408 let ranges = file_ranges
409 .iter()
410 .filter(|file_range| {
411 !matches!(
412 file_range.primary_key_range().and_then(|(min, max)| {
413 filter.overlaps_encoded_bounds(&codec, min, max)
414 }),
415 Some(false)
416 )
417 })
418 .cloned()
419 .collect::<smallvec::SmallVec<[_; 2]>>();
420 if ranges.is_empty() {
421 continue;
422 }
423
424 let stream = scan_series_file_ranges(
425 part_metrics.clone(),
426 ranges,
427 filter.clone(),
428 codec.clone(),
429 );
430 sources.push(Box::pin(stream) as BoxedRecordBatchStream);
431 continue;
432 }
433
434 return UnexpectedSnafu {
435 reason: format!(
436 "series reader received unsupported range index {}",
437 index.index
438 ),
439 }
440 .fail();
441 }
442
443 if split_batch_size.is_some() {
444 sources = sources
445 .into_iter()
446 .map(|stream| Box::pin(SplitRecordBatchStream::new(stream)) as BoxedRecordBatchStream)
447 .collect();
448 }
449 let estimated_batch_size = split_batch_size.unwrap_or(DEFAULT_READ_BATCH_SIZE);
450 let stream = SeqScan::build_flat_reader_from_sources(
451 &stream_ctx,
452 sources,
453 None,
454 Some(&part_metrics),
455 false,
456 compute_parallel_channel_size(estimated_batch_size),
457 )
458 .await?;
459 let stream = match cache_key {
460 Some(key) => cache_flat_range_stream(
461 stream,
462 stream_ctx.input.cache_strategy.clone(),
463 key,
464 part_metrics,
465 ),
466 None => stream,
467 };
468 Ok((stream, estimated_batch_size))
469}
470
471fn scan_series_file_ranges(
472 part_metrics: PartitionMetrics,
473 ranges: smallvec::SmallVec<[crate::sst::parquet::file_range::FileRange; 2]>,
474 filter: MetricSeriesFilter,
475 codec: SparsePrimaryKeyCodec,
476) -> impl futures::Stream<Item = Result<datatypes::arrow::record_batch::RecordBatch>> {
477 try_stream! {
478 let fetch_metrics = part_metrics
479 .explain_verbose()
480 .then(|| Arc::new(ParquetFetchMetrics::default()));
481 let mut reader_metrics = ReaderMetrics {
482 fetch_metrics: fetch_metrics.clone(),
483 ..Default::default()
484 };
485 let mut primary_key_filter = filter.primary_key_filter(codec);
486
487 for range in ranges {
488 let build_start = Instant::now();
489 let Some(mut reader) = range
490 .reader_by_primary_key(
491 primary_key_filter.as_mut(),
492 fetch_metrics.as_deref(),
493 )
494 .await?
495 else {
496 continue;
497 };
498 let build_cost = build_start.elapsed();
499 reader_metrics.build_cost += build_cost;
500 part_metrics.inc_build_reader_cost(build_cost);
501
502 let scan_start = Instant::now();
503 while let Some(record_batch) = reader.next_batch().await? {
504 reader_metrics.num_record_batches += 1;
505 reader_metrics.num_batches += 1;
506 reader_metrics.num_rows += record_batch.num_rows();
507
508 let num_rows_before_filter = record_batch.num_rows();
509 let Some(record_batch) = range.precise_filter_flat(
510 record_batch,
511 range.pre_filter_mode().skip_fields(),
512 true,
513 )? else {
514 reader_metrics.filter_metrics.rows_precise_filtered +=
515 num_rows_before_filter;
516 continue;
517 };
518 reader_metrics.filter_metrics.rows_precise_filtered +=
519 num_rows_before_filter - record_batch.num_rows();
520
521 let record_batch = if let Some(mapper) = range.compaction_projection_mapper() {
522 mapper.project(record_batch)?
523 } else {
524 record_batch
525 };
526 if let Some(compat) = range.compat_batch() {
527 yield compat.compat(record_batch)?;
528 } else {
529 yield record_batch;
530 }
531 }
532 reader_metrics.scan_cost += scan_start.elapsed();
533 }
534
535 reader_metrics.observe_rows("series_data");
536 reader_metrics.filter_metrics.observe();
537 part_metrics.merge_reader_metrics(&reader_metrics, None);
538 }
539}
540
541#[cfg(test)]
542mod tests {
543 use store_api::codec::PrimaryKeyEncoding;
544
545 use super::*;
546 use crate::error::DecodeSnafu;
547 use crate::test_util::sst_util::sst_region_metadata_with_encoding;
548
549 fn series(table_id: u32, tsid: u64) -> MetricSeriesId {
550 MetricSeriesId { table_id, tsid }
551 }
552
553 fn assigned_batch(
554 partitions: usize,
555 partition: usize,
556 series: Vec<MetricSeriesId>,
557 ) -> AssignedSeriesBatch {
558 let mut collector = SeriesBatchCollector::new(partitions).unwrap();
559 collector.push(series);
560 collector.finish().remove(partition)
561 }
562
563 #[test]
564 fn series_range_assignment_is_stable_across_batch_boundaries() {
565 let input = vec![
566 series(1, 0),
567 series(2, 1u64 << 62),
568 series(1, 1u64 << 63),
569 series(2, 3u64 << 62),
570 series(1, u64::MAX),
571 series(3, 0),
572 ];
573 let mut first = SeriesBatchCollector::new(4).unwrap();
574 first.push(input.clone());
575 let first = first.finish();
576
577 let mut second = SeriesBatchCollector::new(4).unwrap();
578 for chunk in input.chunks(2) {
579 second.push(chunk.to_vec());
580 }
581 let second = second.finish();
582
583 assert_eq!(
584 first
585 .iter()
586 .map(AssignedSeriesBatch::series)
587 .collect::<Vec<_>>(),
588 second
589 .iter()
590 .map(AssignedSeriesBatch::series)
591 .collect::<Vec<_>>()
592 );
593 assert_eq!(&[series(1, 0), series(3, 0)], first[0].series());
594 assert_eq!(&[series(2, 1u64 << 62)], first[1].series());
595 assert_eq!(
596 input.len(),
597 first
598 .iter()
599 .map(|batch| batch.series().len())
600 .sum::<usize>()
601 );
602 }
603
604 #[test]
605 fn series_ranges_cover_the_tsid_domain() {
606 let ranges = (0..3)
607 .map(|partition| SeriesRange::new(partition, 3).unwrap())
608 .collect::<Vec<_>>();
609 assert_eq!(0, ranges[0].start);
610 assert_eq!(TSID_DOMAIN_END, ranges[2].end);
611 assert_eq!(ranges[0].end, ranges[1].start);
612 assert_eq!(ranges[1].end, ranges[2].start);
613 assert_eq!(0, SeriesRange::partition_for(0, 3));
614 assert_eq!(2, SeriesRange::partition_for(u64::MAX, 3));
615
616 for (partition, range) in ranges.iter().enumerate() {
617 assert_eq!(partition, SeriesRange::partition_for(range.start as u64, 3));
618 if range.end < TSID_DOMAIN_END {
619 assert_eq!(
620 partition + 1,
621 SeriesRange::partition_for(range.end as u64, 3)
622 );
623 }
624 }
625 }
626
627 #[test]
628 fn metric_series_filter_matches_encoded_primary_key() {
629 let metadata = Arc::new(sst_region_metadata_with_encoding(
630 PrimaryKeyEncoding::Sparse,
631 ));
632 let assigned = assigned_batch(1, 0, vec![series(1, 10), series(2, 20)]);
633 let filter = MetricSeriesFilter::new(&assigned);
634 let codec = SparsePrimaryKeyCodec::new(&metadata);
635 let mut primary_key_filter = filter.primary_key_filter(codec);
636
637 for selected in [series(1, 10), series(2, 20)] {
638 assert!(
639 primary_key_filter
640 .matches(&encode_series(&metadata, selected.table_id, selected.tsid))
641 .context(DecodeSnafu)
642 .unwrap()
643 );
644 }
645 for unselected in [series(2, 10), series(1, 20)] {
646 assert!(
647 !primary_key_filter
648 .matches(&encode_series(
649 &metadata,
650 unselected.table_id,
651 unselected.tsid,
652 ))
653 .context(DecodeSnafu)
654 .unwrap()
655 );
656 }
657 }
658
659 fn encode_series(
660 metadata: &store_api::metadata::RegionMetadataRef,
661 table_id: u32,
662 tsid: u64,
663 ) -> Vec<u8> {
664 let codec = SparsePrimaryKeyCodec::new(metadata);
665 let mut primary_key = Vec::new();
666 codec
667 .encode_internal(table_id, tsid, &mut primary_key)
668 .unwrap();
669 primary_key
670 }
671
672 #[test]
673 fn series_range_overlaps_single_table_primary_key_bounds() {
674 let metadata = Arc::new(sst_region_metadata_with_encoding(
675 PrimaryKeyEncoding::Sparse,
676 ));
677 let assigned = assigned_batch(2, 0, vec![series(1, 10), series(2, 20)]);
678 let filter = MetricSeriesFilter::new(&assigned);
679 let codec = SparsePrimaryKeyCodec::new(&metadata);
680
681 assert_eq!(
683 Some(true),
684 filter.overlaps_encoded_bounds(
685 &codec,
686 &encode_series(&metadata, 1, 100),
687 &encode_series(&metadata, 1, 200),
688 )
689 );
690 assert_eq!(
691 Some(false),
692 filter.overlaps_encoded_bounds(
693 &codec,
694 &encode_series(&metadata, 1, 1u64 << 63),
695 &encode_series(&metadata, 1, u64::MAX),
696 )
697 );
698 assert_eq!(
699 Some(true),
700 filter.overlaps_encoded_bounds(
701 &codec,
702 &encode_series(&metadata, 3, 10),
703 &encode_series(&metadata, 3, 20),
704 )
705 );
706 }
707
708 #[test]
709 fn series_range_keeps_multiple_table_primary_key_bounds() {
710 let metadata = Arc::new(sst_region_metadata_with_encoding(
711 PrimaryKeyEncoding::Sparse,
712 ));
713 let assigned = assigned_batch(2, 0, vec![series(1, 10), series(2, 20)]);
714 let filter = MetricSeriesFilter::new(&assigned);
715 let codec = SparsePrimaryKeyCodec::new(&metadata);
716
717 assert_eq!(
718 Some(true),
719 filter.overlaps_encoded_bounds(
720 &codec,
721 &encode_series(&metadata, 1, 1u64 << 63),
722 &encode_series(&metadata, 2, u64::MAX),
723 )
724 );
725 }
726
727 #[test]
728 fn series_statistics_keep_invalid_or_inverted_bounds() {
729 let metadata = Arc::new(sst_region_metadata_with_encoding(
730 PrimaryKeyEncoding::Sparse,
731 ));
732 let assigned = assigned_batch(2, 0, vec![series(1, 10)]);
733 let filter = MetricSeriesFilter::new(&assigned);
734 let codec = SparsePrimaryKeyCodec::new(&metadata);
735
736 assert_eq!(
737 None,
738 filter.overlaps_encoded_bounds(&codec, b"invalid", b"bounds")
739 );
740 assert_eq!(
741 None,
742 filter.overlaps_encoded_bounds(
743 &codec,
744 &encode_series(&metadata, 2, 0),
745 &encode_series(&metadata, 1, 0),
746 )
747 );
748 }
749}