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