Skip to main content

mito2/memtable/
time_series.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
15use std::collections::btree_map::Entry;
16use std::collections::{BTreeMap, Bound, HashSet};
17use std::fmt::{Debug, Formatter};
18use std::iter;
19use std::sync::atomic::{AtomicI64, AtomicU64, AtomicUsize, Ordering};
20use std::sync::{Arc, RwLock};
21use std::time::{Duration, Instant};
22
23use api::v1::OpType;
24use common_recordbatch::filter::SimpleFilterEvaluator;
25use common_telemetry::{debug, error};
26use common_time::Timestamp;
27use datatypes::arrow;
28use datatypes::arrow::array::ArrayRef;
29use datatypes::arrow_array::StringArray;
30use datatypes::data_type::{ConcreteDataType, DataType};
31use datatypes::prelude::{ScalarVector, Vector, VectorRef};
32use datatypes::schema::ColumnSchema;
33use datatypes::types::TimestampType;
34use datatypes::value::{Value, ValueRef};
35use datatypes::vectors::{
36    Helper, StringVector, TimestampMicrosecondVector, TimestampMillisecondVector,
37    TimestampNanosecondVector, TimestampSecondVector, UInt8Vector, UInt64Vector,
38};
39use mito_codec::key_values::KeyValue;
40use mito_codec::row_converter::{DensePrimaryKeyCodec, PrimaryKeyCodecExt};
41use snafu::{OptionExt, ResultExt, ensure};
42use store_api::metadata::RegionMetadataRef;
43use store_api::storage::{ColumnId, SequenceRange};
44use table::predicate::Predicate;
45
46use crate::error::{
47    self, ComputeArrowSnafu, ConvertVectorSnafu, EncodeSnafu, PrimaryKeyLengthMismatchSnafu, Result,
48};
49use crate::flush::WriteBufferManagerRef;
50use crate::memtable::builder::FieldBuilder;
51use crate::memtable::bulk::part::BulkPart;
52use crate::memtable::simple_bulk_memtable::SimpleBulkMemtable;
53use crate::memtable::stats::WriteMetrics;
54use crate::memtable::{
55    AllocTracker, BatchToRecordBatchContext, BoxedBatchIterator, BoxedRecordBatchIterator,
56    IterBuilder, KeyValues, MemScanMetrics, Memtable, MemtableBuilder, MemtableId, MemtableRange,
57    MemtableRangeContext, MemtableRanges, MemtableRef, MemtableStats, RangesOptions,
58    read_column_ids_from_projection,
59};
60use crate::metrics::{
61    MEMTABLE_ACTIVE_FIELD_BUILDER_COUNT, MEMTABLE_ACTIVE_SERIES_COUNT, READ_ROWS_TOTAL,
62    READ_STAGE_ELAPSED,
63};
64use crate::read::dedup::LastNonNullIter;
65use crate::read::prune::PruneTimeIterator;
66use crate::read::scan_region::PredicateGroup;
67use crate::read::{Batch, BatchBuilder, BatchColumn};
68use crate::region::options::MergeMode;
69
70/// Initial vector builder capacity.
71const INITIAL_BUILDER_CAPACITY: usize = 4;
72
73/// Vector builder capacity.
74const BUILDER_CAPACITY: usize = 512;
75
76fn checked_string_values_size(mut lengths: impl Iterator<Item = usize>) -> Option<i32> {
77    lengths.try_fold(0_i32, |size, len| {
78        size.checked_add(i32::try_from(len).ok()?)
79    })
80}
81
82/// Checks whether a string field builder that currently holds `current_len` bytes can
83/// additionally accommodate `space_needed` bytes of string data, given that the offset of a
84/// single Arrow string array is limited to `limit` (`i32::MAX` in production).
85///
86/// Returns:
87/// - `Ok(true)` if `current_len + space_needed <= limit`;
88/// - `Ok(false)` if `space_needed <= limit` but `current_len + space_needed > limit`, i.e.
89///   the current builder is too full but an empty builder could accommodate the batch (the
90///   caller may freeze the current builder and replay the batch onto an empty one);
91/// - `Err(InvalidBatchSnafu)` if the batch's string data alone (`space_needed`, or the fact
92///   that its total bytes cannot even be represented as `i32`) exceeds `limit`, i.e. the
93///   batch can never be accommodated by any builder.
94fn check_string_capacity(current_len: i32, space_needed: Option<i32>, limit: i32) -> Result<bool> {
95    let Some(space_needed) = space_needed.filter(|&space| space <= limit) else {
96        return error::InvalidBatchSnafu {
97            reason: format!(
98                "String data of the batch exceeds the Arrow string array offset limit ({limit}) and cannot be stored in a single column"
99            ),
100        }
101        .fail();
102    };
103    Ok(current_len
104        .checked_add(space_needed)
105        .is_some_and(|v| v <= limit))
106}
107
108/// Scans all string fields of a batch and checks whether they can be accommodated.
109///
110/// Every string field is checked (no early return on `Ok(false)`) so that an intrinsic
111/// oversize in *any* field surfaces as an error even when an earlier field only overflows
112/// the current builder.
113///
114/// Returns:
115/// - `Ok(true)` if every string field fits into the current builders;
116/// - `Ok(false)` if no field is intrinsically oversized but at least one field only fits
117///   into an empty builder (the caller may freeze the current builder and replay the batch
118///   onto an empty one);
119/// - `Err(InvalidBatchSnafu)` if any single field's string data alone (`space_needed`, or
120///   the fact that its total bytes cannot even be represented as `i32`) exceeds `limit`,
121///   i.e. that field can never be accommodated by any builder.
122fn scan_string_capacity(
123    fields: &[VectorRef],
124    field_builders: &[Option<FieldBuilder>],
125    field_types: &[ConcreteDataType],
126    limit: i32,
127) -> Result<bool> {
128    let mut can_fit_current = true;
129    for ((field_src, field_dest), field_type) in fields
130        .iter()
131        .zip(field_builders.iter())
132        .zip(field_types.iter())
133    {
134        if !matches!(field_type, ConcreteDataType::String(_)) {
135            continue;
136        }
137        let current_size = match field_dest {
138            Some(FieldBuilder::String(builder)) => builder.next_offset(),
139            None => 0,
140            Some(FieldBuilder::Other(_)) => unreachable!(),
141        };
142        let array = field_src.to_arrow_array();
143        let space_needed = if let Some(string_array) = array.as_any().downcast_ref::<StringArray>()
144        {
145            i32::try_from(string_array.value_data().len()).ok()
146        } else {
147            let string_vector = field_src
148                .as_any()
149                .downcast_ref::<StringVector>()
150                .with_context(|| error::InvalidBatchSnafu {
151                    reason: format!(
152                        "Field type mismatch, expecting String, given: {}",
153                        field_src.data_type()
154                    ),
155                })?;
156            checked_string_values_size(string_vector.iter_data().flatten().map(str::len))
157        };
158        can_fit_current &= check_string_capacity(current_size, space_needed, limit)?;
159    }
160    Ok(can_fit_current)
161}
162
163/// Builder to build [TimeSeriesMemtable].
164#[derive(Debug, Default)]
165pub struct TimeSeriesMemtableBuilder {
166    write_buffer_manager: Option<WriteBufferManagerRef>,
167    dedup: bool,
168    merge_mode: MergeMode,
169}
170
171impl TimeSeriesMemtableBuilder {
172    /// Creates a new builder with specific `write_buffer_manager`.
173    pub fn new(
174        write_buffer_manager: Option<WriteBufferManagerRef>,
175        dedup: bool,
176        merge_mode: MergeMode,
177    ) -> Self {
178        Self {
179            write_buffer_manager,
180            dedup,
181            merge_mode,
182        }
183    }
184}
185
186impl MemtableBuilder for TimeSeriesMemtableBuilder {
187    fn build(&self, id: MemtableId, metadata: &RegionMetadataRef) -> MemtableRef {
188        if metadata.primary_key.is_empty() {
189            Arc::new(SimpleBulkMemtable::new(
190                id,
191                metadata.clone(),
192                self.write_buffer_manager.clone(),
193                self.dedup,
194                self.merge_mode,
195            ))
196        } else {
197            Arc::new(TimeSeriesMemtable::new(
198                metadata.clone(),
199                id,
200                self.write_buffer_manager.clone(),
201                self.dedup,
202                self.merge_mode,
203            ))
204        }
205    }
206
207    fn use_bulk_insert(&self, _metadata: &RegionMetadataRef) -> bool {
208        // Now if we can use simple bulk memtable, the input request is already
209        // a bulk write request and won't call this method.
210        false
211    }
212}
213
214/// Memtable implementation that groups rows by their primary key.
215pub struct TimeSeriesMemtable {
216    id: MemtableId,
217    region_metadata: RegionMetadataRef,
218    row_codec: Arc<DensePrimaryKeyCodec>,
219    series_set: SeriesSet,
220    alloc_tracker: AllocTracker,
221    max_timestamp: AtomicI64,
222    min_timestamp: AtomicI64,
223    max_sequence: AtomicU64,
224    dedup: bool,
225    merge_mode: MergeMode,
226    /// Total written rows in memtable. This also includes deleted and duplicated rows.
227    num_rows: AtomicUsize,
228}
229
230impl TimeSeriesMemtable {
231    pub fn new(
232        region_metadata: RegionMetadataRef,
233        id: MemtableId,
234        write_buffer_manager: Option<WriteBufferManagerRef>,
235        dedup: bool,
236        merge_mode: MergeMode,
237    ) -> Self {
238        let row_codec = Arc::new(DensePrimaryKeyCodec::new(&region_metadata));
239        let series_set = SeriesSet::new(region_metadata.clone(), row_codec.clone());
240        Self {
241            id,
242            region_metadata,
243            series_set,
244            row_codec,
245            alloc_tracker: AllocTracker::new(write_buffer_manager),
246            max_timestamp: AtomicI64::new(i64::MIN),
247            min_timestamp: AtomicI64::new(i64::MAX),
248            max_sequence: AtomicU64::new(0),
249            dedup,
250            merge_mode,
251            num_rows: Default::default(),
252        }
253    }
254
255    /// Updates memtable stats.
256    fn update_stats(&self, stats: WriteMetrics) {
257        self.alloc_tracker
258            .on_allocation(stats.key_bytes + stats.value_bytes);
259        self.max_timestamp.fetch_max(stats.max_ts, Ordering::SeqCst);
260        self.min_timestamp.fetch_min(stats.min_ts, Ordering::SeqCst);
261        self.max_sequence
262            .fetch_max(stats.max_sequence, Ordering::SeqCst);
263        self.num_rows.fetch_add(stats.num_rows, Ordering::SeqCst);
264    }
265
266    fn write_key_value(&self, kv: KeyValue, stats: &mut WriteMetrics) -> Result<()> {
267        ensure!(
268            self.row_codec.num_fields() == kv.num_primary_keys(),
269            PrimaryKeyLengthMismatchSnafu {
270                expect: self.row_codec.num_fields(),
271                actual: kv.num_primary_keys(),
272            }
273        );
274
275        let primary_key_encoded = self
276            .row_codec
277            .encode(kv.primary_keys())
278            .context(EncodeSnafu)?;
279
280        let (key_allocated, value_allocated) =
281            self.series_set.push_to_series(primary_key_encoded, &kv);
282        stats.key_bytes += key_allocated;
283        stats.value_bytes += value_allocated;
284
285        // safety: timestamp of kv must be both present and a valid timestamp value.
286        let ts = kv
287            .timestamp()
288            .try_into_timestamp()
289            .unwrap()
290            .unwrap()
291            .value();
292        stats.min_ts = stats.min_ts.min(ts);
293        stats.max_ts = stats.max_ts.max(ts);
294        Ok(())
295    }
296}
297
298impl Debug for TimeSeriesMemtable {
299    fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
300        f.debug_struct("TimeSeriesMemtable").finish()
301    }
302}
303
304impl Memtable for TimeSeriesMemtable {
305    fn id(&self) -> MemtableId {
306        self.id
307    }
308
309    fn write(&self, kvs: &KeyValues) -> Result<()> {
310        if kvs.is_empty() {
311            return Ok(());
312        }
313
314        let mut local_stats = WriteMetrics::default();
315
316        for kv in kvs.iter() {
317            self.write_key_value(kv, &mut local_stats)?;
318        }
319        local_stats.value_bytes += kvs.num_rows() * std::mem::size_of::<Timestamp>();
320        local_stats.value_bytes += kvs.num_rows() * std::mem::size_of::<OpType>();
321        local_stats.max_sequence = kvs.max_sequence();
322        local_stats.num_rows = kvs.num_rows();
323        // TODO(hl): this maybe inaccurate since for-iteration may return early.
324        // We may lift the primary key length check out of Memtable::write
325        // so that we can ensure writing to memtable will succeed.
326        self.update_stats(local_stats);
327        Ok(())
328    }
329
330    fn write_one(&self, key_value: KeyValue) -> Result<()> {
331        let mut metrics = WriteMetrics::default();
332        let res = self.write_key_value(key_value, &mut metrics);
333        metrics.value_bytes += std::mem::size_of::<Timestamp>() + std::mem::size_of::<OpType>();
334        metrics.max_sequence = key_value.sequence();
335        metrics.num_rows = 1;
336
337        if res.is_ok() {
338            self.update_stats(metrics);
339        }
340        res
341    }
342
343    fn write_bulk(&self, part: BulkPart) -> Result<()> {
344        // Default implementation fallback to row iteration.
345        let mutation = part.to_mutation(&self.region_metadata)?;
346        let mut metrics = WriteMetrics::default();
347        if let Some(key_values) = KeyValues::new(&self.region_metadata, mutation) {
348            for kv in key_values.iter() {
349                self.write_key_value(kv, &mut metrics)?
350            }
351        }
352
353        metrics.max_sequence = part.sequence;
354        metrics.max_ts = part.max_timestamp;
355        metrics.min_ts = part.min_timestamp;
356        metrics.num_rows = part.num_rows();
357        self.update_stats(metrics);
358        Ok(())
359    }
360
361    fn ranges(
362        &self,
363        projection: Option<&[ColumnId]>,
364        options: RangesOptions,
365    ) -> Result<MemtableRanges> {
366        let predicate = options.predicate;
367        let sequence = options.sequence;
368        let read_column_ids = read_column_ids_from_projection(&self.region_metadata, projection);
369        let projection = if let Some(projection) = projection {
370            projection.iter().copied().collect()
371        } else {
372            self.region_metadata
373                .field_columns()
374                .map(|c| c.column_id)
375                .collect()
376        };
377        let batch_to_record_batch = Arc::new(BatchToRecordBatchContext::new(
378            self.region_metadata.clone(),
379            read_column_ids,
380        ));
381        let builder = Box::new(TimeSeriesIterBuilder {
382            series_set: self.series_set.clone(),
383            projection,
384            predicate: predicate.clone(),
385            dedup: self.dedup,
386            merge_mode: self.merge_mode,
387            sequence,
388            batch_to_record_batch,
389        });
390        let context = Arc::new(MemtableRangeContext::new(self.id, builder, predicate));
391        let range_stats = self.stats();
392        let range = MemtableRange::new(context, range_stats);
393        Ok(MemtableRanges {
394            ranges: [(0, range)].into(),
395        })
396    }
397
398    fn is_empty(&self) -> bool {
399        self.series_set.series.read().unwrap().0.is_empty()
400    }
401
402    fn freeze(&self) -> Result<()> {
403        self.alloc_tracker.done_allocating();
404
405        Ok(())
406    }
407
408    fn stats(&self) -> MemtableStats {
409        let estimated_bytes = self.alloc_tracker.bytes_allocated();
410
411        if estimated_bytes == 0 {
412            // no rows ever written
413            return MemtableStats {
414                estimated_bytes,
415                time_range: None,
416                num_rows: 0,
417                num_ranges: 0,
418                max_sequence: 0,
419                series_count: 0,
420            };
421        }
422        let ts_type = self
423            .region_metadata
424            .time_index_column()
425            .column_schema
426            .data_type
427            .clone()
428            .as_timestamp()
429            .expect("Timestamp column must have timestamp type");
430        let max_timestamp = ts_type.create_timestamp(self.max_timestamp.load(Ordering::Relaxed));
431        let min_timestamp = ts_type.create_timestamp(self.min_timestamp.load(Ordering::Relaxed));
432        let series_count = self.series_set.series.read().unwrap().0.len();
433        MemtableStats {
434            estimated_bytes,
435            time_range: Some((min_timestamp, max_timestamp)),
436            num_rows: self.num_rows.load(Ordering::Relaxed),
437            num_ranges: 1,
438            max_sequence: self.max_sequence.load(Ordering::Relaxed),
439            series_count,
440        }
441    }
442
443    fn fork(&self, id: MemtableId, metadata: &RegionMetadataRef) -> MemtableRef {
444        Arc::new(TimeSeriesMemtable::new(
445            metadata.clone(),
446            id,
447            self.alloc_tracker.write_buffer_manager(),
448            self.dedup,
449            self.merge_mode,
450        ))
451    }
452}
453
454#[derive(Default)]
455struct SeriesMap(BTreeMap<Vec<u8>, Arc<RwLock<Series>>>);
456
457impl Drop for SeriesMap {
458    fn drop(&mut self) {
459        let num_series = self.0.len();
460        let num_field_builders = self
461            .0
462            .values()
463            .map(|v| v.read().unwrap().active.num_field_builders())
464            .sum::<usize>();
465        MEMTABLE_ACTIVE_SERIES_COUNT.sub(num_series as i64);
466        MEMTABLE_ACTIVE_FIELD_BUILDER_COUNT.sub(num_field_builders as i64);
467    }
468}
469
470#[derive(Clone)]
471pub(crate) struct SeriesSet {
472    region_metadata: RegionMetadataRef,
473    series: Arc<RwLock<SeriesMap>>,
474    codec: Arc<DensePrimaryKeyCodec>,
475}
476
477impl SeriesSet {
478    fn new(region_metadata: RegionMetadataRef, codec: Arc<DensePrimaryKeyCodec>) -> Self {
479        Self {
480            region_metadata,
481            series: Default::default(),
482            codec,
483        }
484    }
485}
486
487impl SeriesSet {
488    /// Push [KeyValue] to SeriesSet with given primary key and return key/value allocated memory size.
489    fn push_to_series(&self, primary_key: Vec<u8>, kv: &KeyValue) -> (usize, usize) {
490        if let Some(series) = self.series.read().unwrap().0.get(&primary_key) {
491            let value_allocated = series.write().unwrap().push(
492                kv.timestamp(),
493                kv.sequence(),
494                kv.op_type(),
495                kv.fields(),
496            );
497            return (0, value_allocated);
498        };
499
500        let mut indices = self.series.write().unwrap();
501        match indices.0.entry(primary_key) {
502            Entry::Vacant(v) => {
503                let key_len = v.key().len();
504                let mut series = Series::new(&self.region_metadata);
505                let value_allocated =
506                    series.push(kv.timestamp(), kv.sequence(), kv.op_type(), kv.fields());
507                v.insert(Arc::new(RwLock::new(series)));
508                (key_len, value_allocated)
509            }
510            // safety: series must exist at given index.
511            Entry::Occupied(v) => {
512                let value_allocated = v.get().write().unwrap().push(
513                    kv.timestamp(),
514                    kv.sequence(),
515                    kv.op_type(),
516                    kv.fields(),
517                );
518                (0, value_allocated)
519            }
520        }
521    }
522
523    #[cfg(test)]
524    fn get_series(&self, primary_key: &[u8]) -> Option<Arc<RwLock<Series>>> {
525        self.series.read().unwrap().0.get(primary_key).cloned()
526    }
527
528    /// Iterates all series in [SeriesSet].
529    fn iter_series(
530        &self,
531        projection: HashSet<ColumnId>,
532        predicate: PredicateGroup,
533        dedup: bool,
534        merge_mode: MergeMode,
535        sequence: Option<SequenceRange>,
536        mem_scan_metrics: Option<MemScanMetrics>,
537    ) -> Result<Iter> {
538        let primary_key_schema = primary_key_schema(&self.region_metadata);
539        let primary_key_datatypes = self
540            .region_metadata
541            .primary_key_columns()
542            .map(|pk| pk.column_schema.data_type.clone())
543            .collect();
544
545        Iter::try_new(
546            self.region_metadata.clone(),
547            self.series.clone(),
548            projection,
549            predicate.predicate().cloned(),
550            primary_key_schema,
551            primary_key_datatypes,
552            self.codec.clone(),
553            dedup,
554            merge_mode,
555            sequence,
556            mem_scan_metrics,
557        )
558    }
559}
560
561/// Creates an arrow [SchemaRef](arrow::datatypes::SchemaRef) that only contains primary keys
562/// of given region schema
563pub(crate) fn primary_key_schema(
564    region_metadata: &RegionMetadataRef,
565) -> arrow::datatypes::SchemaRef {
566    let fields = region_metadata
567        .primary_key_columns()
568        .map(|pk| {
569            arrow::datatypes::Field::new(
570                pk.column_schema.name.clone(),
571                pk.column_schema.data_type.as_arrow_type(),
572                pk.column_schema.is_nullable(),
573            )
574        })
575        .collect::<Vec<_>>();
576    Arc::new(arrow::datatypes::Schema::new(fields))
577}
578
579/// Metrics for reading the memtable.
580#[derive(Debug, Default)]
581struct Metrics {
582    /// Total series in the memtable.
583    total_series: usize,
584    /// Number of series pruned.
585    num_pruned_series: usize,
586    /// Number of rows read.
587    num_rows: usize,
588    /// Number of batch read.
589    num_batches: usize,
590    /// Duration to scan the memtable.
591    scan_cost: Duration,
592}
593
594struct Iter {
595    metadata: RegionMetadataRef,
596    series: Arc<RwLock<SeriesMap>>,
597    projection: HashSet<ColumnId>,
598    last_key: Option<Vec<u8>>,
599    predicate: Vec<SimpleFilterEvaluator>,
600    pk_schema: arrow::datatypes::SchemaRef,
601    pk_datatypes: Vec<ConcreteDataType>,
602    codec: Arc<DensePrimaryKeyCodec>,
603    dedup: bool,
604    merge_mode: MergeMode,
605    sequence: Option<SequenceRange>,
606    metrics: Metrics,
607    mem_scan_metrics: Option<MemScanMetrics>,
608}
609
610impl Iter {
611    #[allow(clippy::too_many_arguments)]
612    pub(crate) fn try_new(
613        metadata: RegionMetadataRef,
614        series: Arc<RwLock<SeriesMap>>,
615        projection: HashSet<ColumnId>,
616        predicate: Option<Predicate>,
617        pk_schema: arrow::datatypes::SchemaRef,
618        pk_datatypes: Vec<ConcreteDataType>,
619        codec: Arc<DensePrimaryKeyCodec>,
620        dedup: bool,
621        merge_mode: MergeMode,
622        sequence: Option<SequenceRange>,
623        mem_scan_metrics: Option<MemScanMetrics>,
624    ) -> Result<Self> {
625        let predicate = predicate
626            .map(|predicate| {
627                predicate
628                    .exprs()
629                    .iter()
630                    .filter_map(SimpleFilterEvaluator::try_new)
631                    .collect::<Vec<_>>()
632            })
633            .unwrap_or_default();
634        Ok(Self {
635            metadata,
636            series,
637            projection,
638            last_key: None,
639            predicate,
640            pk_schema,
641            pk_datatypes,
642            codec,
643            dedup,
644            merge_mode,
645            sequence,
646            metrics: Metrics::default(),
647            mem_scan_metrics,
648        })
649    }
650
651    fn report_mem_scan_metrics(&mut self) {
652        if let Some(mem_scan_metrics) = self.mem_scan_metrics.take() {
653            let inner = crate::memtable::MemScanMetricsData {
654                total_series: self.metrics.total_series,
655                num_rows: self.metrics.num_rows,
656                num_batches: self.metrics.num_batches,
657                scan_cost: self.metrics.scan_cost,
658                ..Default::default()
659            };
660            mem_scan_metrics.merge_inner(&inner);
661        }
662    }
663}
664
665impl Drop for Iter {
666    fn drop(&mut self) {
667        debug!(
668            "Iter {} time series memtable, metrics: {:?}",
669            self.metadata.region_id, self.metrics
670        );
671
672        // Report MemScanMetrics if not already reported
673        self.report_mem_scan_metrics();
674
675        READ_ROWS_TOTAL
676            .with_label_values(&["time_series_memtable"])
677            .inc_by(self.metrics.num_rows as u64);
678        READ_STAGE_ELAPSED
679            .with_label_values(&["scan_memtable"])
680            .observe(self.metrics.scan_cost.as_secs_f64());
681    }
682}
683
684impl Iterator for Iter {
685    type Item = Result<Batch>;
686
687    fn next(&mut self) -> Option<Self::Item> {
688        let start = Instant::now();
689        let map = self.series.read().unwrap();
690        let range = match &self.last_key {
691            None => map.0.range::<Vec<u8>, _>(..),
692            Some(last_key) => map
693                .0
694                .range::<Vec<u8>, _>((Bound::Excluded(last_key), Bound::Unbounded)),
695        };
696
697        // TODO(hl): maybe yield more than one time series to amortize range overhead.
698        for (primary_key, series) in range {
699            self.metrics.total_series += 1;
700
701            let mut series = series.write().unwrap();
702            if !self.predicate.is_empty()
703                && !prune_primary_key(
704                    &self.codec,
705                    primary_key.as_slice(),
706                    &mut series,
707                    &self.pk_datatypes,
708                    self.pk_schema.clone(),
709                    &self.predicate,
710                )
711            {
712                // read next series
713                self.metrics.num_pruned_series += 1;
714                continue;
715            }
716            self.last_key = Some(primary_key.clone());
717
718            let values = series.compact(&self.metadata);
719            let batch = values.and_then(|v| {
720                v.to_batch(
721                    primary_key,
722                    &self.metadata,
723                    &self.projection,
724                    self.sequence,
725                    self.dedup,
726                    self.merge_mode,
727                )
728            });
729
730            // Update metrics.
731            self.metrics.num_batches += 1;
732            self.metrics.num_rows += batch.as_ref().map(|b| b.num_rows()).unwrap_or(0);
733            self.metrics.scan_cost += start.elapsed();
734
735            return Some(batch);
736        }
737        drop(map); // Explicitly drop the read lock
738        self.metrics.scan_cost += start.elapsed();
739
740        // Report MemScanMetrics before returning None
741        self.report_mem_scan_metrics();
742
743        None
744    }
745}
746
747fn prune_primary_key(
748    codec: &Arc<DensePrimaryKeyCodec>,
749    pk: &[u8],
750    series: &mut Series,
751    datatypes: &[ConcreteDataType],
752    pk_schema: arrow::datatypes::SchemaRef,
753    predicates: &[SimpleFilterEvaluator],
754) -> bool {
755    // no primary key, we simply return true.
756    if pk_schema.fields().is_empty() {
757        return true;
758    }
759
760    // retrieve primary key values from cache or decode from bytes.
761    let pk_values = if let Some(pk_values) = series.pk_cache.as_ref() {
762        pk_values
763    } else {
764        let pk_values = codec.decode_dense_without_column_id(pk);
765        if let Err(e) = pk_values {
766            error!(e; "Failed to decode primary key");
767            return true;
768        }
769        series.update_pk_cache(pk_values.unwrap());
770        series.pk_cache.as_ref().unwrap()
771    };
772
773    // evaluate predicates against primary key values
774    let mut result = true;
775    for predicate in predicates {
776        // ignore predicates that are not referencing primary key columns
777        let Ok(index) = pk_schema.index_of(predicate.column_name()) else {
778            continue;
779        };
780        // Safety: arrow schema and datatypes are constructed from the same source.
781        let scalar_value = pk_values[index]
782            .try_to_scalar_value(&datatypes[index])
783            .unwrap();
784        result &= predicate.evaluate_scalar(&scalar_value).unwrap_or(true);
785    }
786
787    result
788}
789
790/// A `Series` holds a list of field values of some given primary key.
791pub struct Series {
792    pk_cache: Option<Vec<Value>>,
793    active: ValueBuilder,
794    frozen: Vec<Values>,
795    region_metadata: RegionMetadataRef,
796    capacity: usize,
797}
798
799impl Series {
800    pub(crate) fn with_capacity(
801        region_metadata: &RegionMetadataRef,
802        init_capacity: usize,
803        capacity: usize,
804    ) -> Self {
805        MEMTABLE_ACTIVE_SERIES_COUNT.inc();
806        Self {
807            pk_cache: None,
808            active: ValueBuilder::new(region_metadata, init_capacity),
809            frozen: vec![],
810            region_metadata: region_metadata.clone(),
811            capacity,
812        }
813    }
814
815    pub(crate) fn new(region_metadata: &RegionMetadataRef) -> Self {
816        Self::with_capacity(region_metadata, INITIAL_BUILDER_CAPACITY, BUILDER_CAPACITY)
817    }
818
819    pub fn is_empty(&self) -> bool {
820        self.active.len() == 0 && self.frozen.is_empty()
821    }
822
823    /// Pushes a row of values into Series. Return the size of values.
824    pub(crate) fn push<'a>(
825        &mut self,
826        ts: ValueRef<'a>,
827        sequence: u64,
828        op_type: OpType,
829        values: impl Iterator<Item = ValueRef<'a>>,
830    ) -> usize {
831        // + 10 to avoid potential reallocation.
832        if self.active.len() + 10 > self.capacity {
833            let region_metadata = self.region_metadata.clone();
834            self.freeze(&region_metadata);
835        }
836        self.active.push(ts, sequence, op_type as u8, values)
837    }
838
839    fn update_pk_cache(&mut self, pk_values: Vec<Value>) {
840        self.pk_cache = Some(pk_values);
841    }
842
843    /// Freezes the active part and push it to `frozen`.
844    pub(crate) fn freeze(&mut self, region_metadata: &RegionMetadataRef) {
845        if self.active.len() != 0 {
846            let mut builder = ValueBuilder::new(region_metadata, INITIAL_BUILDER_CAPACITY);
847            std::mem::swap(&mut self.active, &mut builder);
848            self.frozen.push(Values::from(builder));
849        }
850    }
851
852    pub(crate) fn extend(
853        &mut self,
854        ts_v: VectorRef,
855        op_type_v: u8,
856        sequence_v: u64,
857        fields: Vec<VectorRef>,
858    ) -> Result<()> {
859        if !self.active.can_accommodate(&fields)? {
860            let region_metadata = self.region_metadata.clone();
861            self.freeze(&region_metadata);
862        }
863        self.active.extend(ts_v, op_type_v, sequence_v, fields)
864    }
865
866    /// Freezes active part to frozen part and compact frozen part to reduce memory fragmentation.
867    /// Returns the frozen and compacted values.
868    pub(crate) fn compact(&mut self, region_metadata: &RegionMetadataRef) -> Result<&Values> {
869        self.freeze(region_metadata);
870
871        let frozen = &self.frozen;
872
873        // Each series must contain at least one row
874        debug_assert!(!frozen.is_empty());
875
876        if frozen.len() > 1 {
877            // TODO(hl): We should keep track of min/max timestamps for each values and avoid
878            // cloning and sorting when values do not overlap with each other.
879
880            let column_size = frozen[0].fields.len() + 3;
881
882            if cfg!(debug_assertions) {
883                debug_assert!(
884                    frozen
885                        .iter()
886                        .zip(frozen.iter().skip(1))
887                        .all(|(prev, next)| { prev.fields.len() == next.fields.len() })
888                );
889            }
890
891            let arrays = frozen.iter().map(|v| v.columns()).collect::<Vec<_>>();
892            let concatenated = (0..column_size)
893                .map(|i| {
894                    let to_concat = arrays.iter().map(|a| a[i].as_ref()).collect::<Vec<_>>();
895                    arrow::compute::concat(&to_concat)
896                })
897                .collect::<std::result::Result<Vec<_>, _>>()
898                .context(ComputeArrowSnafu)?;
899
900            debug_assert_eq!(concatenated.len(), column_size);
901            let values = Values::from_columns(&concatenated)?;
902            self.frozen = vec![values];
903        };
904        Ok(&self.frozen[0])
905    }
906
907    pub fn read_to_values(&self) -> Vec<Values> {
908        let mut res = Vec::with_capacity(self.frozen.len() + 1);
909        res.extend(self.frozen.iter().cloned());
910        res.push(self.active.finish_cloned());
911        res
912    }
913}
914
915/// `ValueBuilder` holds all the vector builders for field columns.
916pub(crate) struct ValueBuilder {
917    timestamp: Vec<i64>,
918    timestamp_type: ConcreteDataType,
919    sequence: Vec<u64>,
920    op_type: Vec<u8>,
921    fields: Vec<Option<FieldBuilder>>,
922    field_schemas: Vec<ColumnSchema>,
923}
924
925impl ValueBuilder {
926    pub(crate) fn new(region_metadata: &RegionMetadataRef, capacity: usize) -> Self {
927        let timestamp_type = region_metadata
928            .time_index_column()
929            .column_schema
930            .data_type
931            .clone();
932        let sequence = Vec::with_capacity(capacity);
933        let op_type = Vec::with_capacity(capacity);
934
935        let field_schemas = region_metadata
936            .field_columns()
937            .map(|c| c.column_schema.clone())
938            .collect::<Vec<_>>();
939        let fields = (0..field_schemas.len()).map(|_| None).collect();
940        Self {
941            timestamp: Vec::with_capacity(capacity),
942            timestamp_type,
943            sequence,
944            op_type,
945            fields,
946            field_schemas,
947        }
948    }
949
950    /// Returns number of field builders.
951    pub fn num_field_builders(&self) -> usize {
952        self.fields.iter().flatten().count()
953    }
954
955    /// Pushes a new row to `ValueBuilder`.
956    /// We don't need primary keys since they've already be encoded.
957    /// Returns the size of field values.
958    ///
959    /// In this method, we don't check the data type of the value, because it is already checked in the caller.
960    pub(crate) fn push<'a>(
961        &mut self,
962        ts: ValueRef,
963        sequence: u64,
964        op_type: u8,
965        fields: impl Iterator<Item = ValueRef<'a>>,
966    ) -> usize {
967        #[cfg(debug_assertions)]
968        let fields = {
969            let field_vec = fields.collect::<Vec<_>>();
970            debug_assert_eq!(field_vec.len(), self.fields.len());
971            field_vec.into_iter()
972        };
973
974        self.timestamp
975            .push(ts.try_into_timestamp().unwrap().unwrap().value());
976        self.sequence.push(sequence);
977        self.op_type.push(op_type);
978        let num_rows = self.timestamp.len();
979        let mut size = 0;
980        for (idx, field_value) in fields.enumerate() {
981            size += field_value.data_size();
982            if !field_value.is_null() || self.fields[idx].is_some() {
983                if let Some(field) = self.fields[idx].as_mut() {
984                    field
985                        .push(field_value)
986                        .unwrap_or_else(|e| panic!("Failed to push field value: {e:?}"));
987                } else {
988                    let mut mutable_vector = FieldBuilder::create(
989                        &self.field_schemas[idx],
990                        num_rows.max(INITIAL_BUILDER_CAPACITY),
991                    );
992                    mutable_vector.push_nulls(num_rows - 1);
993                    mutable_vector
994                        .push(field_value)
995                        .unwrap_or_else(|e| panic!("unexpected field value: {e:?}"));
996                    self.fields[idx] = Some(mutable_vector);
997                    MEMTABLE_ACTIVE_FIELD_BUILDER_COUNT.inc();
998                }
999            }
1000        }
1001
1002        size
1003    }
1004
1005    /// Checks if current value builder have sufficient space to accommodate `fields`.
1006    ///
1007    /// Returns `Ok(false)` if the current builder lacks the remaining space to accommodate
1008    /// the fields due to offset overflow, but an empty builder would be able to accommodate
1009    /// the batch (the caller may freeze the current builder and replay the batch onto an
1010    /// empty one).
1011    ///
1012    /// Returns `Err(InvalidBatchSnafu)` if the string data of a single batch itself exceeds
1013    /// the Arrow string array offset limit and thus can never be accommodated, not even by an
1014    /// empty builder.
1015    pub(crate) fn can_accommodate(&self, fields: &[VectorRef]) -> Result<bool> {
1016        let data_types = self
1017            .field_schemas
1018            .iter()
1019            .map(|x| x.data_type.clone())
1020            .collect::<Vec<_>>();
1021        scan_string_capacity(fields, &self.fields, &data_types, i32::MAX)
1022    }
1023
1024    pub(crate) fn extend(
1025        &mut self,
1026        ts_v: VectorRef,
1027        op_type: u8,
1028        sequence: u64,
1029        fields: Vec<VectorRef>,
1030    ) -> Result<()> {
1031        let num_rows_before = self.timestamp.len();
1032        let num_rows_to_write = ts_v.len();
1033        self.timestamp.reserve(num_rows_to_write);
1034        match self.timestamp_type {
1035            ConcreteDataType::Timestamp(TimestampType::Second(_)) => {
1036                self.timestamp.extend(
1037                    ts_v.as_any()
1038                        .downcast_ref::<TimestampSecondVector>()
1039                        .unwrap()
1040                        .iter_data()
1041                        .map(|v| v.unwrap().0.value()),
1042                );
1043            }
1044            ConcreteDataType::Timestamp(TimestampType::Millisecond(_)) => {
1045                self.timestamp.extend(
1046                    ts_v.as_any()
1047                        .downcast_ref::<TimestampMillisecondVector>()
1048                        .unwrap()
1049                        .iter_data()
1050                        .map(|v| v.unwrap().0.value()),
1051                );
1052            }
1053            ConcreteDataType::Timestamp(TimestampType::Microsecond(_)) => {
1054                self.timestamp.extend(
1055                    ts_v.as_any()
1056                        .downcast_ref::<TimestampMicrosecondVector>()
1057                        .unwrap()
1058                        .iter_data()
1059                        .map(|v| v.unwrap().0.value()),
1060                );
1061            }
1062            ConcreteDataType::Timestamp(TimestampType::Nanosecond(_)) => {
1063                self.timestamp.extend(
1064                    ts_v.as_any()
1065                        .downcast_ref::<TimestampNanosecondVector>()
1066                        .unwrap()
1067                        .iter_data()
1068                        .map(|v| v.unwrap().0.value()),
1069                );
1070            }
1071            _ => unreachable!(),
1072        };
1073
1074        self.op_type.reserve(num_rows_to_write);
1075        self.op_type
1076            .extend(iter::repeat_n(op_type, num_rows_to_write));
1077        self.sequence.reserve(num_rows_to_write);
1078        self.sequence
1079            .extend(iter::repeat_n(sequence, num_rows_to_write));
1080
1081        for (field_idx, (field_src, field_dest)) in
1082            fields.into_iter().zip(self.fields.iter_mut()).enumerate()
1083        {
1084            let builder = field_dest.get_or_insert_with(|| {
1085                let mut field_builder =
1086                    FieldBuilder::create(&self.field_schemas[field_idx], INITIAL_BUILDER_CAPACITY);
1087                field_builder.push_nulls(num_rows_before);
1088                field_builder
1089            });
1090            match builder {
1091                FieldBuilder::String(builder) => {
1092                    let array = field_src.to_arrow_array();
1093                    if let Some(string_array) = array.as_any().downcast_ref::<StringArray>() {
1094                        builder.append_array(string_array);
1095                    } else {
1096                        let string_vector = field_src
1097                            .as_any()
1098                            .downcast_ref::<StringVector>()
1099                            .with_context(|| error::InvalidBatchSnafu {
1100                                reason: format!(
1101                                    "Field type mismatch, expecting String, given: {}",
1102                                    field_src.data_type()
1103                                ),
1104                            })?;
1105                        for value in string_vector.iter_data() {
1106                            if let Some(value) = value {
1107                                builder.append(value);
1108                            } else {
1109                                builder.append_null();
1110                            }
1111                        }
1112                    }
1113                }
1114                FieldBuilder::Other(builder) => {
1115                    let len = field_src.len();
1116                    builder
1117                        .extend_slice_of(&*field_src, 0, len)
1118                        .context(error::ComputeVectorSnafu)?;
1119                }
1120            }
1121        }
1122        Ok(())
1123    }
1124
1125    /// Returns the length of [ValueBuilder]
1126    fn len(&self) -> usize {
1127        let sequence_len = self.sequence.len();
1128        debug_assert_eq!(sequence_len, self.op_type.len());
1129        debug_assert_eq!(sequence_len, self.timestamp.len());
1130        sequence_len
1131    }
1132
1133    fn finish_cloned(&self) -> Values {
1134        let num_rows = self.sequence.len();
1135        let fields = self
1136            .fields
1137            .iter()
1138            .enumerate()
1139            .map(|(i, v)| {
1140                if let Some(v) = v {
1141                    MEMTABLE_ACTIVE_FIELD_BUILDER_COUNT.dec();
1142                    v.finish_cloned()
1143                } else {
1144                    let mut builder = FieldBuilder::create(&self.field_schemas[i], num_rows);
1145                    builder.push_nulls(num_rows);
1146                    builder.finish()
1147                }
1148            })
1149            .collect::<Vec<_>>();
1150
1151        let sequence = Arc::new(UInt64Vector::from_vec(self.sequence.clone()));
1152        let op_type = Arc::new(UInt8Vector::from_vec(self.op_type.clone()));
1153        let timestamp: VectorRef = match self.timestamp_type {
1154            ConcreteDataType::Timestamp(TimestampType::Second(_)) => {
1155                Arc::new(TimestampSecondVector::from_vec(self.timestamp.clone()))
1156            }
1157            ConcreteDataType::Timestamp(TimestampType::Millisecond(_)) => {
1158                Arc::new(TimestampMillisecondVector::from_vec(self.timestamp.clone()))
1159            }
1160            ConcreteDataType::Timestamp(TimestampType::Microsecond(_)) => {
1161                Arc::new(TimestampMicrosecondVector::from_vec(self.timestamp.clone()))
1162            }
1163            ConcreteDataType::Timestamp(TimestampType::Nanosecond(_)) => {
1164                Arc::new(TimestampNanosecondVector::from_vec(self.timestamp.clone()))
1165            }
1166            _ => unreachable!(),
1167        };
1168
1169        if cfg!(debug_assertions) {
1170            debug_assert_eq!(timestamp.len(), sequence.len());
1171            debug_assert_eq!(timestamp.len(), op_type.len());
1172            for field in &fields {
1173                debug_assert_eq!(timestamp.len(), field.len());
1174            }
1175        }
1176
1177        Values {
1178            timestamp,
1179            sequence,
1180            op_type,
1181            fields,
1182        }
1183    }
1184}
1185
1186/// [Values] holds an immutable vectors of field columns, including `sequence` and `op_type`.
1187#[derive(Clone)]
1188pub struct Values {
1189    pub(crate) timestamp: VectorRef,
1190    pub(crate) sequence: Arc<UInt64Vector>,
1191    pub(crate) op_type: Arc<UInt8Vector>,
1192    pub(crate) fields: Vec<VectorRef>,
1193}
1194
1195impl Values {
1196    /// Converts [Values] to `Batch`, applies the optional sequence filter, sorts the batch
1197    /// according to `timestamp, sequence` desc, and applies dedup/merge according to `merge_mode`.
1198    pub fn to_batch(
1199        &self,
1200        primary_key: &[u8],
1201        metadata: &RegionMetadataRef,
1202        projection: &HashSet<ColumnId>,
1203        sequence: Option<SequenceRange>,
1204        dedup: bool,
1205        merge_mode: MergeMode,
1206    ) -> Result<Batch> {
1207        let builder = BatchBuilder::with_required_columns(
1208            primary_key.to_vec(),
1209            self.timestamp.clone(),
1210            self.sequence.clone(),
1211            self.op_type.clone(),
1212        );
1213
1214        let fields = metadata
1215            .field_columns()
1216            .zip(self.fields.iter())
1217            .filter_map(|(c, f)| {
1218                projection.get(&c.column_id).map(|c| BatchColumn {
1219                    column_id: *c,
1220                    data: f.clone(),
1221                })
1222            })
1223            .collect();
1224
1225        let mut batch = builder.with_fields(fields).build()?;
1226        // The sequence filter must be applied before dedup/merge to:
1227        // - avoid dropping a timestamp when the newest row is out of range
1228        // - avoid filling null fields from rows that should be excluded by the sequence filter.
1229        batch.filter_by_sequence(sequence)?;
1230
1231        match (dedup, merge_mode) {
1232            // append-only, keep duplicate rows.
1233            (false, _) => batch.sort(false)?,
1234            // keep the last row for each timestamp.
1235            (true, MergeMode::LastRow) => batch.sort(true)?,
1236            // keep the last non-null value for each field.
1237            (true, MergeMode::LastNonNull) => {
1238                batch.sort(false)?;
1239                batch.merge_last_non_null()?;
1240            }
1241        }
1242        Ok(batch)
1243    }
1244
1245    /// Returns a vector of all columns converted to arrow [Array](datatypes::arrow::array::Array) in [Values].
1246    fn columns(&self) -> Vec<ArrayRef> {
1247        let mut res = Vec::with_capacity(3 + self.fields.len());
1248        res.push(self.timestamp.to_arrow_array());
1249        res.push(self.sequence.to_arrow_array());
1250        res.push(self.op_type.to_arrow_array());
1251        res.extend(self.fields.iter().map(|f| f.to_arrow_array()));
1252        res
1253    }
1254
1255    /// Builds a new [Values] instance from columns.
1256    fn from_columns(cols: &[ArrayRef]) -> Result<Self> {
1257        debug_assert!(cols.len() >= 3);
1258        let timestamp = Helper::try_into_vector(&cols[0]).context(ConvertVectorSnafu)?;
1259        let sequence =
1260            Arc::new(UInt64Vector::try_from_arrow_array(&cols[1]).context(ConvertVectorSnafu)?);
1261        let op_type =
1262            Arc::new(UInt8Vector::try_from_arrow_array(&cols[2]).context(ConvertVectorSnafu)?);
1263        let fields = Helper::try_into_vectors(&cols[3..]).context(ConvertVectorSnafu)?;
1264
1265        Ok(Self {
1266            timestamp,
1267            sequence,
1268            op_type,
1269            fields,
1270        })
1271    }
1272}
1273
1274impl From<ValueBuilder> for Values {
1275    fn from(mut value: ValueBuilder) -> Self {
1276        let num_rows = value.len();
1277        let fields = value
1278            .fields
1279            .iter_mut()
1280            .enumerate()
1281            .map(|(i, v)| {
1282                if let Some(v) = v {
1283                    MEMTABLE_ACTIVE_FIELD_BUILDER_COUNT.dec();
1284                    v.finish()
1285                } else {
1286                    let mut builder = FieldBuilder::create(&value.field_schemas[i], num_rows);
1287                    builder.push_nulls(num_rows);
1288                    builder.finish()
1289                }
1290            })
1291            .collect::<Vec<_>>();
1292
1293        let sequence = Arc::new(UInt64Vector::from_vec(value.sequence));
1294        let op_type = Arc::new(UInt8Vector::from_vec(value.op_type));
1295        let timestamp: VectorRef = match value.timestamp_type {
1296            ConcreteDataType::Timestamp(TimestampType::Second(_)) => {
1297                Arc::new(TimestampSecondVector::from_vec(value.timestamp))
1298            }
1299            ConcreteDataType::Timestamp(TimestampType::Millisecond(_)) => {
1300                Arc::new(TimestampMillisecondVector::from_vec(value.timestamp))
1301            }
1302            ConcreteDataType::Timestamp(TimestampType::Microsecond(_)) => {
1303                Arc::new(TimestampMicrosecondVector::from_vec(value.timestamp))
1304            }
1305            ConcreteDataType::Timestamp(TimestampType::Nanosecond(_)) => {
1306                Arc::new(TimestampNanosecondVector::from_vec(value.timestamp))
1307            }
1308            _ => unreachable!(),
1309        };
1310
1311        if cfg!(debug_assertions) {
1312            debug_assert_eq!(timestamp.len(), sequence.len());
1313            debug_assert_eq!(timestamp.len(), op_type.len());
1314            for field in &fields {
1315                debug_assert_eq!(timestamp.len(), field.len());
1316            }
1317        }
1318
1319        Self {
1320            timestamp,
1321            sequence,
1322            op_type,
1323            fields,
1324        }
1325    }
1326}
1327
1328struct TimeSeriesIterBuilder {
1329    series_set: SeriesSet,
1330    projection: HashSet<ColumnId>,
1331    predicate: PredicateGroup,
1332    dedup: bool,
1333    sequence: Option<SequenceRange>,
1334    merge_mode: MergeMode,
1335    batch_to_record_batch: Arc<BatchToRecordBatchContext>,
1336}
1337
1338impl IterBuilder for TimeSeriesIterBuilder {
1339    fn build(&self, metrics: Option<MemScanMetrics>) -> Result<BoxedBatchIterator> {
1340        let iter = self.series_set.iter_series(
1341            self.projection.clone(),
1342            self.predicate.clone(),
1343            self.dedup,
1344            self.merge_mode,
1345            self.sequence,
1346            metrics,
1347        )?;
1348        if self.merge_mode == MergeMode::LastNonNull {
1349            let iter = LastNonNullIter::new(iter);
1350            Ok(Box::new(iter))
1351        } else {
1352            Ok(Box::new(iter))
1353        }
1354    }
1355
1356    fn is_record_batch(&self) -> bool {
1357        true
1358    }
1359
1360    fn build_record_batch(
1361        &self,
1362        time_range: Option<(Timestamp, Timestamp)>,
1363        metrics: Option<MemScanMetrics>,
1364    ) -> Result<BoxedRecordBatchIterator> {
1365        let iter = self.build(metrics)?;
1366        let iter: BoxedBatchIterator = if let Some(time_range) = time_range {
1367            let time_filters = self.predicate.time_filters();
1368            Box::new(PruneTimeIterator::new(iter, time_range, time_filters))
1369        } else {
1370            iter
1371        };
1372        Ok(self.batch_to_record_batch.adapt_iter(iter))
1373    }
1374}
1375
1376#[cfg(test)]
1377mod tests {
1378    use std::collections::{HashMap, HashSet};
1379
1380    use api::helper::ColumnDataTypeWrapper;
1381    use api::v1::helper::row;
1382    use api::v1::value::ValueData;
1383    use api::v1::{Mutation, Rows, SemanticType};
1384    use common_time::Timestamp;
1385    use datatypes::prelude::{ConcreteDataType, ScalarVector};
1386    use datatypes::schema::ColumnSchema;
1387    use datatypes::value::{OrderedFloat, Value};
1388    use datatypes::vectors::{Float64Vector, Int64Vector, TimestampMillisecondVector};
1389    use mito_codec::row_converter::SortField;
1390    use store_api::metadata::{ColumnMetadata, RegionMetadataBuilder};
1391    use store_api::storage::RegionId;
1392
1393    use super::*;
1394    use crate::memtable::builder::StringBuilder;
1395    use crate::test_util::column_metadata_to_column_schema;
1396
1397    fn schema_for_test() -> RegionMetadataRef {
1398        let mut builder = RegionMetadataBuilder::new(RegionId::new(123, 456));
1399        builder
1400            .push_column_metadata(ColumnMetadata {
1401                column_schema: ColumnSchema::new("k0", ConcreteDataType::string_datatype(), false),
1402                semantic_type: SemanticType::Tag,
1403                column_id: 0,
1404            })
1405            .push_column_metadata(ColumnMetadata {
1406                column_schema: ColumnSchema::new("k1", ConcreteDataType::int64_datatype(), false),
1407                semantic_type: SemanticType::Tag,
1408                column_id: 1,
1409            })
1410            .push_column_metadata(ColumnMetadata {
1411                column_schema: ColumnSchema::new(
1412                    "ts",
1413                    ConcreteDataType::timestamp_millisecond_datatype(),
1414                    false,
1415                ),
1416                semantic_type: SemanticType::Timestamp,
1417                column_id: 2,
1418            })
1419            .push_column_metadata(ColumnMetadata {
1420                column_schema: ColumnSchema::new("v0", ConcreteDataType::int64_datatype(), true),
1421                semantic_type: SemanticType::Field,
1422                column_id: 3,
1423            })
1424            .push_column_metadata(ColumnMetadata {
1425                column_schema: ColumnSchema::new("v1", ConcreteDataType::float64_datatype(), true),
1426                semantic_type: SemanticType::Field,
1427                column_id: 4,
1428            })
1429            .primary_key(vec![0, 1]);
1430        let region_metadata = builder.build().unwrap();
1431        Arc::new(region_metadata)
1432    }
1433
1434    fn ts_value_ref(val: i64) -> ValueRef<'static> {
1435        ValueRef::Timestamp(Timestamp::new_millisecond(val))
1436    }
1437
1438    fn field_value_ref(v0: i64, v1: f64) -> impl Iterator<Item = ValueRef<'static>> {
1439        vec![ValueRef::Int64(v0), ValueRef::Float64(OrderedFloat(v1))].into_iter()
1440    }
1441
1442    fn check_values(values: &Values, expect: &[(i64, u64, u8, i64, f64)]) {
1443        let ts = values
1444            .timestamp
1445            .as_any()
1446            .downcast_ref::<TimestampMillisecondVector>()
1447            .unwrap();
1448
1449        let v0 = values.fields[0]
1450            .as_any()
1451            .downcast_ref::<Int64Vector>()
1452            .unwrap();
1453        let v1 = values.fields[1]
1454            .as_any()
1455            .downcast_ref::<Float64Vector>()
1456            .unwrap();
1457        let read = ts
1458            .iter_data()
1459            .zip(values.sequence.iter_data())
1460            .zip(values.op_type.iter_data())
1461            .zip(v0.iter_data())
1462            .zip(v1.iter_data())
1463            .map(|((((ts, sequence), op_type), v0), v1)| {
1464                (
1465                    ts.unwrap().0.value(),
1466                    sequence.unwrap(),
1467                    op_type.unwrap(),
1468                    v0.unwrap(),
1469                    v1.unwrap(),
1470                )
1471            })
1472            .collect::<Vec<_>>();
1473        assert_eq!(expect, &read);
1474    }
1475
1476    #[test]
1477    fn test_series() {
1478        let region_metadata = schema_for_test();
1479        let mut series = Series::new(&region_metadata);
1480        series.push(ts_value_ref(1), 0, OpType::Put, field_value_ref(1, 10.1));
1481        series.push(ts_value_ref(2), 0, OpType::Put, field_value_ref(2, 10.2));
1482        assert_eq!(2, series.active.timestamp.len());
1483        assert_eq!(0, series.frozen.len());
1484
1485        let values = series.compact(&region_metadata).unwrap();
1486        check_values(values, &[(1, 0, 1, 1, 10.1), (2, 0, 1, 2, 10.2)]);
1487        assert_eq!(0, series.active.timestamp.len());
1488        assert_eq!(1, series.frozen.len());
1489    }
1490
1491    #[test]
1492    fn test_series_with_nulls() {
1493        let region_metadata = schema_for_test();
1494        let mut series = Series::new(&region_metadata);
1495        // col1: NULL 1 2 3
1496        // col2: NULL NULL 10.2 NULL
1497        series.push(
1498            ts_value_ref(1),
1499            0,
1500            OpType::Put,
1501            vec![ValueRef::Null, ValueRef::Null].into_iter(),
1502        );
1503        series.push(
1504            ts_value_ref(1),
1505            0,
1506            OpType::Put,
1507            vec![ValueRef::Int64(1), ValueRef::Null].into_iter(),
1508        );
1509        series.push(ts_value_ref(1), 2, OpType::Put, field_value_ref(2, 10.2));
1510        series.push(
1511            ts_value_ref(1),
1512            3,
1513            OpType::Put,
1514            vec![ValueRef::Int64(2), ValueRef::Null].into_iter(),
1515        );
1516        assert_eq!(4, series.active.timestamp.len());
1517        assert_eq!(0, series.frozen.len());
1518
1519        let values = series.compact(&region_metadata).unwrap();
1520        assert_eq!(values.fields[0].null_count(), 1);
1521        assert_eq!(values.fields[1].null_count(), 3);
1522        assert_eq!(0, series.active.timestamp.len());
1523        assert_eq!(1, series.frozen.len());
1524    }
1525
1526    fn check_value(batch: &Batch, expect: Vec<Vec<Value>>) {
1527        let ts_len = batch.timestamps().len();
1528        assert_eq!(batch.sequences().len(), ts_len);
1529        assert_eq!(batch.op_types().len(), ts_len);
1530        for f in batch.fields() {
1531            assert_eq!(f.data.len(), ts_len);
1532        }
1533
1534        let mut rows = vec![];
1535        for idx in 0..ts_len {
1536            let mut row = Vec::with_capacity(batch.fields().len() + 3);
1537            row.push(batch.timestamps().get(idx));
1538            row.push(batch.sequences().get(idx));
1539            row.push(batch.op_types().get(idx));
1540            row.extend(batch.fields().iter().map(|f| f.data.get(idx)));
1541            rows.push(row);
1542        }
1543
1544        assert_eq!(expect.len(), rows.len());
1545        for (idx, row) in rows.iter().enumerate() {
1546            assert_eq!(&expect[idx], row);
1547        }
1548    }
1549
1550    #[test]
1551    fn test_values_sort() {
1552        let schema = schema_for_test();
1553        let timestamp = Arc::new(TimestampMillisecondVector::from_vec(vec![1, 2, 3, 4, 3]));
1554        let sequence = Arc::new(UInt64Vector::from_vec(vec![1, 1, 1, 1, 2]));
1555        let op_type = Arc::new(UInt8Vector::from_vec(vec![1, 1, 1, 1, 0]));
1556
1557        let fields = vec![
1558            Arc::new(Int64Vector::from_vec(vec![4, 3, 2, 1, 2])) as Arc<_>,
1559            Arc::new(Float64Vector::from_vec(vec![1.1, 2.1, 4.2, 3.3, 4.2])) as Arc<_>,
1560        ];
1561        let values = Values {
1562            timestamp: timestamp as Arc<_>,
1563            sequence,
1564            op_type,
1565            fields,
1566        };
1567
1568        let batch = values
1569            .to_batch(
1570                b"test",
1571                &schema,
1572                &[0, 1, 2, 3, 4].into_iter().collect(),
1573                None,
1574                true,
1575                MergeMode::LastRow,
1576            )
1577            .unwrap();
1578        check_value(
1579            &batch,
1580            vec![
1581                vec![
1582                    Value::Timestamp(Timestamp::new_millisecond(1)),
1583                    Value::UInt64(1),
1584                    Value::UInt8(1),
1585                    Value::Int64(4),
1586                    Value::Float64(OrderedFloat(1.1)),
1587                ],
1588                vec![
1589                    Value::Timestamp(Timestamp::new_millisecond(2)),
1590                    Value::UInt64(1),
1591                    Value::UInt8(1),
1592                    Value::Int64(3),
1593                    Value::Float64(OrderedFloat(2.1)),
1594                ],
1595                vec![
1596                    Value::Timestamp(Timestamp::new_millisecond(3)),
1597                    Value::UInt64(2),
1598                    Value::UInt8(0),
1599                    Value::Int64(2),
1600                    Value::Float64(OrderedFloat(4.2)),
1601                ],
1602                vec![
1603                    Value::Timestamp(Timestamp::new_millisecond(4)),
1604                    Value::UInt64(1),
1605                    Value::UInt8(1),
1606                    Value::Int64(1),
1607                    Value::Float64(OrderedFloat(3.3)),
1608                ],
1609            ],
1610        )
1611    }
1612
1613    #[test]
1614    fn test_last_non_null_should_filter_by_sequence_before_merge_drop_ts() {
1615        let schema = schema_for_test();
1616        let projection: HashSet<_> = [0, 1, 2, 3, 4].into_iter().collect();
1617
1618        // Same timestamp, newest sequence is out of range. We should still keep the timestamp by
1619        // using the latest row *within* the sequence range as the base row.
1620        //
1621        // Expect after filtering seq<=2:
1622        // - base row: seq=2
1623        // - v0 from seq=2, v1 filled from seq=1
1624        let timestamp = Arc::new(TimestampMillisecondVector::from_vec(vec![1, 1, 1]));
1625        let sequence = Arc::new(UInt64Vector::from_vec(vec![1, 2, 3]));
1626        let op_type = Arc::new(UInt8Vector::from_vec(vec![OpType::Put as u8; 3]));
1627        let fields = vec![
1628            Arc::new(Int64Vector::from(vec![None, Some(10), None])) as Arc<_>,
1629            Arc::new(Float64Vector::from(vec![Some(1.5), None, None])) as Arc<_>,
1630        ];
1631        let values = Values {
1632            timestamp: timestamp as Arc<_>,
1633            sequence,
1634            op_type,
1635            fields,
1636        };
1637
1638        let batch = values
1639            .to_batch(
1640                b"test",
1641                &schema,
1642                &projection,
1643                Some(SequenceRange::LtEq { max: 2 }),
1644                true,
1645                MergeMode::LastNonNull,
1646            )
1647            .unwrap();
1648
1649        check_value(
1650            &batch,
1651            vec![vec![
1652                Value::Timestamp(Timestamp::new_millisecond(1)),
1653                Value::UInt64(2),
1654                Value::UInt8(OpType::Put as u8),
1655                Value::Int64(10),
1656                Value::Float64(OrderedFloat(1.5)),
1657            ]],
1658        );
1659    }
1660
1661    #[test]
1662    fn test_last_non_null_should_filter_by_sequence_before_merge_no_fill_from_out_of_range_row() {
1663        let schema = schema_for_test();
1664        let projection: HashSet<_> = [0, 1, 2, 3, 4].into_iter().collect();
1665
1666        // Same timestamp, older sequence is out of range. We must not fill null fields using rows
1667        // that should be excluded by the sequence filter.
1668        //
1669        // Expect after filtering seq>1:
1670        // - keep only seq=2 row, v0 stays NULL.
1671        let timestamp = Arc::new(TimestampMillisecondVector::from_vec(vec![1, 1]));
1672        let sequence = Arc::new(UInt64Vector::from_vec(vec![1, 2]));
1673        let op_type = Arc::new(UInt8Vector::from_vec(vec![OpType::Put as u8; 2]));
1674        let fields = vec![
1675            Arc::new(Int64Vector::from(vec![Some(10), None])) as Arc<_>,
1676            Arc::new(Float64Vector::from(vec![Some(1.0), Some(1.0)])) as Arc<_>,
1677        ];
1678        let values = Values {
1679            timestamp: timestamp as Arc<_>,
1680            sequence,
1681            op_type,
1682            fields,
1683        };
1684
1685        let batch = values
1686            .to_batch(
1687                b"test",
1688                &schema,
1689                &projection,
1690                Some(SequenceRange::Gt { min: 1 }),
1691                true,
1692                MergeMode::LastNonNull,
1693            )
1694            .unwrap();
1695
1696        check_value(
1697            &batch,
1698            vec![vec![
1699                Value::Timestamp(Timestamp::new_millisecond(1)),
1700                Value::UInt64(2),
1701                Value::UInt8(OpType::Put as u8),
1702                Value::Null,
1703                Value::Float64(OrderedFloat(1.0)),
1704            ]],
1705        );
1706    }
1707
1708    fn build_key_values(schema: &RegionMetadataRef, k0: String, k1: i64, len: usize) -> KeyValues {
1709        let column_schema = schema
1710            .column_metadatas
1711            .iter()
1712            .map(|c| api::v1::ColumnSchema {
1713                column_name: c.column_schema.name.clone(),
1714                datatype: ColumnDataTypeWrapper::try_from(c.column_schema.data_type.clone())
1715                    .unwrap()
1716                    .datatype() as i32,
1717                semantic_type: c.semantic_type as i32,
1718                ..Default::default()
1719            })
1720            .collect();
1721
1722        let rows = (0..len)
1723            .map(|i| {
1724                row(vec![
1725                    ValueData::StringValue(k0.clone()),
1726                    ValueData::I64Value(k1),
1727                    ValueData::TimestampMillisecondValue(i as i64),
1728                    ValueData::I64Value(i as i64),
1729                    ValueData::F64Value(i as f64),
1730                ])
1731            })
1732            .collect();
1733        let mutation = api::v1::Mutation {
1734            op_type: 1,
1735            sequence: 0,
1736            rows: Some(Rows {
1737                schema: column_schema,
1738                rows,
1739            }),
1740            write_hint: None,
1741        };
1742        KeyValues::new(schema.as_ref(), mutation).unwrap()
1743    }
1744
1745    #[test]
1746    fn test_series_set_concurrency() {
1747        let schema = schema_for_test();
1748        let row_codec = Arc::new(DensePrimaryKeyCodec::with_fields(
1749            schema
1750                .primary_key_columns()
1751                .map(|c| {
1752                    (
1753                        c.column_id,
1754                        SortField::new(c.column_schema.data_type.clone()),
1755                    )
1756                })
1757                .collect(),
1758        ));
1759        let set = Arc::new(SeriesSet::new(schema.clone(), row_codec));
1760
1761        let concurrency = 32;
1762        let pk_num = concurrency * 2;
1763        let mut handles = Vec::with_capacity(concurrency);
1764        for i in 0..concurrency {
1765            let set = set.clone();
1766            let schema = schema.clone();
1767            let column_schemas = schema
1768                .column_metadatas
1769                .iter()
1770                .map(column_metadata_to_column_schema)
1771                .collect::<Vec<_>>();
1772            let handle = std::thread::spawn(move || {
1773                for j in i * 100..(i + 1) * 100 {
1774                    let pk = j % pk_num;
1775                    let primary_key = format!("pk-{}", pk).as_bytes().to_vec();
1776
1777                    let kvs = KeyValues::new(
1778                        &schema,
1779                        Mutation {
1780                            op_type: OpType::Put as i32,
1781                            sequence: j as u64,
1782                            rows: Some(Rows {
1783                                schema: column_schemas.clone(),
1784                                rows: vec![row(vec![
1785                                    ValueData::StringValue(format!("{}", j)),
1786                                    ValueData::I64Value(j as i64),
1787                                    ValueData::TimestampMillisecondValue(j as i64),
1788                                    ValueData::I64Value(j as i64),
1789                                    ValueData::F64Value(j as f64),
1790                                ])],
1791                            }),
1792                            write_hint: None,
1793                        },
1794                    )
1795                    .unwrap();
1796                    set.push_to_series(primary_key, &kvs.iter().next().unwrap());
1797                }
1798            });
1799            handles.push(handle);
1800        }
1801        for h in handles {
1802            h.join().unwrap();
1803        }
1804
1805        let mut timestamps = Vec::with_capacity(concurrency * 100);
1806        let mut sequences = Vec::with_capacity(concurrency * 100);
1807        let mut op_types = Vec::with_capacity(concurrency * 100);
1808        let mut v0 = Vec::with_capacity(concurrency * 100);
1809
1810        for i in 0..pk_num {
1811            let pk = format!("pk-{}", i).as_bytes().to_vec();
1812            let series = set.get_series(&pk).unwrap();
1813            let mut guard = series.write().unwrap();
1814            let values = guard.compact(&schema).unwrap();
1815            timestamps.extend(values.sequence.iter_data().map(|v| v.unwrap() as i64));
1816            sequences.extend(values.sequence.iter_data().map(|v| v.unwrap() as i64));
1817            op_types.extend(values.op_type.iter_data().map(|v| v.unwrap()));
1818            v0.extend(
1819                values
1820                    .fields
1821                    .first()
1822                    .unwrap()
1823                    .as_any()
1824                    .downcast_ref::<Int64Vector>()
1825                    .unwrap()
1826                    .iter_data()
1827                    .map(|v| v.unwrap()),
1828            );
1829        }
1830
1831        let expected_sequence = (0..(concurrency * 100) as i64).collect::<HashSet<_>>();
1832        assert_eq!(
1833            expected_sequence,
1834            sequences.iter().copied().collect::<HashSet<_>>()
1835        );
1836
1837        op_types.iter().all(|op| *op == OpType::Put as u8);
1838        assert_eq!(
1839            expected_sequence,
1840            timestamps.iter().copied().collect::<HashSet<_>>()
1841        );
1842
1843        assert_eq!(timestamps, sequences);
1844        assert_eq!(v0, timestamps);
1845    }
1846
1847    #[test]
1848    fn test_memtable() {
1849        common_telemetry::init_default_ut_logging();
1850        check_memtable_dedup(true);
1851        check_memtable_dedup(false);
1852    }
1853
1854    fn check_memtable_dedup(dedup: bool) {
1855        let schema = schema_for_test();
1856        let kvs = build_key_values(&schema, "hello".to_string(), 42, 100);
1857        let memtable = TimeSeriesMemtable::new(schema, 42, None, dedup, MergeMode::LastRow);
1858        memtable.write(&kvs).unwrap();
1859        memtable.write(&kvs).unwrap();
1860
1861        let mut expected_ts: HashMap<i64, usize> = HashMap::new();
1862        for ts in kvs.iter().map(|kv| {
1863            kv.timestamp()
1864                .try_into_timestamp()
1865                .unwrap()
1866                .unwrap()
1867                .value()
1868        }) {
1869            *expected_ts.entry(ts).or_default() += if dedup { 1 } else { 2 };
1870        }
1871
1872        let ranges = memtable.ranges(None, RangesOptions::default()).unwrap();
1873        let range = ranges.ranges.into_values().next().unwrap();
1874        let iter = range.build_iter().unwrap();
1875        let mut read = HashMap::new();
1876
1877        for ts in iter
1878            .flat_map(|batch| {
1879                batch
1880                    .unwrap()
1881                    .timestamps()
1882                    .as_any()
1883                    .downcast_ref::<TimestampMillisecondVector>()
1884                    .unwrap()
1885                    .iter_data()
1886                    .collect::<Vec<_>>()
1887                    .into_iter()
1888            })
1889            .map(|v| v.unwrap().0.value())
1890        {
1891            *read.entry(ts).or_default() += 1;
1892        }
1893        assert_eq!(expected_ts, read);
1894
1895        let stats = memtable.stats();
1896        assert!(stats.bytes_allocated() > 0);
1897        assert_eq!(
1898            Some((
1899                Timestamp::new_millisecond(0),
1900                Timestamp::new_millisecond(99)
1901            )),
1902            stats.time_range()
1903        );
1904    }
1905
1906    #[test]
1907    fn test_memtable_projection() {
1908        common_telemetry::init_default_ut_logging();
1909        let schema = schema_for_test();
1910        let kvs = build_key_values(&schema, "hello".to_string(), 42, 100);
1911        let memtable = TimeSeriesMemtable::new(schema, 42, None, true, MergeMode::LastRow);
1912        memtable.write(&kvs).unwrap();
1913
1914        let iter = memtable
1915            .ranges(Some(&[3]), RangesOptions::default())
1916            .unwrap()
1917            .build(None)
1918            .unwrap();
1919
1920        let mut v0_all = vec![];
1921
1922        for res in iter {
1923            let batch = res.unwrap();
1924            assert_eq!(1, batch.fields().len());
1925            let v0 = batch
1926                .fields()
1927                .first()
1928                .unwrap()
1929                .data
1930                .as_any()
1931                .downcast_ref::<Int64Vector>()
1932                .unwrap();
1933            v0_all.extend(v0.iter_data().map(|v| v.unwrap()));
1934        }
1935        assert_eq!((0..100i64).collect::<Vec<_>>(), v0_all);
1936    }
1937
1938    #[test]
1939    fn test_memtable_concurrent_write_read() {
1940        common_telemetry::init_default_ut_logging();
1941        let schema = schema_for_test();
1942        let memtable = Arc::new(TimeSeriesMemtable::new(
1943            schema.clone(),
1944            42,
1945            None,
1946            true,
1947            MergeMode::LastRow,
1948        ));
1949
1950        // Number of writer threads
1951        let num_writers = 10;
1952        // Number of reader threads
1953        let num_readers = 5;
1954        // Number of series per writer
1955        let series_per_writer = 100;
1956        // Number of rows per series
1957        let rows_per_series = 10;
1958        // Total number of series
1959        let total_series = num_writers * series_per_writer;
1960
1961        // Create a barrier to synchronize the start of all threads
1962        let barrier = Arc::new(std::sync::Barrier::new(num_writers + num_readers + 1));
1963
1964        // Spawn writer threads
1965        let mut writer_handles = Vec::with_capacity(num_writers);
1966        for writer_id in 0..num_writers {
1967            let memtable = memtable.clone();
1968            let schema = schema.clone();
1969            let barrier = barrier.clone();
1970
1971            let handle = std::thread::spawn(move || {
1972                // Wait for all threads to be ready
1973                barrier.wait();
1974
1975                // Create and write series
1976                for series_id in 0..series_per_writer {
1977                    let series_key = format!("writer-{}-series-{}", writer_id, series_id);
1978                    let kvs =
1979                        build_key_values(&schema, series_key, series_id as i64, rows_per_series);
1980                    memtable.write(&kvs).unwrap();
1981                }
1982            });
1983
1984            writer_handles.push(handle);
1985        }
1986
1987        // Spawn reader threads
1988        let mut reader_handles = Vec::with_capacity(num_readers);
1989        for _ in 0..num_readers {
1990            let memtable = memtable.clone();
1991            let barrier = barrier.clone();
1992
1993            let handle = std::thread::spawn(move || {
1994                barrier.wait();
1995
1996                for _ in 0..10 {
1997                    let iter = memtable
1998                        .ranges(None, RangesOptions::default())
1999                        .unwrap()
2000                        .build(None)
2001                        .unwrap();
2002                    for batch_result in iter {
2003                        let _ = batch_result.unwrap();
2004                    }
2005                }
2006            });
2007
2008            reader_handles.push(handle);
2009        }
2010
2011        barrier.wait();
2012
2013        for handle in writer_handles {
2014            handle.join().unwrap();
2015        }
2016        for handle in reader_handles {
2017            handle.join().unwrap();
2018        }
2019
2020        let iter = memtable
2021            .ranges(None, RangesOptions::default())
2022            .unwrap()
2023            .build(None)
2024            .unwrap();
2025        let mut series_count = 0;
2026        let mut row_count = 0;
2027
2028        for batch_result in iter {
2029            let batch = batch_result.unwrap();
2030            series_count += 1;
2031            row_count += batch.num_rows();
2032        }
2033        assert_eq!(total_series, series_count);
2034        assert_eq!(total_series * rows_per_series, row_count);
2035    }
2036
2037    #[test]
2038    fn test_build_record_batch_iter_from_memtable() {
2039        let schema = schema_for_test();
2040        let memtable = TimeSeriesMemtable::new(schema.clone(), 1, None, true, MergeMode::LastRow);
2041
2042        let kvs = build_key_values(&schema, "test".to_string(), 1, 10);
2043        memtable.write(&kvs).unwrap();
2044
2045        let read_column_ids: Vec<ColumnId> = schema
2046            .column_metadatas
2047            .iter()
2048            .map(|c| c.column_id)
2049            .collect();
2050        let ranges = memtable
2051            .ranges(Some(&read_column_ids), RangesOptions::default())
2052            .unwrap();
2053        assert_eq!(1, ranges.ranges.len());
2054
2055        let range = ranges.ranges.into_values().next().unwrap();
2056        let mut iter = range.build_record_batch_iter(None, None).unwrap();
2057        let rb = iter.next().transpose().unwrap().unwrap();
2058        assert_eq!(10, rb.num_rows());
2059        // k0, k1 (pk columns), v0, v1 (field columns), ts, __primary_key, __sequence, __op_type
2060        let schema = rb.schema();
2061        let column_names: Vec<_> = schema.fields().iter().map(|f| f.name().as_str()).collect();
2062        assert_eq!(
2063            column_names,
2064            vec![
2065                "k0",
2066                "k1",
2067                "v0",
2068                "v1",
2069                "ts",
2070                "__primary_key",
2071                "__sequence",
2072                "__op_type",
2073            ]
2074        );
2075        assert!(iter.next().is_none());
2076    }
2077
2078    #[test]
2079    fn test_build_record_batch_iter_with_time_range() {
2080        let schema = schema_for_test();
2081        let memtable = TimeSeriesMemtable::new(schema.clone(), 1, None, true, MergeMode::LastRow);
2082
2083        let kvs = build_key_values(&schema, "test".to_string(), 1, 10);
2084        memtable.write(&kvs).unwrap();
2085
2086        let read_column_ids: Vec<ColumnId> = schema
2087            .column_metadatas
2088            .iter()
2089            .map(|c| c.column_id)
2090            .collect();
2091        let ranges = memtable
2092            .ranges(Some(&read_column_ids), RangesOptions::default())
2093            .unwrap();
2094        assert_eq!(1, ranges.ranges.len());
2095
2096        let time_range = (Timestamp::new_millisecond(3), Timestamp::new_millisecond(7));
2097
2098        let range = ranges.ranges.into_values().next().unwrap();
2099        let mut iter = range
2100            .build_record_batch_iter(Some(time_range), None)
2101            .unwrap();
2102
2103        let mut total_rows = 0;
2104        let mut all_timestamps = Vec::new();
2105        while let Some(rb) = iter.next().transpose().unwrap() {
2106            total_rows += rb.num_rows();
2107            let ts_col = rb
2108                .column_by_name("ts")
2109                .unwrap()
2110                .as_any()
2111                .downcast_ref::<datatypes::arrow::array::TimestampMillisecondArray>()
2112                .unwrap();
2113            for i in 0..ts_col.len() {
2114                all_timestamps.push(ts_col.value(i));
2115            }
2116        }
2117        assert_eq!(5, total_rows);
2118        all_timestamps.sort();
2119        assert_eq!(vec![3, 4, 5, 6, 7], all_timestamps);
2120    }
2121
2122    /// Helper to create a TimeSeriesIterBuilder from a memtable and schema.
2123    fn build_iter_builder(
2124        schema: &RegionMetadataRef,
2125        memtable: &TimeSeriesMemtable,
2126        projection: Option<&[ColumnId]>,
2127        dedup: bool,
2128        merge_mode: MergeMode,
2129        sequence: Option<SequenceRange>,
2130    ) -> TimeSeriesIterBuilder {
2131        let read_column_ids = read_column_ids_from_projection(schema, projection);
2132        let field_projection = if let Some(projection) = projection {
2133            projection.iter().copied().collect()
2134        } else {
2135            schema.field_columns().map(|c| c.column_id).collect()
2136        };
2137        let adapter_context = Arc::new(BatchToRecordBatchContext::new(
2138            schema.clone(),
2139            read_column_ids,
2140        ));
2141        TimeSeriesIterBuilder {
2142            series_set: memtable.series_set.clone(),
2143            projection: field_projection,
2144            predicate: PredicateGroup::default(),
2145            dedup,
2146            merge_mode,
2147            sequence,
2148            batch_to_record_batch: adapter_context,
2149        }
2150    }
2151
2152    #[test]
2153    fn test_iter_builder_build_record_batch_basic() {
2154        let schema = schema_for_test();
2155        let memtable = TimeSeriesMemtable::new(schema.clone(), 1, None, true, MergeMode::LastRow);
2156
2157        let kvs = build_key_values(&schema, "hello".to_string(), 42, 10);
2158        memtable.write(&kvs).unwrap();
2159
2160        let builder = build_iter_builder(&schema, &memtable, None, true, MergeMode::LastRow, None);
2161
2162        let mut iter = builder.build_record_batch(None, None).unwrap();
2163        let rb = iter.next().transpose().unwrap().unwrap();
2164        assert_eq!(10, rb.num_rows());
2165
2166        let rb_schema = rb.schema();
2167        let col_names: Vec<_> = rb_schema
2168            .fields()
2169            .iter()
2170            .map(|f| f.name().as_str())
2171            .collect();
2172        assert_eq!(
2173            col_names,
2174            vec![
2175                "k0",
2176                "k1",
2177                "v0",
2178                "v1",
2179                "ts",
2180                "__primary_key",
2181                "__sequence",
2182                "__op_type",
2183            ]
2184        );
2185
2186        assert!(iter.next().is_none());
2187    }
2188
2189    #[test]
2190    fn test_iter_builder_build_record_batch_with_projection() {
2191        let schema = schema_for_test();
2192        let memtable = TimeSeriesMemtable::new(schema.clone(), 1, None, true, MergeMode::LastRow);
2193
2194        let kvs = build_key_values(&schema, "test".to_string(), 1, 5);
2195        memtable.write(&kvs).unwrap();
2196
2197        // Project only field v0 (column_id=3) and ts (column_id=2).
2198        let projection = vec![2, 3];
2199        let builder = build_iter_builder(
2200            &schema,
2201            &memtable,
2202            Some(&projection),
2203            true,
2204            MergeMode::LastRow,
2205            None,
2206        );
2207
2208        let mut iter = builder.build_record_batch(None, None).unwrap();
2209        let rb = iter.next().transpose().unwrap().unwrap();
2210        assert_eq!(5, rb.num_rows());
2211
2212        let rb_schema = rb.schema();
2213        let col_names: Vec<_> = rb_schema
2214            .fields()
2215            .iter()
2216            .map(|f| f.name().as_str())
2217            .collect();
2218        // Only projected columns + internal columns.
2219        assert_eq!(
2220            col_names,
2221            vec!["v0", "ts", "__primary_key", "__sequence", "__op_type",]
2222        );
2223
2224        assert!(iter.next().is_none());
2225    }
2226
2227    #[test]
2228    fn test_iter_builder_build_record_batch_multiple_series() {
2229        let schema = schema_for_test();
2230        let memtable = TimeSeriesMemtable::new(schema.clone(), 1, None, true, MergeMode::LastRow);
2231
2232        let kvs_a = build_key_values(&schema, "aaa".to_string(), 1, 3);
2233        let kvs_b = build_key_values(&schema, "bbb".to_string(), 2, 4);
2234        memtable.write(&kvs_a).unwrap();
2235        memtable.write(&kvs_b).unwrap();
2236
2237        let builder = build_iter_builder(&schema, &memtable, None, true, MergeMode::LastRow, None);
2238
2239        let iter = builder.build_record_batch(None, None).unwrap();
2240        let mut total_rows = 0;
2241        for rb in iter {
2242            let rb = rb.unwrap();
2243            total_rows += rb.num_rows();
2244            assert_eq!(8, rb.num_columns());
2245        }
2246        assert_eq!(7, total_rows);
2247    }
2248
2249    #[test]
2250    fn test_iter_builder_build_record_batch_dedup() {
2251        let schema = schema_for_test();
2252        let memtable = TimeSeriesMemtable::new(schema.clone(), 1, None, true, MergeMode::LastRow);
2253
2254        // Write same data twice — dedup should keep only one copy per timestamp.
2255        let kvs = build_key_values(&schema, "dup".to_string(), 10, 5);
2256        memtable.write(&kvs).unwrap();
2257        memtable.write(&kvs).unwrap();
2258
2259        let builder = build_iter_builder(&schema, &memtable, None, true, MergeMode::LastRow, None);
2260
2261        let iter = builder.build_record_batch(None, None).unwrap();
2262        let total_rows: usize = iter.map(|rb| rb.unwrap().num_rows()).sum();
2263        assert_eq!(5, total_rows);
2264    }
2265
2266    #[test]
2267    fn test_iter_builder_build_record_batch_no_dedup() {
2268        let schema = schema_for_test();
2269        let memtable = TimeSeriesMemtable::new(schema.clone(), 1, None, false, MergeMode::LastRow);
2270
2271        let kvs = build_key_values(&schema, "dup".to_string(), 10, 5);
2272        memtable.write(&kvs).unwrap();
2273        memtable.write(&kvs).unwrap();
2274
2275        let builder = build_iter_builder(&schema, &memtable, None, false, MergeMode::LastRow, None);
2276
2277        let iter = builder.build_record_batch(None, None).unwrap();
2278        let total_rows: usize = iter.map(|rb| rb.unwrap().num_rows()).sum();
2279        assert_eq!(10, total_rows);
2280    }
2281
2282    #[test]
2283    fn test_iter_builder_build_record_batch_with_sequence_filter() {
2284        let schema = schema_for_test();
2285        let memtable = TimeSeriesMemtable::new(schema.clone(), 1, None, true, MergeMode::LastRow);
2286
2287        // build_key_values creates a mutation with base sequence=0.
2288        // Each row gets sequence = base + row_index, so 5 rows get sequences 0,1,2,3,4.
2289        let kvs = build_key_values(&schema, "seq".to_string(), 1, 5);
2290        memtable.write(&kvs).unwrap();
2291
2292        // Filter to sequence > 4 — should yield no rows.
2293        let builder = build_iter_builder(
2294            &schema,
2295            &memtable,
2296            None,
2297            true,
2298            MergeMode::LastRow,
2299            Some(SequenceRange::Gt { min: 4 }),
2300        );
2301
2302        let iter = builder.build_record_batch(None, None).unwrap();
2303        let total_rows: usize = iter.map(|rb| rb.unwrap().num_rows()).sum();
2304        assert_eq!(0, total_rows);
2305
2306        // Filter to sequence <= 2 — should yield 3 rows (sequences 0, 1, 2).
2307        let builder = build_iter_builder(
2308            &schema,
2309            &memtable,
2310            None,
2311            true,
2312            MergeMode::LastRow,
2313            Some(SequenceRange::LtEq { max: 2 }),
2314        );
2315
2316        let iter = builder.build_record_batch(None, None).unwrap();
2317        let total_rows: usize = iter.map(|rb| rb.unwrap().num_rows()).sum();
2318        assert_eq!(3, total_rows);
2319    }
2320
2321    #[test]
2322    fn test_iter_builder_build_record_batch_data_correctness() {
2323        use datatypes::arrow::array::{
2324            Float64Array, Int64Array, TimestampMillisecondArray, UInt8Array,
2325        };
2326
2327        let schema = schema_for_test();
2328        let memtable = TimeSeriesMemtable::new(schema.clone(), 1, None, true, MergeMode::LastRow);
2329
2330        let kvs = build_key_values(&schema, "check".to_string(), 7, 3);
2331        memtable.write(&kvs).unwrap();
2332
2333        let builder = build_iter_builder(&schema, &memtable, None, true, MergeMode::LastRow, None);
2334
2335        let mut iter = builder.build_record_batch(None, None).unwrap();
2336        let rb = iter.next().transpose().unwrap().unwrap();
2337        assert_eq!(3, rb.num_rows());
2338
2339        // Verify timestamp values.
2340        let ts_col = rb
2341            .column_by_name("ts")
2342            .unwrap()
2343            .as_any()
2344            .downcast_ref::<TimestampMillisecondArray>()
2345            .unwrap();
2346        let timestamps: Vec<_> = (0..ts_col.len()).map(|i| ts_col.value(i)).collect();
2347        assert_eq!(vec![0, 1, 2], timestamps);
2348
2349        // Verify field v0 values.
2350        let v0_col = rb
2351            .column_by_name("v0")
2352            .unwrap()
2353            .as_any()
2354            .downcast_ref::<Int64Array>()
2355            .unwrap();
2356        let v0_values: Vec<_> = (0..v0_col.len()).map(|i| v0_col.value(i)).collect();
2357        assert_eq!(vec![0, 1, 2], v0_values);
2358
2359        // Verify field v1 values.
2360        let v1_col = rb
2361            .column_by_name("v1")
2362            .unwrap()
2363            .as_any()
2364            .downcast_ref::<Float64Array>()
2365            .unwrap();
2366        let v1_values: Vec<_> = (0..v1_col.len()).map(|i| v1_col.value(i)).collect();
2367        assert_eq!(vec![0.0, 1.0, 2.0], v1_values);
2368
2369        // Verify op_type is all Put (1).
2370        let op_col = rb
2371            .column_by_name("__op_type")
2372            .unwrap()
2373            .as_any()
2374            .downcast_ref::<UInt8Array>()
2375            .unwrap();
2376        for i in 0..op_col.len() {
2377            assert_eq!(OpType::Put as u8, op_col.value(i));
2378        }
2379
2380        assert!(iter.next().is_none());
2381    }
2382
2383    #[test]
2384    fn test_can_accommodate_string_offset_overflow() {
2385        // A batch whose string data alone exceeds the Arrow string array offset limit can
2386        // never be accommodated, not even by an empty builder. `can_accommodate` must return
2387        // a structured `InvalidBatch` error instead of `Ok(false)`: the latter would make
2388        // `Series::extend` freeze the current builder and replay the same oversized batch
2389        // onto an empty builder, which then panics on offset overflow in `StringBuilder`.
2390        assert!(check_string_capacity(0, Some(101), 100).is_err());
2391        assert!(check_string_capacity(0, None, i32::MAX).is_err());
2392
2393        // A batch that fits in an empty builder but not in the current one reports
2394        // `Ok(false)` so the caller can freeze and replay onto an empty builder.
2395        assert!(!check_string_capacity(95, Some(10), 100).unwrap());
2396        assert!(!check_string_capacity(91, Some(10), 100).unwrap());
2397
2398        // The current builder can accommodate the batch when the combined size fits,
2399        // including the case where it exactly reaches the limit.
2400        assert!(check_string_capacity(0, Some(100), 100).unwrap());
2401        assert!(check_string_capacity(90, Some(10), 100).unwrap());
2402    }
2403
2404    #[test]
2405    fn test_can_accommodate_checks_all_string_fields() {
2406        // Two string fields: the first only overflows the *current* builder (an empty
2407        // builder would fit it), the second is intrinsically oversized. The scan must not
2408        // early-return `Ok(false)` on the first field: it has to check every string field so
2409        // the intrinsic oversize of the second field still surfaces as `InvalidBatch`.
2410        let limit = 10;
2411        // 8 bytes already in the current builder: 8 + 4 > 10, but 4 <= 10, so this field
2412        // alone fits an empty builder.
2413        let mut current = StringBuilder::with_capacity(1, 8);
2414        current.append("12345678");
2415        let field_builders = [Some(FieldBuilder::String(current)), None];
2416        let field_types = [
2417            ConcreteDataType::string_datatype(),
2418            ConcreteDataType::string_datatype(),
2419        ];
2420
2421        let first = Arc::new(StringVector::from(StringArray::from(vec!["abcd"]))) as VectorRef;
2422        // 15 bytes > limit: intrinsically oversized, can never fit any builder.
2423        let second = Arc::new(StringVector::from(StringArray::from(vec![
2424            "123456789012345",
2425        ]))) as VectorRef;
2426
2427        let err = scan_string_capacity(&[first, second], &field_builders, &field_types, limit)
2428            .unwrap_err();
2429        assert!(
2430            matches!(err, error::Error::InvalidBatch { .. }),
2431            "expected InvalidBatch, got {err:?}"
2432        );
2433
2434        // Sanity: without the oversized field the result is just `Ok(false)` (current full,
2435        // empty builder fits), which lets the caller freeze and replay onto an empty builder.
2436        let first_only = Arc::new(StringVector::from(StringArray::from(vec!["abcd"]))) as VectorRef;
2437        assert!(
2438            !scan_string_capacity(&[first_only], &field_builders, &field_types, limit).unwrap()
2439        );
2440    }
2441}