1pub(crate) mod chunk_reader;
18pub mod context;
19pub(crate) mod json_align;
20pub mod part;
21pub mod part_reader;
22mod row_group_reader;
23
24use std::collections::{BTreeMap, HashSet};
25use std::sync::atomic::{AtomicI64, AtomicU64, AtomicUsize, Ordering};
26use std::sync::{Arc, LazyLock, Mutex, RwLock};
27use std::time::Instant;
28
29fn env_usize(name: &str, default: usize) -> usize {
31 std::env::var(name)
32 .ok()
33 .and_then(|v| v.parse().ok())
34 .unwrap_or(default)
35}
36
37use common_time::Timestamp;
38use datatypes::arrow::datatypes::SchemaRef;
39use mito_codec::key_values::KeyValue;
40use rayon::prelude::*;
41use serde::{Deserialize, Serialize};
42use serde_with::{DisplayFromStr, serde_as};
43use store_api::metadata::RegionMetadataRef;
44use store_api::storage::{ColumnId, FileId, RegionId, SequenceRange};
45use tokio::sync::Semaphore;
46
47use crate::error::{Result, UnsupportedOperationSnafu};
48use crate::flush::WriteBufferManagerRef;
49use crate::memtable::bulk::context::BulkIterContext;
50use crate::memtable::bulk::json_align::Json2Aligner;
51use crate::memtable::bulk::part::{
52 BulkPart, BulkPartEncodeMetrics, BulkPartEncoder, MultiBulkPart, UnorderedPart,
53 should_prune_bulk_part,
54};
55use crate::memtable::bulk::part_reader::BulkPartBatchIter;
56use crate::memtable::stats::WriteMetrics;
57use crate::memtable::{
58 AllocTracker, BoxedBatchIterator, BoxedRecordBatchIterator, EncodedBulkPart, EncodedRange,
59 IterBuilder, KeyValues, MemScanMetrics, Memtable, MemtableBuilder, MemtableId, MemtableRange,
60 MemtableRangeContext, MemtableRanges, MemtableRef, MemtableStats, RangesOptions,
61};
62use crate::read::flat_dedup::{FlatDedupIterator, FlatLastNonNull, FlatLastRow};
63use crate::read::flat_merge::FlatMergeIterator;
64use crate::region::options::MergeMode;
65use crate::sst::parquet::DEFAULT_ROW_GROUP_SIZE;
66use crate::sst::parquet::flat_format::field_column_start;
67
68const DEFAULT_MERGE_THRESHOLD: usize = 16;
70
71static MERGE_THRESHOLD: LazyLock<usize> =
73 LazyLock::new(|| env_usize("GREPTIME_BULK_MERGE_THRESHOLD", DEFAULT_MERGE_THRESHOLD));
74
75const DEFAULT_MAX_MERGE_GROUPS: usize = 32;
77
78static MAX_MERGE_GROUPS: LazyLock<usize> =
80 LazyLock::new(|| env_usize("GREPTIME_BULK_MAX_MERGE_GROUPS", DEFAULT_MAX_MERGE_GROUPS));
81
82pub(crate) static ENCODE_ROW_THRESHOLD: LazyLock<usize> = LazyLock::new(|| {
85 env_usize(
86 "GREPTIME_BULK_ENCODE_ROW_THRESHOLD",
87 10 * DEFAULT_ROW_GROUP_SIZE,
88 )
89});
90
91const DEFAULT_ENCODE_BYTES_THRESHOLD: usize = 64 * 1024 * 1024;
93
94static ENCODE_BYTES_THRESHOLD: LazyLock<usize> = LazyLock::new(|| {
97 env_usize(
98 "GREPTIME_BULK_ENCODE_BYTES_THRESHOLD",
99 DEFAULT_ENCODE_BYTES_THRESHOLD,
100 )
101});
102
103#[serde_as]
105#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
106#[serde(default)]
107pub struct BulkMemtableConfig {
108 #[serde_as(as = "DisplayFromStr")]
110 pub merge_threshold: usize,
111 #[serde_as(as = "DisplayFromStr")]
113 pub encode_row_threshold: usize,
114 #[serde_as(as = "DisplayFromStr")]
116 pub encode_bytes_threshold: usize,
117 #[serde_as(as = "DisplayFromStr")]
119 pub max_merge_groups: usize,
120}
121
122impl Default for BulkMemtableConfig {
123 fn default() -> Self {
124 Self {
125 merge_threshold: *MERGE_THRESHOLD,
126 encode_row_threshold: *ENCODE_ROW_THRESHOLD,
127 encode_bytes_threshold: *ENCODE_BYTES_THRESHOLD,
128 max_merge_groups: *MAX_MERGE_GROUPS,
129 }
130 .sanitize()
131 }
132}
133
134impl BulkMemtableConfig {
135 fn sanitize(mut self) -> Self {
136 if self.merge_threshold == 0 {
137 self.merge_threshold = DEFAULT_MERGE_THRESHOLD;
138 }
139 self
140 }
141}
142
143enum MergedPart {
145 Multi(MultiBulkPart),
147 Encoded(EncodedBulkPart),
149}
150
151struct CollectedParts {
153 groups: Vec<Vec<PartToMerge>>,
155}
156
157#[derive(Default)]
159struct BulkParts {
160 unordered_part: UnorderedPart,
162 parts: Vec<BulkPartWrapper>,
164}
165
166impl BulkParts {
167 fn num_parts(&self) -> usize {
169 let unordered_count = if self.unordered_part.is_empty() { 0 } else { 1 };
170 self.parts.len() + unordered_count
171 }
172
173 fn is_empty(&self) -> bool {
175 self.unordered_part.is_empty() && self.parts.is_empty()
176 }
177
178 fn should_merge_parts(&self, merge_threshold: usize, encode_bytes_threshold: usize) -> bool {
180 let mut bulk_count = 0;
181 let mut encoded_count = 0;
182
183 for wrapper in &self.parts {
184 if !Self::is_merge_candidate(wrapper, encode_bytes_threshold) {
185 continue;
186 }
187
188 if wrapper.part.is_encoded() {
189 encoded_count += 1;
190 } else {
191 bulk_count += 1;
192 }
193
194 if bulk_count >= merge_threshold || encoded_count >= merge_threshold {
195 return true;
196 }
197 }
198
199 false
200 }
201
202 fn is_merge_candidate(wrapper: &BulkPartWrapper, encode_bytes_threshold: usize) -> bool {
204 !wrapper.merging
205 && Self::is_merge_candidate_by_size(
206 wrapper.part.is_encoded(),
207 wrapper.part.estimated_size(),
208 encode_bytes_threshold,
209 )
210 }
211
212 fn is_merge_candidate_by_size(
213 is_encoded: bool,
214 estimated_size: usize,
215 encode_bytes_threshold: usize,
216 ) -> bool {
217 !is_encoded || estimated_size <= encode_bytes_threshold
218 }
219
220 fn should_compact_unordered_part(&self, encode_bytes_threshold: usize) -> bool {
222 self.unordered_part.should_compact()
223 || self.unordered_part.estimated_bytes() > encode_bytes_threshold
224 }
225
226 fn collect_parts_to_merge(
230 &mut self,
231 merge_threshold: usize,
232 max_merge_groups: usize,
233 encode_bytes_threshold: usize,
234 ) -> CollectedParts {
235 let mut bulk_indices = Vec::new();
236 let mut encoded_indices = Vec::new();
237
238 for (idx, wrapper) in self.parts.iter().enumerate() {
239 if !Self::is_merge_candidate(wrapper, encode_bytes_threshold) {
240 continue;
241 }
242 let num_rows = wrapper.part.num_rows();
243 if wrapper.part.is_encoded() {
244 encoded_indices.push((idx, num_rows));
245 } else {
246 bulk_indices.push((idx, num_rows));
247 }
248 }
249
250 let mut groups = Vec::new();
251
252 if bulk_indices.len() >= merge_threshold {
254 groups.extend(self.collect_and_group_parts(
255 bulk_indices,
256 merge_threshold,
257 max_merge_groups,
258 ));
259 }
260
261 if encoded_indices.len() >= merge_threshold {
263 groups.extend(self.collect_and_group_parts(
264 encoded_indices,
265 merge_threshold,
266 max_merge_groups,
267 ));
268 }
269
270 CollectedParts { groups }
271 }
272
273 fn collect_and_group_parts(
275 &mut self,
276 mut indices: Vec<(usize, usize)>,
277 merge_threshold: usize,
278 max_merge_groups: usize,
279 ) -> Vec<Vec<PartToMerge>> {
280 if indices.is_empty() {
281 return Vec::new();
282 }
283
284 indices.sort_unstable_by_key(|(_, num_rows)| *num_rows);
285 indices
286 .chunks(merge_threshold)
287 .take(max_merge_groups)
288 .map(|chunk| {
289 chunk
290 .iter()
291 .map(|(idx, _)| {
292 let wrapper = &mut self.parts[*idx];
293 wrapper.merging = true;
294 wrapper.part.clone()
295 })
296 .collect()
297 })
298 .collect()
299 }
300
301 fn install_merged_parts<I>(
304 &mut self,
305 merged_parts: I,
306 merged_file_ids: &HashSet<FileId>,
307 ) -> usize
308 where
309 I: IntoIterator<Item = MergedPart>,
310 {
311 let mut total_output_rows = 0;
312
313 for merged_part in merged_parts {
314 match merged_part {
315 MergedPart::Encoded(encoded_part) => {
316 total_output_rows += encoded_part.metadata().num_rows;
317 self.parts.push(BulkPartWrapper {
318 part: PartToMerge::Encoded {
319 part: encoded_part,
320 file_id: FileId::random(),
321 },
322 merging: false,
323 });
324 }
325 MergedPart::Multi(multi_part) => {
326 total_output_rows += multi_part.num_rows();
327 self.parts.push(BulkPartWrapper {
328 part: PartToMerge::Multi {
329 part: multi_part,
330 file_id: FileId::random(),
331 },
332 merging: false,
333 });
334 }
335 }
336 }
337
338 self.parts
339 .retain(|wrapper| !merged_file_ids.contains(&wrapper.file_id()));
340
341 total_output_rows
342 }
343
344 fn reset_merging_flags(&mut self, file_ids: &HashSet<FileId>) {
347 for wrapper in &mut self.parts {
348 if file_ids.contains(&wrapper.file_id()) {
349 wrapper.merging = false;
350 }
351 }
352 }
353}
354
355struct MergingFlagsGuard<'a> {
358 bulk_parts: &'a RwLock<BulkParts>,
359 file_ids: &'a HashSet<FileId>,
360 success: bool,
361}
362
363impl<'a> MergingFlagsGuard<'a> {
364 fn new(bulk_parts: &'a RwLock<BulkParts>, file_ids: &'a HashSet<FileId>) -> Self {
366 Self {
367 bulk_parts,
368 file_ids,
369 success: false,
370 }
371 }
372
373 fn mark_success(&mut self) {
376 self.success = true;
377 }
378}
379
380impl<'a> Drop for MergingFlagsGuard<'a> {
381 fn drop(&mut self) {
382 if !self.success
383 && let Ok(mut parts) = self.bulk_parts.write()
384 {
385 parts.reset_merging_flags(self.file_ids);
386 }
387 }
388}
389
390pub struct BulkMemtable {
392 id: MemtableId,
393 config: BulkMemtableConfig,
395 parts: Arc<RwLock<BulkParts>>,
396 metadata: RegionMetadataRef,
397 alloc_tracker: AllocTracker,
398 max_timestamp: AtomicI64,
399 min_timestamp: AtomicI64,
400 max_sequence: AtomicU64,
401 num_rows: AtomicUsize,
402 compactor: Arc<Mutex<MemtableCompactor>>,
404 compact_dispatcher: Option<Arc<CompactDispatcher>>,
406 append_mode: bool,
408 merge_mode: MergeMode,
410 row_group_size: usize,
412}
413
414impl std::fmt::Debug for BulkMemtable {
415 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
416 f.debug_struct("BulkMemtable")
417 .field("id", &self.id)
418 .field("num_rows", &self.num_rows.load(Ordering::Relaxed))
419 .field("min_timestamp", &self.min_timestamp.load(Ordering::Relaxed))
420 .field("max_timestamp", &self.max_timestamp.load(Ordering::Relaxed))
421 .field("max_sequence", &self.max_sequence.load(Ordering::Relaxed))
422 .finish()
423 }
424}
425
426impl Memtable for BulkMemtable {
427 fn id(&self) -> MemtableId {
428 self.id
429 }
430
431 fn write(&self, _kvs: &KeyValues) -> Result<()> {
432 UnsupportedOperationSnafu {
433 err_msg: "write() is not supported for bulk memtable",
434 }
435 .fail()
436 }
437
438 fn write_one(&self, _key_value: KeyValue) -> Result<()> {
439 UnsupportedOperationSnafu {
440 err_msg: "write_one() is not supported for bulk memtable",
441 }
442 .fail()
443 }
444
445 fn write_bulk(&self, fragment: BulkPart) -> Result<()> {
446 let local_metrics = WriteMetrics {
447 key_bytes: 0,
448 value_bytes: fragment.estimated_size(),
449 min_ts: fragment.min_timestamp,
450 max_ts: fragment.max_timestamp,
451 num_rows: fragment.num_rows(),
452 max_sequence: fragment.sequence,
453 };
454
455 {
456 let mut bulk_parts = self.parts.write().unwrap();
457
458 let fragment_size = fragment.estimated_size();
459 if bulk_parts.unordered_part.should_accept(fragment.num_rows())
461 && fragment_size <= self.config.encode_bytes_threshold
462 {
463 bulk_parts.unordered_part.push(fragment);
464
465 if bulk_parts.should_compact_unordered_part(self.config.encode_bytes_threshold)
467 && let Some(bulk_part) = bulk_parts.unordered_part.to_bulk_part()?
468 {
469 bulk_parts.parts.push(BulkPartWrapper {
470 part: PartToMerge::Bulk {
471 part: bulk_part,
472 file_id: FileId::random(),
473 },
474 merging: false,
475 });
476 bulk_parts.unordered_part.clear();
477 }
478 } else {
479 bulk_parts.parts.push(BulkPartWrapper {
480 part: PartToMerge::Bulk {
481 part: fragment,
482 file_id: FileId::random(),
483 },
484 merging: false,
485 });
486 }
487
488 self.update_stats(local_metrics);
493 }
494
495 if self.should_compact() {
496 self.schedule_compact();
497 }
498
499 Ok(())
500 }
501
502 fn ranges(
503 &self,
504 projection: Option<&[ColumnId]>,
505 options: RangesOptions,
506 ) -> Result<MemtableRanges> {
507 let predicate = options.predicate;
508 let sequence = options.sequence;
509 let mut ranges = BTreeMap::new();
510 let mut range_id = 0;
511
512 let context = Arc::new(BulkIterContext::new_with_pre_filter_mode(
514 self.metadata.clone(),
515 projection,
516 predicate.predicate().cloned(),
517 options.for_flush,
518 options.pre_filter_mode,
519 options.batch_size,
520 )?);
521
522 {
524 let bulk_parts = self.parts.read().unwrap();
525
526 if !bulk_parts.unordered_part.is_empty()
528 && let Some(unordered_bulk_part) = bulk_parts.unordered_part.to_bulk_part()?
529 {
530 let part_stats = unordered_bulk_part.to_memtable_stats(&self.metadata);
531 let range = MemtableRange::new(
532 Arc::new(MemtableRangeContext::new(
533 self.id,
534 Box::new(BulkRangeIterBuilder {
535 part: unordered_bulk_part,
536 context: context.clone(),
537 sequence,
538 }),
539 predicate.clone(),
540 )),
541 part_stats,
542 );
543 ranges.insert(range_id, range);
544 range_id += 1;
545 }
546
547 for part_wrapper in bulk_parts.parts.iter() {
549 if part_wrapper.part.num_rows() == 0 {
551 continue;
552 }
553
554 let part_stats = part_wrapper.part.to_memtable_stats(&self.metadata);
555 let iter_builder: Box<dyn IterBuilder> = match &part_wrapper.part {
556 PartToMerge::Bulk { part, .. } => Box::new(BulkRangeIterBuilder {
557 part: part.clone(),
558 context: context.clone(),
559 sequence,
560 }),
561 PartToMerge::Multi { part, .. } => Box::new(MultiBulkRangeIterBuilder {
562 part: part.clone(),
563 context: context.clone(),
564 sequence,
565 }),
566 PartToMerge::Encoded { part, file_id } => {
567 Box::new(EncodedBulkRangeIterBuilder {
568 file_id: *file_id,
569 part: part.clone(),
570 context: context.clone(),
571 sequence,
572 })
573 }
574 };
575
576 let range = MemtableRange::new(
577 Arc::new(MemtableRangeContext::new(
578 self.id,
579 iter_builder,
580 predicate.clone(),
581 )),
582 part_stats,
583 );
584 ranges.insert(range_id, range);
585 range_id += 1;
586 }
587 }
588
589 Ok(MemtableRanges { ranges })
590 }
591
592 fn is_empty(&self) -> bool {
593 let bulk_parts = self.parts.read().unwrap();
594 bulk_parts.is_empty()
595 }
596
597 fn freeze(&self) -> Result<()> {
598 self.alloc_tracker.done_allocating();
599 Ok(())
600 }
601
602 fn stats(&self) -> MemtableStats {
603 let estimated_bytes = self.alloc_tracker.bytes_allocated();
604
605 if estimated_bytes == 0 || self.num_rows.load(Ordering::Relaxed) == 0 {
606 return MemtableStats {
607 estimated_bytes,
608 time_range: None,
609 num_rows: 0,
610 num_ranges: 0,
611 max_sequence: 0,
612 series_count: 0,
613 };
614 }
615
616 let ts_type = self
617 .metadata
618 .time_index_column()
619 .column_schema
620 .data_type
621 .clone()
622 .as_timestamp()
623 .expect("Timestamp column must have timestamp type");
624 let max_timestamp = ts_type.create_timestamp(self.max_timestamp.load(Ordering::Relaxed));
625 let min_timestamp = ts_type.create_timestamp(self.min_timestamp.load(Ordering::Relaxed));
626
627 let num_ranges = self.parts.read().unwrap().num_parts();
628
629 MemtableStats {
630 estimated_bytes,
631 time_range: Some((min_timestamp, max_timestamp)),
632 num_rows: self.num_rows.load(Ordering::Relaxed),
633 num_ranges,
634 max_sequence: self.max_sequence.load(Ordering::Relaxed),
635 series_count: self.estimated_series_count(),
636 }
637 }
638
639 fn fork(&self, id: MemtableId, metadata: &RegionMetadataRef) -> MemtableRef {
640 Arc::new(Self {
641 id,
642 config: self.config.clone(),
643 parts: Arc::new(RwLock::new(BulkParts::default())),
644 metadata: metadata.clone(),
645 alloc_tracker: AllocTracker::new(self.alloc_tracker.write_buffer_manager()),
646 max_timestamp: AtomicI64::new(i64::MIN),
647 min_timestamp: AtomicI64::new(i64::MAX),
648 max_sequence: AtomicU64::new(0),
649 num_rows: AtomicUsize::new(0),
650 compactor: Arc::new(Mutex::new(MemtableCompactor::new(
651 metadata.region_id,
652 id,
653 self.config.clone(),
654 self.row_group_size,
655 ))),
656 compact_dispatcher: self.compact_dispatcher.clone(),
657 append_mode: self.append_mode,
658 merge_mode: self.merge_mode,
659 row_group_size: self.row_group_size,
660 })
661 }
662
663 fn compact(&self, for_flush: bool) -> Result<()> {
664 let mut compactor = self.compactor.lock().unwrap();
665
666 if for_flush {
667 return Ok(());
668 }
669
670 let should_merge = self.parts.read().unwrap().should_merge_parts(
672 self.config.merge_threshold,
673 self.config.encode_bytes_threshold,
674 );
675 if should_merge {
676 compactor.merge_parts(
677 &self.parts,
678 &self.metadata,
679 !self.append_mode,
680 self.merge_mode,
681 )?;
682 }
683
684 Ok(())
685 }
686}
687
688impl BulkMemtable {
689 pub fn new(
691 id: MemtableId,
692 config: BulkMemtableConfig,
693 metadata: RegionMetadataRef,
694 write_buffer_manager: Option<WriteBufferManagerRef>,
695 compact_dispatcher: Option<Arc<CompactDispatcher>>,
696 append_mode: bool,
697 merge_mode: MergeMode,
698 ) -> Self {
699 Self::new_with_row_group_size(
700 id,
701 config,
702 metadata,
703 write_buffer_manager,
704 compact_dispatcher,
705 append_mode,
706 merge_mode,
707 DEFAULT_ROW_GROUP_SIZE,
708 )
709 }
710
711 #[allow(clippy::too_many_arguments)]
713 pub fn new_with_row_group_size(
714 id: MemtableId,
715 config: BulkMemtableConfig,
716 metadata: RegionMetadataRef,
717 write_buffer_manager: Option<WriteBufferManagerRef>,
718 compact_dispatcher: Option<Arc<CompactDispatcher>>,
719 append_mode: bool,
720 merge_mode: MergeMode,
721 row_group_size: usize,
722 ) -> Self {
723 let config = config.sanitize();
724 let region_id = metadata.region_id;
725 Self {
726 id,
727 config: config.clone(),
728 parts: Arc::new(RwLock::new(BulkParts::default())),
729 metadata,
730 alloc_tracker: AllocTracker::new(write_buffer_manager),
731 max_timestamp: AtomicI64::new(i64::MIN),
732 min_timestamp: AtomicI64::new(i64::MAX),
733 max_sequence: AtomicU64::new(0),
734 num_rows: AtomicUsize::new(0),
735 compactor: Arc::new(Mutex::new(MemtableCompactor::new(
736 region_id,
737 id,
738 config,
739 row_group_size,
740 ))),
741 compact_dispatcher,
742 append_mode,
743 merge_mode,
744 row_group_size,
745 }
746 }
747
748 #[cfg(test)]
750 pub fn set_unordered_part_threshold(&self, threshold: usize) {
751 self.parts
752 .write()
753 .unwrap()
754 .unordered_part
755 .set_threshold(threshold);
756 }
757
758 #[cfg(test)]
760 pub fn set_unordered_part_compact_threshold(&self, compact_threshold: usize) {
761 self.parts
762 .write()
763 .unwrap()
764 .unordered_part
765 .set_compact_threshold(compact_threshold);
766 }
767
768 fn update_stats(&self, stats: WriteMetrics) {
772 self.alloc_tracker
773 .on_allocation(stats.key_bytes + stats.value_bytes);
774
775 self.max_timestamp
776 .fetch_max(stats.max_ts, Ordering::Relaxed);
777 self.min_timestamp
778 .fetch_min(stats.min_ts, Ordering::Relaxed);
779 self.max_sequence
780 .fetch_max(stats.max_sequence, Ordering::Relaxed);
781 self.num_rows.fetch_add(stats.num_rows, Ordering::Relaxed);
782 }
783
784 fn estimated_series_count(&self) -> usize {
786 let bulk_parts = self.parts.read().unwrap();
787 bulk_parts
788 .parts
789 .iter()
790 .map(|part_wrapper| part_wrapper.part.series_count())
791 .sum()
792 }
793
794 fn should_compact(&self) -> bool {
796 let parts = self.parts.read().unwrap();
797 parts.should_merge_parts(
798 self.config.merge_threshold,
799 self.config.encode_bytes_threshold,
800 )
801 }
802
803 fn schedule_compact(&self) {
805 if let Some(dispatcher) = &self.compact_dispatcher {
806 let task = MemCompactTask {
807 metadata: self.metadata.clone(),
808 parts: self.parts.clone(),
809 config: self.config.clone(),
810 compactor: self.compactor.clone(),
811 append_mode: self.append_mode,
812 merge_mode: self.merge_mode,
813 };
814
815 dispatcher.dispatch_compact(task);
816 } else {
817 if let Err(e) = self.compact(false) {
819 common_telemetry::error!(e; "Failed to compact table");
820 }
821 }
822 }
823}
824
825pub struct BulkRangeIterBuilder {
827 pub part: BulkPart,
828 pub context: Arc<BulkIterContext>,
829 pub sequence: Option<SequenceRange>,
830}
831
832struct MultiBulkRangeIterBuilder {
834 part: MultiBulkPart,
835 context: Arc<BulkIterContext>,
836 sequence: Option<SequenceRange>,
837}
838
839impl IterBuilder for BulkRangeIterBuilder {
840 fn build(&self, _metrics: Option<MemScanMetrics>) -> Result<BoxedBatchIterator> {
841 UnsupportedOperationSnafu {
842 err_msg: "BatchIterator is not supported for bulk memtable",
843 }
844 .fail()
845 }
846
847 fn is_record_batch(&self) -> bool {
848 true
849 }
850
851 fn build_record_batch(
852 &self,
853 _time_range: Option<(Timestamp, Timestamp)>,
854 metrics: Option<MemScanMetrics>,
855 ) -> Result<BoxedRecordBatchIterator> {
856 let metadata = self.context.read_format().metadata();
857 if should_prune_bulk_part(&self.part.batch, &self.context, metadata) {
858 return Ok(Box::new(std::iter::empty()));
859 }
860
861 let series_count = self.part.estimated_series_count();
862 let iter = BulkPartBatchIter::from_single(
863 self.part.batch.clone(),
864 self.context.clone(),
865 self.sequence,
866 series_count,
867 metrics,
868 );
869
870 Ok(Box::new(iter))
871 }
872
873 fn record_batch_schema_hint(&self) -> Option<SchemaRef> {
874 Some(self.part.schema())
875 }
876
877 fn encoded_range(&self) -> Option<EncodedRange> {
878 None
879 }
880}
881
882impl IterBuilder for MultiBulkRangeIterBuilder {
883 fn build(&self, _metrics: Option<MemScanMetrics>) -> Result<BoxedBatchIterator> {
884 UnsupportedOperationSnafu {
885 err_msg: "BatchIterator is not supported for multi bulk memtable",
886 }
887 .fail()
888 }
889
890 fn is_record_batch(&self) -> bool {
891 true
892 }
893
894 fn build_record_batch(
895 &self,
896 _time_range: Option<(Timestamp, Timestamp)>,
897 metrics: Option<MemScanMetrics>,
898 ) -> Result<BoxedRecordBatchIterator> {
899 match self
900 .part
901 .read(self.context.clone(), self.sequence, metrics)?
902 {
903 Some(iter) => Ok(iter),
904 None => Ok(Box::new(std::iter::empty())),
906 }
907 }
908
909 fn record_batch_schema_hint(&self) -> Option<SchemaRef> {
910 self.part.schemas().next()
911 }
912
913 fn encoded_range(&self) -> Option<EncodedRange> {
914 None
915 }
916}
917
918struct EncodedBulkRangeIterBuilder {
920 file_id: FileId,
921 part: EncodedBulkPart,
922 context: Arc<BulkIterContext>,
923 sequence: Option<SequenceRange>,
924}
925
926impl IterBuilder for EncodedBulkRangeIterBuilder {
927 fn build(&self, _metrics: Option<MemScanMetrics>) -> Result<BoxedBatchIterator> {
928 UnsupportedOperationSnafu {
929 err_msg: "BatchIterator is not supported for encoded bulk memtable",
930 }
931 .fail()
932 }
933
934 fn is_record_batch(&self) -> bool {
935 true
936 }
937
938 fn build_record_batch(
939 &self,
940 _time_range: Option<(Timestamp, Timestamp)>,
941 metrics: Option<MemScanMetrics>,
942 ) -> Result<BoxedRecordBatchIterator> {
943 if let Some(iter) = self
944 .part
945 .read(self.context.clone(), self.sequence, metrics)?
946 {
947 Ok(iter)
948 } else {
949 Ok(Box::new(std::iter::empty()))
951 }
952 }
953
954 fn record_batch_schema_hint(&self) -> Option<SchemaRef> {
955 Some(self.part.schema())
956 }
957
958 fn encoded_range(&self) -> Option<EncodedRange> {
959 Some(EncodedRange {
960 data: self.part.data().clone(),
961 sst_info: self.part.to_sst_info(self.file_id),
962 })
963 }
964}
965
966struct BulkPartWrapper {
967 part: PartToMerge,
969 merging: bool,
971}
972
973impl BulkPartWrapper {
974 fn file_id(&self) -> FileId {
976 self.part.file_id()
977 }
978}
979
980#[derive(Clone)]
982enum PartToMerge {
983 Bulk { part: BulkPart, file_id: FileId },
985 Multi {
987 part: MultiBulkPart,
988 file_id: FileId,
989 },
990 Encoded {
992 part: EncodedBulkPart,
993 file_id: FileId,
994 },
995}
996
997impl PartToMerge {
998 fn file_id(&self) -> FileId {
1000 match self {
1001 PartToMerge::Bulk { file_id, .. } => *file_id,
1002 PartToMerge::Multi { file_id, .. } => *file_id,
1003 PartToMerge::Encoded { file_id, .. } => *file_id,
1004 }
1005 }
1006
1007 fn min_timestamp(&self) -> i64 {
1009 match self {
1010 PartToMerge::Bulk { part, .. } => part.min_timestamp,
1011 PartToMerge::Multi { part, .. } => part.min_timestamp(),
1012 PartToMerge::Encoded { part, .. } => part.metadata().min_timestamp,
1013 }
1014 }
1015
1016 fn max_timestamp(&self) -> i64 {
1018 match self {
1019 PartToMerge::Bulk { part, .. } => part.max_timestamp,
1020 PartToMerge::Multi { part, .. } => part.max_timestamp(),
1021 PartToMerge::Encoded { part, .. } => part.metadata().max_timestamp,
1022 }
1023 }
1024
1025 fn num_rows(&self) -> usize {
1027 match self {
1028 PartToMerge::Bulk { part, .. } => part.num_rows(),
1029 PartToMerge::Multi { part, .. } => part.num_rows(),
1030 PartToMerge::Encoded { part, .. } => part.metadata().num_rows,
1031 }
1032 }
1033
1034 fn max_sequence(&self) -> u64 {
1036 match self {
1037 PartToMerge::Bulk { part, .. } => part.sequence,
1038 PartToMerge::Multi { part, .. } => part.max_sequence(),
1039 PartToMerge::Encoded { part, .. } => part.metadata().max_sequence,
1040 }
1041 }
1042
1043 fn series_count(&self) -> usize {
1045 match self {
1046 PartToMerge::Bulk { part, .. } => part.estimated_series_count(),
1047 PartToMerge::Multi { part, .. } => part.series_count(),
1048 PartToMerge::Encoded { part, .. } => part.metadata().num_series as usize,
1049 }
1050 }
1051
1052 fn is_encoded(&self) -> bool {
1054 matches!(self, PartToMerge::Encoded { .. })
1055 }
1056
1057 fn estimated_size(&self) -> usize {
1059 match self {
1060 PartToMerge::Bulk { part, .. } => part.estimated_size(),
1061 PartToMerge::Multi { part, .. } => part.estimated_size(),
1062 PartToMerge::Encoded { part, .. } => part.size_bytes(),
1063 }
1064 }
1065
1066 fn batch_size_statistic(&self) -> Option<(u64, u64)> {
1068 match self {
1069 PartToMerge::Bulk { part, .. } => {
1070 Some((part.num_rows() as u64, part.estimated_size() as u64))
1071 }
1072 PartToMerge::Multi { part, .. } => {
1073 Some((part.num_rows() as u64, part.estimated_size() as u64))
1074 }
1075 PartToMerge::Encoded { part, .. } => part
1076 .metadata()
1077 .parquet_metadata
1078 .row_groups()
1079 .iter()
1080 .map(|row_group| {
1081 let uncompressed_bytes = row_group
1082 .columns()
1083 .iter()
1084 .map(|column| column.uncompressed_size() as u64)
1085 .sum();
1086 (row_group.num_rows() as u64, uncompressed_bytes)
1087 })
1088 .max_by_key(|(_, uncompressed_bytes)| *uncompressed_bytes),
1089 }
1090 }
1091
1092 fn to_memtable_stats(&self, region_metadata: &RegionMetadataRef) -> MemtableStats {
1094 match self {
1095 PartToMerge::Bulk { part, .. } => part.to_memtable_stats(region_metadata),
1096 PartToMerge::Multi { part, .. } => part.to_memtable_stats(region_metadata),
1097 PartToMerge::Encoded { part, .. } => part.to_memtable_stats(),
1098 }
1099 }
1100
1101 fn arrow_schema(&self) -> SchemaRef {
1103 match self {
1104 PartToMerge::Bulk { part, .. } => part.schema(),
1105 PartToMerge::Multi { part, .. } => part
1108 .schemas()
1109 .next()
1110 .expect("MultiBulkPart must contain at least one record batch"),
1111 PartToMerge::Encoded { part, .. } => part.schema(),
1112 }
1113 }
1114
1115 fn create_iterator(
1117 self,
1118 context: Arc<BulkIterContext>,
1119 ) -> Result<Option<BoxedRecordBatchIterator>> {
1120 match self {
1121 PartToMerge::Bulk { part, .. } => {
1122 let series_count = part.estimated_series_count();
1123 let iter = BulkPartBatchIter::from_single(
1124 part.batch,
1125 context,
1126 None, series_count,
1128 None, );
1130 Ok(Some(Box::new(iter) as BoxedRecordBatchIterator))
1131 }
1132 PartToMerge::Multi { part, .. } => part.read(context, None, None),
1133 PartToMerge::Encoded { part, .. } => part.read(context, None, None),
1134 }
1135 }
1136}
1137
1138struct MemtableCompactor {
1139 region_id: RegionId,
1140 memtable_id: MemtableId,
1141 config: BulkMemtableConfig,
1143 row_group_size: usize,
1145}
1146
1147impl MemtableCompactor {
1148 fn new(
1150 region_id: RegionId,
1151 memtable_id: MemtableId,
1152 config: BulkMemtableConfig,
1153 row_group_size: usize,
1154 ) -> Self {
1155 Self {
1156 region_id,
1157 memtable_id,
1158 config,
1159 row_group_size,
1160 }
1161 }
1162
1163 fn merge_parts(
1165 &mut self,
1166 bulk_parts: &RwLock<BulkParts>,
1167 metadata: &RegionMetadataRef,
1168 dedup: bool,
1169 merge_mode: MergeMode,
1170 ) -> Result<()> {
1171 let start = Instant::now();
1172
1173 let collected = bulk_parts.write().unwrap().collect_parts_to_merge(
1175 self.config.merge_threshold,
1176 self.config.max_merge_groups,
1177 self.config.encode_bytes_threshold,
1178 );
1179
1180 if collected.groups.is_empty() {
1181 return Ok(());
1182 }
1183
1184 let merged_file_ids: HashSet<FileId> = collected
1186 .groups
1187 .iter()
1188 .flatten()
1189 .map(|part| part.file_id())
1190 .collect();
1191 let mut guard = MergingFlagsGuard::new(bulk_parts, &merged_file_ids);
1192
1193 let num_groups = collected.groups.len();
1194 let num_parts: usize = collected.groups.iter().map(|g| g.len()).sum();
1195
1196 let encode_row_threshold = self.config.encode_row_threshold;
1197 let encode_bytes_threshold = self.config.encode_bytes_threshold;
1198 let row_group_size = self.row_group_size;
1199
1200 let merged_parts = collected
1202 .groups
1203 .into_par_iter()
1204 .map(|group| {
1205 Self::merge_parts_group(
1206 group,
1207 metadata,
1208 dedup,
1209 merge_mode,
1210 encode_row_threshold,
1211 encode_bytes_threshold,
1212 row_group_size,
1213 )
1214 })
1215 .collect::<Result<Vec<Option<MergedPart>>>>()?;
1216
1217 let total_output_rows = {
1219 let mut parts = bulk_parts.write().unwrap();
1220 parts.install_merged_parts(merged_parts.into_iter().flatten(), &merged_file_ids)
1221 };
1222
1223 guard.mark_success();
1224
1225 common_telemetry::debug!(
1226 "BulkMemtable {} {} concurrent compact {} groups, {} parts, {} rows, cost: {:?}",
1227 self.region_id,
1228 self.memtable_id,
1229 num_groups,
1230 num_parts,
1231 total_output_rows,
1232 start.elapsed()
1233 );
1234
1235 Ok(())
1236 }
1237
1238 #[allow(clippy::too_many_arguments)]
1240 fn merge_parts_group(
1241 parts_to_merge: Vec<PartToMerge>,
1242 metadata: &RegionMetadataRef,
1243 dedup: bool,
1244 merge_mode: MergeMode,
1245 encode_row_threshold: usize,
1246 encode_bytes_threshold: usize,
1247 row_group_size: usize,
1248 ) -> Result<Option<MergedPart>> {
1249 if parts_to_merge.is_empty() {
1250 return Ok(None);
1251 }
1252
1253 let min_timestamp = parts_to_merge
1255 .iter()
1256 .map(|p| p.min_timestamp())
1257 .min()
1258 .unwrap_or(i64::MAX);
1259 let max_timestamp = parts_to_merge
1260 .iter()
1261 .map(|p| p.max_timestamp())
1262 .max()
1263 .unwrap_or(i64::MIN);
1264 let max_sequence = parts_to_merge
1265 .iter()
1266 .map(|p| p.max_sequence())
1267 .max()
1268 .unwrap_or(0);
1269
1270 let estimated_total_rows: usize = parts_to_merge.iter().map(|p| p.num_rows()).sum();
1272 let estimated_total_bytes: usize = parts_to_merge.iter().map(|p| p.estimated_size()).sum();
1273 let estimated_series_count = parts_to_merge
1274 .iter()
1275 .map(|p| p.series_count())
1276 .max()
1277 .unwrap_or(0);
1278
1279 let batch_size = crate::batch_size::estimate_batch_size(
1280 parts_to_merge
1281 .iter()
1282 .filter_map(PartToMerge::batch_size_statistic),
1283 );
1284 let context = Arc::new(BulkIterContext::new(
1285 metadata.clone(),
1286 None, None, true,
1289 batch_size,
1290 )?);
1291
1292 let aligner = Json2Aligner::try_new(parts_to_merge.iter().map(PartToMerge::arrow_schema))?;
1293
1294 let iterators: Vec<BoxedRecordBatchIterator> = parts_to_merge
1295 .into_iter()
1296 .filter_map(|part| part.create_iterator(context.clone()).ok().flatten())
1297 .map(|iter| aligner.wrap_iter(iter))
1298 .collect();
1299
1300 if iterators.is_empty() {
1301 return Ok(None);
1302 }
1303
1304 let merged_iter = FlatMergeIterator::new(aligner.schema().clone(), iterators, batch_size)?;
1305
1306 let boxed_iter: BoxedRecordBatchIterator = if dedup {
1307 match merge_mode {
1308 MergeMode::LastRow => {
1309 let dedup_iter = FlatDedupIterator::new(merged_iter, FlatLastRow::new(false));
1310 Box::new(dedup_iter)
1311 }
1312 MergeMode::LastNonNull => {
1313 let field_column_start =
1314 field_column_start(metadata, aligner.schema().fields().len());
1315
1316 let dedup_iter = FlatDedupIterator::new(
1317 merged_iter,
1318 FlatLastNonNull::new(field_column_start, false),
1319 );
1320 Box::new(dedup_iter)
1321 }
1322 }
1323 } else {
1324 Box::new(merged_iter)
1325 };
1326
1327 if estimated_total_rows > encode_row_threshold
1329 || estimated_total_bytes > encode_bytes_threshold
1330 {
1331 let encoder = BulkPartEncoder::new(metadata.clone(), row_group_size)?;
1332 let mut metrics = BulkPartEncodeMetrics::default();
1333 let encoded_part = encoder.encode_record_batch_iter(
1334 boxed_iter,
1335 aligner.schema().clone(),
1336 min_timestamp,
1337 max_timestamp,
1338 max_sequence,
1339 &mut metrics,
1340 )?;
1341
1342 common_telemetry::trace!("merge_parts_group metrics: {:?}", metrics);
1343
1344 Ok(encoded_part.map(MergedPart::Encoded))
1345 } else {
1346 let mut batches = Vec::new();
1348 let mut actual_total_rows = 0;
1349
1350 for batch_result in boxed_iter {
1351 let batch = batch_result?;
1352 actual_total_rows += batch.num_rows();
1353 batches.push(batch);
1354 }
1355
1356 if actual_total_rows == 0 {
1357 return Ok(None);
1358 }
1359
1360 let multi_part = MultiBulkPart::new(
1361 batches,
1362 min_timestamp,
1363 max_timestamp,
1364 max_sequence,
1365 estimated_series_count,
1366 metadata,
1367 );
1368
1369 common_telemetry::trace!(
1370 "merge_parts_group created MultiBulkPart: rows={}, batches={}",
1371 actual_total_rows,
1372 multi_part.num_batches()
1373 );
1374
1375 Ok(Some(MergedPart::Multi(multi_part)))
1376 }
1377 }
1378}
1379
1380struct MemCompactTask {
1382 metadata: RegionMetadataRef,
1383 parts: Arc<RwLock<BulkParts>>,
1384 config: BulkMemtableConfig,
1386 compactor: Arc<Mutex<MemtableCompactor>>,
1388 append_mode: bool,
1390 merge_mode: MergeMode,
1392}
1393
1394impl MemCompactTask {
1395 fn compact(&self) -> Result<()> {
1396 let mut compactor = self.compactor.lock().unwrap();
1397
1398 let should_merge = self.parts.read().unwrap().should_merge_parts(
1399 self.config.merge_threshold,
1400 self.config.encode_bytes_threshold,
1401 );
1402 if should_merge {
1403 compactor.merge_parts(
1404 &self.parts,
1405 &self.metadata,
1406 !self.append_mode,
1407 self.merge_mode,
1408 )?;
1409 }
1410
1411 Ok(())
1412 }
1413}
1414
1415#[derive(Debug)]
1417pub struct CompactDispatcher {
1418 semaphore: Arc<Semaphore>,
1419}
1420
1421impl CompactDispatcher {
1422 pub fn new(permits: usize) -> Self {
1424 Self {
1425 semaphore: Arc::new(Semaphore::new(permits)),
1426 }
1427 }
1428
1429 fn dispatch_compact(&self, task: MemCompactTask) {
1431 let semaphore = self.semaphore.clone();
1432 common_runtime::spawn_global(async move {
1433 let Ok(_permit) = semaphore.acquire().await else {
1434 return;
1435 };
1436
1437 common_runtime::spawn_blocking_global(move || {
1438 if let Err(e) = task.compact() {
1439 common_telemetry::error!(e; "Failed to compact memtable, region: {}", task.metadata.region_id);
1440 }
1441 });
1442 });
1443 }
1444}
1445
1446#[derive(Debug)]
1448pub struct BulkMemtableBuilder {
1449 config: BulkMemtableConfig,
1451 write_buffer_manager: Option<WriteBufferManagerRef>,
1452 compact_dispatcher: Option<Arc<CompactDispatcher>>,
1453 append_mode: bool,
1454 merge_mode: MergeMode,
1455 row_group_size: usize,
1457}
1458
1459impl Default for BulkMemtableBuilder {
1460 fn default() -> Self {
1461 Self {
1462 config: BulkMemtableConfig::default(),
1463 write_buffer_manager: None,
1464 compact_dispatcher: None,
1465 append_mode: false,
1466 merge_mode: MergeMode::default(),
1467 row_group_size: DEFAULT_ROW_GROUP_SIZE,
1468 }
1469 }
1470}
1471
1472impl BulkMemtableBuilder {
1473 pub fn new(
1475 write_buffer_manager: Option<WriteBufferManagerRef>,
1476 append_mode: bool,
1477 merge_mode: MergeMode,
1478 ) -> Self {
1479 Self {
1480 write_buffer_manager,
1481 append_mode,
1482 merge_mode,
1483 ..Default::default()
1484 }
1485 }
1486
1487 pub fn with_config(mut self, config: BulkMemtableConfig) -> Self {
1489 self.config = config;
1490 self
1491 }
1492
1493 pub fn with_row_group_size(mut self, row_group_size: usize) -> Self {
1495 self.row_group_size = row_group_size;
1496 self
1497 }
1498
1499 pub fn with_compact_dispatcher(mut self, compact_dispatcher: Arc<CompactDispatcher>) -> Self {
1501 self.compact_dispatcher = Some(compact_dispatcher);
1502 self
1503 }
1504
1505 #[cfg(test)]
1506 pub(crate) fn config(&self) -> &BulkMemtableConfig {
1507 &self.config
1508 }
1509}
1510
1511impl MemtableBuilder for BulkMemtableBuilder {
1512 fn build(&self, id: MemtableId, metadata: &RegionMetadataRef) -> MemtableRef {
1513 Arc::new(BulkMemtable::new_with_row_group_size(
1514 id,
1515 self.config.clone(),
1516 metadata.clone(),
1517 self.write_buffer_manager.clone(),
1518 self.compact_dispatcher.clone(),
1519 self.append_mode,
1520 self.merge_mode,
1521 self.row_group_size,
1522 ))
1523 }
1524
1525 fn use_bulk_insert(&self, _metadata: &RegionMetadataRef) -> bool {
1526 true
1527 }
1528}
1529
1530#[cfg(test)]
1531mod tests {
1532 use api::helper::encode_json_value;
1533 use api::v1::value::ValueData;
1534 use api::v1::{Mutation, Row, Rows, SemanticType};
1535 use datatypes::data_type::ConcreteDataType;
1536 use datatypes::extension::json::{JsonExtensionType, JsonMetadata};
1537 use datatypes::json::value::JsonValue;
1538 use datatypes::schema::ColumnSchema;
1539 use datatypes::types::json_type::{JsonNativeType, JsonObjectType};
1540 use mito_codec::row_converter::build_primary_key_codec;
1541 use serde_json::json;
1542 use store_api::metadata::{ColumnMetadata, RegionMetadataBuilder, RegionMetadataRef};
1543
1544 use super::*;
1545 use crate::memtable::bulk::part::BulkPartConverter;
1546 use crate::read::scan_region::PredicateGroup;
1547 use crate::sst::{FlatSchemaOptions, to_flat_sst_arrow_schema};
1548 use crate::test_util::memtable_util::{
1549 build_key_values_with_ts_seq_values, metadata_for_test, region_metadata_to_row_schema,
1550 };
1551
1552 fn create_bulk_part_with_converter(
1553 k0: &str,
1554 k1: u32,
1555 timestamps: Vec<i64>,
1556 values: Vec<Option<f64>>,
1557 sequence: u64,
1558 ) -> Result<BulkPart> {
1559 let metadata = metadata_for_test();
1560 let capacity = 100;
1561 let primary_key_codec = build_primary_key_codec(&metadata);
1562 let schema = to_flat_sst_arrow_schema(
1563 &metadata,
1564 &FlatSchemaOptions::from_encoding(metadata.primary_key_encoding),
1565 );
1566
1567 let mut converter =
1568 BulkPartConverter::new(&metadata, schema, capacity, primary_key_codec, true);
1569
1570 let key_values = build_key_values_with_ts_seq_values(
1571 &metadata,
1572 k0.to_string(),
1573 k1,
1574 timestamps.into_iter(),
1575 values.into_iter(),
1576 sequence,
1577 );
1578
1579 converter.append_key_values(&key_values)?;
1580 converter.convert()
1581 }
1582
1583 #[test]
1584 fn test_bulk_memtable_sanitizes_zero_merge_threshold() {
1585 let metadata = metadata_for_test();
1586 let config = BulkMemtableConfig {
1587 merge_threshold: 0,
1588 ..Default::default()
1589 };
1590
1591 let memtable =
1592 BulkMemtable::new(999, config, metadata, None, None, false, MergeMode::LastRow);
1593
1594 assert_eq!(DEFAULT_MERGE_THRESHOLD, memtable.config.merge_threshold);
1595 assert_eq!(
1596 DEFAULT_MERGE_THRESHOLD,
1597 memtable.compactor.lock().unwrap().config.merge_threshold
1598 );
1599 }
1600
1601 #[test]
1602 fn test_bulk_memtable_write_read() {
1603 let metadata = metadata_for_test();
1604 let memtable = BulkMemtable::new(
1605 999,
1606 BulkMemtableConfig::default(),
1607 metadata.clone(),
1608 None,
1609 None,
1610 false,
1611 MergeMode::LastRow,
1612 );
1613 memtable.set_unordered_part_threshold(0);
1615
1616 let test_data = [
1617 (
1618 "key_a",
1619 1u32,
1620 vec![1000i64, 2000i64],
1621 vec![Some(10.5), Some(20.5)],
1622 100u64,
1623 ),
1624 (
1625 "key_b",
1626 2u32,
1627 vec![1500i64, 2500i64],
1628 vec![Some(15.5), Some(25.5)],
1629 200u64,
1630 ),
1631 ("key_c", 3u32, vec![3000i64], vec![Some(30.5)], 300u64),
1632 ];
1633
1634 for (k0, k1, timestamps, values, seq) in test_data.iter() {
1635 let part =
1636 create_bulk_part_with_converter(k0, *k1, timestamps.clone(), values.clone(), *seq)
1637 .unwrap();
1638 memtable.write_bulk(part).unwrap();
1639 }
1640
1641 let stats = memtable.stats();
1642 assert_eq!(5, stats.num_rows);
1643 assert_eq!(3, stats.num_ranges);
1644 assert_eq!(300, stats.max_sequence);
1645
1646 let (min_ts, max_ts) = stats.time_range.unwrap();
1647 assert_eq!(1000, min_ts.value());
1648 assert_eq!(3000, max_ts.value());
1649
1650 let predicate_group = PredicateGroup::new(&metadata, &[]).unwrap();
1651 let ranges = memtable
1652 .ranges(
1653 None,
1654 RangesOptions::default().with_predicate(predicate_group),
1655 )
1656 .unwrap();
1657
1658 assert_eq!(3, ranges.ranges.len());
1659 let total_rows: usize = ranges.ranges.values().map(|r| r.stats().num_rows()).sum();
1660 assert_eq!(5, total_rows);
1661
1662 for (_range_id, range) in ranges.ranges.iter() {
1663 assert!(range.num_rows() > 0);
1664 assert!(range.is_record_batch());
1665
1666 let record_batch_iter = range.build_record_batch_iter(None, None).unwrap();
1667
1668 let mut total_rows = 0;
1669 for batch_result in record_batch_iter {
1670 let batch = batch_result.unwrap();
1671 total_rows += batch.num_rows();
1672 assert!(batch.num_rows() > 0);
1673 assert_eq!(8, batch.num_columns());
1674 }
1675 assert_eq!(total_rows, range.num_rows());
1676 }
1677 }
1678
1679 #[test]
1680 fn test_bulk_memtable_compact_parts_with_json2() {
1681 let metadata = mock_metadata_with_json2();
1682
1683 let config = BulkMemtableConfig {
1684 merge_threshold: 2,
1685 encode_row_threshold: 1,
1686 encode_bytes_threshold: 1,
1687 ..Default::default()
1688 };
1689 let memtable = BulkMemtable::new(
1690 999,
1691 config,
1692 metadata.clone(),
1693 None,
1694 None,
1695 true,
1696 MergeMode::LastRow,
1697 );
1698 memtable.set_unordered_part_threshold(0);
1699
1700 let part1 = mock_bulk_part_with_json2(&metadata, vec![1000, 2000], 100).unwrap();
1701 let part2 = mock_bulk_part_with_json2(&metadata, vec![3000, 4000], 200).unwrap();
1702
1703 memtable.write_bulk(part1).unwrap();
1704 memtable.write_bulk(part2).unwrap();
1705 memtable.compact(false).unwrap();
1706
1707 let stats = memtable.stats();
1708 assert_eq!(4, stats.num_rows);
1709 assert_eq!(201, stats.max_sequence);
1710
1711 let predicate_group = PredicateGroup::new(&metadata, &[]).unwrap();
1712 let opts = RangesOptions::default().with_predicate(predicate_group);
1713 let ranges = memtable.ranges(None, opts).unwrap();
1714
1715 assert_eq!(1, ranges.ranges.len());
1716 let total_rows: usize = ranges.ranges.values().map(|r| r.stats().num_rows()).sum();
1717 assert_eq!(4, total_rows);
1718 }
1719
1720 fn mock_metadata_with_json2() -> RegionMetadataRef {
1721 let col_meta_1 = ColumnMetadata {
1722 column_schema: ColumnSchema::new(
1723 "ts",
1724 ConcreteDataType::timestamp_millisecond_datatype(),
1725 false,
1726 ),
1727 semantic_type: SemanticType::Timestamp,
1728 column_id: 0,
1729 };
1730
1731 let data_type = ConcreteDataType::json2(JsonNativeType::Object(JsonObjectType::new()));
1732 let mut col_schema = ColumnSchema::new("data", data_type, true);
1733 let extension = JsonExtensionType::new(Arc::new(JsonMetadata::default()));
1734 col_schema.with_extension_type(&extension).unwrap();
1735
1736 let col_meta_2 = ColumnMetadata {
1737 column_schema: col_schema,
1738 semantic_type: SemanticType::Field,
1739 column_id: 1,
1740 };
1741 let mut builder = RegionMetadataBuilder::new(RegionId::new(123, 789));
1742 builder
1743 .push_column_metadata(col_meta_1)
1744 .push_column_metadata(col_meta_2)
1745 .primary_key(vec![]);
1746 Arc::new(builder.build().unwrap())
1747 }
1748
1749 fn mock_bulk_part_with_json2(
1750 metadata: &RegionMetadataRef,
1751 timestamps: Vec<i64>,
1752 sequence: u64,
1753 ) -> Result<BulkPart> {
1754 let capacity = timestamps.len();
1755 let primary_key_codec = build_primary_key_codec(metadata);
1756 let json_type = JsonNativeType::Object(JsonObjectType::from([
1757 ("id".to_string(), JsonNativeType::i64()),
1758 (
1759 "payload".to_string(),
1760 JsonNativeType::Object(JsonObjectType::from([(
1761 "message".to_string(),
1762 JsonNativeType::String,
1763 )])),
1764 ),
1765 ]));
1766 let mut options = FlatSchemaOptions::from_encoding(metadata.primary_key_encoding);
1767 options
1768 .concretized_json_types
1769 .insert("data".to_string(), json_type.as_arrow_type());
1770 let schema = to_flat_sst_arrow_schema(metadata, &options);
1771
1772 let mut converter =
1773 BulkPartConverter::new(metadata, schema, capacity, primary_key_codec, true);
1774
1775 let rows = timestamps
1776 .into_iter()
1777 .map(|ts| {
1778 let val1 = api::v1::Value {
1779 value_data: Some(ValueData::TimestampMillisecondValue(ts)),
1780 };
1781 let value_data = ValueData::JsonValue(encode_json_value(JsonValue::from(json!({
1782 "id": ts,
1783 "payload": {
1784 "message": format!("row-{ts}"),
1785 },
1786 }))));
1787 let val2 = api::v1::Value {
1788 value_data: Some(value_data),
1789 };
1790 Row {
1791 values: vec![val1, val2],
1792 }
1793 })
1794 .collect();
1795
1796 let mutation = Mutation {
1797 op_type: 1,
1798 sequence,
1799 rows: Some(Rows {
1800 schema: region_metadata_to_row_schema(metadata),
1801 rows,
1802 }),
1803 write_hint: None,
1804 };
1805 let key_values = KeyValues::new(metadata.as_ref(), mutation).unwrap();
1806
1807 converter.append_key_values(&key_values)?;
1808 converter.convert()
1809 }
1810
1811 #[test]
1812 fn test_bulk_memtable_ranges_with_projection() {
1813 let metadata = metadata_for_test();
1814 let memtable = BulkMemtable::new(
1815 111,
1816 BulkMemtableConfig::default(),
1817 metadata.clone(),
1818 None,
1819 None,
1820 false,
1821 MergeMode::LastRow,
1822 );
1823
1824 let bulk_part = create_bulk_part_with_converter(
1825 "projection_test",
1826 5,
1827 vec![5000, 6000, 7000],
1828 vec![Some(50.0), Some(60.0), Some(70.0)],
1829 500,
1830 )
1831 .unwrap();
1832
1833 memtable.write_bulk(bulk_part).unwrap();
1834
1835 let projection = vec![4u32];
1836 let predicate_group = PredicateGroup::new(&metadata, &[]).unwrap();
1837 let ranges = memtable
1838 .ranges(
1839 Some(&projection),
1840 RangesOptions::default().with_predicate(predicate_group),
1841 )
1842 .unwrap();
1843
1844 assert_eq!(1, ranges.ranges.len());
1845 let range = ranges.ranges.get(&0).unwrap();
1846
1847 assert!(range.is_record_batch());
1848 let record_batch_iter = range.build_record_batch_iter(None, None).unwrap();
1849
1850 let mut total_rows = 0;
1851 for batch_result in record_batch_iter {
1852 let batch = batch_result.unwrap();
1853 assert!(batch.num_rows() > 0);
1854 assert_eq!(5, batch.num_columns());
1855 total_rows += batch.num_rows();
1856 }
1857 assert_eq!(3, total_rows);
1858 }
1859
1860 #[test]
1861 fn test_bulk_memtable_unsupported_operations() {
1862 let metadata = metadata_for_test();
1863 let memtable = BulkMemtable::new(
1864 111,
1865 BulkMemtableConfig::default(),
1866 metadata.clone(),
1867 None,
1868 None,
1869 false,
1870 MergeMode::LastRow,
1871 );
1872
1873 let key_values = build_key_values_with_ts_seq_values(
1874 &metadata,
1875 "test".to_string(),
1876 1,
1877 vec![1000].into_iter(),
1878 vec![Some(1.0)].into_iter(),
1879 1,
1880 );
1881
1882 let err = memtable.write(&key_values).unwrap_err();
1883 assert!(err.to_string().contains("not supported"));
1884
1885 let kv = key_values.iter().next().unwrap();
1886 let err = memtable.write_one(kv).unwrap_err();
1887 assert!(err.to_string().contains("not supported"));
1888 }
1889
1890 #[test]
1891 fn test_bulk_memtable_freeze() {
1892 let metadata = metadata_for_test();
1893 let memtable = BulkMemtable::new(
1894 222,
1895 BulkMemtableConfig::default(),
1896 metadata.clone(),
1897 None,
1898 None,
1899 false,
1900 MergeMode::LastRow,
1901 );
1902
1903 let bulk_part = create_bulk_part_with_converter(
1904 "freeze_test",
1905 10,
1906 vec![10000],
1907 vec![Some(100.0)],
1908 1000,
1909 )
1910 .unwrap();
1911
1912 memtable.write_bulk(bulk_part).unwrap();
1913 memtable.freeze().unwrap();
1914
1915 let stats_after_freeze = memtable.stats();
1916 assert_eq!(1, stats_after_freeze.num_rows);
1917 }
1918
1919 #[test]
1920 fn test_bulk_memtable_fork() {
1921 let metadata = metadata_for_test();
1922 let original_memtable = BulkMemtable::new(
1923 333,
1924 BulkMemtableConfig::default(),
1925 metadata.clone(),
1926 None,
1927 None,
1928 false,
1929 MergeMode::LastRow,
1930 );
1931
1932 let bulk_part =
1933 create_bulk_part_with_converter("fork_test", 15, vec![15000], vec![Some(150.0)], 1500)
1934 .unwrap();
1935
1936 original_memtable.write_bulk(bulk_part).unwrap();
1937
1938 let forked_memtable = original_memtable.fork(444, &metadata);
1939
1940 assert_eq!(forked_memtable.id(), 444);
1941 assert!(forked_memtable.is_empty());
1942 assert_eq!(0, forked_memtable.stats().num_rows);
1943
1944 assert!(!original_memtable.is_empty());
1945 assert_eq!(1, original_memtable.stats().num_rows);
1946 }
1947
1948 #[test]
1949 fn test_bulk_memtable_ranges_multiple_parts() {
1950 let metadata = metadata_for_test();
1951 let memtable = BulkMemtable::new(
1952 777,
1953 BulkMemtableConfig::default(),
1954 metadata.clone(),
1955 None,
1956 None,
1957 false,
1958 MergeMode::LastRow,
1959 );
1960 memtable.set_unordered_part_threshold(0);
1962
1963 let parts_data = vec![
1964 (
1965 "part1",
1966 1u32,
1967 vec![1000i64, 1100i64],
1968 vec![Some(10.0), Some(11.0)],
1969 100u64,
1970 ),
1971 (
1972 "part2",
1973 2u32,
1974 vec![2000i64, 2100i64],
1975 vec![Some(20.0), Some(21.0)],
1976 200u64,
1977 ),
1978 ("part3", 3u32, vec![3000i64], vec![Some(30.0)], 300u64),
1979 ];
1980
1981 for (k0, k1, timestamps, values, seq) in parts_data {
1982 let part = create_bulk_part_with_converter(k0, k1, timestamps, values, seq).unwrap();
1983 memtable.write_bulk(part).unwrap();
1984 }
1985
1986 let predicate_group = PredicateGroup::new(&metadata, &[]).unwrap();
1987 let ranges = memtable
1988 .ranges(
1989 None,
1990 RangesOptions::default().with_predicate(predicate_group),
1991 )
1992 .unwrap();
1993
1994 assert_eq!(3, ranges.ranges.len());
1995 let total_rows: usize = ranges.ranges.values().map(|r| r.stats().num_rows()).sum();
1996 assert_eq!(5, total_rows);
1997 assert_eq!(3, ranges.ranges.len());
1998
1999 for (range_id, range) in ranges.ranges.iter() {
2000 assert!(*range_id < 3);
2001 assert!(range.num_rows() > 0);
2002 assert!(range.is_record_batch());
2003 }
2004 }
2005
2006 #[test]
2007 fn test_bulk_memtable_ranges_with_sequence_filter() {
2008 let metadata = metadata_for_test();
2009 let memtable = BulkMemtable::new(
2010 888,
2011 BulkMemtableConfig::default(),
2012 metadata.clone(),
2013 None,
2014 None,
2015 false,
2016 MergeMode::LastRow,
2017 );
2018
2019 let part = create_bulk_part_with_converter(
2020 "seq_test",
2021 1,
2022 vec![1000, 2000, 3000],
2023 vec![Some(10.0), Some(20.0), Some(30.0)],
2024 500,
2025 )
2026 .unwrap();
2027
2028 memtable.write_bulk(part).unwrap();
2029
2030 let predicate_group = PredicateGroup::new(&metadata, &[]).unwrap();
2031 let sequence_filter = Some(SequenceRange::LtEq { max: 400 }); let ranges = memtable
2033 .ranges(
2034 None,
2035 RangesOptions::default()
2036 .with_predicate(predicate_group)
2037 .with_sequence(sequence_filter),
2038 )
2039 .unwrap();
2040
2041 assert_eq!(1, ranges.ranges.len());
2042 let range = ranges.ranges.get(&0).unwrap();
2043
2044 let mut record_batch_iter = range.build_record_batch_iter(None, None).unwrap();
2045 assert!(record_batch_iter.next().is_none());
2046 }
2047
2048 #[test]
2049 fn test_bulk_memtable_ranges_with_encoded_parts() {
2050 let metadata = metadata_for_test();
2051 let config = BulkMemtableConfig {
2052 merge_threshold: 8,
2053 ..Default::default()
2054 };
2055 let memtable = BulkMemtable::new(
2056 999,
2057 config,
2058 metadata.clone(),
2059 None,
2060 None,
2061 false,
2062 MergeMode::LastRow,
2063 );
2064 memtable.set_unordered_part_threshold(0);
2066
2067 for i in 0..10 {
2069 let part = create_bulk_part_with_converter(
2070 &format!("key_{}", i),
2071 i,
2072 vec![1000 + i as i64 * 100],
2073 vec![Some(i as f64 * 10.0)],
2074 100 + i as u64,
2075 )
2076 .unwrap();
2077 memtable.write_bulk(part).unwrap();
2078 }
2079
2080 memtable.compact(false).unwrap();
2081
2082 let predicate_group = PredicateGroup::new(&metadata, &[]).unwrap();
2083 let ranges = memtable
2084 .ranges(
2085 None,
2086 RangesOptions::default().with_predicate(predicate_group),
2087 )
2088 .unwrap();
2089
2090 assert_eq!(3, ranges.ranges.len());
2092 let total_rows: usize = ranges.ranges.values().map(|r| r.stats().num_rows()).sum();
2093 assert_eq!(10, total_rows);
2094
2095 for (_range_id, range) in ranges.ranges.iter() {
2096 assert!(range.num_rows() > 0);
2097 assert!(range.is_record_batch());
2098
2099 let record_batch_iter = range.build_record_batch_iter(None, None).unwrap();
2100 let mut total_rows = 0;
2101 for batch_result in record_batch_iter {
2102 let batch = batch_result.unwrap();
2103 total_rows += batch.num_rows();
2104 assert!(batch.num_rows() > 0);
2105 }
2106 assert_eq!(total_rows, range.num_rows());
2107 }
2108 }
2109
2110 #[test]
2111 fn test_bulk_memtable_unordered_part() {
2112 let metadata = metadata_for_test();
2113 let memtable = BulkMemtable::new(
2114 1001,
2115 BulkMemtableConfig::default(),
2116 metadata.clone(),
2117 None,
2118 None,
2119 false,
2120 MergeMode::LastRow,
2121 );
2122
2123 memtable.set_unordered_part_threshold(5);
2126 memtable.set_unordered_part_compact_threshold(10);
2128
2129 for i in 0..3 {
2131 let part = create_bulk_part_with_converter(
2132 &format!("key_{}", i),
2133 i,
2134 vec![1000 + i as i64 * 100, 1100 + i as i64 * 100],
2135 vec![Some(i as f64 * 10.0), Some(i as f64 * 10.0 + 1.0)],
2136 100 + i as u64,
2137 )
2138 .unwrap();
2139 assert_eq!(2, part.num_rows());
2140 memtable.write_bulk(part).unwrap();
2141 }
2142
2143 let stats = memtable.stats();
2145 assert_eq!(6, stats.num_rows);
2146
2147 for i in 3..5 {
2150 let part = create_bulk_part_with_converter(
2151 &format!("key_{}", i),
2152 i,
2153 vec![1000 + i as i64 * 100, 1100 + i as i64 * 100],
2154 vec![Some(i as f64 * 10.0), Some(i as f64 * 10.0 + 1.0)],
2155 100 + i as u64,
2156 )
2157 .unwrap();
2158 memtable.write_bulk(part).unwrap();
2159 }
2160
2161 let stats = memtable.stats();
2163 assert_eq!(10, stats.num_rows);
2164
2165 let predicate_group = PredicateGroup::new(&metadata, &[]).unwrap();
2167 let ranges = memtable
2168 .ranges(
2169 None,
2170 RangesOptions::default().with_predicate(predicate_group),
2171 )
2172 .unwrap();
2173
2174 assert!(!ranges.ranges.is_empty());
2176 let total_rows: usize = ranges.ranges.values().map(|r| r.stats().num_rows()).sum();
2177 assert_eq!(10, total_rows);
2178
2179 let mut total_rows_read = 0;
2181 for (_range_id, range) in ranges.ranges.iter() {
2182 assert!(range.is_record_batch());
2183 let record_batch_iter = range.build_record_batch_iter(None, None).unwrap();
2184
2185 for batch_result in record_batch_iter {
2186 let batch = batch_result.unwrap();
2187 total_rows_read += batch.num_rows();
2188 }
2189 }
2190 assert_eq!(10, total_rows_read);
2191 }
2192
2193 #[test]
2194 fn test_bulk_memtable_unordered_part_mixed_sizes() {
2195 let metadata = metadata_for_test();
2196 let memtable = BulkMemtable::new(
2197 1002,
2198 BulkMemtableConfig::default(),
2199 metadata.clone(),
2200 None,
2201 None,
2202 false,
2203 MergeMode::LastRow,
2204 );
2205
2206 memtable.set_unordered_part_threshold(4);
2208 memtable.set_unordered_part_compact_threshold(8);
2209
2210 for i in 0..2 {
2212 let part = create_bulk_part_with_converter(
2213 &format!("small_{}", i),
2214 i,
2215 vec![1000 + i as i64, 2000 + i as i64, 3000 + i as i64],
2216 vec![Some(i as f64), Some(i as f64 + 1.0), Some(i as f64 + 2.0)],
2217 10 + i as u64,
2218 )
2219 .unwrap();
2220 assert_eq!(3, part.num_rows());
2221 memtable.write_bulk(part).unwrap();
2222 }
2223
2224 let large_part = create_bulk_part_with_converter(
2226 "large_key",
2227 100,
2228 vec![5000, 6000, 7000, 8000, 9000],
2229 vec![
2230 Some(100.0),
2231 Some(101.0),
2232 Some(102.0),
2233 Some(103.0),
2234 Some(104.0),
2235 ],
2236 50,
2237 )
2238 .unwrap();
2239 assert_eq!(5, large_part.num_rows());
2240 memtable.write_bulk(large_part).unwrap();
2241
2242 let part = create_bulk_part_with_converter(
2244 "small_2",
2245 2,
2246 vec![4000, 4100],
2247 vec![Some(20.0), Some(21.0)],
2248 30,
2249 )
2250 .unwrap();
2251 memtable.write_bulk(part).unwrap();
2252
2253 let stats = memtable.stats();
2254 assert_eq!(13, stats.num_rows); let predicate_group = PredicateGroup::new(&metadata, &[]).unwrap();
2258 let ranges = memtable
2259 .ranges(
2260 None,
2261 RangesOptions::default().with_predicate(predicate_group),
2262 )
2263 .unwrap();
2264
2265 let total_rows: usize = ranges.ranges.values().map(|r| r.stats().num_rows()).sum();
2266 assert_eq!(13, total_rows);
2267
2268 let mut total_rows_read = 0;
2269 for (_range_id, range) in ranges.ranges.iter() {
2270 let record_batch_iter = range.build_record_batch_iter(None, None).unwrap();
2271 for batch_result in record_batch_iter {
2272 let batch = batch_result.unwrap();
2273 total_rows_read += batch.num_rows();
2274 }
2275 }
2276 assert_eq!(13, total_rows_read);
2277 }
2278
2279 #[test]
2280 fn test_bulk_memtable_unordered_part_byte_threshold() {
2281 let metadata = metadata_for_test();
2282 let first =
2283 create_bulk_part_with_converter("first", 1, vec![1000], vec![Some(1.0)], 1).unwrap();
2284 let second =
2285 create_bulk_part_with_converter("second", 2, vec![2000], vec![Some(2.0)], 2).unwrap();
2286 let threshold = first.estimated_size().max(second.estimated_size());
2287 let memtable = BulkMemtable::new(
2288 1004,
2289 BulkMemtableConfig {
2290 encode_bytes_threshold: threshold,
2291 ..Default::default()
2292 },
2293 metadata.clone(),
2294 None,
2295 None,
2296 false,
2297 MergeMode::LastRow,
2298 );
2299 memtable.set_unordered_part_compact_threshold(usize::MAX);
2300
2301 memtable.write_bulk(first).unwrap();
2303 assert_eq!(1, memtable.parts.read().unwrap().unordered_part.num_parts());
2304
2305 memtable.write_bulk(second).unwrap();
2307 let parts = memtable.parts.read().unwrap();
2308 assert!(parts.unordered_part.is_empty());
2309 assert_eq!(1, parts.parts.len());
2310 drop(parts);
2311
2312 let oversized =
2313 create_bulk_part_with_converter("oversized", 3, vec![3000], vec![Some(3.0)], 3)
2314 .unwrap();
2315 let oversized_memtable = BulkMemtable::new(
2316 1005,
2317 BulkMemtableConfig {
2318 encode_bytes_threshold: oversized.estimated_size() - 1,
2319 ..Default::default()
2320 },
2321 metadata,
2322 None,
2323 None,
2324 false,
2325 MergeMode::LastRow,
2326 );
2327 oversized_memtable.write_bulk(oversized).unwrap();
2328 let parts = oversized_memtable.parts.read().unwrap();
2329 assert!(parts.unordered_part.is_empty());
2330 assert_eq!(1, parts.parts.len());
2331 }
2332
2333 #[test]
2334 fn test_bulk_memtable_unordered_part_with_ranges() {
2335 let metadata = metadata_for_test();
2336 let memtable = BulkMemtable::new(
2337 1003,
2338 BulkMemtableConfig::default(),
2339 metadata.clone(),
2340 None,
2341 None,
2342 false,
2343 MergeMode::LastRow,
2344 );
2345
2346 memtable.set_unordered_part_threshold(3);
2348 memtable.set_unordered_part_compact_threshold(100); for i in 0..3 {
2352 let part = create_bulk_part_with_converter(
2353 &format!("key_{}", i),
2354 i,
2355 vec![1000 + i as i64 * 100],
2356 vec![Some(i as f64 * 10.0)],
2357 100 + i as u64,
2358 )
2359 .unwrap();
2360 assert_eq!(1, part.num_rows());
2361 memtable.write_bulk(part).unwrap();
2362 }
2363
2364 let stats = memtable.stats();
2365 assert_eq!(3, stats.num_rows);
2366
2367 let predicate_group = PredicateGroup::new(&metadata, &[]).unwrap();
2369 let ranges = memtable
2370 .ranges(
2371 None,
2372 RangesOptions::default().with_predicate(predicate_group),
2373 )
2374 .unwrap();
2375
2376 assert_eq!(1, ranges.ranges.len());
2378 let total_rows: usize = ranges.ranges.values().map(|r| r.stats().num_rows()).sum();
2379 assert_eq!(3, total_rows);
2380
2381 let range = ranges.ranges.get(&0).unwrap();
2383 let record_batch_iter = range.build_record_batch_iter(None, None).unwrap();
2384
2385 let mut total_rows = 0;
2386 for batch_result in record_batch_iter {
2387 let batch = batch_result.unwrap();
2388 total_rows += batch.num_rows();
2389 assert!(batch.num_rows() > 0);
2391 }
2392 assert_eq!(3, total_rows);
2393 }
2394
2395 fn create_bulk_part_wrapper(part: BulkPart) -> BulkPartWrapper {
2397 BulkPartWrapper {
2398 part: PartToMerge::Bulk {
2399 part,
2400 file_id: FileId::random(),
2401 },
2402 merging: false,
2403 }
2404 }
2405
2406 #[test]
2407 fn test_should_merge_parts_below_threshold() {
2408 let mut bulk_parts = BulkParts::default();
2409
2410 for i in 0..DEFAULT_MERGE_THRESHOLD - 1 {
2412 let part = create_bulk_part_with_converter(
2413 &format!("key_{}", i),
2414 i as u32,
2415 vec![1000 + i as i64 * 100],
2416 vec![Some(i as f64 * 10.0)],
2417 100 + i as u64,
2418 )
2419 .unwrap();
2420 bulk_parts.parts.push(create_bulk_part_wrapper(part));
2421 }
2422
2423 assert!(!bulk_parts.should_merge_parts(DEFAULT_MERGE_THRESHOLD, usize::MAX));
2425 }
2426
2427 #[test]
2428 fn test_should_merge_parts_at_threshold() {
2429 let mut bulk_parts = BulkParts::default();
2430 let merge_threshold = 8;
2431
2432 for i in 0..merge_threshold {
2434 let part = create_bulk_part_with_converter(
2435 &format!("key_{}", i),
2436 i as u32,
2437 vec![1000 + i as i64 * 100],
2438 vec![Some(i as f64 * 10.0)],
2439 100 + i as u64,
2440 )
2441 .unwrap();
2442 bulk_parts.parts.push(create_bulk_part_wrapper(part));
2443 }
2444
2445 assert!(bulk_parts.should_merge_parts(merge_threshold, usize::MAX));
2447 }
2448
2449 #[test]
2450 fn test_raw_parts_do_not_use_encoded_merge_limit() {
2451 let mut bulk_parts = BulkParts::default();
2452 let merge_threshold = 4;
2453 for i in 0..merge_threshold {
2454 let part = create_bulk_part_with_converter(
2455 &format!("key_{}", i),
2456 i as u32,
2457 vec![1000 + i as i64],
2458 vec![Some(i as f64)],
2459 100 + i as u64,
2460 )
2461 .unwrap();
2462 bulk_parts.parts.push(create_bulk_part_wrapper(part));
2463 }
2464
2465 let max_size = bulk_parts.parts[0].part.estimated_size();
2466 assert!(bulk_parts.should_merge_parts(merge_threshold, max_size));
2467 assert_eq!(
2468 1,
2469 bulk_parts
2470 .collect_parts_to_merge(merge_threshold, 1, max_size)
2471 .groups
2472 .len()
2473 );
2474 }
2475
2476 #[test]
2477 fn test_should_merge_parts_with_merging_flag() {
2478 let mut bulk_parts = BulkParts::default();
2479 let merge_threshold = 8;
2480
2481 for i in 0..10 {
2483 let part = create_bulk_part_with_converter(
2484 &format!("key_{}", i),
2485 i as u32,
2486 vec![1000 + i as i64 * 100],
2487 vec![Some(i as f64 * 10.0)],
2488 100 + i as u64,
2489 )
2490 .unwrap();
2491 bulk_parts.parts.push(create_bulk_part_wrapper(part));
2492 }
2493
2494 assert!(bulk_parts.should_merge_parts(merge_threshold, usize::MAX));
2496
2497 for wrapper in bulk_parts.parts.iter_mut().take(3) {
2499 wrapper.merging = true;
2500 }
2501
2502 assert!(!bulk_parts.should_merge_parts(merge_threshold, usize::MAX));
2504 }
2505
2506 #[test]
2507 fn test_collect_parts_to_merge_grouping() {
2508 let mut bulk_parts = BulkParts::default();
2509
2510 for i in 0..16 {
2512 let num_rows = (i % 4) + 1; let timestamps: Vec<i64> = (0..num_rows)
2514 .map(|j| 1000 + i as i64 * 100 + j as i64)
2515 .collect();
2516 let values: Vec<Option<f64>> =
2517 (0..num_rows).map(|j| Some((i * 10 + j) as f64)).collect();
2518 let part = create_bulk_part_with_converter(
2519 &format!("key_{}", i),
2520 i as u32,
2521 timestamps,
2522 values,
2523 100 + i as u64,
2524 )
2525 .unwrap();
2526 bulk_parts.parts.push(create_bulk_part_wrapper(part));
2527 }
2528
2529 assert!(bulk_parts.should_merge_parts(DEFAULT_MERGE_THRESHOLD, usize::MAX));
2531
2532 let collected = bulk_parts.collect_parts_to_merge(
2534 DEFAULT_MERGE_THRESHOLD,
2535 DEFAULT_MAX_MERGE_GROUPS,
2536 usize::MAX,
2537 );
2538
2539 assert!(!collected.groups.is_empty());
2541
2542 for group in &collected.groups {
2544 assert!(!group.is_empty());
2545 }
2546
2547 let total_parts: usize = collected.groups.iter().map(|g| g.len()).sum();
2549 assert_eq!(16, total_parts);
2550 }
2551
2552 #[test]
2553 fn test_encoded_merge_candidate_size_limit() {
2554 assert!(BulkParts::is_merge_candidate_by_size(false, usize::MAX, 8));
2555 assert!(BulkParts::is_merge_candidate_by_size(true, 8, 8));
2556 assert!(!BulkParts::is_merge_candidate_by_size(true, 9, 8));
2557 }
2558
2559 #[test]
2560 fn test_encoded_part_batch_size_uses_largest_uncompressed_row_group() {
2561 const NUM_ROWS: usize = 14;
2562 const ROW_GROUP_SIZE: usize = 4;
2563
2564 let metadata = metadata_for_test();
2565 let timestamps = (0..NUM_ROWS as i64).collect::<Vec<_>>();
2566 let field_values = (0..NUM_ROWS)
2567 .map(|value| Some(value as f64))
2568 .collect::<Vec<_>>();
2569 let bulk_part =
2570 create_bulk_part_with_converter("key", 0, timestamps, field_values, 0).unwrap();
2571 let encoder = BulkPartEncoder::new(metadata, ROW_GROUP_SIZE).unwrap();
2572 let encoded_part = encoder.encode_part(&bulk_part).unwrap().unwrap();
2573 let max_row_group = encoded_part
2574 .metadata()
2575 .parquet_metadata
2576 .row_groups()
2577 .iter()
2578 .max_by_key(|row_group| {
2579 row_group
2580 .columns()
2581 .iter()
2582 .map(|column| column.uncompressed_size() as u64)
2583 .sum::<u64>()
2584 })
2585 .unwrap();
2586 let max_uncompressed_size = max_row_group
2587 .columns()
2588 .iter()
2589 .map(|column| column.uncompressed_size() as u64)
2590 .sum();
2591 let expected = (max_row_group.num_rows() as u64, max_uncompressed_size);
2592 let part = PartToMerge::Encoded {
2593 part: encoded_part,
2594 file_id: FileId::random(),
2595 };
2596
2597 assert_eq!(Some(expected), part.batch_size_statistic());
2598 }
2599
2600 #[test]
2601 fn test_bulk_memtable_ranges_with_multi_bulk_part() {
2602 let metadata = metadata_for_test();
2603 let merge_threshold = 8;
2604 let config = BulkMemtableConfig {
2605 merge_threshold,
2606 ..Default::default()
2607 };
2608 let memtable = BulkMemtable::new(
2609 2005,
2610 config,
2611 metadata.clone(),
2612 None,
2613 None,
2614 false,
2615 MergeMode::LastRow,
2616 );
2617 memtable.set_unordered_part_threshold(0);
2619
2620 for i in 0..merge_threshold {
2624 let part = create_bulk_part_with_converter(
2625 &format!("key_{}", i),
2626 i as u32,
2627 vec![1000 + i as i64 * 100, 2000 + i as i64 * 100],
2628 vec![Some(i as f64 * 10.0), Some(i as f64 * 10.0 + 1.0)],
2629 100 + i as u64,
2630 )
2631 .unwrap();
2632 memtable.write_bulk(part).unwrap();
2633 }
2634
2635 memtable.compact(false).unwrap();
2637
2638 let predicate_group = PredicateGroup::new(&metadata, &[]).unwrap();
2640 let ranges = memtable
2641 .ranges(
2642 None,
2643 RangesOptions::default().with_predicate(predicate_group),
2644 )
2645 .unwrap();
2646
2647 assert_eq!(1, ranges.ranges.len());
2648 let expected_rows = merge_threshold * 2; let total_rows: usize = ranges.ranges.values().map(|r| r.stats().num_rows()).sum();
2650 assert_eq!(expected_rows, total_rows);
2651
2652 let mut total_rows_read = 0;
2654 for (_range_id, range) in ranges.ranges.iter() {
2655 assert!(range.is_record_batch());
2656 let record_batch_iter = range.build_record_batch_iter(None, None).unwrap();
2657
2658 for batch_result in record_batch_iter {
2659 let batch = batch_result.unwrap();
2660 total_rows_read += batch.num_rows();
2661 }
2662 }
2663 assert_eq!(expected_rows, total_rows_read);
2664 }
2665
2666 #[test]
2667 fn test_multi_bulk_range_iter_builder_all_pruned() {
2668 let metadata = metadata_for_test();
2669 let merge_threshold = 8;
2670 let config = BulkMemtableConfig {
2671 merge_threshold,
2672 ..Default::default()
2673 };
2674 let memtable = BulkMemtable::new(
2675 2006,
2676 config,
2677 metadata.clone(),
2678 None,
2679 None,
2680 false,
2681 MergeMode::LastRow,
2682 );
2683 memtable.set_unordered_part_threshold(0);
2684
2685 for i in 0..merge_threshold {
2687 let part = create_bulk_part_with_converter(
2688 &format!("key_{}", i),
2689 i as u32,
2690 vec![1000 + i as i64 * 100, 2000 + i as i64 * 100],
2691 vec![Some(i as f64 * 10.0), Some(i as f64 * 10.0 + 1.0)],
2692 100 + i as u64,
2693 )
2694 .unwrap();
2695 memtable.write_bulk(part).unwrap();
2696 }
2697 memtable.compact(false).unwrap();
2698
2699 let filter = datafusion_expr::col("k0").eq(datafusion_expr::lit("nonexistent"));
2701 let predicate_group = PredicateGroup::new(&metadata, &[filter]).unwrap();
2702 let ranges = memtable
2703 .ranges(
2704 None,
2705 RangesOptions::default().with_predicate(predicate_group),
2706 )
2707 .unwrap();
2708
2709 for (_range_id, range) in ranges.ranges.iter() {
2712 assert!(range.is_record_batch());
2713 let record_batch_iter = range.build_record_batch_iter(None, None).unwrap();
2714 let total_rows: usize = record_batch_iter.map(|r| r.unwrap().num_rows()).sum();
2715 assert_eq!(0, total_rows);
2716 }
2717 }
2718}