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