Skip to main content

mito2/memtable/
time_series.rs

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