1use 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
67pub type MemtableId = u32;
71
72#[derive(Clone)]
74pub struct RangesOptions {
75 pub for_flush: bool,
77 pub pre_filter_mode: PreFilterMode,
79 pub predicate: PredicateGroup,
81 pub sequence: Option<SequenceRange>,
83 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 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 #[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 #[must_use]
120 pub fn with_predicate(mut self, predicate: PredicateGroup) -> Self {
121 self.predicate = predicate;
122 self
123 }
124
125 #[must_use]
127 pub fn with_sequence(mut self, sequence: Option<SequenceRange>) -> Self {
128 self.sequence = sequence;
129 self
130 }
131
132 #[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 pub estimated_bytes: usize,
144 pub time_range: Option<(Timestamp, Timestamp)>,
147 pub num_rows: usize,
149 pub num_ranges: usize,
151 pub max_sequence: SequenceNumber,
153 pub min_sequence: SequenceNumber,
157 pub series_count: usize,
159}
160
161impl MemtableStats {
162 #[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 pub fn bytes_allocated(&self) -> usize {
177 self.estimated_bytes
178 }
179
180 pub fn time_range(&self) -> Option<(Timestamp, Timestamp)> {
182 self.time_range
183 }
184
185 pub fn num_rows(&self) -> usize {
187 self.num_rows
188 }
189
190 pub fn num_ranges(&self) -> usize {
192 self.num_ranges
193 }
194
195 pub fn max_sequence(&self) -> SequenceNumber {
197 self.max_sequence
198 }
199
200 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#[derive(Default)]
212pub struct MemtableRanges {
213 pub ranges: BTreeMap<usize, MemtableRange>,
215}
216
217impl MemtableRanges {
218 pub fn num_rows(&self) -> usize {
220 self.ranges.values().map(|r| r.stats().num_rows()).sum()
221 }
222
223 pub fn series_count(&self) -> usize {
225 self.ranges.values().map(|r| r.stats().series_count()).sum()
226 }
227
228 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
258pub trait Memtable: Send + Sync + fmt::Debug {
260 fn id(&self) -> MemtableId;
262
263 fn write(&self, kvs: &KeyValues) -> Result<()>;
265
266 fn write_one(&self, key_value: KeyValue) -> Result<()>;
268
269 fn write_bulk(&self, part: crate::memtable::bulk::part::BulkPart) -> Result<()>;
271
272 fn ranges(
276 &self,
277 projection: Option<&[ColumnId]>,
278 options: RangesOptions,
279 ) -> Result<MemtableRanges>;
280
281 fn is_empty(&self) -> bool;
283
284 fn freeze(&self) -> Result<()>;
286
287 fn stats(&self) -> MemtableStats;
289
290 fn min_sequence(&self) -> SequenceNumber;
294
295 fn fork(&self, id: MemtableId, metadata: &RegionMetadataRef) -> MemtableRef;
299
300 fn compact(&self, for_flush: bool) -> Result<()> {
304 let _ = for_flush;
305 Ok(())
306 }
307}
308
309pub type MemtableRef = Arc<dyn Memtable>;
310
311pub trait MemtableBuilder: Send + Sync + fmt::Debug {
313 fn build(&self, id: MemtableId, metadata: &RegionMetadataRef) -> MemtableRef;
315
316 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#[derive(Default)]
327pub struct AllocTracker {
328 write_buffer_manager: Option<WriteBufferManagerRef>,
329 bytes_allocated: AtomicUsize,
331 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 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 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 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 pub(crate) fn bytes_allocated(&self) -> usize {
380 self.bytes_allocated.load(Ordering::Relaxed)
381 }
382
383 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 if let Some(write_buffer_manager) = &self.write_buffer_manager {
400 write_buffer_manager.free_mem(bytes_allocated);
401 }
402 }
403}
404
405#[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
414pub(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 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 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, 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#[derive(Clone, Default)]
553pub struct MemScanMetrics(Arc<Mutex<MemScanMetricsData>>);
554
555impl MemScanMetrics {
556 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 pub(crate) fn data(&self) -> MemScanMetricsData {
569 self.0.lock().unwrap().clone()
570 }
571}
572
573#[derive(Clone, Default)]
574pub(crate) struct MemScanMetricsData {
575 pub(crate) total_series: usize,
577 pub(crate) num_rows: usize,
579 pub(crate) num_batches: usize,
581 pub(crate) scan_cost: Duration,
583 pub(crate) prefilter_cost: Duration,
585 pub(crate) prefilter_rows_filtered: usize,
587}
588
589pub struct EncodedRange {
591 pub data: Bytes,
593 pub sst_info: SstInfo,
595}
596
597pub trait IterBuilder: Send + Sync {
600 fn build(&self, metrics: Option<MemScanMetrics>) -> Result<BoxedBatchIterator>;
602
603 fn is_record_batch(&self) -> bool {
605 false
606 }
607
608 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 fn record_batch_schema_hint(&self) -> Option<SchemaRef> {
626 None
627 }
628
629 fn encoded_range(&self) -> Option<EncodedRange> {
631 None
632 }
633}
634
635pub type BoxedIterBuilder = Box<dyn IterBuilder>;
636
637pub 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
656pub struct BatchToRecordBatchContext {
658 metadata: RegionMetadataRef,
659 codec: Arc<dyn PrimaryKeyCodec>,
660 read_column_ids: Vec<ColumnId>,
661}
662
663impl BatchToRecordBatchContext {
664 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
688pub struct MemtableRangeContext {
690 id: MemtableId,
692 builder: BoxedIterBuilder,
694 predicate: PredicateGroup,
696 batch_to_record_batch: Option<Arc<BatchToRecordBatchContext>>,
698}
699
700pub type MemtableRangeContextRef = Arc<MemtableRangeContext>;
701
702impl MemtableRangeContext {
703 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 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#[derive(Clone)]
726pub struct MemtableRange {
727 context: MemtableRangeContextRef,
729 stats: MemtableStats,
731}
732
733impl MemtableRange {
734 pub fn new(context: MemtableRangeContextRef, stats: MemtableStats) -> Self {
736 Self { context, stats }
737 }
738
739 pub fn stats(&self) -> &MemtableStats {
741 &self.stats
742 }
743
744 pub fn id(&self) -> MemtableId {
746 self.context.id
747 }
748
749 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 pub fn build_iter(&self) -> Result<BoxedBatchIterator> {
768 self.context.builder.build(None)
769 }
770
771 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 pub fn record_batch_schema_hint(&self) -> Option<SchemaRef> {
803 self.context.builder.record_batch_schema_hint()
804 }
805
806 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 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 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}