Skip to main content

mito2/memtable/
time_partition.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15//! Partitions memtables by time.
16
17use 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
45/// Initial time window if not specified.
46const INITIAL_TIME_WINDOW: Duration = Duration::from_days(1);
47
48/// A partition holds rows with timestamps between `[min, max)`.
49#[derive(Debug, Clone)]
50pub struct TimePartition {
51    /// Memtable of the partition.
52    memtable: MemtableRef,
53    /// Time range of the partition. `min` is inclusive and `max` is exclusive.
54    time_range: PartTimeRange,
55}
56
57impl TimePartition {
58    /// Returns whether the `ts` belongs to the partition.
59    fn contains_timestamp(&self, ts: Timestamp) -> bool {
60        self.time_range.contains_timestamp(ts)
61    }
62
63    /// Write rows to the part.
64    fn write(&self, kvs: &KeyValues) -> Result<()> {
65        self.memtable.write(kvs)
66    }
67
68    /// Writes a record batch to memtable.
69    fn write_record_batch(&self, rb: BulkPart) -> Result<()> {
70        self.memtable.write_bulk(rb)
71    }
72
73    /// Write a partial [BulkPart] according to [TimePartition::time_range].
74    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            // SAFETY: we only iterate the array within index bound.
94            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            // SAFETY: Already allocated sufficient capacity
110            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            // SAFETY: Already allocated sufficient capacity
120            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        // safety: we've checked res is not empty
139        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
146/// Filters the given part according to min (inclusive) and max (exclusive) timestamp range.
147/// Returns [None] if no matching rows.
148pub 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/// Partitions.
204#[derive(Debug)]
205pub struct TimePartitions {
206    /// Mutable data of partitions.
207    inner: Mutex<PartitionsInner>,
208    /// Duration of a partition.
209    part_duration: Duration,
210    /// Metadata of the region.
211    metadata: RegionMetadataRef,
212    /// Builder of memtables.
213    builder: MemtableBuilderRef,
214    /// Primary key encoder.
215    primary_key_codec: Arc<dyn PrimaryKeyCodec>,
216
217    /// Cached schema for bulk insert.
218    /// This field is Some if the memtable uses bulk insert.
219    bulk_schema: Option<SchemaRef>,
220}
221
222pub type TimePartitionsRef = Arc<TimePartitions>;
223
224impl TimePartitions {
225    /// Returns a new empty partition list with optional duration.
226    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    /// Write key values to memtables.
250    ///
251    /// It creates new partitions if necessary.
252    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                // Always store primary keys for bulk mode.
260                true,
261            );
262            converter.append_key_values(kvs)?;
263            let part = converter.convert()?;
264
265            return self.write_bulk_inner(part);
266        }
267
268        // Get all parts.
269        let parts = self.list_partitions();
270
271        // Checks whether all rows belongs to a single part. Checks in reverse order as we usually
272        // put to latest part.
273        for part in parts.iter().rev() {
274            let mut all_in_partition = true;
275            for kv in kvs.iter() {
276                // Safety: We checked the schema in the write request.
277                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            // We can write all rows to this part.
288            return part.write(kvs);
289        }
290
291        // Slow path: We have to split kvs by partitions.
292        self.write_multi_parts(kvs, &parts)
293    }
294
295    /// Writes a bulk part.
296    pub fn write_bulk(&self, part: BulkPart) -> Result<()> {
297        // Convert the bulk part if bulk_schema is Some
298        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                // Always store primary keys for bulk mode.
305                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    /// Writes a bulk part without converting.
319    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        // Get all parts.
329        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            // fast path: all timestamps fall in one time partition.
339            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    // Creates or gets parts with given start timestamp.
357    // Acquires the lock to avoid others create the same partition.
358    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    /// Append memtables in partitions to `memtables`.
401    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    /// Returns the number of partitions.
407    pub fn num_partitions(&self) -> usize {
408        let inner = self.inner.lock().unwrap();
409        inner.parts.len()
410    }
411
412    /// Returns true if all memtables are empty.
413    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    /// Freezes all memtables.
419    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    /// Forks latest partition and updates the partition duration if `part_duration` is Some.
428    pub fn fork(&self, metadata: &RegionMetadataRef, part_duration: Option<Duration>) -> Self {
429        // Fall back to the existing partition duration.
430        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            // If there is no partition, then we create a new partition with the new duration.
441            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        // Use the max timestamp to compute the new time range for the memtable.
451        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                // Forks the latest partition, but compute the time range based on the new duration.
459                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    /// Returns partition duration.
479    pub(crate) fn part_duration(&self) -> Duration {
480        self.part_duration
481    }
482
483    /// Returns the memtable builder.
484    pub(crate) fn memtable_builder(&self) -> &MemtableBuilderRef {
485        &self.builder
486    }
487
488    /// Returns memory usage.
489    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    /// Returns the number of rows.
499    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    /// Append memtables in partitions to small vec.
509    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    /// Returns the next memtable id.
515    pub(crate) fn next_memtable_id(&self) -> MemtableId {
516        let inner = self.inner.lock().unwrap();
517        inner.next_memtable_id
518    }
519
520    /// Creates a new empty partition list from this list and a `part_duration`.
521    /// It falls back to the old partition duration if `part_duration` is `None`.
522    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    /// Returns all partitions.
538    fn list_partitions(&self) -> PartitionVec {
539        let inner = self.inner.lock().unwrap();
540        inner.parts.clone()
541    }
542
543    /// Find existing partitions that match the bulk data's time range and identify
544    /// any new partitions that need to be created
545    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        // First find any existing partitions that overlap
556        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        // safety: self.part_duration can only be present when reach here.
565        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        // SAFETY: Timestamps won't overflow when converting to Second.
570        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    /// Returns partition duration, or use default 1day duration is not present.
661    fn part_duration_or_default(&self) -> Duration {
662        self.part_duration
663    }
664
665    /// Write to multiple partitions.
666    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            // Safety: We used the timestamp before.
672            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                // We need to write it to a new part.
690                // Safety: `new()` ensures duration is always Some if we do to this method.
691                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        // Writes rows to existing parts.
709        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        // Creates new parts and writes to them. Acquires the lock to avoid others create
716        // the same partition.
717        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    /// Timeseries count in all time partitions.
729    pub(crate) fn series_count(&self) -> usize {
730        self.inner.lock().unwrap().series_count()
731    }
732}
733
734/// Computes the start timestamp of the partition for `ts`.
735///
736/// It always use bucket size in seconds which should fit all timestamp resolution.
737fn partition_start_timestamp(ts: Timestamp, bucket: Duration) -> Option<Timestamp> {
738    // Safety: We convert it to seconds so it never returns `None`.
739    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    /// All partitions.
748    parts: PartitionVec,
749    /// Next memtable id.
750    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/// Time range of a partition.
783#[derive(Debug, Clone, Copy)]
784struct PartTimeRange {
785    /// Inclusive min timestamp of rows in the partition.
786    min_timestamp: Timestamp,
787    /// Exclusive max timestamp of rows in the partition.
788    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    /// Returns whether the `ts` belongs to the partition.
805    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, // sequence 0, 1, 2, 3, 4
847        );
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], &timestamps[..]);
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, // sequence 0, 1, 2
879        );
880        // It should creates a new partition.
881        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, // sequence 3, 4, 5
891        );
892        // Still writes to the same partition.
893        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], &timestamps[..]);
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, // sequence 0, 1
929        );
930        // It should creates a new partition.
931        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, // sequence 2, 3, 4, 5
941        );
942        // Writes 2 rows to the old partition and 1 row to a new partition.
943        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], &timestamps[..]);
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], &timestamps[..]);
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        // Won't update the duration if it's None.
1002        let new_parts = new_parts.new_with_part_duration(None, None);
1003        assert_eq!(Duration::from_secs(5), new_parts.part_duration());
1004        // Don't need to create new memtables.
1005        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        // Don't need to create new memtables.
1010        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        // Need to build a new memtable as duration is still None.
1015        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        // Won't update the duration.
1040        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        // Won't update the duration.
1058        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        // Although we don't fork a memtable multiple times, we still add a test for it.
1065        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        // Case 1: No time range partitioning
1078        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        // Case 2: With time range partitioning
1093        let partitions = TimePartitions::new(
1094            metadata.clone(),
1095            builder.clone(),
1096            0,
1097            Some(Duration::from_secs(5)),
1098        );
1099
1100        // Create two existing partitions: [0, 5000) and [5000, 10000)
1101        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        // Test case 2a: Query fully within existing partition
1112        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        // Test case 2b: Query spanning multiple existing partitions
1125        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        // Test case 2c: Query requiring new partition
1139        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        // Test case 2d: Query partially overlapping existing partition
1152        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        // Test case 2e: Corner case
1166        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        // Test case 2f: Corner case with
1180        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        // Test case 2g: Cross 0
1194        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        // Test case 3: sparse data
1208        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        // Test case 1: Write to single partition
1290        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], &timestamps[..]);
1304
1305        // Test case 2: Write across multiple existing partitions
1306        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        // Check first partition [0, 5000)
1312        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], &timestamps[..]);
1320        // Check second partition [5000, 10000)
1321        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], &timestamps[..]);
1329
1330        // Test case 3: Write requiring new partition
1331        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        // Check new partition [10000, 15000)
1339        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], &timestamps[..]);
1347
1348        // Test case 4: Write with no time range partitioning
1349        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], &timestamps[..]);
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        // Test splitting with range [3000, 6000)
1406        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        // Test splitting with range that includes no points
1414        let result = filter_record_batch(&part, 3000, 4000).unwrap();
1415        assert!(result.is_none());
1416
1417        // Test splitting with range that includes all points
1418        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}