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