Skip to main content

mito2/memtable/
simple_bulk_memtable.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15#[cfg(any(test, feature = "test"))]
16mod test_only;
17
18use std::collections::HashSet;
19use std::fmt::{Debug, Formatter};
20use std::sync::atomic::{AtomicI64, AtomicU64, AtomicUsize, Ordering};
21use std::sync::{Arc, RwLock};
22use std::time::{Duration, Instant};
23
24use api::v1::OpType;
25use datatypes::vectors::Helper;
26use mito_codec::key_values::KeyValue;
27use rayon::prelude::*;
28use snafu::{OptionExt, ResultExt};
29use store_api::metadata::RegionMetadataRef;
30use store_api::storage::{ColumnId, SequenceNumber};
31
32use crate::flush::WriteBufferManagerRef;
33use crate::memtable::bulk::part::BulkPart;
34use crate::memtable::stats::WriteMetrics;
35use crate::memtable::time_series::Series;
36use crate::memtable::{
37    AllocTracker, BatchToRecordBatchContext, BoxedBatchIterator, IterBuilder, KeyValues,
38    MemScanMetrics, Memtable, MemtableId, MemtableRange, MemtableRangeContext, MemtableRanges,
39    MemtableRef, MemtableStats, RangesOptions, read_column_ids_from_projection,
40};
41use crate::metrics::MEMTABLE_ACTIVE_SERIES_COUNT;
42use crate::read::Batch;
43use crate::read::dedup::LastNonNullIter;
44use crate::region::options::MergeMode;
45use crate::{error, metrics};
46
47pub struct SimpleBulkMemtable {
48    id: MemtableId,
49    region_metadata: RegionMetadataRef,
50    alloc_tracker: AllocTracker,
51    max_timestamp: AtomicI64,
52    min_timestamp: AtomicI64,
53    max_sequence: AtomicU64,
54    min_sequence: AtomicU64,
55    dedup: bool,
56    merge_mode: MergeMode,
57    num_rows: AtomicUsize,
58    series: RwLock<Series>,
59}
60
61impl Drop for SimpleBulkMemtable {
62    fn drop(&mut self) {
63        MEMTABLE_ACTIVE_SERIES_COUNT.dec();
64    }
65}
66
67impl SimpleBulkMemtable {
68    pub fn new(
69        id: MemtableId,
70        region_metadata: RegionMetadataRef,
71        write_buffer_manager: Option<WriteBufferManagerRef>,
72        dedup: bool,
73        merge_mode: MergeMode,
74    ) -> Self {
75        let series = RwLock::new(Series::with_capacity(&region_metadata, 1024, 8192));
76
77        Self {
78            id,
79            region_metadata,
80            alloc_tracker: AllocTracker::new(write_buffer_manager),
81            max_timestamp: AtomicI64::new(i64::MIN),
82            min_timestamp: AtomicI64::new(i64::MAX),
83            max_sequence: AtomicU64::new(0),
84            min_sequence: AtomicU64::new(u64::MAX),
85            dedup,
86            merge_mode,
87            num_rows: AtomicUsize::new(0),
88            series,
89        }
90    }
91
92    fn build_projection(&self, projection: Option<&[ColumnId]>) -> HashSet<ColumnId> {
93        if let Some(projection) = projection {
94            projection.iter().copied().collect()
95        } else {
96            self.region_metadata
97                .field_columns()
98                .map(|c| c.column_id)
99                .collect()
100        }
101    }
102
103    fn write_key_value(&self, kv: KeyValue, stats: &mut WriteMetrics) {
104        let ts = kv.timestamp();
105        let sequence = kv.sequence();
106        let op_type = kv.op_type();
107        let mut series = self.series.write().unwrap();
108        let size = series.push(ts, sequence, op_type, kv.fields());
109        stats.value_bytes += size;
110        // safety: timestamp of kv must be both present and a valid timestamp value.
111        let ts = kv
112            .timestamp()
113            .try_into_timestamp()
114            .unwrap()
115            .unwrap()
116            .value();
117        stats.min_ts = stats.min_ts.min(ts);
118        stats.max_ts = stats.max_ts.max(ts);
119        stats.min_sequence = stats.min_sequence.min(sequence);
120    }
121
122    /// Updates memtable stats.
123    fn update_stats(&self, stats: WriteMetrics) {
124        self.alloc_tracker
125            .on_allocation(stats.key_bytes + stats.value_bytes);
126        self.num_rows.fetch_add(stats.num_rows, Ordering::SeqCst);
127        self.max_timestamp.fetch_max(stats.max_ts, Ordering::SeqCst);
128        self.min_timestamp.fetch_min(stats.min_ts, Ordering::SeqCst);
129        self.max_sequence
130            .fetch_max(stats.max_sequence, Ordering::SeqCst);
131        self.min_sequence
132            .fetch_min(stats.min_sequence, Ordering::SeqCst);
133    }
134
135    #[cfg(test)]
136    fn schema(&self) -> &RegionMetadataRef {
137        &self.region_metadata
138    }
139}
140
141impl Debug for SimpleBulkMemtable {
142    fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
143        f.debug_struct("SimpleBulkMemtable").finish()
144    }
145}
146
147impl Memtable for SimpleBulkMemtable {
148    fn id(&self) -> MemtableId {
149        self.id
150    }
151
152    fn write(&self, kvs: &KeyValues) -> error::Result<()> {
153        let mut stats = WriteMetrics::default();
154        let max_sequence = kvs.max_sequence();
155        for kv in kvs.iter() {
156            self.write_key_value(kv, &mut stats);
157        }
158        stats.max_sequence = max_sequence;
159        stats.num_rows = kvs.num_rows();
160        self.update_stats(stats);
161        Ok(())
162    }
163
164    fn write_one(&self, kv: KeyValue) -> error::Result<()> {
165        debug_assert_eq!(0, kv.num_primary_keys());
166        let mut stats = WriteMetrics::default();
167        self.write_key_value(kv, &mut stats);
168        stats.num_rows = 1;
169        stats.max_sequence = kv.sequence();
170        self.update_stats(stats);
171        Ok(())
172    }
173
174    fn write_bulk(&self, part: BulkPart) -> error::Result<()> {
175        let rb = &part.batch;
176
177        let ts = Helper::try_into_vector(
178            rb.column_by_name(&self.region_metadata.time_index_column().column_schema.name)
179                .with_context(|| error::InvalidRequestSnafu {
180                    region_id: self.region_metadata.region_id,
181                    reason: "Timestamp not found",
182                })?,
183        )
184        .context(error::ConvertVectorSnafu)?;
185
186        let sequence = part.sequence;
187
188        let fields: Vec<_> = self
189            .region_metadata
190            .field_columns()
191            .map(|f| {
192                let array = rb.column_by_name(&f.column_schema.name).ok_or_else(|| {
193                    error::InvalidRequestSnafu {
194                        region_id: self.region_metadata.region_id,
195                        reason: format!("Column {} not found", f.column_schema.name),
196                    }
197                    .build()
198                })?;
199                Helper::try_into_vector(array).context(error::ConvertVectorSnafu)
200            })
201            .collect::<error::Result<Vec<_>>>()?;
202
203        let mut series = self.series.write().unwrap();
204        let extend_timer = metrics::REGION_WORKER_HANDLE_WRITE_ELAPSED
205            .with_label_values(&["bulk_extend"])
206            .start_timer();
207        series.extend(ts, OpType::Put as u8, sequence, fields)?;
208        extend_timer.observe_duration();
209
210        self.update_stats(WriteMetrics {
211            key_bytes: 0,
212            value_bytes: part.estimated_size(),
213            min_ts: part.min_timestamp,
214            max_ts: part.max_timestamp,
215            num_rows: part.num_rows(),
216            max_sequence: sequence,
217            min_sequence: sequence,
218        });
219        Ok(())
220    }
221
222    fn ranges(
223        &self,
224        projection: Option<&[ColumnId]>,
225        options: RangesOptions,
226    ) -> error::Result<MemtableRanges> {
227        let predicate = options.predicate;
228        let sequence = options.sequence;
229        let start_time = Instant::now();
230        let read_column_ids = read_column_ids_from_projection(&self.region_metadata, projection);
231        let projection = Arc::new(self.build_projection(projection));
232
233        // Use the memtable's overall time range and max sequence for all ranges
234        let max_sequence = self.max_sequence.load(Ordering::Relaxed);
235        let min_sequence = self.min_sequence.load(Ordering::Relaxed);
236        let time_range = {
237            let num_rows = self.num_rows.load(Ordering::Relaxed);
238            if num_rows > 0 {
239                let ts_type = self.region_metadata.time_index_type();
240                let max_timestamp =
241                    ts_type.create_timestamp(self.max_timestamp.load(Ordering::Relaxed));
242                let min_timestamp =
243                    ts_type.create_timestamp(self.min_timestamp.load(Ordering::Relaxed));
244                Some((min_timestamp, max_timestamp))
245            } else {
246                None
247            }
248        };
249
250        let values = self.series.read().unwrap().read_to_values();
251        let batch_to_record_batch = Arc::new(BatchToRecordBatchContext::new(
252            self.region_metadata.clone(),
253            read_column_ids.clone(),
254        ));
255
256        let contexts = values
257            .into_par_iter()
258            .filter_map(|v| {
259                let filtered = match v.to_batch(
260                    &[],
261                    &self.region_metadata,
262                    &projection,
263                    sequence,
264                    self.dedup,
265                    self.merge_mode,
266                ) {
267                    Ok(filtered) => filtered,
268                    Err(e) => {
269                        return Some(Err(e));
270                    }
271                };
272                if filtered.is_empty() {
273                    None
274                } else {
275                    Some(Ok(filtered))
276                }
277            })
278            .map(|result| {
279                result.map(|batch| {
280                    let num_rows = batch.num_rows();
281                    let estimated_bytes = batch.memory_size();
282
283                    let range_stats = MemtableStats {
284                        estimated_bytes,
285                        time_range,
286                        num_rows,
287                        num_ranges: 1,
288                        max_sequence,
289                        min_sequence,
290                        series_count: 1,
291                    };
292
293                    let builder = BatchRangeBuilder {
294                        batch,
295                        merge_mode: self.merge_mode,
296                        scan_cost: start_time.elapsed(),
297                    };
298                    (
299                        range_stats,
300                        Arc::new(MemtableRangeContext::new_with_batch_to_record_batch(
301                            self.id,
302                            Box::new(builder),
303                            predicate.clone(),
304                            Some(batch_to_record_batch.clone()),
305                        )),
306                    )
307                })
308            })
309            .collect::<error::Result<Vec<_>>>()?;
310
311        let ranges = contexts
312            .into_iter()
313            .enumerate()
314            .map(|(idx, (range_stats, context))| (idx, MemtableRange::new(context, range_stats)))
315            .collect();
316
317        Ok(MemtableRanges { ranges })
318    }
319
320    fn is_empty(&self) -> bool {
321        self.series.read().unwrap().is_empty()
322    }
323
324    fn freeze(&self) -> error::Result<()> {
325        self.series.write().unwrap().freeze(&self.region_metadata);
326        Ok(())
327    }
328
329    fn stats(&self) -> MemtableStats {
330        let estimated_bytes = self.alloc_tracker.bytes_allocated();
331        let num_rows = self.num_rows.load(Ordering::Relaxed);
332        if num_rows == 0 {
333            // no rows ever written
334            return MemtableStats {
335                estimated_bytes,
336                time_range: None,
337                num_rows: 0,
338                num_ranges: 0,
339                max_sequence: 0,
340                min_sequence: 0,
341                series_count: 0,
342            };
343        }
344        let ts_type = self.region_metadata.time_index_type();
345        let max_timestamp = ts_type.create_timestamp(self.max_timestamp.load(Ordering::Relaxed));
346        let min_timestamp = ts_type.create_timestamp(self.min_timestamp.load(Ordering::Relaxed));
347        MemtableStats {
348            estimated_bytes,
349            time_range: Some((min_timestamp, max_timestamp)),
350            num_rows,
351            num_ranges: 1,
352            max_sequence: self.max_sequence.load(Ordering::Relaxed),
353            min_sequence: self.min_sequence.load(Ordering::Relaxed),
354            series_count: 1,
355        }
356    }
357
358    fn min_sequence(&self) -> SequenceNumber {
359        if self.num_rows.load(Ordering::Relaxed) == 0 {
360            return 0;
361        }
362        self.min_sequence.load(Ordering::Relaxed)
363    }
364
365    fn fork(&self, id: MemtableId, metadata: &RegionMetadataRef) -> MemtableRef {
366        Arc::new(Self::new(
367            id,
368            metadata.clone(),
369            self.alloc_tracker.write_buffer_manager(),
370            self.dedup,
371            self.merge_mode,
372        ))
373    }
374}
375
376#[derive(Clone)]
377pub struct BatchRangeBuilder {
378    pub batch: Batch,
379    pub merge_mode: MergeMode,
380    scan_cost: Duration,
381}
382
383impl IterBuilder for BatchRangeBuilder {
384    fn build(&self, metrics: Option<MemScanMetrics>) -> error::Result<BoxedBatchIterator> {
385        let batch = self.batch.clone();
386        if let Some(metrics) = metrics {
387            let inner = crate::memtable::MemScanMetricsData {
388                total_series: 1,
389                num_rows: batch.num_rows(),
390                num_batches: 1,
391                scan_cost: self.scan_cost,
392                ..Default::default()
393            };
394            metrics.merge_inner(&inner);
395        }
396
397        let iter = Iter {
398            batch: Some(Ok(batch)),
399        };
400
401        if self.merge_mode == MergeMode::LastNonNull {
402            Ok(Box::new(LastNonNullIter::new(iter)))
403        } else {
404            Ok(Box::new(iter))
405        }
406    }
407}
408
409struct Iter {
410    batch: Option<error::Result<Batch>>,
411}
412
413impl Iterator for Iter {
414    type Item = error::Result<Batch>;
415
416    fn next(&mut self) -> Option<Self::Item> {
417        self.batch.take()
418    }
419}
420
421#[cfg(test)]
422mod tests {
423    use std::sync::Arc;
424
425    use api::v1::helper::row;
426    use api::v1::value::ValueData;
427    use api::v1::{Mutation, OpType, Rows, SemanticType};
428    use common_recordbatch::DfRecordBatch;
429    use common_time::Timestamp;
430    use datatypes::arrow::array::{ArrayRef, Float64Array, RecordBatch, TimestampMillisecondArray};
431    use datatypes::arrow_array::StringArray;
432    use datatypes::data_type::ConcreteDataType;
433    use datatypes::prelude::{ScalarVector, Vector};
434    use datatypes::schema::ColumnSchema;
435    use datatypes::value::Value;
436    use datatypes::vectors::TimestampMillisecondVector;
437    use store_api::metadata::{ColumnMetadata, RegionMetadataBuilder};
438    use store_api::storage::{RegionId, SequenceNumber, SequenceRange};
439
440    use super::*;
441    use crate::region::options::MergeMode;
442    use crate::test_util::column_metadata_to_column_schema;
443
444    fn new_test_metadata() -> RegionMetadataRef {
445        let mut builder = RegionMetadataBuilder::new(1.into());
446        builder
447            .push_column_metadata(ColumnMetadata {
448                column_schema: ColumnSchema::new(
449                    "ts",
450                    ConcreteDataType::timestamp_millisecond_datatype(),
451                    false,
452                ),
453                semantic_type: SemanticType::Timestamp,
454                column_id: 1,
455            })
456            .push_column_metadata(ColumnMetadata {
457                column_schema: ColumnSchema::new("f1", ConcreteDataType::float64_datatype(), true),
458                semantic_type: SemanticType::Field,
459                column_id: 2,
460            })
461            .push_column_metadata(ColumnMetadata {
462                column_schema: ColumnSchema::new("f2", ConcreteDataType::string_datatype(), true),
463                semantic_type: SemanticType::Field,
464                column_id: 3,
465            });
466        Arc::new(builder.build().unwrap())
467    }
468
469    fn new_test_memtable(dedup: bool, merge_mode: MergeMode) -> SimpleBulkMemtable {
470        SimpleBulkMemtable::new(1, new_test_metadata(), None, dedup, merge_mode)
471    }
472
473    fn build_key_values(
474        metadata: &RegionMetadataRef,
475        sequence: SequenceNumber,
476        row_values: &[(i64, f64, String)],
477        op_type: OpType,
478    ) -> KeyValues {
479        let column_schemas: Vec<_> = metadata
480            .column_metadatas
481            .iter()
482            .map(column_metadata_to_column_schema)
483            .collect();
484
485        let rows: Vec<_> = row_values
486            .iter()
487            .map(|(ts, f1, f2)| {
488                row(vec![
489                    ValueData::TimestampMillisecondValue(*ts),
490                    ValueData::F64Value(*f1),
491                    ValueData::StringValue(f2.clone()),
492                ])
493            })
494            .collect();
495        let mutation = Mutation {
496            op_type: op_type as i32,
497            sequence,
498            rows: Some(Rows {
499                schema: column_schemas,
500                rows,
501            }),
502            write_hint: None,
503        };
504        KeyValues::new(metadata, mutation).unwrap()
505    }
506
507    #[test]
508    fn test_min_sequence_covers_single_and_batch_writes_and_fork() {
509        let memtable = new_test_memtable(true, MergeMode::LastNonNull);
510        let schema = memtable.schema();
511        let rows = [(1, 1.0, "a".to_string()), (2, 2.0, "b".to_string())];
512        let newer = build_key_values(schema, 100, &rows, OpType::Put);
513        assert_eq!(0, memtable.min_sequence());
514        memtable.write(&newer).unwrap();
515        assert_eq!(100, memtable.stats().min_sequence);
516        assert_eq!(100, memtable.min_sequence());
517        let older = build_key_values(schema, 10, &rows, OpType::Put);
518        memtable.write_one(older.iter().next().unwrap()).unwrap();
519        assert_eq!(10, memtable.stats().min_sequence);
520        assert_eq!(10, memtable.min_sequence());
521        assert_eq!(101, memtable.stats().max_sequence);
522        let fork = memtable.fork(2, schema);
523        assert!(fork.is_empty());
524        assert_eq!(0, fork.min_sequence());
525        fork.write(&newer).unwrap();
526        assert_eq!(100, fork.stats().min_sequence);
527        assert_eq!(100, fork.min_sequence());
528    }
529
530    #[test]
531    fn test_write_and_iter() {
532        let memtable = new_test_memtable(false, MergeMode::LastRow);
533        memtable
534            .write(&build_key_values(
535                &memtable.region_metadata,
536                0,
537                &[(1, 1.0, "a".to_string())],
538                OpType::Put,
539            ))
540            .unwrap();
541        memtable
542            .write(&build_key_values(
543                &memtable.region_metadata,
544                1,
545                &[(2, 2.0, "b".to_string())],
546                OpType::Put,
547            ))
548            .unwrap();
549
550        let mut iter = memtable
551            .ranges(None, RangesOptions::default())
552            .unwrap()
553            .build(None)
554            .unwrap();
555        let batch = iter.next().unwrap().unwrap();
556        assert_eq!(2, batch.num_rows());
557        assert_eq!(2, batch.fields().len());
558        let ts_v = batch
559            .timestamps()
560            .as_any()
561            .downcast_ref::<TimestampMillisecondVector>()
562            .unwrap();
563        assert_eq!(Value::Timestamp(Timestamp::new_millisecond(1)), ts_v.get(0));
564        assert_eq!(Value::Timestamp(Timestamp::new_millisecond(2)), ts_v.get(1));
565    }
566
567    #[test]
568    fn test_projection() {
569        let memtable = new_test_memtable(false, MergeMode::LastRow);
570        memtable
571            .write(&build_key_values(
572                &memtable.region_metadata,
573                0,
574                &[(1, 1.0, "a".to_string())],
575                OpType::Put,
576            ))
577            .unwrap();
578
579        let mut iter = memtable
580            .ranges(None, RangesOptions::default())
581            .unwrap()
582            .build(None)
583            .unwrap();
584        let batch = iter.next().unwrap().unwrap();
585        assert_eq!(1, batch.num_rows());
586        assert_eq!(2, batch.fields().len());
587
588        let ts_v = batch
589            .timestamps()
590            .as_any()
591            .downcast_ref::<TimestampMillisecondVector>()
592            .unwrap();
593        assert_eq!(Value::Timestamp(Timestamp::new_millisecond(1)), ts_v.get(0));
594
595        // Only project column 2 (f1)
596        let projection = vec![2];
597        let mut iter = memtable
598            .ranges(Some(&projection), RangesOptions::default())
599            .unwrap()
600            .build(None)
601            .unwrap();
602        let batch = iter.next().unwrap().unwrap();
603
604        assert_eq!(1, batch.num_rows());
605        assert_eq!(1, batch.fields().len()); // only f1
606        assert_eq!(2, batch.fields()[0].column_id);
607    }
608
609    #[test]
610    fn test_dedup() {
611        let memtable = new_test_memtable(true, MergeMode::LastRow);
612        memtable
613            .write(&build_key_values(
614                &memtable.region_metadata,
615                0,
616                &[(1, 1.0, "a".to_string())],
617                OpType::Put,
618            ))
619            .unwrap();
620        memtable
621            .write(&build_key_values(
622                &memtable.region_metadata,
623                1,
624                &[(1, 2.0, "b".to_string())],
625                OpType::Put,
626            ))
627            .unwrap();
628        let mut iter = memtable
629            .ranges(None, RangesOptions::default())
630            .unwrap()
631            .build(None)
632            .unwrap();
633        let batch = iter.next().unwrap().unwrap();
634
635        assert_eq!(1, batch.num_rows()); // deduped to 1 row
636        assert_eq!(2.0, batch.fields()[0].data.get(0).as_f64_lossy().unwrap()); // last write wins
637    }
638
639    #[test]
640    fn test_write_one() {
641        let memtable = new_test_memtable(false, MergeMode::LastRow);
642        let kvs = build_key_values(
643            &memtable.region_metadata,
644            0,
645            &[(1, 1.0, "a".to_string())],
646            OpType::Put,
647        );
648        let kv = kvs.iter().next().unwrap();
649        memtable.write_one(kv).unwrap();
650
651        let mut iter = memtable
652            .ranges(None, RangesOptions::default())
653            .unwrap()
654            .build(None)
655            .unwrap();
656        let batch = iter.next().unwrap().unwrap();
657        assert_eq!(1, batch.num_rows());
658    }
659
660    #[tokio::test]
661    async fn test_single_range() {
662        let memtable = new_test_memtable(true, MergeMode::LastRow);
663        let kvs = build_key_values(
664            &memtable.region_metadata,
665            0,
666            &[(1, 1.0, "a".to_string())],
667            OpType::Put,
668        );
669        memtable.write_one(kvs.iter().next().unwrap()).unwrap();
670
671        let kvs = build_key_values(
672            &memtable.region_metadata,
673            1,
674            &[(1, 2.0, "b".to_string())],
675            OpType::Put,
676        );
677        memtable.write_one(kvs.iter().next().unwrap()).unwrap();
678        memtable.freeze().unwrap();
679
680        let ranges = memtable.ranges(None, RangesOptions::default()).unwrap();
681        assert_eq!(ranges.ranges.len(), 1);
682        let range = ranges.ranges.into_values().next().unwrap();
683        let mut reader = range.context.builder.build(None).unwrap();
684
685        let mut num_rows = 0;
686        while let Some(b) = reader.next().transpose().unwrap() {
687            num_rows += b.num_rows();
688            assert_eq!(b.fields()[1].data.get(0).as_string(), Some("b".to_string()));
689        }
690        assert_eq!(num_rows, 1);
691    }
692
693    #[test]
694    fn test_write_bulk() {
695        let memtable = new_test_memtable(false, MergeMode::LastRow);
696        let arrow_schema = memtable.schema().schema.arrow_schema().clone();
697        let arrays = vec![
698            Arc::new(TimestampMillisecondArray::from(vec![1, 2])) as ArrayRef,
699            Arc::new(Float64Array::from(vec![1.0, 2.0])) as ArrayRef,
700            Arc::new(StringArray::from(vec!["a", "b"])) as ArrayRef,
701        ];
702        let rb = DfRecordBatch::try_new(arrow_schema, arrays).unwrap();
703
704        let part = BulkPart {
705            batch: rb,
706            sequence: 1,
707            min_sequence: 1,
708            min_timestamp: 1,
709            max_timestamp: 2,
710            timestamp_index: 0,
711            raw_data: None,
712        };
713        memtable.write_bulk(part).unwrap();
714
715        let mut iter = memtable
716            .ranges(None, RangesOptions::default())
717            .unwrap()
718            .build(None)
719            .unwrap();
720        let batch = iter.next().unwrap().unwrap();
721        assert_eq!(2, batch.num_rows());
722
723        let stats = memtable.stats();
724        assert_eq!(1, stats.max_sequence);
725        assert_eq!(2, stats.num_rows);
726        assert_eq!(
727            Some((Timestamp::new_millisecond(1), Timestamp::new_millisecond(2))),
728            stats.time_range
729        );
730
731        let kvs = build_key_values(
732            &memtable.region_metadata,
733            2,
734            &[(3, 3.0, "c".to_string())],
735            OpType::Put,
736        );
737        memtable.write(&kvs).unwrap();
738        let mut iter = memtable
739            .ranges(None, RangesOptions::default())
740            .unwrap()
741            .build(None)
742            .unwrap();
743        let batch = iter.next().unwrap().unwrap();
744        assert_eq!(3, batch.num_rows());
745        assert_eq!(
746            vec![1, 2, 3],
747            batch
748                .timestamps()
749                .as_any()
750                .downcast_ref::<TimestampMillisecondVector>()
751                .unwrap()
752                .iter_data()
753                .map(|t| { t.unwrap().0.value() })
754                .collect::<Vec<_>>()
755        );
756    }
757
758    #[test]
759    fn test_is_empty() {
760        let memtable = new_test_memtable(false, MergeMode::LastRow);
761        assert!(memtable.is_empty());
762
763        memtable
764            .write(&build_key_values(
765                &memtable.region_metadata,
766                0,
767                &[(1, 1.0, "a".to_string())],
768                OpType::Put,
769            ))
770            .unwrap();
771        assert!(!memtable.is_empty());
772    }
773
774    #[test]
775    fn test_stats() {
776        let memtable = new_test_memtable(false, MergeMode::LastRow);
777        let stats = memtable.stats();
778        assert_eq!(0, stats.num_rows);
779        assert!(stats.time_range.is_none());
780
781        memtable
782            .write(&build_key_values(
783                &memtable.region_metadata,
784                0,
785                &[(1, 1.0, "a".to_string())],
786                OpType::Put,
787            ))
788            .unwrap();
789        let stats = memtable.stats();
790        assert_eq!(1, stats.num_rows);
791        assert!(stats.time_range.is_some());
792    }
793
794    #[test]
795    fn test_fork() {
796        let memtable = new_test_memtable(false, MergeMode::LastRow);
797        memtable
798            .write(&build_key_values(
799                &memtable.region_metadata,
800                0,
801                &[(1, 1.0, "a".to_string())],
802                OpType::Put,
803            ))
804            .unwrap();
805
806        let forked = memtable.fork(2, &memtable.region_metadata);
807        assert!(forked.is_empty());
808    }
809
810    #[test]
811    fn test_sequence_filter() {
812        let memtable = new_test_memtable(false, MergeMode::LastRow);
813        memtable
814            .write(&build_key_values(
815                &memtable.region_metadata,
816                0,
817                &[(1, 1.0, "a".to_string())],
818                OpType::Put,
819            ))
820            .unwrap();
821        memtable
822            .write(&build_key_values(
823                &memtable.region_metadata,
824                1,
825                &[(2, 2.0, "b".to_string())],
826                OpType::Put,
827            ))
828            .unwrap();
829
830        // Filter with sequence 0 should only return first write
831        let mut iter = memtable
832            .ranges(
833                None,
834                RangesOptions {
835                    sequence: Some(SequenceRange::LtEq { max: 0 }),
836                    ..Default::default()
837                },
838            )
839            .unwrap()
840            .build(None)
841            .unwrap();
842        let batch = iter.next().unwrap().unwrap();
843        assert_eq!(1, batch.num_rows());
844        assert_eq!(1.0, batch.fields()[0].data.get(0).as_f64_lossy().unwrap());
845    }
846
847    fn rb_with_large_string(
848        ts: i64,
849        string_len: i32,
850        region_meta: &RegionMetadataRef,
851    ) -> RecordBatch {
852        let schema = region_meta.schema.arrow_schema().clone();
853        RecordBatch::try_new(
854            schema,
855            vec![
856                Arc::new(StringArray::from_iter_values(
857                    ["a".repeat(string_len as usize).clone()].into_iter(),
858                )) as ArrayRef,
859                Arc::new(TimestampMillisecondArray::from_iter_values(
860                    [ts].into_iter(),
861                )) as ArrayRef,
862            ],
863        )
864        .unwrap()
865    }
866
867    #[test]
868    fn test_write_read_large_string() {
869        let mut builder = RegionMetadataBuilder::new(RegionId::new(123, 456));
870        builder
871            .push_column_metadata(ColumnMetadata {
872                column_schema: ColumnSchema::new("k0", ConcreteDataType::string_datatype(), false),
873                semantic_type: SemanticType::Field,
874                column_id: 0,
875            })
876            .push_column_metadata(ColumnMetadata {
877                column_schema: ColumnSchema::new(
878                    "ts",
879                    ConcreteDataType::timestamp_millisecond_datatype(),
880                    false,
881                ),
882                semantic_type: SemanticType::Timestamp,
883                column_id: 1,
884            })
885            .primary_key(vec![]);
886        let region_meta = Arc::new(builder.build().unwrap());
887        let memtable =
888            SimpleBulkMemtable::new(0, region_meta.clone(), None, true, MergeMode::LastRow);
889        memtable
890            .write_bulk(BulkPart {
891                batch: rb_with_large_string(0, i32::MAX, &region_meta),
892                max_timestamp: 0,
893                min_timestamp: 0,
894                sequence: 0,
895                min_sequence: 0,
896                timestamp_index: 1,
897                raw_data: None,
898            })
899            .unwrap();
900
901        memtable.freeze().unwrap();
902        memtable
903            .write_bulk(BulkPart {
904                batch: rb_with_large_string(1, 3, &region_meta),
905                max_timestamp: 1,
906                min_timestamp: 1,
907                sequence: 1,
908                min_sequence: 1,
909                timestamp_index: 1,
910                raw_data: None,
911            })
912            .unwrap();
913        let MemtableRanges { ranges, .. } =
914            memtable.ranges(None, RangesOptions::default()).unwrap();
915        let mut rows = 0;
916        for range in ranges.into_values() {
917            let iter = range.build_iter().unwrap();
918            for batch in iter {
919                rows += batch.unwrap().num_rows();
920            }
921        }
922        assert_eq!(rows, 2);
923    }
924
925    #[test]
926    fn test_build_record_batch_iter_from_memtable() {
927        let memtable = new_test_memtable(false, MergeMode::LastRow);
928
929        let kvs = build_key_values(
930            &memtable.region_metadata,
931            0,
932            &[(1, 1.0, "a".to_string()), (2, 2.0, "b".to_string())],
933            OpType::Put,
934        );
935        memtable.write(&kvs).unwrap();
936
937        let read_column_ids: Vec<ColumnId> = memtable
938            .region_metadata
939            .column_metadatas
940            .iter()
941            .map(|c| c.column_id)
942            .collect();
943        let ranges = memtable
944            .ranges(Some(&read_column_ids), RangesOptions::default())
945            .unwrap();
946        assert!(!ranges.ranges.is_empty());
947
948        let mut total_rows = 0;
949        for range in ranges.ranges.into_values() {
950            let mut iter = range.build_record_batch_iter(None, None).unwrap();
951            while let Some(rb) = iter.next().transpose().unwrap() {
952                total_rows += rb.num_rows();
953                let schema = rb.schema();
954                let column_names: Vec<_> =
955                    schema.fields().iter().map(|f| f.name().as_str()).collect();
956                assert_eq!(
957                    column_names,
958                    vec!["f1", "f2", "ts", "__primary_key", "__sequence", "__op_type"]
959                );
960            }
961        }
962        assert_eq!(2, total_rows);
963    }
964}