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