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