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