Skip to main content

mito2/read/
series_scan.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15//! Per-series scan implementation.
16
17use 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
62/// Timeout to send a batch to a sender.
63const SEND_TIMEOUT: Duration = Duration::from_micros(100);
64
65/// Maximum number of candidate series retained before distributing assignments.
66const CANDIDATE_SERIES_ASSIGNMENT_THRESHOLD: usize = 1_000_000;
67
68/// Legacy data receivers for output partitions.
69type LegacyReceiverList = Vec<Option<Receiver<Result<SeriesBatch>>>>;
70
71/// Candidate assignment receivers for output partitions.
72type CandidateReceiverList = Vec<Option<CandidateReceiver>>;
73
74/// Candidate assignment senders for output partitions.
75type 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
95/// Input required by a partition-local series reader.
96struct 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
117/// Scans a region and returns sorted rows of a series in the same partition.
118///
119/// The output order is always order by `(primary key, time index)` inside every
120/// partition.
121/// Always returns the same series (primary key) to the same partition.
122pub struct SeriesScan {
123    /// Implementation used by this scan.
124    mode: SeriesScanMode,
125    /// Properties of the scanner.
126    properties: ScannerProperties,
127    /// Context of streams.
128    stream_ctx: Arc<StreamContext>,
129    /// Shared pruner for file range building.
130    pruner: Arc<Pruner>,
131    /// Legacy data receivers for each partition.
132    legacy_receivers: Mutex<LegacyReceiverList>,
133    /// Candidate assignment receivers for each partition.
134    candidate_receivers: Mutex<CandidateReceiverList>,
135    /// Metrics for each partition.
136    /// The scanner only sets in query and keeps it empty during compaction.
137    metrics_list: Arc<PartitionMetricsList>,
138}
139
140impl SeriesScan {
141    /// Creates a new [SeriesScan].
142    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        // Create the shared pruner with number of workers equal to CPU cores.
158        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        // Safety: `scan_batch_in_partition` validates the index, and the receiver
320        // list always contains one entry for each output partition.
321        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    /// Scans the region and returns a stream.
442    #[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    /// Scan [`Batch`] in all partitions one by one.
457    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    /// Checks resource limit for the scanner.
483    pub(crate) fn check_scan_limit(&self) -> Result<()> {
484        // Sum the total number of files across all partitions
485        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            // Partition streams can be consumed sequentially. A bounded channel for an unpolled
526            // partition could block the distributor and prevent it from sending more assignments
527            // to the partition currently being consumed, causing a deadlock.
528            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    /// Returns the input.
657    pub(crate) fn input(&self) -> &ScanInput {
658        &self.stream_ctx.input
659    }
660
661    /// Returns the scan mode for tests.
662    pub(crate) fn mode(&self) -> &'static str {
663        self.mode.as_str()
664    }
665}
666
667/// The distributor scans series and distributes them to different partitions.
668struct SeriesDistributor {
669    /// Context for the scan stream.
670    stream_ctx: Arc<StreamContext>,
671    /// Semaphore for file scanning and range-level merging.
672    range_semaphore: Option<Arc<Semaphore>>,
673    /// Semaphore for the final merge across all range streams.
674    /// Must be separate from `range_semaphore` to avoid deadlock: final merge tasks
675    /// hold a permit while waiting for data from range-level merge tasks, which also
676    /// need permits to produce data.
677    final_merge_semaphore: Option<Arc<Semaphore>>,
678    /// Partition ranges to scan.
679    partitions: Vec<Vec<PartitionRange>>,
680    /// Shared pruner for file range building.
681    pruner: Arc<Pruner>,
682    /// Senders of all partitions.
683    senders: SenderList,
684    /// Metrics set to report.
685    /// The distributor report the metrics as an additional partition.
686    /// This may double the scan cost of the [SeriesScan] metrics. We can
687    /// get per-partition metrics in verbose mode to see the metrics of the
688    /// distributor.
689    metrics_set: ExecutionPlanMetricsSet,
690    metrics_list: Arc<PartitionMetricsList>,
691    /// Whether to use verbose logging and collect detailed metrics.
692    explain_verbose: bool,
693}
694
695impl SeriesDistributor {
696    /// Executes the distributor.
697    #[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    /// Scans all parts in flat format using FlatSeriesBatchDivider.
708    #[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        // Initialize reference counts for all partition ranges.
714        for partition_ranges in &self.partitions {
715            self.pruner.add_partition_ranges(partition_ranges);
716        }
717
718        // Create PartitionPruner covering all partitions
719        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        // Start fetch time before building sources so scan cost contains
734        // build part cost.
735        let mut fetch_start = Instant::now();
736
737        // Builds one deduped stream per partition range, then merges across ranges.
738        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        // Each partition range stream is already deduped, so skip dedup here.
780        // Use a separate semaphore for the final merge to avoid deadlock with
781        // range-level merge tasks that share the range_semaphore.
782        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            // Use divider to split series
806            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        // Send any remaining batch in the divider
820        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
842/// Discovers metric series and sends one reader assignment to each partition.
843struct 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/// Batches of the same series.
950#[derive(Debug)]
951pub enum SeriesBatch {
952    Flat(FlatSeriesBatch),
953}
954
955impl SeriesBatch {
956    /// Returns the number of batches.
957    pub fn num_batches(&self) -> usize {
958        match self {
959            SeriesBatch::Flat(flat_batch) => flat_batch.batches.len(),
960        }
961    }
962
963    /// Returns the total number of rows across all batches.
964    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/// Batches of the same series in flat format.
972#[derive(Default, Debug)]
973pub struct FlatSeriesBatch {
974    pub batches: SmallVec<[RecordBatch; 4]>,
975}
976
977/// List of senders.
978struct SenderList {
979    senders: Vec<Option<Sender<Result<SeriesBatch>>>>,
980    /// Number of None senders.
981    num_nones: usize,
982    /// Index of the current partition to send.
983    sender_idx: usize,
984    /// Number of timeout.
985    num_timeout: usize,
986    /// Number of full senders.
987    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    /// Finds a partition and tries to send the batch to the partition.
1003    /// Returns None if it sends successfully.
1004    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                    // Safety: we send Ok.
1018                    batch = res.unwrap();
1019                }
1020                Err(TrySendError::Closed(res)) => {
1021                    self.senders[sender_idx] = None;
1022                    self.num_nones += 1;
1023                    // Safety: we send Ok.
1024                    batch = res.unwrap();
1025                }
1026            }
1027        }
1028
1029        Ok(Some(batch))
1030    }
1031
1032    /// Finds a partition and sends the batch to the partition.
1033    async fn send_batch(&mut self, mut batch: SeriesBatch) -> Result<()> {
1034        // Sends the batch without blocking first.
1035        match self.try_send_batch(batch)? {
1036            Some(b) => {
1037                // Unable to send batch to partition.
1038                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            // Adds a timeout to avoid blocking indefinitely and sending
1053            // the batch in a round-robin fashion when some partitions
1054            // don't poll their inputs. This may happen if we have a
1055            // node like sort merging. But it is rare when we are using SeriesScan.
1056            match sender.send_timeout(Ok(batch), SEND_TIMEOUT).await {
1057                Ok(()) => break,
1058                Err(SendTimeoutError::Timeout(res)) => {
1059                    self.num_timeout += 1;
1060                    // Safety: we send Ok.
1061                    batch = res.unwrap();
1062                }
1063                Err(SendTimeoutError::Closed(res)) => {
1064                    self.senders[sender_idx] = None;
1065                    self.num_nones += 1;
1066                    // Safety: we send Ok.
1067                    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/// A divider to split flat record batches by time series.
1111///
1112/// It only ensures rows of the same series are returned in the same [FlatSeriesBatch].
1113/// However, a [FlatSeriesBatch] may contain rows from multiple series.
1114#[derive(Default)]
1115struct FlatSeriesBatchDivider {
1116    buffer: FlatSeriesBatch,
1117}
1118
1119impl FlatSeriesBatchDivider {
1120    /// Pushes a record batch into the divider.
1121    ///
1122    /// Returns a [FlatSeriesBatch] if we ensure the batch contains all rows of the series in it.
1123    fn push(&mut self, batch: RecordBatch) -> Option<FlatSeriesBatch> {
1124        // If buffer is empty
1125        if self.buffer.batches.is_empty() {
1126            self.buffer.batches.push(batch);
1127            return None;
1128        }
1129
1130        // Gets the primary key column from the incoming batch.
1131        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        // Gets the last primary key of the incoming batch.
1143        let batch_last_pk =
1144            primary_key_at(batch_pk_array, batch_pk_values, batch_pk_array.len() - 1);
1145        // Gets the last primary key of the buffer.
1146        // Safety: the buffer is not empty.
1147        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 last primary key in the batch is the same as last primary key in the buffer.
1162        if batch_last_pk == buffer_last_pk {
1163            self.buffer.batches.push(batch);
1164            return None;
1165        }
1166        // Otherwise, the batch must have a different primary key, we find the first offset of the
1167        // changed primary key.
1168        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        // Splits the batch at the change offset
1181        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        // Creates the result from current buffer + first part of new batch
1190        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        // Pushes remaining part to the buffer if it exists
1196        if let Some(remaining_part) = remaining_part {
1197            self.buffer.batches.push(remaining_part);
1198        }
1199
1200        Some(result)
1201    }
1202
1203    /// Returns the final [FlatSeriesBatch].
1204    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
1213/// Helper function to extract primary key bytes at a specific index from [PrimaryKeyArray].
1214fn 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}