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, HashMap};
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, RegionId, SequenceNumber, SequenceRange};
35
36use crate::config::MitoConfig;
37use crate::error::{InvalidRegionOptionsSnafu, Result, UnsupportedOperationSnafu};
38use crate::flush::WriteBufferManagerRef;
39use crate::memtable::bulk::{BulkMemtableBuilder, BulkMemtableConfig, 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    /// Lower bound of the row sequences written to this memtable or range.
154    /// May remain below the actual minimum after deduplication. Zero is conservative
155    /// when a range does not track its minimum; the value is unused for empty memtables.
156    pub min_sequence: SequenceNumber,
157    /// Number of estimated timeseries in memtable.
158    pub series_count: usize,
159}
160
161impl MemtableStats {
162    /// Attaches the time range to the stats.
163    #[cfg(any(test, feature = "test"))]
164    pub fn with_time_range(mut self, time_range: Option<(Timestamp, Timestamp)>) -> Self {
165        self.time_range = time_range;
166        self
167    }
168
169    #[cfg(feature = "test")]
170    pub fn with_max_sequence(mut self, max_sequence: SequenceNumber) -> Self {
171        self.max_sequence = max_sequence;
172        self
173    }
174
175    /// Returns the estimated bytes allocated by this memtable.
176    pub fn bytes_allocated(&self) -> usize {
177        self.estimated_bytes
178    }
179
180    /// Returns the time range of the memtable.
181    pub fn time_range(&self) -> Option<(Timestamp, Timestamp)> {
182        self.time_range
183    }
184
185    /// Returns the num of total rows in memtable.
186    pub fn num_rows(&self) -> usize {
187        self.num_rows
188    }
189
190    /// Returns the number of ranges in the memtable.
191    pub fn num_ranges(&self) -> usize {
192        self.num_ranges
193    }
194
195    /// Returns the maximum sequence number in the memtable.
196    pub fn max_sequence(&self) -> SequenceNumber {
197        self.max_sequence
198    }
199
200    /// Series count in memtable.
201    pub fn series_count(&self) -> usize {
202        self.series_count
203    }
204}
205
206pub type BoxedBatchIterator = Box<dyn Iterator<Item = Result<Batch>> + Send>;
207
208pub type BoxedRecordBatchIterator = Box<dyn Iterator<Item = Result<RecordBatch>> + Send>;
209
210/// Ranges in a memtable.
211#[derive(Default)]
212pub struct MemtableRanges {
213    /// Range IDs and ranges.
214    pub ranges: BTreeMap<usize, MemtableRange>,
215}
216
217impl MemtableRanges {
218    /// Returns the total number of rows across all ranges.
219    pub fn num_rows(&self) -> usize {
220        self.ranges.values().map(|r| r.stats().num_rows()).sum()
221    }
222
223    /// Returns the total series count across all ranges.
224    pub fn series_count(&self) -> usize {
225        self.ranges.values().map(|r| r.stats().series_count()).sum()
226    }
227
228    /// Returns the maximum sequence number across all ranges.
229    pub fn max_sequence(&self) -> SequenceNumber {
230        self.ranges
231            .values()
232            .map(|r| r.stats().max_sequence())
233            .max()
234            .unwrap_or(0)
235    }
236}
237
238impl IterBuilder for MemtableRanges {
239    fn build(&self, _metrics: Option<MemScanMetrics>) -> Result<BoxedBatchIterator> {
240        ensure!(
241            self.ranges.len() == 1,
242            UnsupportedOperationSnafu {
243                err_msg: format!(
244                    "Building an iterator from MemtableRanges expects 1 range, but got {}",
245                    self.ranges.len()
246                ),
247            }
248        );
249
250        self.ranges.values().next().unwrap().build_iter()
251    }
252
253    fn is_record_batch(&self) -> bool {
254        self.ranges.values().all(|range| range.is_record_batch())
255    }
256}
257
258/// In memory write buffer.
259pub trait Memtable: Send + Sync + fmt::Debug {
260    /// Returns the id of this memtable.
261    fn id(&self) -> MemtableId;
262
263    /// Writes key values into the memtable.
264    fn write(&self, kvs: &KeyValues) -> Result<()>;
265
266    /// Writes one key value pair into the memtable.
267    fn write_one(&self, key_value: KeyValue) -> Result<()>;
268
269    /// Writes an encoded batch of into memtable.
270    fn write_bulk(&self, part: crate::memtable::bulk::part::BulkPart) -> Result<()>;
271
272    /// Returns the ranges in the memtable.
273    ///
274    /// The returned map contains the range id and the range after applying the predicate.
275    fn ranges(
276        &self,
277        projection: Option<&[ColumnId]>,
278        options: RangesOptions,
279    ) -> Result<MemtableRanges>;
280
281    /// Returns true if the memtable is empty.
282    fn is_empty(&self) -> bool;
283
284    /// Turns a mutable memtable into an immutable memtable.
285    fn freeze(&self) -> Result<()>;
286
287    /// Returns the [MemtableStats] info of Memtable.
288    fn stats(&self) -> MemtableStats;
289
290    /// Returns a conservative row-sequence lower bound without computing full statistics.
291    /// Returns zero when write statistics are unavailable. Callers must ignore
292    /// empty memtables when using this bound to constrain compaction.
293    fn min_sequence(&self) -> SequenceNumber;
294
295    /// Forks this (immutable) memtable and returns a new mutable memtable with specific memtable `id`.
296    ///
297    /// A region must freeze the memtable before invoking this method.
298    fn fork(&self, id: MemtableId, metadata: &RegionMetadataRef) -> MemtableRef;
299
300    /// Compacts the memtable.
301    ///
302    /// The `for_flush` is true when the flush job calls this method.
303    fn compact(&self, for_flush: bool) -> Result<()> {
304        let _ = for_flush;
305        Ok(())
306    }
307}
308
309pub type MemtableRef = Arc<dyn Memtable>;
310
311/// Builder to build a new [Memtable].
312pub trait MemtableBuilder: Send + Sync + fmt::Debug {
313    /// Builds a new memtable instance.
314    fn build(&self, id: MemtableId, metadata: &RegionMetadataRef) -> MemtableRef;
315
316    /// Returns true if the memtable supports bulk insert and benefits from it.
317    fn use_bulk_insert(&self, metadata: &RegionMetadataRef) -> bool {
318        let _metadata = metadata;
319        false
320    }
321}
322
323pub type MemtableBuilderRef = Arc<dyn MemtableBuilder>;
324
325/// Memtable memory allocation tracker.
326#[derive(Default)]
327pub struct AllocTracker {
328    write_buffer_manager: Option<WriteBufferManagerRef>,
329    /// Bytes allocated by the tracker.
330    bytes_allocated: AtomicUsize,
331    /// Whether allocating is done.
332    is_done_allocating: AtomicBool,
333}
334
335impl fmt::Debug for AllocTracker {
336    fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
337        f.debug_struct("AllocTracker")
338            .field("bytes_allocated", &self.bytes_allocated)
339            .field("is_done_allocating", &self.is_done_allocating)
340            .finish()
341    }
342}
343
344impl AllocTracker {
345    /// Returns a new [AllocTracker].
346    pub fn new(write_buffer_manager: Option<WriteBufferManagerRef>) -> AllocTracker {
347        AllocTracker {
348            write_buffer_manager,
349            bytes_allocated: AtomicUsize::new(0),
350            is_done_allocating: AtomicBool::new(false),
351        }
352    }
353
354    /// Tracks `bytes` memory is allocated.
355    pub(crate) fn on_allocation(&self, bytes: usize) {
356        self.bytes_allocated.fetch_add(bytes, Ordering::Relaxed);
357        WRITE_BUFFER_BYTES.add(bytes as i64);
358        if let Some(write_buffer_manager) = &self.write_buffer_manager {
359            write_buffer_manager.reserve_mem(bytes);
360        }
361    }
362
363    /// Marks we have finished allocating memory so we can free it from
364    /// the write buffer's limit.
365    ///
366    /// The region MUST ensure that it calls this method inside the region writer's write lock.
367    pub(crate) fn done_allocating(&self) {
368        if let Some(write_buffer_manager) = &self.write_buffer_manager
369            && self
370                .is_done_allocating
371                .compare_exchange(false, true, Ordering::Relaxed, Ordering::Relaxed)
372                .is_ok()
373        {
374            write_buffer_manager.schedule_free_mem(self.bytes_allocated.load(Ordering::Relaxed));
375        }
376    }
377
378    /// Returns bytes allocated.
379    pub(crate) fn bytes_allocated(&self) -> usize {
380        self.bytes_allocated.load(Ordering::Relaxed)
381    }
382
383    /// Returns the write buffer manager.
384    pub(crate) fn write_buffer_manager(&self) -> Option<WriteBufferManagerRef> {
385        self.write_buffer_manager.clone()
386    }
387}
388
389impl Drop for AllocTracker {
390    fn drop(&mut self) {
391        if !self.is_done_allocating.load(Ordering::Relaxed) {
392            self.done_allocating();
393        }
394
395        let bytes_allocated = self.bytes_allocated.load(Ordering::Relaxed);
396        WRITE_BUFFER_BYTES.sub(bytes_allocated as i64);
397
398        // Memory tracked by this tracker is freed.
399        if let Some(write_buffer_manager) = &self.write_buffer_manager {
400            write_buffer_manager.free_mem(bytes_allocated);
401        }
402    }
403}
404
405/// Provider of memtable builders for regions.
406#[derive(Clone)]
407pub(crate) struct MemtableBuilderProvider {
408    write_buffer_manager: Option<WriteBufferManagerRef>,
409    config: Arc<MitoConfig>,
410    default_bulk_memtable_config: BulkMemtableConfig,
411    compact_dispatcher: Arc<CompactDispatcher>,
412}
413
414/// Ensures JSON2 columns are not used with [`TimeSeriesMemtable`].
415pub(crate) fn ensure_json2_not_use_time_series_memtable(
416    metadata: &RegionMetadata,
417    options: &RegionOptions,
418) -> Result<()> {
419    if metadata
420        .column_metadatas
421        .iter()
422        .any(|x| x.column_schema.data_type.is_json2())
423    {
424        ensure!(
425            !matches!(&options.memtable, Some(MemtableOptions::TimeSeries)),
426            InvalidRegionOptionsSnafu {
427                reason: "JSON2 columns only support BulkMemtable",
428            }
429        );
430    }
431    Ok(())
432}
433
434impl MemtableBuilderProvider {
435    pub(crate) fn new(
436        write_buffer_manager: Option<WriteBufferManagerRef>,
437        config: Arc<MitoConfig>,
438    ) -> Self {
439        let compact_dispatcher =
440            Arc::new(CompactDispatcher::new(config.max_background_compactions));
441        let default_bulk_memtable_config = BulkMemtableConfig::default_for_write_buffer_size(
442            config.global_write_buffer_size.as_bytes() as usize,
443        );
444
445        Self {
446            write_buffer_manager,
447            config,
448            default_bulk_memtable_config,
449            compact_dispatcher,
450        }
451    }
452
453    /// Parses region options with this provider's default bulk memtable config.
454    pub(crate) fn parse_options(
455        &self,
456        region_id: RegionId,
457        options: &HashMap<String, String>,
458    ) -> Result<RegionOptions> {
459        RegionOptions::try_from_options_with_bulk_config(
460            region_id,
461            options,
462            &self.default_bulk_memtable_config,
463        )
464    }
465
466    pub(crate) fn builder_for_options(&self, options: &RegionOptions) -> MemtableBuilderRef {
467        let dedup = options.need_dedup();
468        let merge_mode = options.merge_mode();
469        let primary_key_encoding = options.primary_key_encoding();
470        let flat_format = options
471            .sst_format
472            .map(|format| format == FormatType::Flat)
473            .unwrap_or(self.config.default_flat_format);
474        if flat_format {
475            if options.memtable.is_some()
476                && !matches!(&options.memtable, Some(MemtableOptions::Bulk(_)))
477            {
478                common_telemetry::info!(
479                    "Overriding memtable config, use BulkMemtable under flat format"
480                );
481            }
482
483            return Arc::new(self.bulk_memtable_builder(dedup, merge_mode, options));
484        }
485
486        if primary_key_encoding == PrimaryKeyEncoding::Sparse {
487            if options.memtable.is_some()
488                && !matches!(&options.memtable, Some(MemtableOptions::Bulk(_)))
489            {
490                common_telemetry::info!(
491                    "Overriding memtable config, use BulkMemtable for sparse primary key encoding"
492                );
493            }
494            return Arc::new(self.bulk_memtable_builder(dedup, merge_mode, options));
495        }
496
497        // The format is not flat.
498        match &options.memtable {
499            Some(MemtableOptions::Bulk(config)) => Arc::new(
500                BulkMemtableBuilder::new(self.write_buffer_manager.clone(), !dedup, merge_mode)
501                    .with_config(config.clone())
502                    .with_row_group_size(options.row_group_size())
503                    .with_float_field_encoding(options.float_field_encoding)
504                    .with_compact_dispatcher(self.compact_dispatcher.clone()),
505            ),
506            Some(MemtableOptions::TimeSeries) => Arc::new(TimeSeriesMemtableBuilder::new(
507                self.write_buffer_manager.clone(),
508                dedup,
509                merge_mode,
510            )),
511            None => self.default_primary_key_memtable_builder(dedup, merge_mode),
512        }
513    }
514
515    fn bulk_memtable_builder(
516        &self,
517        dedup: bool,
518        merge_mode: MergeMode,
519        options: &RegionOptions,
520    ) -> BulkMemtableBuilder {
521        let mut builder = BulkMemtableBuilder::new(
522            self.write_buffer_manager.clone(),
523            !dedup, // append_mode: true if not dedup, false if dedup
524            merge_mode,
525        )
526        .with_config(self.default_bulk_memtable_config.clone())
527        .with_row_group_size(options.row_group_size())
528        .with_float_field_encoding(options.float_field_encoding)
529        .with_compact_dispatcher(self.compact_dispatcher.clone());
530
531        if let Some(MemtableOptions::Bulk(config)) = &options.memtable {
532            builder = builder.with_config(config.clone());
533        }
534
535        builder
536    }
537
538    fn default_primary_key_memtable_builder(
539        &self,
540        dedup: bool,
541        merge_mode: MergeMode,
542    ) -> MemtableBuilderRef {
543        Arc::new(TimeSeriesMemtableBuilder::new(
544            self.write_buffer_manager.clone(),
545            dedup,
546            merge_mode,
547        ))
548    }
549}
550
551/// Metrics for scanning a memtable.
552#[derive(Clone, Default)]
553pub struct MemScanMetrics(Arc<Mutex<MemScanMetricsData>>);
554
555impl MemScanMetrics {
556    /// Merges the metrics.
557    pub(crate) fn merge_inner(&self, inner: &MemScanMetricsData) {
558        let mut metrics = self.0.lock().unwrap();
559        metrics.total_series += inner.total_series;
560        metrics.num_rows += inner.num_rows;
561        metrics.num_batches += inner.num_batches;
562        metrics.scan_cost += inner.scan_cost;
563        metrics.prefilter_cost += inner.prefilter_cost;
564        metrics.prefilter_rows_filtered += inner.prefilter_rows_filtered;
565    }
566
567    /// Gets the metrics data.
568    pub(crate) fn data(&self) -> MemScanMetricsData {
569        self.0.lock().unwrap().clone()
570    }
571}
572
573#[derive(Clone, Default)]
574pub(crate) struct MemScanMetricsData {
575    /// Total series in the memtable.
576    pub(crate) total_series: usize,
577    /// Number of rows read.
578    pub(crate) num_rows: usize,
579    /// Number of batch read.
580    pub(crate) num_batches: usize,
581    /// Duration to scan the memtable.
582    pub(crate) scan_cost: Duration,
583    /// Duration of prefilter in memtable scan.
584    pub(crate) prefilter_cost: Duration,
585    /// Number of rows filtered by prefilter in memtable scan.
586    pub(crate) prefilter_rows_filtered: usize,
587}
588
589/// Encoded range in the memtable.
590pub struct EncodedRange {
591    /// Encoded file data.
592    pub data: Bytes,
593    /// Metadata of the encoded range.
594    pub sst_info: SstInfo,
595}
596
597/// Builder to build an iterator to read the range.
598/// The builder should know the projection and the predicate to build the iterator.
599pub trait IterBuilder: Send + Sync {
600    /// Returns the iterator to read the range.
601    fn build(&self, metrics: Option<MemScanMetrics>) -> Result<BoxedBatchIterator>;
602
603    /// Returns whether the iterator is a record batch iterator.
604    fn is_record_batch(&self) -> bool {
605        false
606    }
607
608    /// Returns the record batch iterator to read the range.
609    /// ## Note
610    /// Implementations should ensure the iterator yields data within given time range.
611    fn build_record_batch(
612        &self,
613        time_range: Option<(Timestamp, Timestamp)>,
614        metrics: Option<MemScanMetrics>,
615    ) -> Result<BoxedRecordBatchIterator> {
616        let _metrics = metrics;
617        let _ = time_range;
618        UnsupportedOperationSnafu {
619            err_msg: "Record batch iterator is not supported by this memtable",
620        }
621        .fail()
622    }
623
624    /// Returns a cheap schema hint for record batches yielded by this builder.
625    fn record_batch_schema_hint(&self) -> Option<SchemaRef> {
626        None
627    }
628
629    /// Returns the [EncodedRange] if the range is already encoded into SST.
630    fn encoded_range(&self) -> Option<EncodedRange> {
631        None
632    }
633}
634
635pub type BoxedIterBuilder = Box<dyn IterBuilder>;
636
637/// Computes the column IDs to read based on the projection.
638///
639/// If `projection` is `Some`, returns those column IDs. If `None`, returns all column IDs
640/// from the metadata.
641pub fn read_column_ids_from_projection(
642    metadata: &RegionMetadataRef,
643    projection: Option<&[ColumnId]>,
644) -> Vec<ColumnId> {
645    if let Some(projection) = projection {
646        projection.to_vec()
647    } else {
648        metadata
649            .column_metadatas
650            .iter()
651            .map(|c| c.column_id)
652            .collect()
653    }
654}
655
656/// Context to adapt batch iterators to record batch iterators for flat scan.
657pub struct BatchToRecordBatchContext {
658    metadata: RegionMetadataRef,
659    codec: Arc<dyn PrimaryKeyCodec>,
660    read_column_ids: Vec<ColumnId>,
661}
662
663impl BatchToRecordBatchContext {
664    /// Creates a new context for adapting batch iterators.
665    pub fn new(metadata: RegionMetadataRef, mut read_column_ids: Vec<ColumnId>) -> Self {
666        if read_column_ids.is_empty() {
667            read_column_ids.push(metadata.time_index_column().column_id);
668        }
669
670        let codec = build_primary_key_codec(&metadata);
671        Self {
672            metadata,
673            codec,
674            read_column_ids,
675        }
676    }
677
678    fn adapt_iter(&self, iter: BoxedBatchIterator) -> BoxedRecordBatchIterator {
679        Box::new(BatchToRecordBatchAdapter::new(
680            iter,
681            self.metadata.clone(),
682            self.codec.clone(),
683            &self.read_column_ids,
684        ))
685    }
686}
687
688/// Context shared by ranges of the same memtable.
689pub struct MemtableRangeContext {
690    /// Id of the memtable.
691    id: MemtableId,
692    /// Iterator builder.
693    builder: BoxedIterBuilder,
694    /// All filters.
695    predicate: PredicateGroup,
696    /// Optional context to adapt batch iterators for flat scans.
697    batch_to_record_batch: Option<Arc<BatchToRecordBatchContext>>,
698}
699
700pub type MemtableRangeContextRef = Arc<MemtableRangeContext>;
701
702impl MemtableRangeContext {
703    /// Creates a new [MemtableRangeContext].
704    pub fn new(id: MemtableId, builder: BoxedIterBuilder, predicate: PredicateGroup) -> Self {
705        Self::new_with_batch_to_record_batch(id, builder, predicate, None)
706    }
707
708    /// Creates a new [MemtableRangeContext] with optional adapter context.
709    pub fn new_with_batch_to_record_batch(
710        id: MemtableId,
711        builder: BoxedIterBuilder,
712        predicate: PredicateGroup,
713        batch_to_record_batch: Option<Arc<BatchToRecordBatchContext>>,
714    ) -> Self {
715        Self {
716            id,
717            builder,
718            predicate,
719            batch_to_record_batch,
720        }
721    }
722}
723
724/// A range in the memtable.
725#[derive(Clone)]
726pub struct MemtableRange {
727    /// Shared context.
728    context: MemtableRangeContextRef,
729    /// Statistics for this memtable range.
730    stats: MemtableStats,
731}
732
733impl MemtableRange {
734    /// Creates a new range from context and stats.
735    pub fn new(context: MemtableRangeContextRef, stats: MemtableStats) -> Self {
736        Self { context, stats }
737    }
738
739    /// Returns the statistics for this range.
740    pub fn stats(&self) -> &MemtableStats {
741        &self.stats
742    }
743
744    /// Returns the id of the memtable to read.
745    pub fn id(&self) -> MemtableId {
746        self.context.id
747    }
748
749    /// Builds an iterator to read the range.
750    /// Filters the result by the specific time range, this ensures memtable won't return
751    /// rows out of the time range when new rows are inserted.
752    pub fn build_prune_iter(
753        &self,
754        time_range: FileTimeRange,
755        metrics: Option<MemScanMetrics>,
756    ) -> Result<BoxedBatchIterator> {
757        let iter = self.context.builder.build(metrics)?;
758        let time_filters = self.context.predicate.time_filters();
759        Ok(Box::new(PruneTimeIterator::new(
760            iter,
761            time_range,
762            time_filters,
763        )))
764    }
765
766    /// Builds an iterator to read all rows in range.
767    pub fn build_iter(&self) -> Result<BoxedBatchIterator> {
768        self.context.builder.build(None)
769    }
770
771    /// Builds a record batch iterator to read rows in range.
772    ///
773    /// For mutable memtables (adapter path), applies time-range pruning to ensure rows
774    /// outside the time range are filtered, matching the behavior of `build_prune_iter`.
775    pub fn build_record_batch_iter(
776        &self,
777        time_range: Option<FileTimeRange>,
778        metrics: Option<MemScanMetrics>,
779    ) -> Result<BoxedRecordBatchIterator> {
780        if self.context.builder.is_record_batch() {
781            return self.context.builder.build_record_batch(time_range, metrics);
782        }
783
784        if let Some(context) = self.context.batch_to_record_batch.as_ref() {
785            let iter = self.context.builder.build(metrics)?;
786            let iter: BoxedBatchIterator = if let Some(time_range) = time_range {
787                let time_filters = self.context.predicate.time_filters();
788                Box::new(PruneTimeIterator::new(iter, time_range, time_filters))
789            } else {
790                iter
791            };
792            return Ok(context.adapt_iter(iter));
793        }
794
795        UnsupportedOperationSnafu {
796            err_msg: "Record batch iterator is not supported by this memtable",
797        }
798        .fail()
799    }
800
801    /// Returns a cheap schema hint for record batches yielded by this range.
802    pub fn record_batch_schema_hint(&self) -> Option<SchemaRef> {
803        self.context.builder.record_batch_schema_hint()
804    }
805
806    /// Returns whether the iterator is a record batch iterator.
807    pub fn is_record_batch(&self) -> bool {
808        self.context.builder.is_record_batch()
809    }
810
811    pub fn num_rows(&self) -> usize {
812        self.stats.num_rows
813    }
814
815    /// Returns the encoded range if available.
816    pub fn encoded(&self) -> Option<EncodedRange> {
817        self.context.builder.encoded_range()
818    }
819}
820
821#[cfg(test)]
822mod tests {
823    use std::collections::HashMap;
824    use std::sync::Arc;
825
826    use common_base::readable_size::ReadableSize;
827    use common_error::ext::WhateverResult;
828    use datatypes::prelude::ConcreteDataType;
829    use datatypes::types::json_type::{JsonNativeType, JsonObjectType};
830    use store_api::metadata::RegionMetadataBuilder;
831    use store_api::storage::RegionId;
832
833    use super::*;
834    use crate::flush::{WriteBufferManager, WriteBufferManagerImpl};
835    use crate::memtable::bulk::BulkMemtableConfig;
836    use crate::test_util::sst_util::sst_region_metadata;
837
838    #[test]
839    fn test_alloc_tracker_without_manager() {
840        let tracker = AllocTracker::new(None);
841        assert_eq!(0, tracker.bytes_allocated());
842        tracker.on_allocation(100);
843        assert_eq!(100, tracker.bytes_allocated());
844        tracker.on_allocation(200);
845        assert_eq!(300, tracker.bytes_allocated());
846
847        tracker.done_allocating();
848        assert_eq!(300, tracker.bytes_allocated());
849    }
850
851    #[test]
852    fn test_alloc_tracker_with_manager() {
853        let manager = Arc::new(WriteBufferManagerImpl::new(1000));
854        {
855            let tracker = AllocTracker::new(Some(manager.clone() as WriteBufferManagerRef));
856
857            tracker.on_allocation(100);
858            assert_eq!(100, tracker.bytes_allocated());
859            assert_eq!(100, manager.memory_usage());
860            assert_eq!(100, manager.mutable_usage());
861
862            for _ in 0..2 {
863                // Done allocating won't free the same memory multiple times.
864                tracker.done_allocating();
865                assert_eq!(100, manager.memory_usage());
866                assert_eq!(0, manager.mutable_usage());
867            }
868        }
869
870        assert_eq!(0, manager.memory_usage());
871        assert_eq!(0, manager.mutable_usage());
872    }
873
874    #[test]
875    fn test_alloc_tracker_without_done_allocating() {
876        let manager = Arc::new(WriteBufferManagerImpl::new(1000));
877        {
878            let tracker = AllocTracker::new(Some(manager.clone() as WriteBufferManagerRef));
879
880            tracker.on_allocation(100);
881            assert_eq!(100, tracker.bytes_allocated());
882            assert_eq!(100, manager.memory_usage());
883            assert_eq!(100, manager.mutable_usage());
884        }
885
886        assert_eq!(0, manager.memory_usage());
887        assert_eq!(0, manager.mutable_usage());
888    }
889
890    #[test]
891    fn test_forced_bulk_memtable_preserves_bulk_config() {
892        let provider = MemtableBuilderProvider::new(None, Arc::new(MitoConfig::default()));
893        let config = BulkMemtableConfig {
894            merge_threshold: 7,
895            encode_row_threshold: 11,
896            encode_bytes_threshold: 13,
897            max_merge_groups: 17,
898        };
899        let options = RegionOptions {
900            memtable: Some(MemtableOptions::Bulk(config.clone())),
901            primary_key_encoding: Some(PrimaryKeyEncoding::Sparse),
902            ..Default::default()
903        };
904
905        let builder =
906            provider.bulk_memtable_builder(options.need_dedup(), options.merge_mode(), &options);
907
908        assert_eq!(&config, builder.config());
909    }
910
911    #[test]
912    fn test_provider_uses_adaptive_config_for_implicit_bulk_builder() {
913        let config = MitoConfig {
914            global_write_buffer_size: ReadableSize::gb(8),
915            ..Default::default()
916        };
917        let provider = MemtableBuilderProvider::new(None, Arc::new(config));
918        let options = RegionOptions::default();
919
920        let builder =
921            provider.bulk_memtable_builder(options.need_dedup(), options.merge_mode(), &options);
922
923        assert_eq!(256 * 1024 * 1024, builder.config().encode_bytes_threshold);
924    }
925
926    #[test]
927    fn test_provider_parses_bulk_memtable_with_adaptive_config() {
928        let config = MitoConfig {
929            global_write_buffer_size: ReadableSize::gb(8),
930            ..Default::default()
931        };
932        let provider = MemtableBuilderProvider::new(None, Arc::new(config));
933        let options = HashMap::from([("memtable.type".to_string(), "bulk".to_string())]);
934
935        let options = provider
936            .parse_options(RegionId::new(0, 0), &options)
937            .unwrap();
938
939        let Some(MemtableOptions::Bulk(config)) = options.memtable else {
940            panic!("expected bulk memtable options");
941        };
942        assert_eq!(256 * 1024 * 1024, config.encode_bytes_threshold);
943    }
944
945    #[test]
946    fn test_provider_preserves_explicit_bulk_encode_bytes_threshold() {
947        let config = MitoConfig {
948            global_write_buffer_size: ReadableSize::gb(8),
949            ..Default::default()
950        };
951        let provider = MemtableBuilderProvider::new(None, Arc::new(config));
952        let options = HashMap::from([
953            ("memtable.type".to_string(), "bulk".to_string()),
954            (
955                "memtable.bulk.encode_bytes_threshold".to_string(),
956                "13".to_string(),
957            ),
958        ]);
959
960        let options = provider
961            .parse_options(RegionId::new(0, 0), &options)
962            .unwrap();
963
964        let Some(MemtableOptions::Bulk(config)) = options.memtable else {
965            panic!("expected bulk memtable options");
966        };
967        assert_eq!(13, config.encode_bytes_threshold);
968    }
969
970    #[test]
971    fn test_json2_requires_bulk_memtable() -> WhateverResult<()> {
972        let mut metadata = sst_region_metadata();
973        metadata.column_metadatas[2].column_schema.data_type =
974            ConcreteDataType::json2(JsonNativeType::Object(JsonObjectType::new()));
975        let metadata = RegionMetadataBuilder::from_existing(metadata).build()?;
976        let mut options = RegionOptions {
977            sst_format: Some(FormatType::PrimaryKey),
978            memtable: Some(MemtableOptions::TimeSeries),
979            ..Default::default()
980        };
981
982        let err = ensure_json2_not_use_time_series_memtable(&metadata, &options).unwrap_err();
983        assert!(
984            err.to_string()
985                .contains("JSON2 columns only support BulkMemtable")
986        );
987
988        options.memtable = Some(MemtableOptions::Bulk(BulkMemtableConfig::default()));
989        ensure_json2_not_use_time_series_memtable(&metadata, &options)?;
990        Ok(())
991    }
992}