1use std::collections::{HashMap, HashSet};
18use std::sync::{Arc, Mutex, MutexGuard};
19use std::time::Duration;
20
21use common_telemetry::debug;
22use common_time::Timestamp;
23use common_time::timestamp::TimeUnit;
24use common_time::timestamp_millis::BucketAligned;
25use datatypes::arrow;
26use datatypes::arrow::array::{
27 ArrayRef, BooleanArray, RecordBatch, RecordBatchOptions, TimestampMicrosecondArray,
28 TimestampMillisecondArray, TimestampNanosecondArray, TimestampSecondArray,
29};
30use datatypes::arrow::buffer::{BooleanBuffer, MutableBuffer};
31use datatypes::arrow::datatypes::{DataType, Int64Type, SchemaRef};
32use mito_codec::key_values::KeyValue;
33use mito_codec::row_converter::{PrimaryKeyCodec, build_primary_key_codec};
34use smallvec::{SmallVec, smallvec};
35use snafu::{OptionExt, ResultExt};
36use store_api::metadata::RegionMetadataRef;
37
38use crate::error;
39use crate::error::{InvalidRequestSnafu, Result};
40use crate::memtable::bulk::part::{BulkPart, BulkPartConverter};
41use crate::memtable::version::SmallMemtableVec;
42use crate::memtable::{KeyValues, MemtableBuilderRef, MemtableId, MemtableRef};
43use crate::sst::{FlatSchemaOptions, to_flat_sst_arrow_schema};
44
45const INITIAL_TIME_WINDOW: Duration = Duration::from_days(1);
47
48#[derive(Debug, Clone)]
50pub struct TimePartition {
51 memtable: MemtableRef,
53 time_range: PartTimeRange,
55}
56
57impl TimePartition {
58 fn contains_timestamp(&self, ts: Timestamp) -> bool {
60 self.time_range.contains_timestamp(ts)
61 }
62
63 fn write(&self, kvs: &KeyValues) -> Result<()> {
65 self.memtable.write(kvs)
66 }
67
68 fn write_record_batch(&self, rb: BulkPart) -> Result<()> {
70 self.memtable.write_bulk(rb)
71 }
72
73 fn write_record_batch_partial(&self, part: &BulkPart) -> Result<()> {
75 let Some(filtered) = filter_record_batch(
76 part,
77 self.time_range.min_timestamp.value(),
78 self.time_range.max_timestamp.value(),
79 )?
80 else {
81 return Ok(());
82 };
83 self.write_record_batch(filtered)
84 }
85}
86
87macro_rules! create_filter_buffer {
88 ($ts_array:expr, $min:expr, $max:expr) => {{
89 let len = $ts_array.len();
90 let mut buffer = MutableBuffer::new(len.div_ceil(64) * 8);
91
92 let f = |idx: usize| -> bool {
93 unsafe {
95 let val = $ts_array.value_unchecked(idx);
96 val >= $min && val < $max
97 }
98 };
99
100 let chunks = len / 64;
101 let remainder = len % 64;
102
103 for chunk in 0..chunks {
104 let mut packed = 0;
105 for bit_idx in 0..64 {
106 let i = bit_idx + chunk * 64;
107 packed |= (f(i) as u64) << bit_idx;
108 }
109 unsafe { buffer.push_unchecked(packed) }
111 }
112
113 if remainder != 0 {
114 let mut packed = 0;
115 for bit_idx in 0..remainder {
116 let i = bit_idx + chunks * 64;
117 packed |= (f(i) as u64) << bit_idx;
118 }
119 unsafe { buffer.push_unchecked(packed) }
121 }
122
123 BooleanArray::new(BooleanBuffer::new(buffer.into(), 0, len), None)
124 }};
125}
126
127macro_rules! handle_timestamp_array {
128 ($ts_array:expr, $array_type:ty, $min:expr, $max:expr) => {{
129 let ts_array = $ts_array.as_any().downcast_ref::<$array_type>().unwrap();
130 let filter = create_filter_buffer!(ts_array, $min, $max);
131
132 let res = arrow::compute::filter(ts_array, &filter).context(error::ComputeArrowSnafu)?;
133 if res.is_empty() {
134 return Ok(None);
135 }
136
137 let i64array = res.as_any().downcast_ref::<$array_type>().unwrap();
138 let max_ts = arrow::compute::max(i64array).unwrap();
140 let min_ts = arrow::compute::min(i64array).unwrap();
141
142 (res, filter, min_ts, max_ts)
143 }};
144}
145
146pub fn filter_record_batch(part: &BulkPart, min: i64, max: i64) -> Result<Option<BulkPart>> {
149 let ts_array = part.timestamps();
150 let (ts_array, filter, min_ts, max_ts) = match ts_array.data_type() {
151 DataType::Timestamp(unit, _) => match unit {
152 arrow::datatypes::TimeUnit::Second => {
153 handle_timestamp_array!(ts_array, TimestampSecondArray, min, max)
154 }
155 arrow::datatypes::TimeUnit::Millisecond => {
156 handle_timestamp_array!(ts_array, TimestampMillisecondArray, min, max)
157 }
158 arrow::datatypes::TimeUnit::Microsecond => {
159 handle_timestamp_array!(ts_array, TimestampMicrosecondArray, min, max)
160 }
161 arrow::datatypes::TimeUnit::Nanosecond => {
162 handle_timestamp_array!(ts_array, TimestampNanosecondArray, min, max)
163 }
164 },
165 _ => {
166 unreachable!("Got data type: {:?}", ts_array.data_type());
167 }
168 };
169
170 let num_rows = ts_array.len();
171 let arrays = part
172 .batch
173 .columns()
174 .iter()
175 .enumerate()
176 .map(|(index, array)| {
177 if index == part.timestamp_index {
178 Ok(ts_array.clone())
179 } else {
180 arrow::compute::filter(&array, &filter).context(error::ComputeArrowSnafu)
181 }
182 })
183 .collect::<Result<Vec<_>>>()?;
184 let batch = RecordBatch::try_new_with_options(
185 part.batch.schema(),
186 arrays,
187 &RecordBatchOptions::default().with_row_count(Some(num_rows)),
188 )
189 .context(error::NewRecordBatchSnafu)?;
190 Ok(Some(BulkPart {
191 batch,
192 max_timestamp: max_ts,
193 min_timestamp: min_ts,
194 sequence: part.sequence,
195 min_sequence: part.min_sequence,
196 timestamp_index: part.timestamp_index,
197 raw_data: None,
198 }))
199}
200
201type PartitionVec = SmallVec<[TimePartition; 2]>;
202
203#[derive(Debug)]
205pub struct TimePartitions {
206 inner: Mutex<PartitionsInner>,
208 part_duration: Duration,
210 metadata: RegionMetadataRef,
212 builder: MemtableBuilderRef,
214 primary_key_codec: Arc<dyn PrimaryKeyCodec>,
216
217 bulk_schema: Option<SchemaRef>,
220}
221
222pub type TimePartitionsRef = Arc<TimePartitions>;
223
224impl TimePartitions {
225 pub fn new(
227 metadata: RegionMetadataRef,
228 builder: MemtableBuilderRef,
229 next_memtable_id: MemtableId,
230 part_duration: Option<Duration>,
231 ) -> Self {
232 let inner = PartitionsInner::new(next_memtable_id);
233 let primary_key_codec = build_primary_key_codec(&metadata);
234 let bulk_schema = builder.use_bulk_insert(&metadata).then(|| {
235 let opts = FlatSchemaOptions::from_encoding(metadata.primary_key_encoding);
236 to_flat_sst_arrow_schema(&metadata, &opts)
237 });
238
239 Self {
240 inner: Mutex::new(inner),
241 part_duration: part_duration.unwrap_or(INITIAL_TIME_WINDOW),
242 metadata,
243 builder,
244 primary_key_codec,
245 bulk_schema,
246 }
247 }
248
249 pub fn write(&self, kvs: &KeyValues) -> Result<()> {
253 if let Some(bulk_schema) = &self.bulk_schema {
254 let mut converter = BulkPartConverter::new(
255 &self.metadata,
256 bulk_schema.clone(),
257 kvs.num_rows(),
258 self.primary_key_codec.clone(),
259 true,
261 );
262 converter.append_key_values(kvs)?;
263 let part = converter.convert()?;
264
265 return self.write_bulk_inner(part);
266 }
267
268 let parts = self.list_partitions();
270
271 for part in parts.iter().rev() {
274 let mut all_in_partition = true;
275 for kv in kvs.iter() {
276 let ts = kv.timestamp().try_into_timestamp().unwrap().unwrap();
278 if !part.contains_timestamp(ts) {
279 all_in_partition = false;
280 break;
281 }
282 }
283 if !all_in_partition {
284 continue;
285 }
286
287 return part.write(kvs);
289 }
290
291 self.write_multi_parts(kvs, &parts)
293 }
294
295 pub fn write_bulk(&self, part: BulkPart) -> Result<()> {
297 let part = if let Some(bulk_schema) = &self.bulk_schema {
299 let converted = crate::memtable::bulk::part::convert_bulk_part(
300 part,
301 &self.metadata,
302 self.primary_key_codec.clone(),
303 bulk_schema.clone(),
304 true,
306 )?;
307 match converted {
308 Some(p) => p,
309 None => return Ok(()),
310 }
311 } else {
312 part
313 };
314
315 self.write_bulk_inner(part)
316 }
317
318 fn write_bulk_inner(&self, part: BulkPart) -> Result<()> {
320 let time_type = self
321 .metadata
322 .time_index_column()
323 .column_schema
324 .data_type
325 .as_timestamp()
326 .unwrap();
327
328 let parts = self.list_partitions();
330 let (matching_parts, missing_parts) = self.find_partitions_by_time_range(
331 part.timestamps(),
332 &parts,
333 time_type.create_timestamp(part.min_timestamp),
334 time_type.create_timestamp(part.max_timestamp),
335 )?;
336
337 if matching_parts.len() == 1 && missing_parts.is_empty() {
338 return matching_parts[0].write_record_batch(part);
340 }
341
342 for matching in matching_parts {
343 matching.write_record_batch_partial(&part)?
344 }
345
346 for missing in missing_parts {
347 let new_part = {
348 let mut inner = self.inner.lock().unwrap();
349 self.get_or_create_time_partition(missing, &mut inner)?
350 };
351 new_part.write_record_batch_partial(&part)?;
352 }
353 Ok(())
354 }
355
356 fn get_or_create_time_partition(
359 &self,
360 part_start: Timestamp,
361 inner: &mut MutexGuard<PartitionsInner>,
362 ) -> Result<TimePartition> {
363 let part_pos = match inner
364 .parts
365 .iter()
366 .position(|part| part.time_range.min_timestamp == part_start)
367 {
368 Some(pos) => pos,
369 None => {
370 let range = PartTimeRange::from_start_duration(part_start, self.part_duration)
371 .with_context(|| InvalidRequestSnafu {
372 region_id: self.metadata.region_id,
373 reason: format!(
374 "Partition time range for {part_start:?} is out of bound, bucket size: {:?}", self.part_duration
375 ),
376 })?;
377 let memtable = self
378 .builder
379 .build(inner.alloc_memtable_id(), &self.metadata);
380 debug!(
381 "Create time partition {:?} for region {}, duration: {:?}, memtable_id: {}, parts_total: {}, metadata: {:?}",
382 range,
383 self.metadata.region_id,
384 self.part_duration,
385 memtable.id(),
386 inner.parts.len() + 1,
387 self.metadata,
388 );
389 let pos = inner.parts.len();
390 inner.parts.push(TimePartition {
391 memtable,
392 time_range: range,
393 });
394 pos
395 }
396 };
397 Ok(inner.parts[part_pos].clone())
398 }
399
400 pub fn list_memtables(&self, memtables: &mut Vec<MemtableRef>) {
402 let inner = self.inner.lock().unwrap();
403 memtables.extend(inner.parts.iter().map(|part| part.memtable.clone()));
404 }
405
406 pub fn num_partitions(&self) -> usize {
408 let inner = self.inner.lock().unwrap();
409 inner.parts.len()
410 }
411
412 pub fn is_empty(&self) -> bool {
414 let inner = self.inner.lock().unwrap();
415 inner.parts.iter().all(|part| part.memtable.is_empty())
416 }
417
418 pub fn freeze(&self) -> Result<()> {
420 let inner = self.inner.lock().unwrap();
421 for part in &*inner.parts {
422 part.memtable.freeze()?;
423 }
424 Ok(())
425 }
426
427 pub fn fork(&self, metadata: &RegionMetadataRef, part_duration: Option<Duration>) -> Self {
429 let part_duration = part_duration.unwrap_or(self.part_duration);
431
432 let mut inner = self.inner.lock().unwrap();
433 let latest_part = inner
434 .parts
435 .iter()
436 .max_by_key(|part| part.time_range.min_timestamp)
437 .cloned();
438
439 let Some(old_part) = latest_part else {
440 return Self::new(
442 metadata.clone(),
443 self.builder.clone(),
444 inner.next_memtable_id,
445 Some(part_duration),
446 );
447 };
448
449 let old_stats = old_part.memtable.stats();
450 let partitions_inner = old_stats
452 .time_range()
453 .and_then(|(_, old_stats_end_timestamp)| {
454 partition_start_timestamp(old_stats_end_timestamp, part_duration)
455 .and_then(|start| PartTimeRange::from_start_duration(start, part_duration))
456 })
457 .map(|part_time_range| {
458 let memtable = old_part.memtable.fork(inner.alloc_memtable_id(), metadata);
460 let part = TimePartition {
461 memtable,
462 time_range: part_time_range,
463 };
464 PartitionsInner::with_partition(part, inner.next_memtable_id)
465 })
466 .unwrap_or_else(|| PartitionsInner::new(inner.next_memtable_id));
467
468 Self {
469 inner: Mutex::new(partitions_inner),
470 part_duration,
471 metadata: metadata.clone(),
472 builder: self.builder.clone(),
473 primary_key_codec: self.primary_key_codec.clone(),
474 bulk_schema: self.bulk_schema.clone(),
475 }
476 }
477
478 pub(crate) fn part_duration(&self) -> Duration {
480 self.part_duration
481 }
482
483 pub(crate) fn memtable_builder(&self) -> &MemtableBuilderRef {
485 &self.builder
486 }
487
488 pub(crate) fn memory_usage(&self) -> usize {
490 let inner = self.inner.lock().unwrap();
491 inner
492 .parts
493 .iter()
494 .map(|part| part.memtable.stats().estimated_bytes)
495 .sum()
496 }
497
498 pub(crate) fn num_rows(&self) -> u64 {
500 let inner = self.inner.lock().unwrap();
501 inner
502 .parts
503 .iter()
504 .map(|part| part.memtable.stats().num_rows as u64)
505 .sum()
506 }
507
508 pub(crate) fn list_memtables_to_small_vec(&self, memtables: &mut SmallMemtableVec) {
510 let inner = self.inner.lock().unwrap();
511 memtables.extend(inner.parts.iter().map(|part| part.memtable.clone()));
512 }
513
514 pub(crate) fn next_memtable_id(&self) -> MemtableId {
516 let inner = self.inner.lock().unwrap();
517 inner.next_memtable_id
518 }
519
520 pub(crate) fn new_with_part_duration(
523 &self,
524 part_duration: Option<Duration>,
525 memtable_builder: Option<MemtableBuilderRef>,
526 ) -> Self {
527 debug_assert!(self.is_empty());
528
529 Self::new(
530 self.metadata.clone(),
531 memtable_builder.unwrap_or_else(|| self.builder.clone()),
532 self.next_memtable_id(),
533 Some(part_duration.unwrap_or(self.part_duration)),
534 )
535 }
536
537 fn list_partitions(&self) -> PartitionVec {
539 let inner = self.inner.lock().unwrap();
540 inner.parts.clone()
541 }
542
543 fn find_partitions_by_time_range<'a>(
546 &self,
547 ts_array: &ArrayRef,
548 existing_parts: &'a [TimePartition],
549 min: Timestamp,
550 max: Timestamp,
551 ) -> Result<(Vec<&'a TimePartition>, Vec<Timestamp>)> {
552 let mut matching = Vec::new();
553
554 let mut present = HashSet::new();
555 for part in existing_parts {
557 let part_time_range = &part.time_range;
558 if !(max < part_time_range.min_timestamp || min >= part_time_range.max_timestamp) {
559 matching.push(part);
560 present.insert(part_time_range.min_timestamp.value());
561 }
562 }
563
564 let part_duration = self.part_duration_or_default();
566 let timestamp_unit = self.metadata.time_index_type().unit();
567
568 let part_duration_sec = part_duration.as_secs() as i64;
569 let start_bucket = min
571 .convert_to(TimeUnit::Second)
572 .unwrap()
573 .value()
574 .div_euclid(part_duration_sec);
575 let end_bucket = max
576 .convert_to(TimeUnit::Second)
577 .unwrap()
578 .value()
579 .div_euclid(part_duration_sec);
580 let bucket_num = (end_bucket - start_bucket + 1) as usize;
581
582 let num_timestamps = ts_array.len();
583 let missing = if bucket_num <= num_timestamps {
584 (start_bucket..=end_bucket)
585 .filter_map(|start_sec| {
586 let Some(timestamp) = Timestamp::new_second(start_sec * part_duration_sec)
587 .convert_to(timestamp_unit)
588 else {
589 return Some(
590 InvalidRequestSnafu {
591 region_id: self.metadata.region_id,
592 reason: format!("Timestamp out of range: {}", start_sec),
593 }
594 .fail(),
595 );
596 };
597 if present.insert(timestamp.value()) {
598 Some(Ok(timestamp))
599 } else {
600 None
601 }
602 })
603 .collect::<Result<Vec<_>>>()?
604 } else {
605 let ts_primitive = match ts_array.data_type() {
606 DataType::Timestamp(unit, _) => match unit {
607 arrow::datatypes::TimeUnit::Second => ts_array
608 .as_any()
609 .downcast_ref::<TimestampSecondArray>()
610 .unwrap()
611 .reinterpret_cast::<Int64Type>(),
612 arrow::datatypes::TimeUnit::Millisecond => ts_array
613 .as_any()
614 .downcast_ref::<TimestampMillisecondArray>()
615 .unwrap()
616 .reinterpret_cast::<Int64Type>(),
617 arrow::datatypes::TimeUnit::Microsecond => ts_array
618 .as_any()
619 .downcast_ref::<TimestampMicrosecondArray>()
620 .unwrap()
621 .reinterpret_cast::<Int64Type>(),
622 arrow::datatypes::TimeUnit::Nanosecond => ts_array
623 .as_any()
624 .downcast_ref::<TimestampNanosecondArray>()
625 .unwrap()
626 .reinterpret_cast::<Int64Type>(),
627 },
628 _ => unreachable!(),
629 };
630
631 ts_primitive
632 .values()
633 .iter()
634 .filter_map(|ts| {
635 let ts = self.metadata.time_index_type().create_timestamp(*ts);
636 let Some(bucket_start) = ts
637 .convert_to(TimeUnit::Second)
638 .and_then(|ts| ts.align_by_bucket(part_duration_sec))
639 .and_then(|ts| ts.convert_to(timestamp_unit))
640 else {
641 return Some(
642 InvalidRequestSnafu {
643 region_id: self.metadata.region_id,
644 reason: format!("Timestamp out of range: {:?}", ts),
645 }
646 .fail(),
647 );
648 };
649 if present.insert(bucket_start.value()) {
650 Some(Ok(bucket_start))
651 } else {
652 None
653 }
654 })
655 .collect::<Result<Vec<_>>>()?
656 };
657 Ok((matching, missing))
658 }
659
660 fn part_duration_or_default(&self) -> Duration {
662 self.part_duration
663 }
664
665 fn write_multi_parts(&self, kvs: &KeyValues, parts: &PartitionVec) -> Result<()> {
667 let mut parts_to_write = HashMap::new();
668 let mut missing_parts = HashMap::new();
669 for kv in kvs.iter() {
670 let mut part_found = false;
671 let ts = kv.timestamp().try_into_timestamp().unwrap().unwrap();
673 for part in parts {
674 if part.contains_timestamp(ts) {
675 parts_to_write
676 .entry(part.time_range.min_timestamp)
677 .or_insert_with(|| PartitionToWrite {
678 partition: part.clone(),
679 key_values: Vec::new(),
680 })
681 .key_values
682 .push(kv);
683 part_found = true;
684 break;
685 }
686 }
687
688 if !part_found {
689 let part_duration = self.part_duration_or_default();
692 let part_start =
693 partition_start_timestamp(ts, part_duration).with_context(|| {
694 InvalidRequestSnafu {
695 region_id: self.metadata.region_id,
696 reason: format!(
697 "timestamp {ts:?} and bucket {part_duration:?} are out of range"
698 ),
699 }
700 })?;
701 missing_parts
702 .entry(part_start)
703 .or_insert_with(Vec::new)
704 .push(kv);
705 }
706 }
707
708 for part_to_write in parts_to_write.into_values() {
710 for kv in part_to_write.key_values {
711 part_to_write.partition.memtable.write_one(kv)?;
712 }
713 }
714
715 let mut inner = self.inner.lock().unwrap();
718 for (part_start, key_values) in missing_parts {
719 let partition = self.get_or_create_time_partition(part_start, &mut inner)?;
720 for kv in key_values {
721 partition.memtable.write_one(kv)?;
722 }
723 }
724
725 Ok(())
726 }
727
728 pub(crate) fn series_count(&self) -> usize {
730 self.inner.lock().unwrap().series_count()
731 }
732}
733
734fn partition_start_timestamp(ts: Timestamp, bucket: Duration) -> Option<Timestamp> {
738 let ts_sec = ts.convert_to(TimeUnit::Second).unwrap();
740 let bucket_sec: i64 = bucket.as_secs().try_into().ok()?;
741 let start_sec = ts_sec.align_by_bucket(bucket_sec)?;
742 start_sec.convert_to(ts.unit())
743}
744
745#[derive(Debug)]
746struct PartitionsInner {
747 parts: PartitionVec,
749 next_memtable_id: MemtableId,
751}
752
753impl PartitionsInner {
754 fn new(next_memtable_id: MemtableId) -> Self {
755 Self {
756 parts: Default::default(),
757 next_memtable_id,
758 }
759 }
760
761 fn with_partition(part: TimePartition, next_memtable_id: MemtableId) -> Self {
762 Self {
763 parts: smallvec![part],
764 next_memtable_id,
765 }
766 }
767
768 fn alloc_memtable_id(&mut self) -> MemtableId {
769 let id = self.next_memtable_id;
770 self.next_memtable_id += 1;
771 id
772 }
773
774 pub(crate) fn series_count(&self) -> usize {
775 self.parts
776 .iter()
777 .map(|p| p.memtable.stats().series_count)
778 .sum()
779 }
780}
781
782#[derive(Debug, Clone, Copy)]
784struct PartTimeRange {
785 min_timestamp: Timestamp,
787 max_timestamp: Timestamp,
789}
790
791impl PartTimeRange {
792 fn from_start_duration(start: Timestamp, duration: Duration) -> Option<Self> {
793 let start_sec = start.convert_to(TimeUnit::Second)?;
794 let end_sec = start_sec.add_duration(duration).ok()?;
795 let min_timestamp = start_sec.convert_to(start.unit())?;
796 let max_timestamp = end_sec.convert_to(start.unit())?;
797
798 Some(Self {
799 min_timestamp,
800 max_timestamp,
801 })
802 }
803
804 fn contains_timestamp(&self, ts: Timestamp) -> bool {
806 self.min_timestamp <= ts && ts < self.max_timestamp
807 }
808}
809
810struct PartitionToWrite<'a> {
811 partition: TimePartition,
812 key_values: Vec<KeyValue<'a>>,
813}
814
815#[cfg(test)]
816mod tests {
817 use std::sync::Arc;
818
819 use api::v1::SemanticType;
820 use datatypes::arrow::array::{ArrayRef, StringArray, TimestampMillisecondArray};
821 use datatypes::arrow::datatypes::{DataType, Field, Schema};
822 use datatypes::arrow::record_batch::RecordBatch;
823 use datatypes::prelude::ConcreteDataType;
824 use datatypes::schema::ColumnSchema;
825 use store_api::metadata::{ColumnMetadata, RegionMetadataBuilder};
826 use store_api::storage::SequenceNumber;
827
828 use super::*;
829 use crate::memtable::time_series::TimeSeriesMemtableBuilder;
830 use crate::memtable::{IterBuilder, RangesOptions};
831 use crate::test_util::memtable_util::{self, collect_iter_timestamps};
832
833 #[test]
834 fn test_no_duration() {
835 let metadata = memtable_util::metadata_for_test();
836 let builder = Arc::new(TimeSeriesMemtableBuilder::default());
837 let partitions = TimePartitions::new(metadata.clone(), builder, 0, None);
838 assert_eq!(0, partitions.num_partitions());
839 assert!(partitions.is_empty());
840
841 let kvs = memtable_util::build_key_values(
842 &metadata,
843 "hello".to_string(),
844 0,
845 &[1000, 3000, 7000, 5000, 6000],
846 0, );
848 partitions.write(&kvs).unwrap();
849
850 assert_eq!(1, partitions.num_partitions());
851 assert!(!partitions.is_empty());
852 let mut memtables = Vec::new();
853 partitions.list_memtables(&mut memtables);
854 assert_eq!(0, memtables[0].id());
855
856 let iter = memtables[0]
857 .ranges(None, RangesOptions::default())
858 .unwrap()
859 .build(None)
860 .unwrap();
861 let timestamps = collect_iter_timestamps(iter);
862 assert_eq!(&[1000, 3000, 5000, 6000, 7000], ×tamps[..]);
863 }
864
865 #[test]
866 fn test_write_single_part() {
867 let metadata = memtable_util::metadata_for_test();
868 let builder = Arc::new(TimeSeriesMemtableBuilder::default());
869 let partitions =
870 TimePartitions::new(metadata.clone(), builder, 0, Some(Duration::from_secs(10)));
871 assert_eq!(0, partitions.num_partitions());
872
873 let kvs = memtable_util::build_key_values(
874 &metadata,
875 "hello".to_string(),
876 0,
877 &[5000, 2000, 0],
878 0, );
880 partitions.write(&kvs).unwrap();
882 assert_eq!(1, partitions.num_partitions());
883 assert!(!partitions.is_empty());
884
885 let kvs = memtable_util::build_key_values(
886 &metadata,
887 "hello".to_string(),
888 0,
889 &[3000, 7000, 4000],
890 3, );
892 partitions.write(&kvs).unwrap();
894 assert_eq!(1, partitions.num_partitions());
895
896 let mut memtables = Vec::new();
897 partitions.list_memtables(&mut memtables);
898 let iter = memtables[0]
899 .ranges(None, RangesOptions::default())
900 .unwrap()
901 .build(None)
902 .unwrap();
903 let timestamps = collect_iter_timestamps(iter);
904 assert_eq!(&[0, 2000, 3000, 4000, 5000, 7000], ×tamps[..]);
905 let parts = partitions.list_partitions();
906 assert_eq!(
907 Timestamp::new_millisecond(0),
908 parts[0].time_range.min_timestamp
909 );
910 assert_eq!(
911 Timestamp::new_millisecond(10000),
912 parts[0].time_range.max_timestamp
913 );
914 }
915
916 #[cfg(test)]
917 fn new_multi_partitions(metadata: &RegionMetadataRef) -> TimePartitions {
918 let builder = Arc::new(TimeSeriesMemtableBuilder::default());
919 let partitions =
920 TimePartitions::new(metadata.clone(), builder, 0, Some(Duration::from_secs(5)));
921 assert_eq!(0, partitions.num_partitions());
922
923 let kvs = memtable_util::build_key_values(
924 metadata,
925 "hello".to_string(),
926 0,
927 &[2000, 0],
928 0, );
930 partitions.write(&kvs).unwrap();
932 assert_eq!(1, partitions.num_partitions());
933 assert!(!partitions.is_empty());
934
935 let kvs = memtable_util::build_key_values(
936 metadata,
937 "hello".to_string(),
938 0,
939 &[3000, 7000, 4000, 5000],
940 2, );
942 partitions.write(&kvs).unwrap();
944 assert_eq!(2, partitions.num_partitions());
945
946 partitions
947 }
948
949 #[test]
950 fn test_write_multi_parts() {
951 let metadata = memtable_util::metadata_for_test();
952 let partitions = new_multi_partitions(&metadata);
953
954 let parts = partitions.list_partitions();
955 let iter = parts[0]
956 .memtable
957 .ranges(None, RangesOptions::default())
958 .unwrap()
959 .build(None)
960 .unwrap();
961 let timestamps = collect_iter_timestamps(iter);
962 assert_eq!(0, parts[0].memtable.id());
963 assert_eq!(
964 Timestamp::new_millisecond(0),
965 parts[0].time_range.min_timestamp
966 );
967 assert_eq!(
968 Timestamp::new_millisecond(5000),
969 parts[0].time_range.max_timestamp
970 );
971 assert_eq!(&[0, 2000, 3000, 4000], ×tamps[..]);
972 let iter = parts[1]
973 .memtable
974 .ranges(None, RangesOptions::default())
975 .unwrap()
976 .build(None)
977 .unwrap();
978 assert_eq!(1, parts[1].memtable.id());
979 let timestamps = collect_iter_timestamps(iter);
980 assert_eq!(&[5000, 7000], ×tamps[..]);
981 assert_eq!(
982 Timestamp::new_millisecond(5000),
983 parts[1].time_range.min_timestamp
984 );
985 assert_eq!(
986 Timestamp::new_millisecond(10000),
987 parts[1].time_range.max_timestamp
988 );
989 }
990
991 #[test]
992 fn test_new_with_part_duration() {
993 let metadata = memtable_util::metadata_for_test();
994 let builder = Arc::new(TimeSeriesMemtableBuilder::default());
995 let partitions = TimePartitions::new(metadata.clone(), builder.clone(), 0, None);
996
997 let new_parts = partitions.new_with_part_duration(Some(Duration::from_secs(5)), None);
998 assert_eq!(Duration::from_secs(5), new_parts.part_duration());
999 assert_eq!(0, new_parts.next_memtable_id());
1000
1001 let new_parts = new_parts.new_with_part_duration(None, None);
1003 assert_eq!(Duration::from_secs(5), new_parts.part_duration());
1004 assert_eq!(0, new_parts.next_memtable_id());
1006
1007 let new_parts = new_parts.new_with_part_duration(Some(Duration::from_secs(10)), None);
1008 assert_eq!(Duration::from_secs(10), new_parts.part_duration());
1009 assert_eq!(0, new_parts.next_memtable_id());
1011
1012 let builder = Arc::new(TimeSeriesMemtableBuilder::default());
1013 let partitions = TimePartitions::new(metadata.clone(), builder.clone(), 0, None);
1014 let new_parts = partitions.new_with_part_duration(None, None);
1016 assert_eq!(INITIAL_TIME_WINDOW, new_parts.part_duration());
1017 assert_eq!(0, new_parts.next_memtable_id());
1018 }
1019
1020 #[test]
1021 fn test_fork_empty() {
1022 let metadata = memtable_util::metadata_for_test();
1023 let builder = Arc::new(TimeSeriesMemtableBuilder::default());
1024 let partitions = TimePartitions::new(metadata.clone(), builder, 0, None);
1025 partitions.freeze().unwrap();
1026 let new_parts = partitions.fork(&metadata, None);
1027 assert_eq!(INITIAL_TIME_WINDOW, new_parts.part_duration());
1028 assert!(new_parts.list_partitions().is_empty());
1029 assert_eq!(0, new_parts.next_memtable_id());
1030
1031 new_parts.freeze().unwrap();
1032 let new_parts = new_parts.fork(&metadata, Some(Duration::from_secs(5)));
1033 assert_eq!(Duration::from_secs(5), new_parts.part_duration());
1034 assert!(new_parts.list_partitions().is_empty());
1035 assert_eq!(0, new_parts.next_memtable_id());
1036
1037 new_parts.freeze().unwrap();
1038 let new_parts = new_parts.fork(&metadata, None);
1039 assert_eq!(Duration::from_secs(5), new_parts.part_duration());
1041 assert!(new_parts.list_partitions().is_empty());
1042 assert_eq!(0, new_parts.next_memtable_id());
1043
1044 new_parts.freeze().unwrap();
1045 let new_parts = new_parts.fork(&metadata, Some(Duration::from_secs(10)));
1046 assert_eq!(Duration::from_secs(10), new_parts.part_duration());
1047 assert!(new_parts.list_partitions().is_empty());
1048 assert_eq!(0, new_parts.next_memtable_id());
1049 }
1050
1051 #[test]
1052 fn test_fork_non_empty_none() {
1053 let metadata = memtable_util::metadata_for_test();
1054 let partitions = new_multi_partitions(&metadata);
1055 partitions.freeze().unwrap();
1056
1057 let new_parts = partitions.fork(&metadata, None);
1059 assert!(new_parts.is_empty());
1060 assert_eq!(Duration::from_secs(5), new_parts.part_duration());
1061 assert_eq!(2, new_parts.list_partitions()[0].memtable.id());
1062 assert_eq!(3, new_parts.next_memtable_id());
1063
1064 let new_parts = partitions.fork(&metadata, Some(Duration::from_secs(10)));
1066 assert!(new_parts.is_empty());
1067 assert_eq!(Duration::from_secs(10), new_parts.part_duration());
1068 assert_eq!(3, new_parts.list_partitions()[0].memtable.id());
1069 assert_eq!(4, new_parts.next_memtable_id());
1070 }
1071
1072 #[test]
1073 fn test_find_partitions_by_time_range() {
1074 let metadata = memtable_util::metadata_for_test();
1075 let builder = Arc::new(TimeSeriesMemtableBuilder::default());
1076
1077 let partitions = TimePartitions::new(metadata.clone(), builder.clone(), 0, None);
1079 let parts = partitions.list_partitions();
1080 let (matching, missing) = partitions
1081 .find_partitions_by_time_range(
1082 &(Arc::new(TimestampMillisecondArray::from_iter_values(1000..=2000)) as ArrayRef),
1083 &parts,
1084 Timestamp::new_millisecond(1000),
1085 Timestamp::new_millisecond(2000),
1086 )
1087 .unwrap();
1088 assert_eq!(matching.len(), 0);
1089 assert_eq!(missing.len(), 1);
1090 assert_eq!(missing[0], Timestamp::new_millisecond(0));
1091
1092 let partitions = TimePartitions::new(
1094 metadata.clone(),
1095 builder.clone(),
1096 0,
1097 Some(Duration::from_secs(5)),
1098 );
1099
1100 let kvs =
1102 memtable_util::build_key_values(&metadata, "hello".to_string(), 0, &[2000, 4000], 0);
1103 partitions.write(&kvs).unwrap();
1104 let kvs =
1105 memtable_util::build_key_values(&metadata, "hello".to_string(), 0, &[7000, 8000], 2);
1106 partitions.write(&kvs).unwrap();
1107
1108 let parts = partitions.list_partitions();
1109 assert_eq!(2, parts.len());
1110
1111 let (matching, missing) = partitions
1113 .find_partitions_by_time_range(
1114 &(Arc::new(TimestampMillisecondArray::from_iter_values(2000..=4000)) as ArrayRef),
1115 &parts,
1116 Timestamp::new_millisecond(2000),
1117 Timestamp::new_millisecond(4000),
1118 )
1119 .unwrap();
1120 assert_eq!(matching.len(), 1);
1121 assert!(missing.is_empty());
1122 assert_eq!(matching[0].time_range.min_timestamp.value(), 0);
1123
1124 let (matching, missing) = partitions
1126 .find_partitions_by_time_range(
1127 &(Arc::new(TimestampMillisecondArray::from_iter_values(3000..=8000)) as ArrayRef),
1128 &parts,
1129 Timestamp::new_millisecond(3000),
1130 Timestamp::new_millisecond(8000),
1131 )
1132 .unwrap();
1133 assert_eq!(matching.len(), 2);
1134 assert!(missing.is_empty());
1135 assert_eq!(matching[0].time_range.min_timestamp.value(), 0);
1136 assert_eq!(matching[1].time_range.min_timestamp.value(), 5000);
1137
1138 let (matching, missing) = partitions
1140 .find_partitions_by_time_range(
1141 &(Arc::new(TimestampMillisecondArray::from_iter_values(12000..=13000)) as ArrayRef),
1142 &parts,
1143 Timestamp::new_millisecond(12000),
1144 Timestamp::new_millisecond(13000),
1145 )
1146 .unwrap();
1147 assert!(matching.is_empty());
1148 assert_eq!(missing.len(), 1);
1149 assert_eq!(missing[0].value(), 10000);
1150
1151 let (matching, missing) = partitions
1153 .find_partitions_by_time_range(
1154 &(Arc::new(TimestampMillisecondArray::from_iter_values(4000..=6000)) as ArrayRef),
1155 &parts,
1156 Timestamp::new_millisecond(4000),
1157 Timestamp::new_millisecond(6000),
1158 )
1159 .unwrap();
1160 assert_eq!(matching.len(), 2);
1161 assert!(missing.is_empty());
1162 assert_eq!(matching[0].time_range.min_timestamp.value(), 0);
1163 assert_eq!(matching[1].time_range.min_timestamp.value(), 5000);
1164
1165 let (matching, missing) = partitions
1167 .find_partitions_by_time_range(
1168 &(Arc::new(TimestampMillisecondArray::from_iter_values(4999..=5000)) as ArrayRef),
1169 &parts,
1170 Timestamp::new_millisecond(4999),
1171 Timestamp::new_millisecond(5000),
1172 )
1173 .unwrap();
1174 assert_eq!(matching.len(), 2);
1175 assert!(missing.is_empty());
1176 assert_eq!(matching[0].time_range.min_timestamp.value(), 0);
1177 assert_eq!(matching[1].time_range.min_timestamp.value(), 5000);
1178
1179 let (matching, missing) = partitions
1181 .find_partitions_by_time_range(
1182 &(Arc::new(TimestampMillisecondArray::from_iter_values(9999..=10000)) as ArrayRef),
1183 &parts,
1184 Timestamp::new_millisecond(9999),
1185 Timestamp::new_millisecond(10000),
1186 )
1187 .unwrap();
1188 assert_eq!(matching.len(), 1);
1189 assert_eq!(1, missing.len());
1190 assert_eq!(matching[0].time_range.min_timestamp.value(), 5000);
1191 assert_eq!(missing[0].value(), 10000);
1192
1193 let (matching, missing) = partitions
1195 .find_partitions_by_time_range(
1196 &(Arc::new(TimestampMillisecondArray::from_iter_values(-1000..=1000)) as ArrayRef),
1197 &parts,
1198 Timestamp::new_millisecond(-1000),
1199 Timestamp::new_millisecond(1000),
1200 )
1201 .unwrap();
1202 assert_eq!(matching.len(), 1);
1203 assert_eq!(matching[0].time_range.min_timestamp.value(), 0);
1204 assert_eq!(1, missing.len());
1205 assert_eq!(missing[0].value(), -5000);
1206
1207 let (matching, missing) = partitions
1209 .find_partitions_by_time_range(
1210 &(Arc::new(TimestampMillisecondArray::from(vec![
1211 -100000000000,
1212 0,
1213 100000000000,
1214 ])) as ArrayRef),
1215 &parts,
1216 Timestamp::new_millisecond(-100000000000),
1217 Timestamp::new_millisecond(100000000000),
1218 )
1219 .unwrap();
1220 assert_eq!(2, matching.len());
1221 assert_eq!(matching[0].time_range.min_timestamp.value(), 0);
1222 assert_eq!(matching[1].time_range.min_timestamp.value(), 5000);
1223 assert_eq!(2, missing.len());
1224 assert_eq!(missing[0].value(), -100000000000);
1225 assert_eq!(missing[1].value(), 100000000000);
1226 }
1227
1228 fn build_part(ts: &[i64], sequence: SequenceNumber) -> BulkPart {
1229 let schema = Arc::new(Schema::new(vec![
1230 Field::new(
1231 "ts",
1232 DataType::Timestamp(arrow::datatypes::TimeUnit::Millisecond, None),
1233 false,
1234 ),
1235 Field::new("val", DataType::Utf8, true),
1236 ]));
1237 let ts_data = ts.to_vec();
1238 let ts_array = Arc::new(TimestampMillisecondArray::from(ts_data));
1239 let val_array = Arc::new(StringArray::from_iter_values(
1240 ts.iter().map(|v| v.to_string()),
1241 ));
1242 let batch = RecordBatch::try_new(
1243 schema,
1244 vec![ts_array.clone() as ArrayRef, val_array.clone() as ArrayRef],
1245 )
1246 .unwrap();
1247 let max_ts = ts.iter().max().copied().unwrap();
1248 let min_ts = ts.iter().min().copied().unwrap();
1249 BulkPart {
1250 batch,
1251 max_timestamp: max_ts,
1252 min_timestamp: min_ts,
1253 sequence,
1254 min_sequence: sequence,
1255 timestamp_index: 0,
1256 raw_data: None,
1257 }
1258 }
1259
1260 #[test]
1261 fn test_write_bulk() {
1262 let mut metadata_builder = RegionMetadataBuilder::new(0.into());
1263 metadata_builder
1264 .push_column_metadata(ColumnMetadata {
1265 column_schema: ColumnSchema::new(
1266 "ts",
1267 ConcreteDataType::timestamp_millisecond_datatype(),
1268 false,
1269 ),
1270 semantic_type: SemanticType::Timestamp,
1271 column_id: 0,
1272 })
1273 .push_column_metadata(ColumnMetadata {
1274 column_schema: ColumnSchema::new("val", ConcreteDataType::string_datatype(), false),
1275 semantic_type: SemanticType::Field,
1276 column_id: 1,
1277 })
1278 .primary_key(vec![]);
1279 let metadata = Arc::new(metadata_builder.build().unwrap());
1280
1281 let builder = Arc::new(TimeSeriesMemtableBuilder::default());
1282 let partitions = TimePartitions::new(
1283 metadata.clone(),
1284 builder.clone(),
1285 0,
1286 Some(Duration::from_secs(5)),
1287 );
1288
1289 partitions
1291 .write_bulk(build_part(&[1000, 2000, 3000], 0))
1292 .unwrap();
1293
1294 let parts = partitions.list_partitions();
1295 assert_eq!(1, parts.len());
1296 let iter = parts[0]
1297 .memtable
1298 .ranges(None, RangesOptions::default())
1299 .unwrap()
1300 .build(None)
1301 .unwrap();
1302 let timestamps = collect_iter_timestamps(iter);
1303 assert_eq!(&[1000, 2000, 3000], ×tamps[..]);
1304
1305 partitions
1307 .write_bulk(build_part(&[4000, 5000, 6000], 1))
1308 .unwrap();
1309 let parts = partitions.list_partitions();
1310 assert_eq!(2, parts.len());
1311 let iter = parts[0]
1313 .memtable
1314 .ranges(None, RangesOptions::default())
1315 .unwrap()
1316 .build(None)
1317 .unwrap();
1318 let timestamps = collect_iter_timestamps(iter);
1319 assert_eq!(&[1000, 2000, 3000, 4000], ×tamps[..]);
1320 let iter = parts[1]
1322 .memtable
1323 .ranges(None, RangesOptions::default())
1324 .unwrap()
1325 .build(None)
1326 .unwrap();
1327 let timestamps = collect_iter_timestamps(iter);
1328 assert_eq!(&[5000, 6000], ×tamps[..]);
1329
1330 partitions
1332 .write_bulk(build_part(&[11000, 12000], 3))
1333 .unwrap();
1334
1335 let parts = partitions.list_partitions();
1336 assert_eq!(3, parts.len());
1337
1338 let iter = parts[2]
1340 .memtable
1341 .ranges(None, RangesOptions::default())
1342 .unwrap()
1343 .build(None)
1344 .unwrap();
1345 let timestamps = collect_iter_timestamps(iter);
1346 assert_eq!(&[11000, 12000], ×tamps[..]);
1347
1348 let partitions = TimePartitions::new(metadata.clone(), builder, 3, None);
1350
1351 partitions
1352 .write_bulk(build_part(&[1000, 5000, 9000], 4))
1353 .unwrap();
1354
1355 let parts = partitions.list_partitions();
1356 assert_eq!(1, parts.len());
1357 let iter = parts[0]
1358 .memtable
1359 .ranges(None, RangesOptions::default())
1360 .unwrap()
1361 .build(None)
1362 .unwrap();
1363 let timestamps = collect_iter_timestamps(iter);
1364 assert_eq!(&[1000, 5000, 9000], ×tamps[..]);
1365 }
1366
1367 #[test]
1368 fn test_split_record_batch() {
1369 let schema = Arc::new(Schema::new(vec![
1370 Field::new(
1371 "ts",
1372 DataType::Timestamp(TimeUnit::Millisecond.as_arrow_time_unit(), None),
1373 false,
1374 ),
1375 Field::new("val", DataType::Utf8, true),
1376 ]));
1377
1378 let ts_array = Arc::new(TimestampMillisecondArray::from(vec![
1379 1000, 2000, 5000, 7000, 8000,
1380 ]));
1381 let val_array = Arc::new(StringArray::from(vec!["a", "b", "c", "d", "e"]));
1382 let batch = RecordBatch::try_new(
1383 schema.clone(),
1384 vec![ts_array as ArrayRef, val_array as ArrayRef],
1385 )
1386 .unwrap();
1387
1388 let part = BulkPart {
1389 batch,
1390 max_timestamp: 8000,
1391 min_timestamp: 1000,
1392 sequence: 0,
1393 min_sequence: 0,
1394 timestamp_index: 0,
1395 raw_data: None,
1396 };
1397
1398 let result = filter_record_batch(&part, 1000, 2000).unwrap();
1399 assert!(result.is_some());
1400 let filtered = result.unwrap();
1401 assert_eq!(filtered.num_rows(), 1);
1402 assert_eq!(filtered.min_timestamp, 1000);
1403 assert_eq!(filtered.max_timestamp, 1000);
1404
1405 let result = filter_record_batch(&part, 3000, 6000).unwrap();
1407 assert!(result.is_some());
1408 let filtered = result.unwrap();
1409 assert_eq!(filtered.num_rows(), 1);
1410 assert_eq!(filtered.min_timestamp, 5000);
1411 assert_eq!(filtered.max_timestamp, 5000);
1412
1413 let result = filter_record_batch(&part, 3000, 4000).unwrap();
1415 assert!(result.is_none());
1416
1417 let result = filter_record_batch(&part, 0, 9000).unwrap();
1419 assert!(result.is_some());
1420 let filtered = result.unwrap();
1421 assert_eq!(filtered.num_rows(), 5);
1422 assert_eq!(filtered.min_timestamp, 1000);
1423 assert_eq!(filtered.max_timestamp, 8000);
1424 }
1425}