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 !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 set_logical_region(&mut self, logical_region: bool) {
591 self.properties.set_logical_region(logical_region);
592 }
593
594 fn set_query_load_region_id(&mut self, region_id: store_api::storage::RegionId) {
595 self.properties.set_query_load_region_id(region_id);
596 }
597
598 fn snapshot_sequence(&self) -> Option<u64> {
599 self.stream_ctx.input.snapshot_sequence
600 }
601}
602
603impl DisplayAs for SeriesScan {
604 fn fmt_as(&self, t: DisplayFormatType, f: &mut fmt::Formatter) -> fmt::Result {
605 write!(
606 f,
607 "SeriesScan: region={}, ",
608 self.stream_ctx.input.mapper.metadata().region_id
609 )?;
610 match t {
611 DisplayFormatType::Default | DisplayFormatType::TreeRender => {
612 self.stream_ctx.format_for_explain(false, f)?;
613 }
614 DisplayFormatType::Verbose => {
615 self.stream_ctx.format_for_explain(true, f)?;
616 }
617 }
618 write!(f, ", \"mode\":\"{}\"", self.mode.as_str())?;
619 if matches!(t, DisplayFormatType::Verbose) {
620 self.metrics_list.format_verbose_metrics(f)?;
621 }
622 Ok(())
623 }
624}
625
626impl fmt::Debug for SeriesScan {
627 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
628 f.debug_struct("SeriesScan")
629 .field("mode", &self.mode)
630 .field("num_ranges", &self.stream_ctx.ranges.len())
631 .finish()
632 }
633}
634
635#[cfg(test)]
636impl SeriesScan {
637 pub(crate) fn input(&self) -> &ScanInput {
639 &self.stream_ctx.input
640 }
641
642 pub(crate) fn mode(&self) -> &'static str {
644 self.mode.as_str()
645 }
646}
647
648struct SeriesDistributor {
650 stream_ctx: Arc<StreamContext>,
652 range_semaphore: Option<Arc<Semaphore>>,
654 final_merge_semaphore: Option<Arc<Semaphore>>,
659 partitions: Vec<Vec<PartitionRange>>,
661 pruner: Arc<Pruner>,
663 senders: SenderList,
665 metrics_set: ExecutionPlanMetricsSet,
671 metrics_list: Arc<PartitionMetricsList>,
672 explain_verbose: bool,
674}
675
676impl SeriesDistributor {
677 #[tracing::instrument(
679 skip_all,
680 fields(region_id = %self.stream_ctx.input.mapper.metadata().region_id)
681 )]
682 async fn execute(&mut self) {
683 if let Err(e) = self.scan_partitions_flat().await {
684 self.senders.send_error(e).await;
685 }
686 }
687
688 #[tracing::instrument(
690 skip_all,
691 fields(region_id = %self.stream_ctx.input.mapper.metadata().region_id)
692 )]
693 async fn scan_partitions_flat(&mut self) -> Result<()> {
694 for partition_ranges in &self.partitions {
696 self.pruner.add_partition_ranges(partition_ranges);
697 }
698
699 let all_partition_ranges: Vec<_> = self.partitions.iter().flatten().cloned().collect();
701 let partition_pruner = Arc::new(PartitionPruner::new(
702 self.pruner.clone(),
703 &all_partition_ranges,
704 ));
705
706 let part_metrics = new_partition_metrics(
707 &self.stream_ctx,
708 self.explain_verbose,
709 &self.metrics_set,
710 self.partitions.len(),
711 &self.metrics_list,
712 );
713 part_metrics.on_first_poll();
714 let mut fetch_start = Instant::now();
717
718 let build_start = Instant::now();
720 let mut tasks = Vec::new();
721 for partition in &self.partitions {
722 for part_range in partition {
723 let stream_ctx = self.stream_ctx.clone();
724 let part_range = *part_range;
725 let part_metrics = part_metrics.clone();
726 let partition_pruner = partition_pruner.clone();
727 let file_scan_semaphore = self.range_semaphore.clone();
728 let merge_semaphore = self.range_semaphore.clone();
729 tasks.push(common_runtime::spawn_query(async move {
730 SeqScan::build_flat_partition_range_read(
731 &stream_ctx,
732 &part_range,
733 false,
734 &part_metrics,
735 partition_pruner,
736 file_scan_semaphore,
737 merge_semaphore,
738 )
739 .await
740 }));
741 }
742 }
743 let mut range_streams = Vec::with_capacity(tasks.len());
744 let mut estimated_batch_sizes = Vec::with_capacity(tasks.len());
745 for task in tasks {
746 let (stream, estimated_batch_size) = task.await.context(JoinSnafu)??;
747 range_streams.push(stream);
748 estimated_batch_sizes.push(estimated_batch_size);
749 }
750 let channel_size =
751 compute_parallel_channel_size(compute_average_batch_size(estimated_batch_sizes));
752 common_telemetry::debug!(
753 "SeriesDistributor built {} range_streams, region: {}, build cost: {:?}, channel_size: {}",
754 range_streams.len(),
755 self.stream_ctx.input.region_metadata().region_id,
756 build_start.elapsed(),
757 channel_size,
758 );
759
760 let mut reader = SeqScan::build_flat_reader_from_sources(
764 &self.stream_ctx,
765 range_streams,
766 self.final_merge_semaphore.clone(),
767 Some(&part_metrics),
768 true,
769 channel_size,
770 )
771 .await?;
772 let mut metrics = SeriesDistributorMetrics::default();
773
774 let mut divider = FlatSeriesBatchDivider::default();
775 while let Some(record_batch) = reader.try_next().await? {
776 metrics.scan_cost += fetch_start.elapsed();
777 metrics.num_batches += 1;
778 metrics.num_rows += record_batch.num_rows();
779
780 debug_assert!(record_batch.num_rows() > 0);
781 if record_batch.num_rows() == 0 {
782 fetch_start = Instant::now();
783 continue;
784 }
785
786 let divider_start = Instant::now();
788 let series_batch = divider.push(record_batch);
789 metrics.divider_cost += divider_start.elapsed();
790 if let Some(series_batch) = series_batch {
791 let yield_start = Instant::now();
792 self.senders
793 .send_batch(SeriesBatch::Flat(series_batch))
794 .await?;
795 metrics.yield_cost += yield_start.elapsed();
796 }
797 fetch_start = Instant::now();
798 }
799
800 let divider_start = Instant::now();
802 let series_batch = divider.finish();
803 metrics.divider_cost += divider_start.elapsed();
804 if let Some(series_batch) = series_batch {
805 let yield_start = Instant::now();
806 self.senders
807 .send_batch(SeriesBatch::Flat(series_batch))
808 .await?;
809 metrics.yield_cost += yield_start.elapsed();
810 }
811
812 metrics.scan_cost += fetch_start.elapsed();
813 metrics.num_series_send_timeout = self.senders.num_timeout;
814 metrics.num_series_send_full = self.senders.num_full;
815 part_metrics.set_distributor_metrics(&metrics);
816
817 part_metrics.on_finish();
818
819 Ok(())
820 }
821}
822
823struct SeriesCandidateDistributor {
825 stream_ctx: Arc<StreamContext>,
826 range_semaphore: Arc<Semaphore>,
827 partitions: Vec<Vec<PartitionRange>>,
828 pruner: Arc<Pruner>,
829 senders: CandidateSenderList,
830 active_receivers: Arc<AtomicUsize>,
831 metrics_set: ExecutionPlanMetricsSet,
832 metrics_list: Arc<PartitionMetricsList>,
833 explain_verbose: bool,
834}
835
836impl SeriesCandidateDistributor {
837 async fn execute(&mut self) {
838 if let Err(e) = self.distribute().await {
839 self.send_error(e);
840 }
841 }
842
843 async fn distribute(&mut self) -> Result<()> {
844 let part_metrics = new_partition_metrics(
845 &self.stream_ctx,
846 self.explain_verbose,
847 &self.metrics_set,
848 self.partitions.len(),
849 &self.metrics_list,
850 );
851 part_metrics.on_first_poll();
852
853 let candidate_scanner = SeriesCandidateScanner::try_new(
854 self.stream_ctx.clone(),
855 self.partitions.clone(),
856 self.pruner.clone(),
857 self.range_semaphore.clone(),
858 self.stream_ctx.input.scan_memory_pool.clone(),
859 self.metrics_set.clone(),
860 part_metrics.clone(),
861 )?;
862 let partition_pruner = candidate_scanner.partition_pruner();
863 let mut candidates = candidate_scanner.build_stream().await?;
864 let mut collector =
865 SeriesBatchCollector::new(self.partitions.len()).context(InvalidSenderSnafu)?;
866 let mut chunked = false;
867 while let Some(batch) = candidates.try_next().await? {
868 if !self.should_fetch_candidates() {
869 part_metrics.on_finish();
870 return Ok(());
871 }
872 collector.push(batch);
873 if collector.len() >= CANDIDATE_SERIES_ASSIGNMENT_THRESHOLD {
874 chunked = true;
875 self.send_assignments(collector.finish(false), &partition_pruner);
876 if !self.should_fetch_candidates() {
877 part_metrics.on_finish();
878 return Ok(());
879 }
880 collector =
881 SeriesBatchCollector::new(self.partitions.len()).context(InvalidSenderSnafu)?;
882 }
883 }
884
885 if collector.len() > 0 {
886 self.send_assignments(collector.finish(!chunked), &partition_pruner);
887 }
888 part_metrics.on_finish();
889 Ok(())
890 }
891
892 fn send_assignments(
893 &mut self,
894 assignments: Vec<AssignedSeriesBatch>,
895 partition_pruner: &Arc<PartitionPruner>,
896 ) {
897 for (partition, assigned_series) in assignments.into_iter().enumerate() {
898 if assigned_series.series().is_empty() {
899 continue;
900 }
901 let Some(sender) = self.senders[partition].as_ref() else {
902 continue;
903 };
904 let sent = sender
905 .send(Ok(SeriesReaderInput {
906 assigned_series,
907 partition_pruner: partition_pruner.clone(),
908 range_semaphore: self.range_semaphore.clone(),
909 }))
910 .is_ok();
911 if !sent {
912 self.senders[partition] = None;
913 }
914 }
915 }
916
917 fn should_fetch_candidates(&self) -> bool {
918 self.active_receivers.load(Ordering::Relaxed) > 0
919 }
920
921 fn send_error(&mut self, error: Error) {
922 let error = Arc::new(error);
923 for sender in self.senders.iter_mut().filter_map(Option::take) {
924 let result = Err(error.clone()).context(ScanSeriesSnafu);
925 let _ = sender.send(result);
926 }
927 }
928}
929
930#[derive(Debug)]
932pub enum SeriesBatch {
933 Flat(FlatSeriesBatch),
934}
935
936impl SeriesBatch {
937 pub fn num_batches(&self) -> usize {
939 match self {
940 SeriesBatch::Flat(flat_batch) => flat_batch.batches.len(),
941 }
942 }
943
944 pub fn num_rows(&self) -> usize {
946 match self {
947 SeriesBatch::Flat(flat_batch) => flat_batch.batches.iter().map(|x| x.num_rows()).sum(),
948 }
949 }
950}
951
952#[derive(Default, Debug)]
954pub struct FlatSeriesBatch {
955 pub batches: SmallVec<[RecordBatch; 4]>,
956}
957
958struct SenderList {
960 senders: Vec<Option<Sender<Result<SeriesBatch>>>>,
961 num_nones: usize,
963 sender_idx: usize,
965 num_timeout: usize,
967 num_full: usize,
969}
970
971impl SenderList {
972 fn new(senders: Vec<Option<Sender<Result<SeriesBatch>>>>) -> Self {
973 let num_nones = senders.iter().filter(|sender| sender.is_none()).count();
974 Self {
975 senders,
976 num_nones,
977 sender_idx: 0,
978 num_timeout: 0,
979 num_full: 0,
980 }
981 }
982
983 fn try_send_batch(&mut self, mut batch: SeriesBatch) -> Result<Option<SeriesBatch>> {
986 for _ in 0..self.senders.len() {
987 ensure!(self.num_nones < self.senders.len(), InvalidSenderSnafu);
988
989 let sender_idx = self.fetch_add_sender_idx();
990 let Some(sender) = &self.senders[sender_idx] else {
991 continue;
992 };
993
994 match sender.try_send(Ok(batch)) {
995 Ok(()) => return Ok(None),
996 Err(TrySendError::Full(res)) => {
997 self.num_full += 1;
998 batch = res.unwrap();
1000 }
1001 Err(TrySendError::Closed(res)) => {
1002 self.senders[sender_idx] = None;
1003 self.num_nones += 1;
1004 batch = res.unwrap();
1006 }
1007 }
1008 }
1009
1010 Ok(Some(batch))
1011 }
1012
1013 async fn send_batch(&mut self, mut batch: SeriesBatch) -> Result<()> {
1015 match self.try_send_batch(batch)? {
1017 Some(b) => {
1018 batch = b;
1020 }
1021 None => {
1022 return Ok(());
1023 }
1024 }
1025
1026 loop {
1027 ensure!(self.num_nones < self.senders.len(), InvalidSenderSnafu);
1028
1029 let sender_idx = self.fetch_add_sender_idx();
1030 let Some(sender) = &self.senders[sender_idx] else {
1031 continue;
1032 };
1033 match sender.send_timeout(Ok(batch), SEND_TIMEOUT).await {
1038 Ok(()) => break,
1039 Err(SendTimeoutError::Timeout(res)) => {
1040 self.num_timeout += 1;
1041 batch = res.unwrap();
1043 }
1044 Err(SendTimeoutError::Closed(res)) => {
1045 self.senders[sender_idx] = None;
1046 self.num_nones += 1;
1047 batch = res.unwrap();
1049 }
1050 }
1051 }
1052
1053 Ok(())
1054 }
1055
1056 async fn send_error(&self, error: Error) {
1057 let error = Arc::new(error);
1058 for sender in self.senders.iter().flatten() {
1059 let result = Err(error.clone()).context(ScanSeriesSnafu);
1060 let _ = sender.send(result).await;
1061 }
1062 }
1063
1064 fn fetch_add_sender_idx(&mut self) -> usize {
1065 let sender_idx = self.sender_idx;
1066 self.sender_idx = (self.sender_idx + 1) % self.senders.len();
1067 sender_idx
1068 }
1069}
1070
1071fn new_partition_metrics(
1072 stream_ctx: &StreamContext,
1073 explain_verbose: bool,
1074 metrics_set: &ExecutionPlanMetricsSet,
1075 partition: usize,
1076 metrics_list: &PartitionMetricsList,
1077) -> PartitionMetrics {
1078 let metrics = PartitionMetrics::new(
1079 stream_ctx.input.mapper.metadata().region_id,
1080 partition,
1081 "SeriesScan",
1082 stream_ctx.query_start,
1083 explain_verbose,
1084 metrics_set,
1085 );
1086
1087 metrics_list.set(partition, metrics.clone());
1088 metrics
1089}
1090
1091#[derive(Default)]
1096struct FlatSeriesBatchDivider {
1097 buffer: FlatSeriesBatch,
1098}
1099
1100impl FlatSeriesBatchDivider {
1101 fn push(&mut self, batch: RecordBatch) -> Option<FlatSeriesBatch> {
1105 if self.buffer.batches.is_empty() {
1107 self.buffer.batches.push(batch);
1108 return None;
1109 }
1110
1111 let pk_column_idx = primary_key_column_index(batch.num_columns());
1113 let batch_pk_column = batch.column(pk_column_idx);
1114 let batch_pk_array = batch_pk_column
1115 .as_any()
1116 .downcast_ref::<PrimaryKeyArray>()
1117 .unwrap();
1118 let batch_pk_values = batch_pk_array
1119 .values()
1120 .as_any()
1121 .downcast_ref::<BinaryArray>()
1122 .unwrap();
1123 let batch_last_pk =
1125 primary_key_at(batch_pk_array, batch_pk_values, batch_pk_array.len() - 1);
1126 let buffer_last_batch = self.buffer.batches.last().unwrap();
1129 let buffer_pk_column = buffer_last_batch.column(pk_column_idx);
1130 let buffer_pk_array = buffer_pk_column
1131 .as_any()
1132 .downcast_ref::<PrimaryKeyArray>()
1133 .unwrap();
1134 let buffer_pk_values = buffer_pk_array
1135 .values()
1136 .as_any()
1137 .downcast_ref::<BinaryArray>()
1138 .unwrap();
1139 let buffer_last_pk =
1140 primary_key_at(buffer_pk_array, buffer_pk_values, buffer_pk_array.len() - 1);
1141
1142 if batch_last_pk == buffer_last_pk {
1144 self.buffer.batches.push(batch);
1145 return None;
1146 }
1147 let batch_pk_keys = batch_pk_array.keys();
1150 let pk_indices = batch_pk_keys.values();
1151 let mut change_offset = 0;
1152 for (i, &key) in pk_indices.iter().enumerate() {
1153 let batch_pk = batch_pk_values.value(key as usize);
1154
1155 if buffer_last_pk != batch_pk {
1156 change_offset = i;
1157 break;
1158 }
1159 }
1160
1161 let (first_part, remaining_part) = if change_offset > 0 {
1163 let first_part = batch.slice(0, change_offset);
1164 let remaining_part = batch.slice(change_offset, batch.num_rows() - change_offset);
1165 (Some(first_part), Some(remaining_part))
1166 } else {
1167 (None, Some(batch))
1168 };
1169
1170 let mut result = std::mem::take(&mut self.buffer);
1172 if let Some(first_part) = first_part {
1173 result.batches.push(first_part);
1174 }
1175
1176 if let Some(remaining_part) = remaining_part {
1178 self.buffer.batches.push(remaining_part);
1179 }
1180
1181 Some(result)
1182 }
1183
1184 fn finish(&mut self) -> Option<FlatSeriesBatch> {
1186 if self.buffer.batches.is_empty() {
1187 None
1188 } else {
1189 Some(std::mem::take(&mut self.buffer))
1190 }
1191 }
1192}
1193
1194fn primary_key_at<'a>(
1196 primary_key: &PrimaryKeyArray,
1197 primary_key_values: &'a BinaryArray,
1198 index: usize,
1199) -> &'a [u8] {
1200 let key = primary_key.keys().value(index);
1201 primary_key_values.value(key as usize)
1202}
1203
1204#[cfg(test)]
1205mod tests {
1206 use std::sync::Arc;
1207
1208 use api::v1::OpType;
1209 use datatypes::arrow::array::{
1210 ArrayRef, BinaryDictionaryBuilder, Int64Array, StringDictionaryBuilder,
1211 TimestampMillisecondArray, UInt8Array, UInt64Array,
1212 };
1213 use datatypes::arrow::datatypes::{DataType, Field, Schema, SchemaRef, TimeUnit, UInt32Type};
1214 use datatypes::arrow::record_batch::RecordBatch;
1215
1216 use super::*;
1217
1218 #[test]
1219 fn candidate_distributor_stops_after_all_receivers_close() {
1220 let (_senders, mut receivers, active_receivers) = new_candidate_channel_list(2);
1221 assert_eq!(2, active_receivers.load(Ordering::Relaxed));
1222
1223 drop(receivers[0].take());
1224 assert_eq!(1, active_receivers.load(Ordering::Relaxed));
1225
1226 drop(receivers[1].take());
1227 assert_eq!(0, active_receivers.load(Ordering::Relaxed));
1228 }
1229
1230 fn new_test_record_batch(
1231 primary_keys: &[&[u8]],
1232 timestamps: &[i64],
1233 sequences: &[u64],
1234 op_types: &[OpType],
1235 fields: &[u64],
1236 ) -> RecordBatch {
1237 let num_rows = timestamps.len();
1238 debug_assert_eq!(sequences.len(), num_rows);
1239 debug_assert_eq!(op_types.len(), num_rows);
1240 debug_assert_eq!(fields.len(), num_rows);
1241 debug_assert_eq!(primary_keys.len(), num_rows);
1242
1243 let columns: Vec<ArrayRef> = vec![
1244 build_test_pk_string_dict_array(primary_keys),
1245 Arc::new(Int64Array::from_iter(
1246 fields.iter().map(|v| Some(*v as i64)),
1247 )),
1248 Arc::new(TimestampMillisecondArray::from_iter_values(
1249 timestamps.iter().copied(),
1250 )),
1251 build_test_pk_array(primary_keys),
1252 Arc::new(UInt64Array::from_iter_values(sequences.iter().copied())),
1253 Arc::new(UInt8Array::from_iter_values(
1254 op_types.iter().map(|v| *v as u8),
1255 )),
1256 ];
1257
1258 RecordBatch::try_new(build_test_flat_schema(), columns).unwrap()
1259 }
1260
1261 fn build_test_pk_string_dict_array(primary_keys: &[&[u8]]) -> ArrayRef {
1262 let mut builder = StringDictionaryBuilder::<UInt32Type>::new();
1263 for &pk in primary_keys {
1264 let pk_str = std::str::from_utf8(pk).unwrap();
1265 builder.append(pk_str).unwrap();
1266 }
1267 Arc::new(builder.finish())
1268 }
1269
1270 fn build_test_pk_array(primary_keys: &[&[u8]]) -> ArrayRef {
1271 let mut builder = BinaryDictionaryBuilder::<UInt32Type>::new();
1272 for &pk in primary_keys {
1273 builder.append(pk).unwrap();
1274 }
1275 Arc::new(builder.finish())
1276 }
1277
1278 fn build_test_flat_schema() -> SchemaRef {
1279 let fields = vec![
1280 Field::new(
1281 "k0",
1282 DataType::Dictionary(Box::new(DataType::UInt32), Box::new(DataType::Utf8)),
1283 false,
1284 ),
1285 Field::new("field0", DataType::Int64, true),
1286 Field::new(
1287 "ts",
1288 DataType::Timestamp(TimeUnit::Millisecond, None),
1289 false,
1290 ),
1291 Field::new(
1292 "__primary_key",
1293 DataType::Dictionary(Box::new(DataType::UInt32), Box::new(DataType::Binary)),
1294 false,
1295 ),
1296 Field::new("__sequence", DataType::UInt64, false),
1297 Field::new("__op_type", DataType::UInt8, false),
1298 ];
1299 Arc::new(Schema::new(fields))
1300 }
1301
1302 #[test]
1303 fn test_empty_buffer_first_push() {
1304 let mut divider = FlatSeriesBatchDivider::default();
1305 let result = divider.finish();
1306 assert!(result.is_none());
1307
1308 let mut divider = FlatSeriesBatchDivider::default();
1309 let batch = new_test_record_batch(
1310 &[b"series1", b"series1"],
1311 &[1000, 2000],
1312 &[1, 2],
1313 &[OpType::Put, OpType::Put],
1314 &[10, 20],
1315 );
1316 let result = divider.push(batch);
1317 assert!(result.is_none());
1318 assert_eq!(divider.buffer.batches.len(), 1);
1319 }
1320
1321 #[test]
1322 fn test_same_series_accumulation() {
1323 let mut divider = FlatSeriesBatchDivider::default();
1324
1325 let batch1 = new_test_record_batch(
1326 &[b"series1", b"series1"],
1327 &[1000, 2000],
1328 &[1, 2],
1329 &[OpType::Put, OpType::Put],
1330 &[10, 20],
1331 );
1332
1333 let batch2 = new_test_record_batch(
1334 &[b"series1", b"series1"],
1335 &[3000, 4000],
1336 &[3, 4],
1337 &[OpType::Put, OpType::Put],
1338 &[30, 40],
1339 );
1340
1341 divider.push(batch1);
1342 let result = divider.push(batch2);
1343 assert!(result.is_none());
1344 let series_batch = divider.finish().unwrap();
1345 assert_eq!(series_batch.batches.len(), 2);
1346 }
1347
1348 #[test]
1349 fn test_series_boundary_detection() {
1350 let mut divider = FlatSeriesBatchDivider::default();
1351
1352 let batch1 = new_test_record_batch(
1353 &[b"series1", b"series1"],
1354 &[1000, 2000],
1355 &[1, 2],
1356 &[OpType::Put, OpType::Put],
1357 &[10, 20],
1358 );
1359
1360 let batch2 = new_test_record_batch(
1361 &[b"series2", b"series2"],
1362 &[3000, 4000],
1363 &[3, 4],
1364 &[OpType::Put, OpType::Put],
1365 &[30, 40],
1366 );
1367
1368 divider.push(batch1);
1369 let series_batch = divider.push(batch2).unwrap();
1370 assert_eq!(series_batch.batches.len(), 1);
1371
1372 assert_eq!(divider.buffer.batches.len(), 1);
1373 }
1374
1375 #[test]
1376 fn test_series_boundary_within_batch() {
1377 let mut divider = FlatSeriesBatchDivider::default();
1378
1379 let batch1 = new_test_record_batch(
1380 &[b"series1", b"series1"],
1381 &[1000, 2000],
1382 &[1, 2],
1383 &[OpType::Put, OpType::Put],
1384 &[10, 20],
1385 );
1386
1387 let batch2 = new_test_record_batch(
1388 &[b"series1", b"series2"],
1389 &[3000, 4000],
1390 &[3, 4],
1391 &[OpType::Put, OpType::Put],
1392 &[30, 40],
1393 );
1394
1395 divider.push(batch1);
1396 let series_batch = divider.push(batch2).unwrap();
1397 assert_eq!(series_batch.batches.len(), 2);
1398 assert_eq!(series_batch.batches[0].num_rows(), 2);
1399 assert_eq!(series_batch.batches[1].num_rows(), 1);
1400
1401 assert_eq!(divider.buffer.batches.len(), 1);
1402 assert_eq!(divider.buffer.batches[0].num_rows(), 1);
1403 }
1404
1405 #[test]
1406 fn test_series_splitting() {
1407 let mut divider = FlatSeriesBatchDivider::default();
1408
1409 let batch1 = new_test_record_batch(&[b"series1"], &[1000], &[1], &[OpType::Put], &[10]);
1410
1411 let batch2 = new_test_record_batch(
1412 &[b"series1", b"series2", b"series2", b"series3"],
1413 &[2000, 3000, 4000, 5000],
1414 &[2, 3, 4, 5],
1415 &[OpType::Put, OpType::Put, OpType::Put, OpType::Put],
1416 &[20, 30, 40, 50],
1417 );
1418
1419 divider.push(batch1);
1420 let series_batch = divider.push(batch2).unwrap();
1421 assert_eq!(series_batch.batches.len(), 2);
1422
1423 let total_rows: usize = series_batch.batches.iter().map(|b| b.num_rows()).sum();
1424 assert_eq!(total_rows, 2);
1425
1426 let final_batch = divider.finish().unwrap();
1427 assert_eq!(final_batch.batches.len(), 1);
1428 assert_eq!(final_batch.batches[0].num_rows(), 3);
1429 }
1430}