Skip to main content

mito2/
memtable.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//! Memtables are write buffers for regions.
16
17use std::collections::BTreeMap;
18use std::fmt;
19use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
20use std::sync::{Arc, Mutex};
21use std::time::Duration;
22
23pub use bulk::part::EncodedBulkPart;
24use bytes::Bytes;
25use common_time::Timestamp;
26use datatypes::arrow::datatypes::SchemaRef;
27use datatypes::arrow::record_batch::RecordBatch;
28use mito_codec::key_values::KeyValue;
29pub use mito_codec::key_values::KeyValues;
30use mito_codec::row_converter::{PrimaryKeyCodec, build_primary_key_codec};
31use snafu::ensure;
32use store_api::codec::PrimaryKeyEncoding;
33use store_api::metadata::{RegionMetadata, RegionMetadataRef};
34use store_api::storage::{ColumnId, SequenceNumber, SequenceRange};
35
36use crate::config::MitoConfig;
37use crate::error::{InvalidRegionOptionsSnafu, Result, UnsupportedOperationSnafu};
38use crate::flush::WriteBufferManagerRef;
39use crate::memtable::bulk::{BulkMemtableBuilder, CompactDispatcher};
40use crate::memtable::time_series::TimeSeriesMemtableBuilder;
41use crate::metrics::WRITE_BUFFER_BYTES;
42use crate::read::Batch;
43use crate::read::batch_adapter::BatchToRecordBatchAdapter;
44use crate::read::prune::PruneTimeIterator;
45use crate::read::scan_region::PredicateGroup;
46use crate::region::options::{MemtableOptions, MergeMode, RegionOptions};
47use crate::sst::FormatType;
48use crate::sst::file::FileTimeRange;
49use crate::sst::parquet::SstInfo;
50use crate::sst::parquet::file_range::PreFilterMode;
51
52mod builder;
53pub mod bulk;
54pub mod simple_bulk_memtable;
55mod stats;
56pub mod time_partition;
57pub mod time_series;
58pub(crate) mod version;
59
60pub use bulk::part::{
61    BulkPart, BulkPartEncoder, BulkPartMeta, UnorderedPart, record_batch_estimated_size,
62    sort_primary_key_record_batch,
63};
64#[cfg(any(test, feature = "test"))]
65pub use time_partition::filter_record_batch;
66
67/// Id for memtables.
68///
69/// Should be unique under the same region.
70pub type MemtableId = u32;
71
72/// Options for querying ranges from a memtable.
73#[derive(Clone)]
74pub struct RangesOptions {
75    /// Whether the ranges are being queried for flush.
76    pub for_flush: bool,
77    /// Mode to pre-filter columns in ranges.
78    pub pre_filter_mode: PreFilterMode,
79    /// Predicate to filter the data.
80    pub predicate: PredicateGroup,
81    /// Sequence range to filter the data.
82    pub sequence: Option<SequenceRange>,
83    /// Maximum number of rows readers should produce in one batch.
84    pub batch_size: usize,
85}
86
87impl Default for RangesOptions {
88    fn default() -> Self {
89        Self {
90            for_flush: false,
91            pre_filter_mode: PreFilterMode::All,
92            predicate: PredicateGroup::default(),
93            sequence: None,
94            batch_size: crate::sst::parquet::DEFAULT_READ_BATCH_SIZE,
95        }
96    }
97}
98
99impl RangesOptions {
100    /// Creates a new [RangesOptions] for flushing.
101    pub fn for_flush() -> Self {
102        Self {
103            for_flush: true,
104            pre_filter_mode: PreFilterMode::All,
105            predicate: PredicateGroup::default(),
106            sequence: None,
107            batch_size: crate::sst::parquet::DEFAULT_READ_BATCH_SIZE,
108        }
109    }
110
111    /// Sets the pre-filter mode.
112    #[must_use]
113    pub fn with_pre_filter_mode(mut self, pre_filter_mode: PreFilterMode) -> Self {
114        self.pre_filter_mode = pre_filter_mode;
115        self
116    }
117
118    /// Sets the predicate.
119    #[must_use]
120    pub fn with_predicate(mut self, predicate: PredicateGroup) -> Self {
121        self.predicate = predicate;
122        self
123    }
124
125    /// Sets the sequence range.
126    #[must_use]
127    pub fn with_sequence(mut self, sequence: Option<SequenceRange>) -> Self {
128        self.sequence = sequence;
129        self
130    }
131
132    /// Sets the maximum number of rows readers should produce in one batch.
133    #[must_use]
134    pub fn with_batch_size(mut self, batch_size: usize) -> Self {
135        self.batch_size = batch_size.clamp(1, crate::sst::parquet::DEFAULT_READ_BATCH_SIZE);
136        self
137    }
138}
139
140#[derive(Debug, Default, Clone)]
141pub struct MemtableStats {
142    /// The estimated bytes allocated by this memtable from heap.
143    pub estimated_bytes: usize,
144    /// The inclusive time range that this memtable contains. It is None if
145    /// and only if the memtable is empty.
146    pub time_range: Option<(Timestamp, Timestamp)>,
147    /// Total rows in memtable
148    pub num_rows: usize,
149    /// Total number of ranges in the memtable.
150    pub num_ranges: usize,
151    /// The maximum sequence number in the memtable.
152    pub max_sequence: SequenceNumber,
153    /// Number of estimated timeseries in memtable.
154    pub series_count: usize,
155}
156
157impl MemtableStats {
158    /// Attaches the time range to the stats.
159    #[cfg(any(test, feature = "test"))]
160    pub fn with_time_range(mut self, time_range: Option<(Timestamp, Timestamp)>) -> Self {
161        self.time_range = time_range;
162        self
163    }
164
165    #[cfg(feature = "test")]
166    pub fn with_max_sequence(mut self, max_sequence: SequenceNumber) -> Self {
167        self.max_sequence = max_sequence;
168        self
169    }
170
171    /// Returns the estimated bytes allocated by this memtable.
172    pub fn bytes_allocated(&self) -> usize {
173        self.estimated_bytes
174    }
175
176    /// Returns the time range of the memtable.
177    pub fn time_range(&self) -> Option<(Timestamp, Timestamp)> {
178        self.time_range
179    }
180
181    /// Returns the num of total rows in memtable.
182    pub fn num_rows(&self) -> usize {
183        self.num_rows
184    }
185
186    /// Returns the number of ranges in the memtable.
187    pub fn num_ranges(&self) -> usize {
188        self.num_ranges
189    }
190
191    /// Returns the maximum sequence number in the memtable.
192    pub fn max_sequence(&self) -> SequenceNumber {
193        self.max_sequence
194    }
195
196    /// Series count in memtable.
197    pub fn series_count(&self) -> usize {
198        self.series_count
199    }
200}
201
202pub type BoxedBatchIterator = Box<dyn Iterator<Item = Result<Batch>> + Send>;
203
204pub type BoxedRecordBatchIterator = Box<dyn Iterator<Item = Result<RecordBatch>> + Send>;
205
206/// Ranges in a memtable.
207#[derive(Default)]
208pub struct MemtableRanges {
209    /// Range IDs and ranges.
210    pub ranges: BTreeMap<usize, MemtableRange>,
211}
212
213impl MemtableRanges {
214    /// Returns the total number of rows across all ranges.
215    pub fn num_rows(&self) -> usize {
216        self.ranges.values().map(|r| r.stats().num_rows()).sum()
217    }
218
219    /// Returns the total series count across all ranges.
220    pub fn series_count(&self) -> usize {
221        self.ranges.values().map(|r| r.stats().series_count()).sum()
222    }
223
224    /// Returns the maximum sequence number across all ranges.
225    pub fn max_sequence(&self) -> SequenceNumber {
226        self.ranges
227            .values()
228            .map(|r| r.stats().max_sequence())
229            .max()
230            .unwrap_or(0)
231    }
232}
233
234impl IterBuilder for MemtableRanges {
235    fn build(&self, _metrics: Option<MemScanMetrics>) -> Result<BoxedBatchIterator> {
236        ensure!(
237            self.ranges.len() == 1,
238            UnsupportedOperationSnafu {
239                err_msg: format!(
240                    "Building an iterator from MemtableRanges expects 1 range, but got {}",
241                    self.ranges.len()
242                ),
243            }
244        );
245
246        self.ranges.values().next().unwrap().build_iter()
247    }
248
249    fn is_record_batch(&self) -> bool {
250        self.ranges.values().all(|range| range.is_record_batch())
251    }
252}
253
254/// In memory write buffer.
255pub trait Memtable: Send + Sync + fmt::Debug {
256    /// Returns the id of this memtable.
257    fn id(&self) -> MemtableId;
258
259    /// Writes key values into the memtable.
260    fn write(&self, kvs: &KeyValues) -> Result<()>;
261
262    /// Writes one key value pair into the memtable.
263    fn write_one(&self, key_value: KeyValue) -> Result<()>;
264
265    /// Writes an encoded batch of into memtable.
266    fn write_bulk(&self, part: crate::memtable::bulk::part::BulkPart) -> Result<()>;
267
268    /// Returns the ranges in the memtable.
269    ///
270    /// The returned map contains the range id and the range after applying the predicate.
271    fn ranges(
272        &self,
273        projection: Option<&[ColumnId]>,
274        options: RangesOptions,
275    ) -> Result<MemtableRanges>;
276
277    /// Returns true if the memtable is empty.
278    fn is_empty(&self) -> bool;
279
280    /// Turns a mutable memtable into an immutable memtable.
281    fn freeze(&self) -> Result<()>;
282
283    /// Returns the [MemtableStats] info of Memtable.
284    fn stats(&self) -> MemtableStats;
285
286    /// Forks this (immutable) memtable and returns a new mutable memtable with specific memtable `id`.
287    ///
288    /// A region must freeze the memtable before invoking this method.
289    fn fork(&self, id: MemtableId, metadata: &RegionMetadataRef) -> MemtableRef;
290
291    /// Compacts the memtable.
292    ///
293    /// The `for_flush` is true when the flush job calls this method.
294    fn compact(&self, for_flush: bool) -> Result<()> {
295        let _ = for_flush;
296        Ok(())
297    }
298}
299
300pub type MemtableRef = Arc<dyn Memtable>;
301
302/// Builder to build a new [Memtable].
303pub trait MemtableBuilder: Send + Sync + fmt::Debug {
304    /// Builds a new memtable instance.
305    fn build(&self, id: MemtableId, metadata: &RegionMetadataRef) -> MemtableRef;
306
307    /// Returns true if the memtable supports bulk insert and benefits from it.
308    fn use_bulk_insert(&self, metadata: &RegionMetadataRef) -> bool {
309        let _metadata = metadata;
310        false
311    }
312}
313
314pub type MemtableBuilderRef = Arc<dyn MemtableBuilder>;
315
316/// Memtable memory allocation tracker.
317#[derive(Default)]
318pub struct AllocTracker {
319    write_buffer_manager: Option<WriteBufferManagerRef>,
320    /// Bytes allocated by the tracker.
321    bytes_allocated: AtomicUsize,
322    /// Whether allocating is done.
323    is_done_allocating: AtomicBool,
324}
325
326impl fmt::Debug for AllocTracker {
327    fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
328        f.debug_struct("AllocTracker")
329            .field("bytes_allocated", &self.bytes_allocated)
330            .field("is_done_allocating", &self.is_done_allocating)
331            .finish()
332    }
333}
334
335impl AllocTracker {
336    /// Returns a new [AllocTracker].
337    pub fn new(write_buffer_manager: Option<WriteBufferManagerRef>) -> AllocTracker {
338        AllocTracker {
339            write_buffer_manager,
340            bytes_allocated: AtomicUsize::new(0),
341            is_done_allocating: AtomicBool::new(false),
342        }
343    }
344
345    /// Tracks `bytes` memory is allocated.
346    pub(crate) fn on_allocation(&self, bytes: usize) {
347        self.bytes_allocated.fetch_add(bytes, Ordering::Relaxed);
348        WRITE_BUFFER_BYTES.add(bytes as i64);
349        if let Some(write_buffer_manager) = &self.write_buffer_manager {
350            write_buffer_manager.reserve_mem(bytes);
351        }
352    }
353
354    /// Marks we have finished allocating memory so we can free it from
355    /// the write buffer's limit.
356    ///
357    /// The region MUST ensure that it calls this method inside the region writer's write lock.
358    pub(crate) fn done_allocating(&self) {
359        if let Some(write_buffer_manager) = &self.write_buffer_manager
360            && self
361                .is_done_allocating
362                .compare_exchange(false, true, Ordering::Relaxed, Ordering::Relaxed)
363                .is_ok()
364        {
365            write_buffer_manager.schedule_free_mem(self.bytes_allocated.load(Ordering::Relaxed));
366        }
367    }
368
369    /// Returns bytes allocated.
370    pub(crate) fn bytes_allocated(&self) -> usize {
371        self.bytes_allocated.load(Ordering::Relaxed)
372    }
373
374    /// Returns the write buffer manager.
375    pub(crate) fn write_buffer_manager(&self) -> Option<WriteBufferManagerRef> {
376        self.write_buffer_manager.clone()
377    }
378}
379
380impl Drop for AllocTracker {
381    fn drop(&mut self) {
382        if !self.is_done_allocating.load(Ordering::Relaxed) {
383            self.done_allocating();
384        }
385
386        let bytes_allocated = self.bytes_allocated.load(Ordering::Relaxed);
387        WRITE_BUFFER_BYTES.sub(bytes_allocated as i64);
388
389        // Memory tracked by this tracker is freed.
390        if let Some(write_buffer_manager) = &self.write_buffer_manager {
391            write_buffer_manager.free_mem(bytes_allocated);
392        }
393    }
394}
395
396/// Provider of memtable builders for regions.
397#[derive(Clone)]
398pub(crate) struct MemtableBuilderProvider {
399    write_buffer_manager: Option<WriteBufferManagerRef>,
400    config: Arc<MitoConfig>,
401    compact_dispatcher: Arc<CompactDispatcher>,
402}
403
404/// Ensures JSON2 columns are not used with [`TimeSeriesMemtable`].
405pub(crate) fn ensure_json2_not_use_time_series_memtable(
406    metadata: &RegionMetadata,
407    options: &RegionOptions,
408) -> Result<()> {
409    if metadata
410        .column_metadatas
411        .iter()
412        .any(|x| x.column_schema.data_type.is_json2())
413    {
414        ensure!(
415            !matches!(&options.memtable, Some(MemtableOptions::TimeSeries)),
416            InvalidRegionOptionsSnafu {
417                reason: "JSON2 columns only support BulkMemtable",
418            }
419        );
420    }
421    Ok(())
422}
423
424impl MemtableBuilderProvider {
425    pub(crate) fn new(
426        write_buffer_manager: Option<WriteBufferManagerRef>,
427        config: Arc<MitoConfig>,
428    ) -> Self {
429        let compact_dispatcher =
430            Arc::new(CompactDispatcher::new(config.max_background_compactions));
431
432        Self {
433            write_buffer_manager,
434            config,
435            compact_dispatcher,
436        }
437    }
438
439    pub(crate) fn builder_for_options(&self, options: &RegionOptions) -> MemtableBuilderRef {
440        let dedup = options.need_dedup();
441        let merge_mode = options.merge_mode();
442        let primary_key_encoding = options.primary_key_encoding();
443        let flat_format = options
444            .sst_format
445            .map(|format| format == FormatType::Flat)
446            .unwrap_or(self.config.default_flat_format);
447        if flat_format {
448            if options.memtable.is_some()
449                && !matches!(&options.memtable, Some(MemtableOptions::Bulk(_)))
450            {
451                common_telemetry::info!(
452                    "Overriding memtable config, use BulkMemtable under flat format"
453                );
454            }
455
456            return Arc::new(self.bulk_memtable_builder(dedup, merge_mode, options));
457        }
458
459        if primary_key_encoding == PrimaryKeyEncoding::Sparse {
460            if options.memtable.is_some()
461                && !matches!(&options.memtable, Some(MemtableOptions::Bulk(_)))
462            {
463                common_telemetry::info!(
464                    "Overriding memtable config, use BulkMemtable for sparse primary key encoding"
465                );
466            }
467            return Arc::new(self.bulk_memtable_builder(dedup, merge_mode, options));
468        }
469
470        // The format is not flat.
471        match &options.memtable {
472            Some(MemtableOptions::Bulk(config)) => Arc::new(
473                BulkMemtableBuilder::new(self.write_buffer_manager.clone(), !dedup, merge_mode)
474                    .with_config(config.clone())
475                    .with_row_group_size(options.row_group_size())
476                    .with_compact_dispatcher(self.compact_dispatcher.clone()),
477            ),
478            Some(MemtableOptions::TimeSeries) => Arc::new(TimeSeriesMemtableBuilder::new(
479                self.write_buffer_manager.clone(),
480                dedup,
481                merge_mode,
482            )),
483            None => self.default_primary_key_memtable_builder(dedup, merge_mode),
484        }
485    }
486
487    fn bulk_memtable_builder(
488        &self,
489        dedup: bool,
490        merge_mode: MergeMode,
491        options: &RegionOptions,
492    ) -> BulkMemtableBuilder {
493        let mut builder = BulkMemtableBuilder::new(
494            self.write_buffer_manager.clone(),
495            !dedup, // append_mode: true if not dedup, false if dedup
496            merge_mode,
497        )
498        .with_row_group_size(options.row_group_size())
499        .with_compact_dispatcher(self.compact_dispatcher.clone());
500
501        if let Some(MemtableOptions::Bulk(config)) = &options.memtable {
502            builder = builder.with_config(config.clone());
503        }
504
505        builder
506    }
507
508    fn default_primary_key_memtable_builder(
509        &self,
510        dedup: bool,
511        merge_mode: MergeMode,
512    ) -> MemtableBuilderRef {
513        Arc::new(TimeSeriesMemtableBuilder::new(
514            self.write_buffer_manager.clone(),
515            dedup,
516            merge_mode,
517        ))
518    }
519}
520
521/// Metrics for scanning a memtable.
522#[derive(Clone, Default)]
523pub struct MemScanMetrics(Arc<Mutex<MemScanMetricsData>>);
524
525impl MemScanMetrics {
526    /// Merges the metrics.
527    pub(crate) fn merge_inner(&self, inner: &MemScanMetricsData) {
528        let mut metrics = self.0.lock().unwrap();
529        metrics.total_series += inner.total_series;
530        metrics.num_rows += inner.num_rows;
531        metrics.num_batches += inner.num_batches;
532        metrics.scan_cost += inner.scan_cost;
533        metrics.prefilter_cost += inner.prefilter_cost;
534        metrics.prefilter_rows_filtered += inner.prefilter_rows_filtered;
535    }
536
537    /// Gets the metrics data.
538    pub(crate) fn data(&self) -> MemScanMetricsData {
539        self.0.lock().unwrap().clone()
540    }
541}
542
543#[derive(Clone, Default)]
544pub(crate) struct MemScanMetricsData {
545    /// Total series in the memtable.
546    pub(crate) total_series: usize,
547    /// Number of rows read.
548    pub(crate) num_rows: usize,
549    /// Number of batch read.
550    pub(crate) num_batches: usize,
551    /// Duration to scan the memtable.
552    pub(crate) scan_cost: Duration,
553    /// Duration of prefilter in memtable scan.
554    pub(crate) prefilter_cost: Duration,
555    /// Number of rows filtered by prefilter in memtable scan.
556    pub(crate) prefilter_rows_filtered: usize,
557}
558
559/// Encoded range in the memtable.
560pub struct EncodedRange {
561    /// Encoded file data.
562    pub data: Bytes,
563    /// Metadata of the encoded range.
564    pub sst_info: SstInfo,
565}
566
567/// Builder to build an iterator to read the range.
568/// The builder should know the projection and the predicate to build the iterator.
569pub trait IterBuilder: Send + Sync {
570    /// Returns the iterator to read the range.
571    fn build(&self, metrics: Option<MemScanMetrics>) -> Result<BoxedBatchIterator>;
572
573    /// Returns whether the iterator is a record batch iterator.
574    fn is_record_batch(&self) -> bool {
575        false
576    }
577
578    /// Returns the record batch iterator to read the range.
579    /// ## Note
580    /// Implementations should ensure the iterator yields data within given time range.
581    fn build_record_batch(
582        &self,
583        time_range: Option<(Timestamp, Timestamp)>,
584        metrics: Option<MemScanMetrics>,
585    ) -> Result<BoxedRecordBatchIterator> {
586        let _metrics = metrics;
587        let _ = time_range;
588        UnsupportedOperationSnafu {
589            err_msg: "Record batch iterator is not supported by this memtable",
590        }
591        .fail()
592    }
593
594    /// Returns a cheap schema hint for record batches yielded by this builder.
595    fn record_batch_schema_hint(&self) -> Option<SchemaRef> {
596        None
597    }
598
599    /// Returns the [EncodedRange] if the range is already encoded into SST.
600    fn encoded_range(&self) -> Option<EncodedRange> {
601        None
602    }
603}
604
605pub type BoxedIterBuilder = Box<dyn IterBuilder>;
606
607/// Computes the column IDs to read based on the projection.
608///
609/// If `projection` is `Some`, returns those column IDs. If `None`, returns all column IDs
610/// from the metadata.
611pub fn read_column_ids_from_projection(
612    metadata: &RegionMetadataRef,
613    projection: Option<&[ColumnId]>,
614) -> Vec<ColumnId> {
615    if let Some(projection) = projection {
616        projection.to_vec()
617    } else {
618        metadata
619            .column_metadatas
620            .iter()
621            .map(|c| c.column_id)
622            .collect()
623    }
624}
625
626/// Context to adapt batch iterators to record batch iterators for flat scan.
627pub struct BatchToRecordBatchContext {
628    metadata: RegionMetadataRef,
629    codec: Arc<dyn PrimaryKeyCodec>,
630    read_column_ids: Vec<ColumnId>,
631}
632
633impl BatchToRecordBatchContext {
634    /// Creates a new context for adapting batch iterators.
635    pub fn new(metadata: RegionMetadataRef, mut read_column_ids: Vec<ColumnId>) -> Self {
636        if read_column_ids.is_empty() {
637            read_column_ids.push(metadata.time_index_column().column_id);
638        }
639
640        let codec = build_primary_key_codec(&metadata);
641        Self {
642            metadata,
643            codec,
644            read_column_ids,
645        }
646    }
647
648    fn adapt_iter(&self, iter: BoxedBatchIterator) -> BoxedRecordBatchIterator {
649        Box::new(BatchToRecordBatchAdapter::new(
650            iter,
651            self.metadata.clone(),
652            self.codec.clone(),
653            &self.read_column_ids,
654        ))
655    }
656}
657
658/// Context shared by ranges of the same memtable.
659pub struct MemtableRangeContext {
660    /// Id of the memtable.
661    id: MemtableId,
662    /// Iterator builder.
663    builder: BoxedIterBuilder,
664    /// All filters.
665    predicate: PredicateGroup,
666    /// Optional context to adapt batch iterators for flat scans.
667    batch_to_record_batch: Option<Arc<BatchToRecordBatchContext>>,
668}
669
670pub type MemtableRangeContextRef = Arc<MemtableRangeContext>;
671
672impl MemtableRangeContext {
673    /// Creates a new [MemtableRangeContext].
674    pub fn new(id: MemtableId, builder: BoxedIterBuilder, predicate: PredicateGroup) -> Self {
675        Self::new_with_batch_to_record_batch(id, builder, predicate, None)
676    }
677
678    /// Creates a new [MemtableRangeContext] with optional adapter context.
679    pub fn new_with_batch_to_record_batch(
680        id: MemtableId,
681        builder: BoxedIterBuilder,
682        predicate: PredicateGroup,
683        batch_to_record_batch: Option<Arc<BatchToRecordBatchContext>>,
684    ) -> Self {
685        Self {
686            id,
687            builder,
688            predicate,
689            batch_to_record_batch,
690        }
691    }
692}
693
694/// A range in the memtable.
695#[derive(Clone)]
696pub struct MemtableRange {
697    /// Shared context.
698    context: MemtableRangeContextRef,
699    /// Statistics for this memtable range.
700    stats: MemtableStats,
701}
702
703impl MemtableRange {
704    /// Creates a new range from context and stats.
705    pub fn new(context: MemtableRangeContextRef, stats: MemtableStats) -> Self {
706        Self { context, stats }
707    }
708
709    /// Returns the statistics for this range.
710    pub fn stats(&self) -> &MemtableStats {
711        &self.stats
712    }
713
714    /// Returns the id of the memtable to read.
715    pub fn id(&self) -> MemtableId {
716        self.context.id
717    }
718
719    /// Builds an iterator to read the range.
720    /// Filters the result by the specific time range, this ensures memtable won't return
721    /// rows out of the time range when new rows are inserted.
722    pub fn build_prune_iter(
723        &self,
724        time_range: FileTimeRange,
725        metrics: Option<MemScanMetrics>,
726    ) -> Result<BoxedBatchIterator> {
727        let iter = self.context.builder.build(metrics)?;
728        let time_filters = self.context.predicate.time_filters();
729        Ok(Box::new(PruneTimeIterator::new(
730            iter,
731            time_range,
732            time_filters,
733        )))
734    }
735
736    /// Builds an iterator to read all rows in range.
737    pub fn build_iter(&self) -> Result<BoxedBatchIterator> {
738        self.context.builder.build(None)
739    }
740
741    /// Builds a record batch iterator to read rows in range.
742    ///
743    /// For mutable memtables (adapter path), applies time-range pruning to ensure rows
744    /// outside the time range are filtered, matching the behavior of `build_prune_iter`.
745    pub fn build_record_batch_iter(
746        &self,
747        time_range: Option<FileTimeRange>,
748        metrics: Option<MemScanMetrics>,
749    ) -> Result<BoxedRecordBatchIterator> {
750        if self.context.builder.is_record_batch() {
751            return self.context.builder.build_record_batch(time_range, metrics);
752        }
753
754        if let Some(context) = self.context.batch_to_record_batch.as_ref() {
755            let iter = self.context.builder.build(metrics)?;
756            let iter: BoxedBatchIterator = if let Some(time_range) = time_range {
757                let time_filters = self.context.predicate.time_filters();
758                Box::new(PruneTimeIterator::new(iter, time_range, time_filters))
759            } else {
760                iter
761            };
762            return Ok(context.adapt_iter(iter));
763        }
764
765        UnsupportedOperationSnafu {
766            err_msg: "Record batch iterator is not supported by this memtable",
767        }
768        .fail()
769    }
770
771    /// Returns a cheap schema hint for record batches yielded by this range.
772    pub fn record_batch_schema_hint(&self) -> Option<SchemaRef> {
773        self.context.builder.record_batch_schema_hint()
774    }
775
776    /// Returns whether the iterator is a record batch iterator.
777    pub fn is_record_batch(&self) -> bool {
778        self.context.builder.is_record_batch()
779    }
780
781    pub fn num_rows(&self) -> usize {
782        self.stats.num_rows
783    }
784
785    /// Returns the encoded range if available.
786    pub fn encoded(&self) -> Option<EncodedRange> {
787        self.context.builder.encoded_range()
788    }
789}
790
791#[cfg(test)]
792mod tests {
793    use std::sync::Arc;
794
795    use common_error::ext::WhateverResult;
796    use datatypes::prelude::ConcreteDataType;
797    use datatypes::types::json_type::{JsonNativeType, JsonObjectType};
798    use store_api::metadata::RegionMetadataBuilder;
799
800    use super::*;
801    use crate::flush::{WriteBufferManager, WriteBufferManagerImpl};
802    use crate::memtable::bulk::BulkMemtableConfig;
803    use crate::test_util::sst_util::sst_region_metadata;
804
805    #[test]
806    fn test_alloc_tracker_without_manager() {
807        let tracker = AllocTracker::new(None);
808        assert_eq!(0, tracker.bytes_allocated());
809        tracker.on_allocation(100);
810        assert_eq!(100, tracker.bytes_allocated());
811        tracker.on_allocation(200);
812        assert_eq!(300, tracker.bytes_allocated());
813
814        tracker.done_allocating();
815        assert_eq!(300, tracker.bytes_allocated());
816    }
817
818    #[test]
819    fn test_alloc_tracker_with_manager() {
820        let manager = Arc::new(WriteBufferManagerImpl::new(1000));
821        {
822            let tracker = AllocTracker::new(Some(manager.clone() as WriteBufferManagerRef));
823
824            tracker.on_allocation(100);
825            assert_eq!(100, tracker.bytes_allocated());
826            assert_eq!(100, manager.memory_usage());
827            assert_eq!(100, manager.mutable_usage());
828
829            for _ in 0..2 {
830                // Done allocating won't free the same memory multiple times.
831                tracker.done_allocating();
832                assert_eq!(100, manager.memory_usage());
833                assert_eq!(0, manager.mutable_usage());
834            }
835        }
836
837        assert_eq!(0, manager.memory_usage());
838        assert_eq!(0, manager.mutable_usage());
839    }
840
841    #[test]
842    fn test_alloc_tracker_without_done_allocating() {
843        let manager = Arc::new(WriteBufferManagerImpl::new(1000));
844        {
845            let tracker = AllocTracker::new(Some(manager.clone() as WriteBufferManagerRef));
846
847            tracker.on_allocation(100);
848            assert_eq!(100, tracker.bytes_allocated());
849            assert_eq!(100, manager.memory_usage());
850            assert_eq!(100, manager.mutable_usage());
851        }
852
853        assert_eq!(0, manager.memory_usage());
854        assert_eq!(0, manager.mutable_usage());
855    }
856
857    #[test]
858    fn test_forced_bulk_memtable_preserves_bulk_config() {
859        let provider = MemtableBuilderProvider::new(None, Arc::new(MitoConfig::default()));
860        let config = BulkMemtableConfig {
861            merge_threshold: 7,
862            encode_row_threshold: 11,
863            encode_bytes_threshold: 13,
864            max_merge_groups: 17,
865        };
866        let options = RegionOptions {
867            memtable: Some(MemtableOptions::Bulk(config.clone())),
868            primary_key_encoding: Some(PrimaryKeyEncoding::Sparse),
869            ..Default::default()
870        };
871
872        let builder =
873            provider.bulk_memtable_builder(options.need_dedup(), options.merge_mode(), &options);
874
875        assert_eq!(&config, builder.config());
876    }
877
878    #[test]
879    fn test_json2_requires_bulk_memtable() -> WhateverResult<()> {
880        let mut metadata = sst_region_metadata();
881        metadata.column_metadatas[2].column_schema.data_type =
882            ConcreteDataType::json2(JsonNativeType::Object(JsonObjectType::new()));
883        let metadata = RegionMetadataBuilder::from_existing(metadata).build()?;
884        let mut options = RegionOptions {
885            sst_format: Some(FormatType::PrimaryKey),
886            memtable: Some(MemtableOptions::TimeSeries),
887            ..Default::default()
888        };
889
890        let err = ensure_json2_not_use_time_series_memtable(&metadata, &options).unwrap_err();
891        assert!(
892            err.to_string()
893                .contains("JSON2 columns only support BulkMemtable")
894        );
895
896        options.memtable = Some(MemtableOptions::Bulk(BulkMemtableConfig::default()));
897        ensure_json2_not_use_time_series_memtable(&metadata, &options)?;
898        Ok(())
899    }
900}