1use std::fmt;
18use std::sync::Arc;
19use std::time::Instant;
20
21use async_stream::try_stream;
22use common_error::ext::BoxedError;
23use common_recordbatch::util::ChainedRecordBatchStream;
24use common_recordbatch::{RecordBatchStreamWrapper, SendableRecordBatchStream};
25use common_telemetry::tracing;
26use datafusion::physical_plan::metrics::ExecutionPlanMetricsSet;
27use datafusion::physical_plan::{DisplayAs, DisplayFormatType};
28use datatypes::schema::SchemaRef;
29use futures::{StreamExt, TryStreamExt};
30use snafu::ensure;
31use store_api::metadata::RegionMetadataRef;
32use store_api::region_engine::{
33 PartitionRange, PrepareRequest, QueryScanContext, RegionScanner, ScannerProperties,
34};
35use store_api::storage::TimeSeriesRowSelector;
36use tokio::sync::Semaphore;
37
38use crate::error::{PartitionOutOfRangeSnafu, Result, TooManyFilesToReadSnafu};
39use crate::read::flat_dedup::{FlatDedupReader, FlatLastNonNull, FlatLastRow};
40use crate::read::flat_merge::FlatMergeReader;
41use crate::read::last_row::FlatLastRowReader;
42use crate::read::pruner::{PartitionPruner, Pruner};
43use crate::read::range::RangeMeta;
44use crate::read::range_cache::{
45 build_range_cache_key, cache_flat_range_stream, cached_flat_range_stream,
46};
47use crate::read::scan_region::{ScanInput, StreamContext};
48use crate::read::scan_util::{
49 PartitionMetrics, PartitionMetricsList, SplitRecordBatchStream, compute_parallel_channel_size,
50 scan_flat_file_ranges, scan_flat_mem_ranges, should_split_flat_batches_for_merge,
51};
52use crate::read::stream::{ConvertBatchStream, ScanBatch, ScanBatchStream};
53use crate::read::{BoxedRecordBatchStream, ScannerMetrics, scan_util};
54use crate::region::options::MergeMode;
55use crate::sst::parquet::DEFAULT_READ_BATCH_SIZE;
56
57pub struct SeqScan {
62 properties: ScannerProperties,
64 stream_ctx: Arc<StreamContext>,
66 pruner: Arc<Pruner>,
68 metrics_list: PartitionMetricsList,
71}
72
73impl SeqScan {
74 pub(crate) fn new(input: ScanInput) -> Self {
77 let mut properties = ScannerProperties::default()
78 .with_append_mode(input.append_mode)
79 .with_total_rows(input.total_rows());
80 if let Some(counters) = input.query_stat_counters.clone() {
81 properties.set_query_stat_counters(counters);
82 }
83 let stream_ctx = Arc::new(StreamContext::seq_scan_ctx(input));
84 properties.partitions = vec![stream_ctx.partition_ranges()];
85
86 let num_workers = common_stat::get_total_cpu_cores().max(1);
88 let pruner = Arc::new(Pruner::new(stream_ctx.clone(), num_workers));
89
90 Self {
91 properties,
92 stream_ctx,
93 pruner,
94 metrics_list: PartitionMetricsList::default(),
95 }
96 }
97
98 #[tracing::instrument(
103 skip_all,
104 fields(region_id = %self.stream_ctx.input.mapper.metadata().region_id)
105 )]
106 pub fn build_stream(&self) -> Result<SendableRecordBatchStream, BoxedError> {
107 let metrics_set = ExecutionPlanMetricsSet::new();
108 let streams = (0..self.properties.partitions.len())
109 .map(|partition: usize| {
110 self.scan_partition(&QueryScanContext::default(), &metrics_set, partition)
111 })
112 .collect::<Result<Vec<_>, _>>()?;
113
114 let aggr_stream = ChainedRecordBatchStream::new(streams).map_err(BoxedError::new)?;
115 Ok(Box::pin(aggr_stream))
116 }
117
118 pub(crate) fn scan_all_partitions(&self) -> Result<ScanBatchStream> {
120 let metrics_set = ExecutionPlanMetricsSet::new();
121
122 let streams = (0..self.properties.partitions.len())
123 .map(|partition| {
124 let metrics = self.new_partition_metrics(false, &metrics_set, partition);
125 self.scan_flat_batch_in_partition(partition, metrics)
126 })
127 .collect::<Result<Vec<_>>>()?;
128
129 Ok(Box::pin(futures::stream::iter(streams).flatten()))
130 }
131
132 pub(crate) async fn build_flat_reader_for_compaction(&self) -> Result<BoxedRecordBatchStream> {
137 assert!(self.stream_ctx.input.compaction);
138
139 let metrics_set = ExecutionPlanMetricsSet::new();
140 let part_metrics = self.new_partition_metrics(false, &metrics_set, 0);
141 debug_assert_eq!(1, self.properties.partitions.len());
142 let partition_ranges = &self.properties.partitions[0];
143
144 let reader = Self::merge_all_flat_ranges_for_compaction(
145 &self.stream_ctx,
146 partition_ranges,
147 &part_metrics,
148 self.pruner.clone(),
149 )
150 .await?;
151 Ok(reader)
152 }
153
154 async fn merge_all_flat_ranges_for_compaction(
157 stream_ctx: &Arc<StreamContext>,
158 partition_ranges: &[PartitionRange],
159 part_metrics: &PartitionMetrics,
160 pruner: Arc<Pruner>,
161 ) -> Result<BoxedRecordBatchStream> {
162 pruner.add_partition_ranges(partition_ranges);
163 let partition_pruner = Arc::new(PartitionPruner::new(pruner, partition_ranges));
164
165 let mut sources = Vec::new();
166 for part_range in partition_ranges {
167 build_flat_sources(
168 stream_ctx,
169 part_range,
170 true,
171 part_metrics,
172 partition_pruner.clone(),
173 &mut sources,
174 None,
175 )
176 .await?;
177 }
178
179 common_telemetry::debug!(
180 "Build flat reader to read all parts, region_id: {}, num_part_ranges: {}, num_sources: {}",
181 stream_ctx.input.mapper.metadata().region_id,
182 partition_ranges.len(),
183 sources.len()
184 );
185 Self::build_flat_reader_from_sources(
186 stream_ctx,
187 sources,
188 None,
189 None,
190 false,
191 compute_parallel_channel_size(DEFAULT_READ_BATCH_SIZE),
192 )
193 .await
194 }
195
196 #[tracing::instrument(level = tracing::Level::DEBUG, skip_all)]
200 pub(crate) async fn build_flat_reader_from_sources(
201 stream_ctx: &StreamContext,
202 mut sources: Vec<BoxedRecordBatchStream>,
203 semaphore: Option<Arc<Semaphore>>,
204 part_metrics: Option<&PartitionMetrics>,
205 skip_dedup: bool,
206 channel_size: usize,
207 ) -> Result<BoxedRecordBatchStream> {
208 if let Some(semaphore) = semaphore.as_ref() {
209 if sources.len() > 1 {
211 sources = stream_ctx.input.create_parallel_flat_sources(
212 sources,
213 semaphore.clone(),
214 channel_size,
215 )?;
216 }
217 }
218
219 let mapper = &stream_ctx.input.mapper;
220 let reader: BoxedRecordBatchStream = if sources.len() == 1 {
221 sources.pop().unwrap()
224 } else {
225 let schema = mapper.input_arrow_schema(stream_ctx.input.compaction);
226 let metrics_reporter = part_metrics.map(|m| m.merge_metrics_reporter());
227 let reader = FlatMergeReader::new(
228 schema,
229 sources,
230 stream_ctx.input.batch_size(),
231 metrics_reporter,
232 )
233 .await?;
234 Box::pin(reader.into_stream())
235 };
236
237 let dedup = !skip_dedup && !stream_ctx.input.append_mode;
238 let dedup_metrics_reporter = part_metrics.map(|m| m.dedup_metrics_reporter());
239 let reader = if dedup {
240 match stream_ctx.input.merge_mode {
241 MergeMode::LastRow => Box::pin(
242 FlatDedupReader::new(
243 reader,
244 FlatLastRow::new(stream_ctx.input.filter_deleted),
245 dedup_metrics_reporter,
246 )
247 .into_stream(),
248 ) as _,
249 MergeMode::LastNonNull => Box::pin(
250 FlatDedupReader::new(
251 reader,
252 FlatLastNonNull::new(
253 mapper.field_column_start(),
254 stream_ctx.input.filter_deleted,
255 ),
256 dedup_metrics_reporter,
257 )
258 .into_stream(),
259 ) as _,
260 }
261 } else {
262 reader
263 };
264
265 let reader = match &stream_ctx.input.series_row_selector {
266 Some(TimeSeriesRowSelector::LastRow { .. }) => {
267 Box::pin(FlatLastRowReader::new(reader).into_stream()) as _
268 }
269 None => reader,
270 };
271
272 Ok(reader)
273 }
274
275 pub(crate) async fn build_flat_partition_range_read(
277 stream_ctx: &Arc<StreamContext>,
278 part_range: &PartitionRange,
279 compaction: bool,
280 part_metrics: &PartitionMetrics,
281 partition_pruner: Arc<PartitionPruner>,
282 file_scan_semaphore: Option<Arc<Semaphore>>,
283 merge_semaphore: Option<Arc<Semaphore>>,
284 ) -> Result<(BoxedRecordBatchStream, usize)> {
285 let cache_key = build_range_cache_key(stream_ctx, part_range);
286
287 if let Some(key) = cache_key.as_ref() {
288 if let Some(value) = stream_ctx.input.cache_strategy.get_range_result(key) {
289 part_metrics.inc_range_cache_hit();
290 return Ok((cached_flat_range_stream(value), DEFAULT_READ_BATCH_SIZE));
291 }
292 part_metrics.inc_range_cache_miss();
293 }
294
295 let mut sources = Vec::new();
296 let split_batch_size = build_flat_sources(
297 stream_ctx,
298 part_range,
299 compaction,
300 part_metrics,
301 partition_pruner,
302 &mut sources,
303 file_scan_semaphore,
304 )
305 .await?;
306 let estimated_rows_per_batch = split_batch_size.unwrap_or(DEFAULT_READ_BATCH_SIZE);
307 let channel_size = compute_parallel_channel_size(estimated_rows_per_batch);
308 let stream = Self::build_flat_reader_from_sources(
309 stream_ctx,
310 sources,
311 merge_semaphore,
312 Some(part_metrics),
313 false,
314 channel_size,
315 )
316 .await?;
317
318 let stream = match cache_key {
319 Some(key) => cache_flat_range_stream(
320 stream,
321 stream_ctx.input.cache_strategy.clone(),
322 key,
323 part_metrics.clone(),
324 ),
325 None => stream,
326 };
327
328 Ok((stream, estimated_rows_per_batch))
329 }
330
331 fn scan_partition_impl(
334 &self,
335 ctx: &QueryScanContext,
336 metrics_set: &ExecutionPlanMetricsSet,
337 partition: usize,
338 ) -> Result<SendableRecordBatchStream> {
339 if ctx.explain_verbose {
340 common_telemetry::info!(
341 "SeqScan partition {}, region_id: {}",
342 partition,
343 self.stream_ctx.input.region_metadata().region_id
344 );
345 }
346
347 let metrics = self.new_partition_metrics(ctx.explain_verbose, metrics_set, partition);
348 let input = &self.stream_ctx.input;
349
350 let batch_stream = self.scan_flat_batch_in_partition(partition, metrics.clone())?;
351 let record_batch_stream = ConvertBatchStream::new(
352 batch_stream,
353 input.mapper.clone(),
354 input.cache_strategy.clone(),
355 metrics,
356 );
357
358 Ok(Box::pin(RecordBatchStreamWrapper::new(
359 input.mapper.output_schema(),
360 Box::pin(record_batch_stream),
361 )))
362 }
363
364 #[tracing::instrument(
365 skip_all,
366 fields(
367 region_id = %self.stream_ctx.input.mapper.metadata().region_id,
368 partition = partition
369 )
370 )]
371 fn scan_flat_batch_in_partition(
372 &self,
373 partition: usize,
374 part_metrics: PartitionMetrics,
375 ) -> Result<ScanBatchStream> {
376 ensure!(
377 partition < self.properties.partitions.len(),
378 PartitionOutOfRangeSnafu {
379 given: partition,
380 all: self.properties.partitions.len(),
381 }
382 );
383
384 if self.properties.partitions[partition].is_empty() {
385 return Ok(Box::pin(futures::stream::empty()));
386 }
387
388 let stream_ctx = self.stream_ctx.clone();
389 let semaphore = self.new_semaphore();
390 let partition_ranges = self.properties.partitions[partition].clone();
391 let compaction = self.stream_ctx.input.compaction;
392 let file_scan_semaphore = if compaction { None } else { semaphore.clone() };
393 let pruner = self.pruner.clone();
394 pruner.add_partition_ranges(&partition_ranges);
399 let partition_pruner = Arc::new(PartitionPruner::new(pruner, &partition_ranges));
400
401 let stream = try_stream! {
402 part_metrics.on_first_poll();
403 let mut fetch_start = Instant::now();
406
407 for part_range in partition_ranges {
409 let (mut reader, _) = Self::build_flat_partition_range_read(
410 &stream_ctx,
411 &part_range,
412 compaction,
413 &part_metrics,
414 partition_pruner.clone(),
415 file_scan_semaphore.clone(),
416 semaphore.clone(),
417 )
418 .await?;
419
420 let mut metrics = ScannerMetrics {
421 scan_cost: fetch_start.elapsed(),
422 ..Default::default()
423 };
424 fetch_start = Instant::now();
425
426 while let Some(record_batch) = reader.try_next().await? {
427 metrics.scan_cost += fetch_start.elapsed();
428 metrics.num_batches += 1;
429 metrics.num_rows += record_batch.num_rows();
430
431 debug_assert!(record_batch.num_rows() > 0);
432 if record_batch.num_rows() == 0 {
433 fetch_start = Instant::now();
434 continue;
435 }
436
437 let yield_start = Instant::now();
438 yield ScanBatch::RecordBatch(record_batch);
439 metrics.yield_cost += yield_start.elapsed();
440
441 fetch_start = Instant::now();
442 }
443
444 metrics.scan_cost += fetch_start.elapsed();
445 fetch_start = Instant::now();
446 part_metrics.merge_metrics(&metrics);
447 }
448
449 part_metrics.on_finish();
450 };
451 Ok(Box::pin(stream))
452 }
453
454 fn new_semaphore(&self) -> Option<Arc<Semaphore>> {
455 if self.properties.target_partitions() > self.properties.num_partitions() {
456 Some(Arc::new(Semaphore::new(
462 self.properties.target_partitions() - self.properties.num_partitions() + 1,
463 )))
464 } else {
465 None
466 }
467 }
468
469 fn new_partition_metrics(
472 &self,
473 explain_verbose: bool,
474 metrics_set: &ExecutionPlanMetricsSet,
475 partition: usize,
476 ) -> PartitionMetrics {
477 let metrics = PartitionMetrics::new(
478 self.stream_ctx.input.mapper.metadata().region_id,
479 partition,
480 get_scanner_type(self.stream_ctx.input.compaction),
481 self.stream_ctx.query_start,
482 explain_verbose,
483 metrics_set,
484 );
485
486 if !self.stream_ctx.input.compaction {
487 self.metrics_list.set(partition, metrics.clone());
488 }
489
490 metrics
491 }
492
493 fn max_files_in_partition(ranges: &[RangeMeta], partition_ranges: &[PartitionRange]) -> usize {
495 partition_ranges
496 .iter()
497 .map(|part_range| {
498 let range_meta = &ranges[part_range.identifier];
499 range_meta.indices.len()
500 })
501 .max()
502 .unwrap_or(0)
503 }
504
505 pub(crate) fn check_scan_limit(&self) -> Result<()> {
507 let total_max_files: usize = self
509 .properties
510 .partitions
511 .iter()
512 .map(|partition| Self::max_files_in_partition(&self.stream_ctx.ranges, partition))
513 .sum();
514
515 let max_concurrent_files = self.stream_ctx.input.max_concurrent_scan_files;
516 if total_max_files > max_concurrent_files {
517 return TooManyFilesToReadSnafu {
518 actual: total_max_files,
519 max: max_concurrent_files,
520 }
521 .fail();
522 }
523
524 Ok(())
525 }
526}
527
528impl RegionScanner for SeqScan {
529 fn name(&self) -> &str {
530 "SeqScan"
531 }
532
533 fn properties(&self) -> &ScannerProperties {
534 &self.properties
535 }
536
537 fn schema(&self) -> SchemaRef {
538 self.stream_ctx.input.mapper.output_schema()
539 }
540
541 fn metadata(&self) -> RegionMetadataRef {
542 self.stream_ctx.input.mapper.metadata().clone()
543 }
544
545 fn scan_partition(
546 &self,
547 ctx: &QueryScanContext,
548 metrics_set: &ExecutionPlanMetricsSet,
549 partition: usize,
550 ) -> Result<SendableRecordBatchStream, BoxedError> {
551 self.scan_partition_impl(ctx, metrics_set, partition)
552 .map_err(BoxedError::new)
553 }
554
555 fn prepare(&mut self, request: PrepareRequest) -> Result<(), BoxedError> {
556 self.properties.prepare(request);
557
558 self.check_scan_limit().map_err(BoxedError::new)?;
559
560 Ok(())
561 }
562
563 fn has_predicate_without_region(&self) -> bool {
564 let predicate = self
565 .stream_ctx
566 .input
567 .predicate_group()
568 .predicate_without_region();
569 predicate.is_some()
570 }
571
572 fn add_dyn_filter_to_predicate(
573 &mut self,
574 filter_exprs: Vec<Arc<dyn datafusion::physical_plan::PhysicalExpr>>,
575 ) -> Vec<bool> {
576 self.stream_ctx.add_dyn_filter_to_predicate(filter_exprs)
577 }
578
579 fn reset_state(&mut self) {
580 self.stream_ctx.input.predicate.clear_dyn_filters();
581 let num_workers = common_stat::get_total_cpu_cores().max(1);
582 self.pruner = Arc::new(Pruner::new(self.stream_ctx.clone(), num_workers));
583 self.metrics_list = PartitionMetricsList::default();
584 }
585
586 fn set_logical_region(&mut self, logical_region: bool) {
587 self.properties.set_logical_region(logical_region);
588 }
589
590 fn set_query_load_region_id(&mut self, region_id: store_api::storage::RegionId) {
591 self.properties.set_query_load_region_id(region_id);
592 }
593
594 fn snapshot_sequence(&self) -> Option<u64> {
595 self.stream_ctx.input.snapshot_sequence
596 }
597}
598
599impl DisplayAs for SeqScan {
600 fn fmt_as(&self, t: DisplayFormatType, f: &mut fmt::Formatter) -> fmt::Result {
601 write!(
602 f,
603 "SeqScan: region={}, ",
604 self.stream_ctx.input.mapper.metadata().region_id
605 )?;
606 match t {
607 DisplayFormatType::Default | DisplayFormatType::TreeRender => {
609 self.stream_ctx.format_for_explain(false, f)
610 }
611 DisplayFormatType::Verbose => {
612 self.stream_ctx.format_for_explain(true, f)?;
613 self.metrics_list.format_verbose_metrics(f)
614 }
615 }
616 }
617}
618
619impl fmt::Debug for SeqScan {
620 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
621 f.debug_struct("SeqScan")
622 .field("num_ranges", &self.stream_ctx.ranges.len())
623 .finish()
624 }
625}
626
627pub(crate) async fn build_flat_sources(
630 stream_ctx: &Arc<StreamContext>,
631 part_range: &PartitionRange,
632 compaction: bool,
633 part_metrics: &PartitionMetrics,
634 partition_pruner: Arc<PartitionPruner>,
635 sources: &mut Vec<BoxedRecordBatchStream>,
636 semaphore: Option<Arc<Semaphore>>,
637) -> Result<Option<usize>> {
638 let range_meta = &stream_ctx.ranges[part_range.identifier];
640 #[cfg(debug_assertions)]
641 if compaction {
642 debug_assert_eq!(range_meta.indices.len(), range_meta.row_group_indices.len());
644 for (i, row_group_idx) in range_meta.row_group_indices.iter().enumerate() {
645 debug_assert_eq!(
647 -1, row_group_idx.row_group_index,
648 "Expect {} range scan all row groups, given: {}",
649 i, row_group_idx.row_group_index,
650 );
651 }
652 }
653
654 let read_type = if compaction {
655 "compaction"
656 } else {
657 "seq_scan_files"
658 };
659 let num_indices = range_meta.row_group_indices.len();
660 if num_indices == 0 {
661 return Ok(None);
662 }
663
664 let split_batch_size = should_split_flat_batches_for_merge(stream_ctx, range_meta);
665 let should_split = split_batch_size.is_some();
666 sources.reserve(num_indices);
667 let mut ordered_sources = Vec::with_capacity(num_indices);
668 ordered_sources.resize_with(num_indices, || None);
669 let mut file_scan_tasks = Vec::new();
670 let pre_filter_mode = stream_ctx.range_pre_filter_mode(part_range);
671
672 for (position, index) in range_meta.row_group_indices.iter().enumerate() {
673 if stream_ctx.is_mem_range_index(*index) {
674 let stream = scan_flat_mem_ranges(
675 stream_ctx.clone(),
676 part_metrics.clone(),
677 *index,
678 range_meta.time_range,
679 );
680 ordered_sources[position] = Some(Box::pin(stream) as _);
681 } else if stream_ctx.is_file_range_index(*index) {
682 if !compaction
685 && partition_pruner.try_skip_manifest_pruned_file_range(*index, part_metrics)
686 {
687 continue;
688 }
689 if let Some(semaphore_ref) = semaphore.as_ref() {
690 let stream_ctx = stream_ctx.clone();
692 let part_metrics = part_metrics.clone();
693 let partition_pruner = partition_pruner.clone();
694 let semaphore = Arc::clone(semaphore_ref);
695 let row_group_index = *index;
696 file_scan_tasks.push(async move {
697 let _permit = semaphore.acquire().await.unwrap();
698 let stream = scan_flat_file_ranges(
699 stream_ctx,
700 part_metrics,
701 row_group_index,
702 read_type,
703 partition_pruner,
704 )
705 .await?;
706 Ok((position, Box::pin(stream) as _))
707 });
708 } else {
709 let stream = scan_flat_file_ranges(
711 stream_ctx.clone(),
712 part_metrics.clone(),
713 *index,
714 read_type,
715 partition_pruner.clone(),
716 )
717 .await?;
718 ordered_sources[position] = Some(Box::pin(stream) as _);
719 }
720 } else {
721 let stream = scan_util::maybe_scan_flat_other_ranges(
722 stream_ctx,
723 *index,
724 part_metrics,
725 pre_filter_mode,
726 )
727 .await?;
728 ordered_sources[position] = Some(stream);
729 }
730 }
731
732 if !file_scan_tasks.is_empty() {
733 let results = futures::future::try_join_all(file_scan_tasks).await?;
734 for (position, stream) in results {
735 ordered_sources[position] = Some(stream);
736 }
737 }
738
739 for stream in ordered_sources.into_iter().flatten() {
740 if should_split {
741 sources.push(Box::pin(SplitRecordBatchStream::new(stream)));
742 } else {
743 sources.push(stream);
744 }
745 }
746
747 if should_split {
748 common_telemetry::debug!(
749 "Splitting record batches, region: {}, sources: {}, part_range: {:?}",
750 stream_ctx.input.region_metadata().region_id,
751 sources.len(),
752 part_range,
753 );
754 }
755
756 Ok(split_batch_size)
757}
758
759#[cfg(test)]
760impl SeqScan {
761 pub(crate) fn input(&self) -> &ScanInput {
763 &self.stream_ctx.input
764 }
765}
766
767fn get_scanner_type(compaction: bool) -> &'static str {
769 if compaction {
770 "SeqScan(compaction)"
771 } else {
772 "SeqScan"
773 }
774}