1use std::fmt;
18use std::sync::atomic::{AtomicUsize, Ordering};
19use std::sync::{Arc, Mutex};
20use std::time::{Duration, Instant};
21
22use async_stream::try_stream;
23use common_error::ext::BoxedError;
24use common_recordbatch::util::ChainedRecordBatchStream;
25use common_recordbatch::{RecordBatchStreamWrapper, SendableRecordBatchStream};
26use common_telemetry::tracing::{self, Instrument};
27use common_telemetry::warn;
28use datafusion::physical_plan::metrics::ExecutionPlanMetricsSet;
29use datafusion::physical_plan::{DisplayAs, DisplayFormatType};
30use datatypes::arrow::array::BinaryArray;
31use datatypes::arrow::record_batch::RecordBatch;
32use datatypes::schema::SchemaRef;
33use futures::{StreamExt, TryStreamExt};
34use smallvec::SmallVec;
35use snafu::{OptionExt, ResultExt, ensure};
36use store_api::metadata::RegionMetadataRef;
37use store_api::region_engine::{
38 PartitionRange, PrepareRequest, QueryScanContext, RegionScanner, ScannerProperties,
39};
40use tokio::sync::Semaphore;
41use tokio::sync::mpsc::error::{SendTimeoutError, TrySendError};
42use tokio::sync::mpsc::{self, Receiver, Sender};
43
44use crate::error::{
45 Error, InvalidSenderSnafu, JoinSnafu, PartitionOutOfRangeSnafu, Result, ScanMultiTimesSnafu,
46 ScanSeriesSnafu, TooManyFilesToReadSnafu,
47};
48use crate::read::ScannerMetrics;
49use crate::read::pruner::{PartitionPruner, Pruner, PrunerOptions};
50use crate::read::scan_region::{ScanInput, StreamContext};
51use crate::read::scan_util::{
52 PartitionMetrics, PartitionMetricsList, SeriesDistributorMetrics, compute_average_batch_size,
53 compute_parallel_channel_size,
54};
55use crate::read::seq_scan::SeqScan;
56use crate::read::series_candidate::{SeriesCandidateScanner, is_sparse_metric_metadata};
57use crate::read::series_reader::{AssignedSeriesBatch, SeriesBatchCollector, SeriesReader};
58use crate::read::stream::{ConvertBatchStream, ScanBatch, ScanBatchStream};
59use crate::sst::parquet::flat_format::primary_key_column_index;
60use crate::sst::parquet::format::PrimaryKeyArray;
61
62const SEND_TIMEOUT: Duration = Duration::from_micros(100);
64
65const CANDIDATE_SERIES_ASSIGNMENT_THRESHOLD: usize = 1_000_000;
67
68type LegacyReceiverList = Vec<Option<Receiver<Result<SeriesBatch>>>>;
70
71type CandidateReceiverList = Vec<Option<CandidateReceiver>>;
73
74type CandidateSenderList = Vec<Option<mpsc::UnboundedSender<Result<SeriesReaderInput>>>>;
76
77struct CandidateReceiver {
78 receiver: mpsc::UnboundedReceiver<Result<SeriesReaderInput>>,
79 active_receivers: Arc<AtomicUsize>,
80}
81
82impl CandidateReceiver {
83 async fn recv(&mut self) -> Option<Result<SeriesReaderInput>> {
84 self.receiver.recv().await
85 }
86}
87
88impl Drop for CandidateReceiver {
89 fn drop(&mut self) {
90 let previous = self.active_receivers.fetch_sub(1, Ordering::Relaxed);
91 debug_assert!(previous > 0);
92 }
93}
94
95struct SeriesReaderInput {
97 assigned_series: AssignedSeriesBatch,
98 partition_pruner: Arc<PartitionPruner>,
99 range_semaphore: Arc<Semaphore>,
100}
101
102#[derive(Debug, Clone, Copy, PartialEq, Eq)]
103enum SeriesScanMode {
104 Legacy,
105 TwoPhase,
106}
107
108impl SeriesScanMode {
109 fn as_str(self) -> &'static str {
110 match self {
111 Self::Legacy => "legacy",
112 Self::TwoPhase => "two_phase",
113 }
114 }
115}
116
117pub struct SeriesScan {
123 mode: SeriesScanMode,
125 properties: ScannerProperties,
127 stream_ctx: Arc<StreamContext>,
129 pruner: Arc<Pruner>,
131 legacy_receivers: Mutex<LegacyReceiverList>,
133 candidate_receivers: Mutex<CandidateReceiverList>,
135 metrics_list: Arc<PartitionMetricsList>,
138}
139
140impl SeriesScan {
141 pub(crate) fn new(input: ScanInput, experimental_series_scan_v2: bool) -> Self {
143 let mode = if experimental_series_scan_v2 && Self::supports_two_phase(&input) {
144 SeriesScanMode::TwoPhase
145 } else {
146 SeriesScanMode::Legacy
147 };
148 let mut properties = ScannerProperties::default()
149 .with_append_mode(input.append_mode)
150 .with_total_rows(input.total_rows());
151 if let Some(counters) = input.query_stat_counters.clone() {
152 properties.set_query_stat_counters(counters);
153 }
154 let stream_ctx = Arc::new(StreamContext::seq_scan_ctx(input));
155 properties.partitions = vec![stream_ctx.partition_ranges()];
156
157 let num_workers = common_stat::get_total_cpu_cores().max(1);
159 let pruner = match mode {
160 SeriesScanMode::Legacy => Arc::new(Pruner::new(stream_ctx.clone(), num_workers)),
161 SeriesScanMode::TwoPhase => Arc::new(Pruner::new_with_options(
162 stream_ctx.clone(),
163 num_workers,
164 PrunerOptions {
165 retain_builders: true,
166 enable_predicate_prefilter: false,
167 },
168 )),
169 };
170
171 Self {
172 mode,
173 properties,
174 stream_ctx,
175 pruner,
176 legacy_receivers: Mutex::new(Vec::new()),
177 candidate_receivers: Mutex::new(Vec::new()),
178 metrics_list: Arc::new(PartitionMetricsList::default()),
179 }
180 }
181
182 fn supports_two_phase(input: &ScanInput) -> bool {
183 if input.sequence_range.is_some() || !is_sparse_metric_metadata(input.region_metadata()) {
184 return false;
185 }
186 #[cfg(feature = "enterprise")]
187 if !input.extension_ranges().is_empty() {
188 return false;
189 }
190 true
191 }
192
193 #[tracing::instrument(
194 skip_all,
195 fields(
196 region_id = %self.stream_ctx.input.mapper.metadata().region_id,
197 partition = partition
198 )
199 )]
200 fn scan_partition_impl(
201 &self,
202 ctx: &QueryScanContext,
203 metrics_set: &ExecutionPlanMetricsSet,
204 partition: usize,
205 ) -> Result<SendableRecordBatchStream> {
206 let metrics = new_partition_metrics(
207 &self.stream_ctx,
208 ctx.explain_verbose,
209 metrics_set,
210 partition,
211 &self.metrics_list,
212 );
213
214 let batch_stream =
215 self.scan_batch_in_partition(ctx, partition, metrics.clone(), metrics_set)?;
216
217 let input = &self.stream_ctx.input;
218 let record_batch_stream = ConvertBatchStream::new(
219 batch_stream,
220 input.mapper.clone(),
221 input.cache_strategy.clone(),
222 metrics,
223 );
224
225 Ok(Box::pin(RecordBatchStreamWrapper::new(
226 input.mapper.output_schema(),
227 Box::pin(record_batch_stream),
228 )))
229 }
230
231 #[tracing::instrument(
232 skip_all,
233 fields(
234 region_id = %self.stream_ctx.input.mapper.metadata().region_id,
235 partition = partition
236 )
237 )]
238 fn scan_batch_in_partition(
239 &self,
240 ctx: &QueryScanContext,
241 partition: usize,
242 part_metrics: PartitionMetrics,
243 metrics_set: &ExecutionPlanMetricsSet,
244 ) -> Result<ScanBatchStream> {
245 if ctx.explain_verbose {
246 common_telemetry::info!(
247 "SeriesScan partition {}, region_id: {}",
248 partition,
249 self.stream_ctx.input.region_metadata().region_id
250 );
251 }
252
253 ensure!(
254 partition < self.properties.num_partitions(),
255 PartitionOutOfRangeSnafu {
256 given: partition,
257 all: self.properties.num_partitions(),
258 }
259 );
260
261 match self.mode {
262 SeriesScanMode::Legacy => self.scan_legacy_batch_in_partition(
263 partition,
264 part_metrics,
265 metrics_set,
266 ctx.explain_verbose,
267 ),
268 SeriesScanMode::TwoPhase => self.scan_two_phase_batch_in_partition(
269 partition,
270 part_metrics,
271 metrics_set,
272 ctx.explain_verbose,
273 ),
274 }
275 }
276
277 fn scan_legacy_batch_in_partition(
278 &self,
279 partition: usize,
280 part_metrics: PartitionMetrics,
281 metrics_set: &ExecutionPlanMetricsSet,
282 explain_verbose: bool,
283 ) -> Result<ScanBatchStream> {
284 self.maybe_start_legacy_distributor(metrics_set, explain_verbose);
285 let mut receiver = self.legacy_receivers.lock().unwrap()[partition]
286 .take()
287 .context(ScanMultiTimesSnafu { partition })?;
288 let stream = try_stream! {
289 part_metrics.on_first_poll();
290
291 let mut fetch_start = Instant::now();
292 let mut metrics = ScannerMetrics::default();
293 while let Some(series) = receiver.recv().await {
294 let series = series?;
295 metrics.scan_cost += fetch_start.elapsed();
296 metrics.num_batches += series.num_batches();
297 metrics.num_rows += series.num_rows();
298
299 let yield_start = Instant::now();
300 yield ScanBatch::Series(series);
301 metrics.yield_cost += yield_start.elapsed();
302 fetch_start = Instant::now();
303 }
304
305 part_metrics.merge_metrics(&metrics);
306 part_metrics.on_finish();
307 };
308 Ok(Box::pin(stream))
309 }
310
311 fn scan_two_phase_batch_in_partition(
312 &self,
313 partition: usize,
314 part_metrics: PartitionMetrics,
315 metrics_set: &ExecutionPlanMetricsSet,
316 explain_verbose: bool,
317 ) -> Result<ScanBatchStream> {
318 self.maybe_start_candidate_distributor(metrics_set, explain_verbose);
319 let mut receiver = self.candidate_receivers.lock().unwrap()[partition]
322 .take()
323 .context(ScanMultiTimesSnafu { partition })?;
324 let stream_ctx = self.stream_ctx.clone();
325 let partition_ranges = self
326 .properties
327 .partitions
328 .iter()
329 .flatten()
330 .copied()
331 .collect::<Vec<_>>();
332 let stream = try_stream! {
333 part_metrics.on_first_poll();
334
335 let mut fetch_start = Instant::now();
336 let mut metrics = ScannerMetrics::default();
337 while let Some(input) = receiver.recv().await {
338 let input = input?;
339 metrics.scan_cost += fetch_start.elapsed();
340
341 let build_start = Instant::now();
342 let reader = SeriesReader::try_new(
343 stream_ctx.clone(),
344 partition_ranges.clone(),
345 input.assigned_series,
346 input.partition_pruner,
347 input.range_semaphore,
348 part_metrics.clone(),
349 )?;
350 let mut reader_stream = reader.build_stream().await?;
351 metrics.scan_cost += build_start.elapsed();
352 fetch_start = Instant::now();
353 while let Some(record_batch) = reader_stream.try_next().await? {
354 metrics.scan_cost += fetch_start.elapsed();
355 metrics.num_batches += 1;
356 metrics.num_rows += record_batch.num_rows();
357
358 let yield_start = Instant::now();
359 yield ScanBatch::RecordBatch(record_batch);
360 metrics.yield_cost += yield_start.elapsed();
361 fetch_start = Instant::now();
362 }
363 }
364 metrics.scan_cost += fetch_start.elapsed();
365
366 part_metrics.merge_metrics(&metrics);
367 part_metrics.on_finish();
368 };
369 Ok(Box::pin(stream))
370 }
371
372 fn maybe_start_legacy_distributor(
373 &self,
374 metrics_set: &ExecutionPlanMetricsSet,
375 explain_verbose: bool,
376 ) {
377 let mut rx_list = self.legacy_receivers.lock().unwrap();
378 if !rx_list.is_empty() {
379 return;
380 }
381
382 let (senders, receivers) = new_legacy_channel_list(self.properties.num_partitions());
383 let mut distributor = SeriesDistributor {
384 stream_ctx: self.stream_ctx.clone(),
385 range_semaphore: Some(Arc::new(Semaphore::new(self.properties.num_partitions()))),
386 final_merge_semaphore: Some(Arc::new(Semaphore::new(self.properties.num_partitions()))),
387 partitions: self.properties.partitions.clone(),
388 pruner: self.pruner.clone(),
389 senders,
390 metrics_set: metrics_set.clone(),
391 metrics_list: self.metrics_list.clone(),
392 explain_verbose,
393 };
394 let region_id = distributor.stream_ctx.input.mapper.metadata().region_id;
395 let span = tracing::info_span!("SeriesScan::distributor", region_id = %region_id);
396 common_runtime::spawn_query(
397 async move {
398 distributor.execute().await;
399 }
400 .instrument(span),
401 );
402
403 *rx_list = receivers;
404 }
405
406 fn maybe_start_candidate_distributor(
407 &self,
408 metrics_set: &ExecutionPlanMetricsSet,
409 explain_verbose: bool,
410 ) {
411 let mut rx_list = self.candidate_receivers.lock().unwrap();
412 if !rx_list.is_empty() {
413 return;
414 }
415
416 let (senders, receivers, active_receivers) =
417 new_candidate_channel_list(self.properties.num_partitions());
418 let mut distributor = SeriesCandidateDistributor {
419 stream_ctx: self.stream_ctx.clone(),
420 range_semaphore: Arc::new(Semaphore::new(self.properties.num_partitions())),
421 partitions: self.properties.partitions.clone(),
422 pruner: self.pruner.clone(),
423 senders,
424 active_receivers,
425 metrics_set: metrics_set.clone(),
426 metrics_list: self.metrics_list.clone(),
427 explain_verbose,
428 };
429 let region_id = distributor.stream_ctx.input.mapper.metadata().region_id;
430 let span = tracing::info_span!("SeriesScan::candidate_distributor", region_id = %region_id);
431 common_runtime::spawn_query(
432 async move {
433 distributor.execute().await;
434 }
435 .instrument(span),
436 );
437
438 *rx_list = receivers;
439 }
440
441 #[tracing::instrument(
443 skip_all,
444 fields(region_id = %self.stream_ctx.input.mapper.metadata().region_id)
445 )]
446 pub(crate) async fn build_stream(&self) -> Result<SendableRecordBatchStream, BoxedError> {
447 let part_num = self.properties.num_partitions();
448 let metrics_set = ExecutionPlanMetricsSet::default();
449 let streams = (0..part_num)
450 .map(|i| self.scan_partition(&QueryScanContext::default(), &metrics_set, i))
451 .collect::<Result<Vec<_>, BoxedError>>()?;
452 let chained_stream = ChainedRecordBatchStream::new(streams).map_err(BoxedError::new)?;
453 Ok(Box::pin(chained_stream))
454 }
455
456 pub(crate) fn scan_all_partitions(&self) -> Result<ScanBatchStream> {
458 let metrics_set = ExecutionPlanMetricsSet::new();
459
460 let streams = (0..self.properties.partitions.len())
461 .map(|partition| {
462 let metrics = new_partition_metrics(
463 &self.stream_ctx,
464 false,
465 &metrics_set,
466 partition,
467 &self.metrics_list,
468 );
469
470 self.scan_batch_in_partition(
471 &QueryScanContext::default(),
472 partition,
473 metrics,
474 &metrics_set,
475 )
476 })
477 .collect::<Result<Vec<_>>>()?;
478
479 Ok(Box::pin(futures::stream::iter(streams).flatten()))
480 }
481
482 pub(crate) fn check_scan_limit(&self) -> Result<()> {
484 let total_files: usize = self
486 .properties
487 .partitions
488 .iter()
489 .flat_map(|partition| partition.iter())
490 .map(|part_range| {
491 let range_meta = &self.stream_ctx.ranges[part_range.identifier];
492 range_meta.indices.len()
493 })
494 .sum();
495
496 let max_concurrent_files = self.stream_ctx.input.max_concurrent_scan_files;
497 if total_files > max_concurrent_files {
498 return TooManyFilesToReadSnafu {
499 actual: total_files,
500 max: max_concurrent_files,
501 }
502 .fail();
503 }
504
505 Ok(())
506 }
507}
508
509fn new_legacy_channel_list(num_partitions: usize) -> (SenderList, LegacyReceiverList) {
510 let (senders, receivers): (Vec<_>, Vec<_>) = (0..num_partitions)
511 .map(|_| {
512 let (sender, receiver) = mpsc::channel(1);
513 (Some(sender), Some(receiver))
514 })
515 .unzip();
516 (SenderList::new(senders), receivers)
517}
518
519fn new_candidate_channel_list(
520 num_partitions: usize,
521) -> (CandidateSenderList, CandidateReceiverList, Arc<AtomicUsize>) {
522 let active_receivers = Arc::new(AtomicUsize::new(num_partitions));
523 let (senders, receivers) = (0..num_partitions)
524 .map(|_| {
525 let (sender, receiver) = mpsc::unbounded_channel();
529 let receiver = CandidateReceiver {
530 receiver,
531 active_receivers: active_receivers.clone(),
532 };
533 (Some(sender), Some(receiver))
534 })
535 .unzip();
536 (senders, receivers, active_receivers)
537}
538
539impl RegionScanner for SeriesScan {
540 fn name(&self) -> &str {
541 "SeriesScan"
542 }
543
544 fn properties(&self) -> &ScannerProperties {
545 &self.properties
546 }
547
548 fn schema(&self) -> SchemaRef {
549 self.stream_ctx.input.mapper.output_schema()
550 }
551
552 fn metadata(&self) -> RegionMetadataRef {
553 self.stream_ctx.input.mapper.metadata().clone()
554 }
555
556 fn scan_partition(
557 &self,
558 ctx: &QueryScanContext,
559 metrics_set: &ExecutionPlanMetricsSet,
560 partition: usize,
561 ) -> Result<SendableRecordBatchStream, BoxedError> {
562 self.scan_partition_impl(ctx, metrics_set, partition)
563 .map_err(BoxedError::new)
564 }
565
566 fn prepare(&mut self, request: PrepareRequest) -> Result<(), BoxedError> {
567 self.properties.prepare(request);
568
569 self.check_scan_limit().map_err(BoxedError::new)?;
570
571 Ok(())
572 }
573
574 fn has_predicate_without_region(&self) -> bool {
575 let predicate = self
576 .stream_ctx
577 .input
578 .predicate_group()
579 .predicate_without_region();
580 predicate.is_some()
581 }
582
583 fn add_dyn_filter_to_predicate(
584 &mut self,
585 filter_exprs: Vec<Arc<dyn datafusion::physical_plan::PhysicalExpr>>,
586 ) -> Vec<bool> {
587 self.stream_ctx.add_dyn_filter_to_predicate(filter_exprs)
588 }
589
590 fn reset_state(&mut self) {
591 self.stream_ctx.input.predicate.clear_dyn_filters();
592 let num_workers = common_stat::get_total_cpu_cores().max(1);
593 self.pruner = match self.mode {
594 SeriesScanMode::Legacy => Arc::new(Pruner::new(self.stream_ctx.clone(), num_workers)),
595 SeriesScanMode::TwoPhase => Arc::new(Pruner::new_with_options(
596 self.stream_ctx.clone(),
597 num_workers,
598 PrunerOptions {
599 retain_builders: true,
600 enable_predicate_prefilter: false,
601 },
602 )),
603 };
604 self.legacy_receivers.lock().unwrap().clear();
605 self.candidate_receivers.lock().unwrap().clear();
606 self.metrics_list = Arc::new(PartitionMetricsList::default());
607 }
608
609 fn set_logical_region(&mut self, logical_region: bool) {
610 self.properties.set_logical_region(logical_region);
611 }
612
613 fn set_query_load_region_id(&mut self, region_id: store_api::storage::RegionId) {
614 self.properties.set_query_load_region_id(region_id);
615 }
616
617 fn snapshot_sequence(&self) -> Option<u64> {
618 self.stream_ctx.input.snapshot_sequence
619 }
620}
621
622impl DisplayAs for SeriesScan {
623 fn fmt_as(&self, t: DisplayFormatType, f: &mut fmt::Formatter) -> fmt::Result {
624 write!(
625 f,
626 "SeriesScan: region={}, ",
627 self.stream_ctx.input.mapper.metadata().region_id
628 )?;
629 match t {
630 DisplayFormatType::Default | DisplayFormatType::TreeRender => {
631 self.stream_ctx.format_for_explain(false, f)?;
632 }
633 DisplayFormatType::Verbose => {
634 self.stream_ctx.format_for_explain(true, f)?;
635 }
636 }
637 write!(f, ", \"mode\":\"{}\"", self.mode.as_str())?;
638 if matches!(t, DisplayFormatType::Verbose) {
639 self.metrics_list.format_verbose_metrics(f)?;
640 }
641 Ok(())
642 }
643}
644
645impl fmt::Debug for SeriesScan {
646 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
647 f.debug_struct("SeriesScan")
648 .field("mode", &self.mode)
649 .field("num_ranges", &self.stream_ctx.ranges.len())
650 .finish()
651 }
652}
653
654#[cfg(test)]
655impl SeriesScan {
656 pub(crate) fn input(&self) -> &ScanInput {
658 &self.stream_ctx.input
659 }
660
661 pub(crate) fn mode(&self) -> &'static str {
663 self.mode.as_str()
664 }
665}
666
667struct SeriesDistributor {
669 stream_ctx: Arc<StreamContext>,
671 range_semaphore: Option<Arc<Semaphore>>,
673 final_merge_semaphore: Option<Arc<Semaphore>>,
678 partitions: Vec<Vec<PartitionRange>>,
680 pruner: Arc<Pruner>,
682 senders: SenderList,
684 metrics_set: ExecutionPlanMetricsSet,
690 metrics_list: Arc<PartitionMetricsList>,
691 explain_verbose: bool,
693}
694
695impl SeriesDistributor {
696 #[tracing::instrument(
698 skip_all,
699 fields(region_id = %self.stream_ctx.input.mapper.metadata().region_id)
700 )]
701 async fn execute(&mut self) {
702 if let Err(e) = self.scan_partitions_flat().await {
703 self.senders.send_error(e).await;
704 }
705 }
706
707 #[tracing::instrument(
709 skip_all,
710 fields(region_id = %self.stream_ctx.input.mapper.metadata().region_id)
711 )]
712 async fn scan_partitions_flat(&mut self) -> Result<()> {
713 for partition_ranges in &self.partitions {
715 self.pruner.add_partition_ranges(partition_ranges);
716 }
717
718 let all_partition_ranges: Vec<_> = self.partitions.iter().flatten().cloned().collect();
720 let partition_pruner = Arc::new(PartitionPruner::new(
721 self.pruner.clone(),
722 &all_partition_ranges,
723 ));
724
725 let part_metrics = new_partition_metrics(
726 &self.stream_ctx,
727 self.explain_verbose,
728 &self.metrics_set,
729 self.partitions.len(),
730 &self.metrics_list,
731 );
732 part_metrics.on_first_poll();
733 let mut fetch_start = Instant::now();
736
737 let build_start = Instant::now();
739 let mut tasks = Vec::new();
740 for partition in &self.partitions {
741 for part_range in partition {
742 let stream_ctx = self.stream_ctx.clone();
743 let part_range = *part_range;
744 let part_metrics = part_metrics.clone();
745 let partition_pruner = partition_pruner.clone();
746 let file_scan_semaphore = self.range_semaphore.clone();
747 let merge_semaphore = self.range_semaphore.clone();
748 tasks.push(common_runtime::spawn_query(async move {
749 SeqScan::build_flat_partition_range_read(
750 &stream_ctx,
751 &part_range,
752 false,
753 &part_metrics,
754 partition_pruner,
755 file_scan_semaphore,
756 merge_semaphore,
757 )
758 .await
759 }));
760 }
761 }
762 let mut range_streams = Vec::with_capacity(tasks.len());
763 let mut estimated_batch_sizes = Vec::with_capacity(tasks.len());
764 for task in tasks {
765 let (stream, estimated_batch_size) = task.await.context(JoinSnafu)??;
766 range_streams.push(stream);
767 estimated_batch_sizes.push(estimated_batch_size);
768 }
769 let channel_size =
770 compute_parallel_channel_size(compute_average_batch_size(estimated_batch_sizes));
771 common_telemetry::debug!(
772 "SeriesDistributor built {} range_streams, region: {}, build cost: {:?}, channel_size: {}",
773 range_streams.len(),
774 self.stream_ctx.input.region_metadata().region_id,
775 build_start.elapsed(),
776 channel_size,
777 );
778
779 let mut reader = SeqScan::build_flat_reader_from_sources(
783 &self.stream_ctx,
784 range_streams,
785 self.final_merge_semaphore.clone(),
786 Some(&part_metrics),
787 true,
788 channel_size,
789 )
790 .await?;
791 let mut metrics = SeriesDistributorMetrics::default();
792
793 let mut divider = FlatSeriesBatchDivider::default();
794 while let Some(record_batch) = reader.try_next().await? {
795 metrics.scan_cost += fetch_start.elapsed();
796 metrics.num_batches += 1;
797 metrics.num_rows += record_batch.num_rows();
798
799 debug_assert!(record_batch.num_rows() > 0);
800 if record_batch.num_rows() == 0 {
801 fetch_start = Instant::now();
802 continue;
803 }
804
805 let divider_start = Instant::now();
807 let series_batch = divider.push(record_batch);
808 metrics.divider_cost += divider_start.elapsed();
809 if let Some(series_batch) = series_batch {
810 let yield_start = Instant::now();
811 self.senders
812 .send_batch(SeriesBatch::Flat(series_batch))
813 .await?;
814 metrics.yield_cost += yield_start.elapsed();
815 }
816 fetch_start = Instant::now();
817 }
818
819 let divider_start = Instant::now();
821 let series_batch = divider.finish();
822 metrics.divider_cost += divider_start.elapsed();
823 if let Some(series_batch) = series_batch {
824 let yield_start = Instant::now();
825 self.senders
826 .send_batch(SeriesBatch::Flat(series_batch))
827 .await?;
828 metrics.yield_cost += yield_start.elapsed();
829 }
830
831 metrics.scan_cost += fetch_start.elapsed();
832 metrics.num_series_send_timeout = self.senders.num_timeout;
833 metrics.num_series_send_full = self.senders.num_full;
834 part_metrics.set_distributor_metrics(&metrics);
835
836 part_metrics.on_finish();
837
838 Ok(())
839 }
840}
841
842struct SeriesCandidateDistributor {
844 stream_ctx: Arc<StreamContext>,
845 range_semaphore: Arc<Semaphore>,
846 partitions: Vec<Vec<PartitionRange>>,
847 pruner: Arc<Pruner>,
848 senders: CandidateSenderList,
849 active_receivers: Arc<AtomicUsize>,
850 metrics_set: ExecutionPlanMetricsSet,
851 metrics_list: Arc<PartitionMetricsList>,
852 explain_verbose: bool,
853}
854
855impl SeriesCandidateDistributor {
856 async fn execute(&mut self) {
857 if let Err(e) = self.distribute().await {
858 self.send_error(e);
859 }
860 }
861
862 async fn distribute(&mut self) -> Result<()> {
863 let part_metrics = new_partition_metrics(
864 &self.stream_ctx,
865 self.explain_verbose,
866 &self.metrics_set,
867 self.partitions.len(),
868 &self.metrics_list,
869 );
870 part_metrics.on_first_poll();
871
872 let candidate_scanner = SeriesCandidateScanner::try_new(
873 self.stream_ctx.clone(),
874 self.partitions.clone(),
875 self.pruner.clone(),
876 self.range_semaphore.clone(),
877 self.stream_ctx.input.scan_memory_pool.clone(),
878 self.metrics_set.clone(),
879 part_metrics.clone(),
880 )?;
881 let partition_pruner = candidate_scanner.partition_pruner();
882 let mut candidates = candidate_scanner.build_stream().await?;
883 let mut collector =
884 SeriesBatchCollector::new(self.partitions.len()).context(InvalidSenderSnafu)?;
885 let mut chunked = false;
886 while let Some(batch) = candidates.try_next().await? {
887 if !self.should_fetch_candidates() {
888 part_metrics.on_finish();
889 return Ok(());
890 }
891 collector.push(batch);
892 if collector.len() >= CANDIDATE_SERIES_ASSIGNMENT_THRESHOLD {
893 chunked = true;
894 self.send_assignments(collector.finish(false), &partition_pruner);
895 if !self.should_fetch_candidates() {
896 part_metrics.on_finish();
897 return Ok(());
898 }
899 collector =
900 SeriesBatchCollector::new(self.partitions.len()).context(InvalidSenderSnafu)?;
901 }
902 }
903
904 if collector.len() > 0 {
905 self.send_assignments(collector.finish(!chunked), &partition_pruner);
906 }
907 part_metrics.on_finish();
908 Ok(())
909 }
910
911 fn send_assignments(
912 &mut self,
913 assignments: Vec<AssignedSeriesBatch>,
914 partition_pruner: &Arc<PartitionPruner>,
915 ) {
916 for (partition, assigned_series) in assignments.into_iter().enumerate() {
917 if assigned_series.series().is_empty() {
918 continue;
919 }
920 let Some(sender) = self.senders[partition].as_ref() else {
921 continue;
922 };
923 let sent = sender
924 .send(Ok(SeriesReaderInput {
925 assigned_series,
926 partition_pruner: partition_pruner.clone(),
927 range_semaphore: self.range_semaphore.clone(),
928 }))
929 .is_ok();
930 if !sent {
931 self.senders[partition] = None;
932 }
933 }
934 }
935
936 fn should_fetch_candidates(&self) -> bool {
937 self.active_receivers.load(Ordering::Relaxed) > 0
938 }
939
940 fn send_error(&mut self, error: Error) {
941 let error = Arc::new(error);
942 for sender in self.senders.iter_mut().filter_map(Option::take) {
943 let result = Err(error.clone()).context(ScanSeriesSnafu);
944 let _ = sender.send(result);
945 }
946 }
947}
948
949#[derive(Debug)]
951pub enum SeriesBatch {
952 Flat(FlatSeriesBatch),
953}
954
955impl SeriesBatch {
956 pub fn num_batches(&self) -> usize {
958 match self {
959 SeriesBatch::Flat(flat_batch) => flat_batch.batches.len(),
960 }
961 }
962
963 pub fn num_rows(&self) -> usize {
965 match self {
966 SeriesBatch::Flat(flat_batch) => flat_batch.batches.iter().map(|x| x.num_rows()).sum(),
967 }
968 }
969}
970
971#[derive(Default, Debug)]
973pub struct FlatSeriesBatch {
974 pub batches: SmallVec<[RecordBatch; 4]>,
975}
976
977struct SenderList {
979 senders: Vec<Option<Sender<Result<SeriesBatch>>>>,
980 num_nones: usize,
982 sender_idx: usize,
984 num_timeout: usize,
986 num_full: usize,
988}
989
990impl SenderList {
991 fn new(senders: Vec<Option<Sender<Result<SeriesBatch>>>>) -> Self {
992 let num_nones = senders.iter().filter(|sender| sender.is_none()).count();
993 Self {
994 senders,
995 num_nones,
996 sender_idx: 0,
997 num_timeout: 0,
998 num_full: 0,
999 }
1000 }
1001
1002 fn try_send_batch(&mut self, mut batch: SeriesBatch) -> Result<Option<SeriesBatch>> {
1005 for _ in 0..self.senders.len() {
1006 ensure!(self.num_nones < self.senders.len(), InvalidSenderSnafu);
1007
1008 let sender_idx = self.fetch_add_sender_idx();
1009 let Some(sender) = &self.senders[sender_idx] else {
1010 continue;
1011 };
1012
1013 match sender.try_send(Ok(batch)) {
1014 Ok(()) => return Ok(None),
1015 Err(TrySendError::Full(res)) => {
1016 self.num_full += 1;
1017 batch = res.unwrap();
1019 }
1020 Err(TrySendError::Closed(res)) => {
1021 self.senders[sender_idx] = None;
1022 self.num_nones += 1;
1023 batch = res.unwrap();
1025 }
1026 }
1027 }
1028
1029 Ok(Some(batch))
1030 }
1031
1032 async fn send_batch(&mut self, mut batch: SeriesBatch) -> Result<()> {
1034 match self.try_send_batch(batch)? {
1036 Some(b) => {
1037 batch = b;
1039 }
1040 None => {
1041 return Ok(());
1042 }
1043 }
1044
1045 loop {
1046 ensure!(self.num_nones < self.senders.len(), InvalidSenderSnafu);
1047
1048 let sender_idx = self.fetch_add_sender_idx();
1049 let Some(sender) = &self.senders[sender_idx] else {
1050 continue;
1051 };
1052 match sender.send_timeout(Ok(batch), SEND_TIMEOUT).await {
1057 Ok(()) => break,
1058 Err(SendTimeoutError::Timeout(res)) => {
1059 self.num_timeout += 1;
1060 batch = res.unwrap();
1062 }
1063 Err(SendTimeoutError::Closed(res)) => {
1064 self.senders[sender_idx] = None;
1065 self.num_nones += 1;
1066 batch = res.unwrap();
1068 }
1069 }
1070 }
1071
1072 Ok(())
1073 }
1074
1075 async fn send_error(&self, error: Error) {
1076 let error = Arc::new(error);
1077 for sender in self.senders.iter().flatten() {
1078 let result = Err(error.clone()).context(ScanSeriesSnafu);
1079 let _ = sender.send(result).await;
1080 }
1081 }
1082
1083 fn fetch_add_sender_idx(&mut self) -> usize {
1084 let sender_idx = self.sender_idx;
1085 self.sender_idx = (self.sender_idx + 1) % self.senders.len();
1086 sender_idx
1087 }
1088}
1089
1090fn new_partition_metrics(
1091 stream_ctx: &StreamContext,
1092 explain_verbose: bool,
1093 metrics_set: &ExecutionPlanMetricsSet,
1094 partition: usize,
1095 metrics_list: &PartitionMetricsList,
1096) -> PartitionMetrics {
1097 let metrics = PartitionMetrics::new(
1098 stream_ctx.input.mapper.metadata().region_id,
1099 partition,
1100 "SeriesScan",
1101 stream_ctx.query_start,
1102 explain_verbose,
1103 metrics_set,
1104 );
1105
1106 metrics_list.set(partition, metrics.clone());
1107 metrics
1108}
1109
1110#[derive(Default)]
1115struct FlatSeriesBatchDivider {
1116 buffer: FlatSeriesBatch,
1117}
1118
1119impl FlatSeriesBatchDivider {
1120 fn push(&mut self, batch: RecordBatch) -> Option<FlatSeriesBatch> {
1124 if self.buffer.batches.is_empty() {
1126 self.buffer.batches.push(batch);
1127 return None;
1128 }
1129
1130 let pk_column_idx = primary_key_column_index(batch.num_columns());
1132 let batch_pk_column = batch.column(pk_column_idx);
1133 let batch_pk_array = batch_pk_column
1134 .as_any()
1135 .downcast_ref::<PrimaryKeyArray>()
1136 .unwrap();
1137 let batch_pk_values = batch_pk_array
1138 .values()
1139 .as_any()
1140 .downcast_ref::<BinaryArray>()
1141 .unwrap();
1142 let batch_last_pk =
1144 primary_key_at(batch_pk_array, batch_pk_values, batch_pk_array.len() - 1);
1145 let buffer_last_batch = self.buffer.batches.last().unwrap();
1148 let buffer_pk_column = buffer_last_batch.column(pk_column_idx);
1149 let buffer_pk_array = buffer_pk_column
1150 .as_any()
1151 .downcast_ref::<PrimaryKeyArray>()
1152 .unwrap();
1153 let buffer_pk_values = buffer_pk_array
1154 .values()
1155 .as_any()
1156 .downcast_ref::<BinaryArray>()
1157 .unwrap();
1158 let buffer_last_pk =
1159 primary_key_at(buffer_pk_array, buffer_pk_values, buffer_pk_array.len() - 1);
1160
1161 if batch_last_pk == buffer_last_pk {
1163 self.buffer.batches.push(batch);
1164 return None;
1165 }
1166 let batch_pk_keys = batch_pk_array.keys();
1169 let pk_indices = batch_pk_keys.values();
1170 let mut change_offset = 0;
1171 for (i, &key) in pk_indices.iter().enumerate() {
1172 let batch_pk = batch_pk_values.value(key as usize);
1173
1174 if buffer_last_pk != batch_pk {
1175 change_offset = i;
1176 break;
1177 }
1178 }
1179
1180 let (first_part, remaining_part) = if change_offset > 0 {
1182 let first_part = batch.slice(0, change_offset);
1183 let remaining_part = batch.slice(change_offset, batch.num_rows() - change_offset);
1184 (Some(first_part), Some(remaining_part))
1185 } else {
1186 (None, Some(batch))
1187 };
1188
1189 let mut result = std::mem::take(&mut self.buffer);
1191 if let Some(first_part) = first_part {
1192 result.batches.push(first_part);
1193 }
1194
1195 if let Some(remaining_part) = remaining_part {
1197 self.buffer.batches.push(remaining_part);
1198 }
1199
1200 Some(result)
1201 }
1202
1203 fn finish(&mut self) -> Option<FlatSeriesBatch> {
1205 if self.buffer.batches.is_empty() {
1206 None
1207 } else {
1208 Some(std::mem::take(&mut self.buffer))
1209 }
1210 }
1211}
1212
1213fn primary_key_at<'a>(
1215 primary_key: &PrimaryKeyArray,
1216 primary_key_values: &'a BinaryArray,
1217 index: usize,
1218) -> &'a [u8] {
1219 let key = primary_key.keys().value(index);
1220 primary_key_values.value(key as usize)
1221}
1222
1223#[cfg(test)]
1224mod tests {
1225 use std::sync::Arc;
1226
1227 use api::v1::OpType;
1228 use datatypes::arrow::array::{
1229 ArrayRef, BinaryDictionaryBuilder, Int64Array, StringDictionaryBuilder,
1230 TimestampMillisecondArray, UInt8Array, UInt64Array,
1231 };
1232 use datatypes::arrow::datatypes::{DataType, Field, Schema, SchemaRef, TimeUnit, UInt32Type};
1233 use datatypes::arrow::record_batch::RecordBatch;
1234
1235 use super::*;
1236 use crate::read::flat_projection::FlatProjectionMapper;
1237 use crate::read::scan_region::PredicateGroup;
1238 use crate::test_util::scheduler_util::SchedulerEnv;
1239 use crate::test_util::sst_util::sst_region_metadata_with_encoding;
1240
1241 #[tokio::test]
1242 async fn two_phase_eligibility_rejects_exact_sequence_range() {
1243 let env = SchedulerEnv::new().await;
1244 let metadata = Arc::new(sst_region_metadata_with_encoding(
1245 store_api::codec::PrimaryKeyEncoding::Sparse,
1246 ));
1247 let predicate = PredicateGroup::default();
1248
1249 let eligible = ScanInput::builder(
1250 env.access_layer.clone(),
1251 FlatProjectionMapper::new(&metadata, [0]).unwrap(),
1252 )
1253 .with_predicate(predicate.clone())
1254 .build();
1255 assert!(SeriesScan::supports_two_phase(&eligible));
1256
1257 let exact_sequence = ScanInput::builder(
1258 env.access_layer.clone(),
1259 FlatProjectionMapper::new(&metadata, [0]).unwrap(),
1260 )
1261 .with_predicate(predicate)
1262 .with_sequence_range(Some(store_api::storage::SequenceRange::GtLtEq {
1263 min: 1,
1264 max: 2,
1265 }))
1266 .build();
1267 assert!(!SeriesScan::supports_two_phase(&exact_sequence));
1268 }
1269
1270 #[test]
1271 fn candidate_distributor_stops_after_all_receivers_close() {
1272 let (_senders, mut receivers, active_receivers) = new_candidate_channel_list(2);
1273 assert_eq!(2, active_receivers.load(Ordering::Relaxed));
1274
1275 drop(receivers[0].take());
1276 assert_eq!(1, active_receivers.load(Ordering::Relaxed));
1277
1278 drop(receivers[1].take());
1279 assert_eq!(0, active_receivers.load(Ordering::Relaxed));
1280 }
1281
1282 fn new_test_record_batch(
1283 primary_keys: &[&[u8]],
1284 timestamps: &[i64],
1285 sequences: &[u64],
1286 op_types: &[OpType],
1287 fields: &[u64],
1288 ) -> RecordBatch {
1289 let num_rows = timestamps.len();
1290 debug_assert_eq!(sequences.len(), num_rows);
1291 debug_assert_eq!(op_types.len(), num_rows);
1292 debug_assert_eq!(fields.len(), num_rows);
1293 debug_assert_eq!(primary_keys.len(), num_rows);
1294
1295 let columns: Vec<ArrayRef> = vec![
1296 build_test_pk_string_dict_array(primary_keys),
1297 Arc::new(Int64Array::from_iter(
1298 fields.iter().map(|v| Some(*v as i64)),
1299 )),
1300 Arc::new(TimestampMillisecondArray::from_iter_values(
1301 timestamps.iter().copied(),
1302 )),
1303 build_test_pk_array(primary_keys),
1304 Arc::new(UInt64Array::from_iter_values(sequences.iter().copied())),
1305 Arc::new(UInt8Array::from_iter_values(
1306 op_types.iter().map(|v| *v as u8),
1307 )),
1308 ];
1309
1310 RecordBatch::try_new(build_test_flat_schema(), columns).unwrap()
1311 }
1312
1313 fn build_test_pk_string_dict_array(primary_keys: &[&[u8]]) -> ArrayRef {
1314 let mut builder = StringDictionaryBuilder::<UInt32Type>::new();
1315 for &pk in primary_keys {
1316 let pk_str = std::str::from_utf8(pk).unwrap();
1317 builder.append(pk_str).unwrap();
1318 }
1319 Arc::new(builder.finish())
1320 }
1321
1322 fn build_test_pk_array(primary_keys: &[&[u8]]) -> ArrayRef {
1323 let mut builder = BinaryDictionaryBuilder::<UInt32Type>::new();
1324 for &pk in primary_keys {
1325 builder.append(pk).unwrap();
1326 }
1327 Arc::new(builder.finish())
1328 }
1329
1330 fn build_test_flat_schema() -> SchemaRef {
1331 let fields = vec![
1332 Field::new(
1333 "k0",
1334 DataType::Dictionary(Box::new(DataType::UInt32), Box::new(DataType::Utf8)),
1335 false,
1336 ),
1337 Field::new("field0", DataType::Int64, true),
1338 Field::new(
1339 "ts",
1340 DataType::Timestamp(TimeUnit::Millisecond, None),
1341 false,
1342 ),
1343 Field::new(
1344 "__primary_key",
1345 DataType::Dictionary(Box::new(DataType::UInt32), Box::new(DataType::Binary)),
1346 false,
1347 ),
1348 Field::new("__sequence", DataType::UInt64, false),
1349 Field::new("__op_type", DataType::UInt8, false),
1350 ];
1351 Arc::new(Schema::new(fields))
1352 }
1353
1354 #[test]
1355 fn test_empty_buffer_first_push() {
1356 let mut divider = FlatSeriesBatchDivider::default();
1357 let result = divider.finish();
1358 assert!(result.is_none());
1359
1360 let mut divider = FlatSeriesBatchDivider::default();
1361 let batch = new_test_record_batch(
1362 &[b"series1", b"series1"],
1363 &[1000, 2000],
1364 &[1, 2],
1365 &[OpType::Put, OpType::Put],
1366 &[10, 20],
1367 );
1368 let result = divider.push(batch);
1369 assert!(result.is_none());
1370 assert_eq!(divider.buffer.batches.len(), 1);
1371 }
1372
1373 #[test]
1374 fn test_same_series_accumulation() {
1375 let mut divider = FlatSeriesBatchDivider::default();
1376
1377 let batch1 = new_test_record_batch(
1378 &[b"series1", b"series1"],
1379 &[1000, 2000],
1380 &[1, 2],
1381 &[OpType::Put, OpType::Put],
1382 &[10, 20],
1383 );
1384
1385 let batch2 = new_test_record_batch(
1386 &[b"series1", b"series1"],
1387 &[3000, 4000],
1388 &[3, 4],
1389 &[OpType::Put, OpType::Put],
1390 &[30, 40],
1391 );
1392
1393 divider.push(batch1);
1394 let result = divider.push(batch2);
1395 assert!(result.is_none());
1396 let series_batch = divider.finish().unwrap();
1397 assert_eq!(series_batch.batches.len(), 2);
1398 }
1399
1400 #[test]
1401 fn test_series_boundary_detection() {
1402 let mut divider = FlatSeriesBatchDivider::default();
1403
1404 let batch1 = new_test_record_batch(
1405 &[b"series1", b"series1"],
1406 &[1000, 2000],
1407 &[1, 2],
1408 &[OpType::Put, OpType::Put],
1409 &[10, 20],
1410 );
1411
1412 let batch2 = new_test_record_batch(
1413 &[b"series2", b"series2"],
1414 &[3000, 4000],
1415 &[3, 4],
1416 &[OpType::Put, OpType::Put],
1417 &[30, 40],
1418 );
1419
1420 divider.push(batch1);
1421 let series_batch = divider.push(batch2).unwrap();
1422 assert_eq!(series_batch.batches.len(), 1);
1423
1424 assert_eq!(divider.buffer.batches.len(), 1);
1425 }
1426
1427 #[test]
1428 fn test_series_boundary_within_batch() {
1429 let mut divider = FlatSeriesBatchDivider::default();
1430
1431 let batch1 = new_test_record_batch(
1432 &[b"series1", b"series1"],
1433 &[1000, 2000],
1434 &[1, 2],
1435 &[OpType::Put, OpType::Put],
1436 &[10, 20],
1437 );
1438
1439 let batch2 = new_test_record_batch(
1440 &[b"series1", b"series2"],
1441 &[3000, 4000],
1442 &[3, 4],
1443 &[OpType::Put, OpType::Put],
1444 &[30, 40],
1445 );
1446
1447 divider.push(batch1);
1448 let series_batch = divider.push(batch2).unwrap();
1449 assert_eq!(series_batch.batches.len(), 2);
1450 assert_eq!(series_batch.batches[0].num_rows(), 2);
1451 assert_eq!(series_batch.batches[1].num_rows(), 1);
1452
1453 assert_eq!(divider.buffer.batches.len(), 1);
1454 assert_eq!(divider.buffer.batches[0].num_rows(), 1);
1455 }
1456
1457 #[test]
1458 fn test_series_splitting() {
1459 let mut divider = FlatSeriesBatchDivider::default();
1460
1461 let batch1 = new_test_record_batch(&[b"series1"], &[1000], &[1], &[OpType::Put], &[10]);
1462
1463 let batch2 = new_test_record_batch(
1464 &[b"series1", b"series2", b"series2", b"series3"],
1465 &[2000, 3000, 4000, 5000],
1466 &[2, 3, 4, 5],
1467 &[OpType::Put, OpType::Put, OpType::Put, OpType::Put],
1468 &[20, 30, 40, 50],
1469 );
1470
1471 divider.push(batch1);
1472 let series_batch = divider.push(batch2).unwrap();
1473 assert_eq!(series_batch.batches.len(), 2);
1474
1475 let total_rows: usize = series_batch.batches.iter().map(|b| b.num_rows()).sum();
1476 assert_eq!(total_rows, 2);
1477
1478 let final_batch = divider.finish().unwrap();
1479 assert_eq!(final_batch.batches.len(), 1);
1480 assert_eq!(final_batch.batches[0].num_rows(), 3);
1481 }
1482}