1#[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(®ion_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 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 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 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 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 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()); 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()); assert_eq!(2.0, batch.fields()[0].data.get(0).as_f64_lossy().unwrap()); }
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 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, ®ion_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, ®ion_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}