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::Json2ExtensionType;
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 col_schema.with_extension_type(&Json2ExtensionType::default());
1734
1735 let col_meta_2 = ColumnMetadata {
1736 column_schema: col_schema,
1737 semantic_type: SemanticType::Field,
1738 column_id: 1,
1739 };
1740 let mut builder = RegionMetadataBuilder::new(RegionId::new(123, 789));
1741 builder
1742 .push_column_metadata(col_meta_1)
1743 .push_column_metadata(col_meta_2)
1744 .primary_key(vec![]);
1745 Arc::new(builder.build().unwrap())
1746 }
1747
1748 fn mock_bulk_part_with_json2(
1749 metadata: &RegionMetadataRef,
1750 timestamps: Vec<i64>,
1751 sequence: u64,
1752 ) -> Result<BulkPart> {
1753 let capacity = timestamps.len();
1754 let primary_key_codec = build_primary_key_codec(metadata);
1755 let json_type = JsonNativeType::Object(JsonObjectType::from([
1756 ("id".to_string(), JsonNativeType::i64()),
1757 (
1758 "payload".to_string(),
1759 JsonNativeType::Object(JsonObjectType::from([(
1760 "message".to_string(),
1761 JsonNativeType::String,
1762 )])),
1763 ),
1764 ]));
1765 let mut options = FlatSchemaOptions::from_encoding(metadata.primary_key_encoding);
1766 options
1767 .concretized_json_types
1768 .insert("data".to_string(), json_type.as_arrow_type());
1769 let schema = to_flat_sst_arrow_schema(metadata, &options);
1770
1771 let mut converter =
1772 BulkPartConverter::new(metadata, schema, capacity, primary_key_codec, true);
1773
1774 let rows = timestamps
1775 .into_iter()
1776 .map(|ts| {
1777 let val1 = api::v1::Value {
1778 value_data: Some(ValueData::TimestampMillisecondValue(ts)),
1779 };
1780 let value_data = ValueData::JsonValue(encode_json_value(JsonValue::from(json!({
1781 "id": ts,
1782 "payload": {
1783 "message": format!("row-{ts}"),
1784 },
1785 }))));
1786 let val2 = api::v1::Value {
1787 value_data: Some(value_data),
1788 };
1789 Row {
1790 values: vec![val1, val2],
1791 }
1792 })
1793 .collect();
1794
1795 let mutation = Mutation {
1796 op_type: 1,
1797 sequence,
1798 rows: Some(Rows {
1799 schema: region_metadata_to_row_schema(metadata),
1800 rows,
1801 }),
1802 write_hint: None,
1803 };
1804 let key_values = KeyValues::new(metadata.as_ref(), mutation).unwrap();
1805
1806 converter.append_key_values(&key_values)?;
1807 converter.convert()
1808 }
1809
1810 #[test]
1811 fn test_bulk_memtable_ranges_with_projection() {
1812 let metadata = metadata_for_test();
1813 let memtable = BulkMemtable::new(
1814 111,
1815 BulkMemtableConfig::default(),
1816 metadata.clone(),
1817 None,
1818 None,
1819 false,
1820 MergeMode::LastRow,
1821 );
1822
1823 let bulk_part = create_bulk_part_with_converter(
1824 "projection_test",
1825 5,
1826 vec![5000, 6000, 7000],
1827 vec![Some(50.0), Some(60.0), Some(70.0)],
1828 500,
1829 )
1830 .unwrap();
1831
1832 memtable.write_bulk(bulk_part).unwrap();
1833
1834 let projection = vec![4u32];
1835 let predicate_group = PredicateGroup::new(&metadata, &[]).unwrap();
1836 let ranges = memtable
1837 .ranges(
1838 Some(&projection),
1839 RangesOptions::default().with_predicate(predicate_group),
1840 )
1841 .unwrap();
1842
1843 assert_eq!(1, ranges.ranges.len());
1844 let range = ranges.ranges.get(&0).unwrap();
1845
1846 assert!(range.is_record_batch());
1847 let record_batch_iter = range.build_record_batch_iter(None, None).unwrap();
1848
1849 let mut total_rows = 0;
1850 for batch_result in record_batch_iter {
1851 let batch = batch_result.unwrap();
1852 assert!(batch.num_rows() > 0);
1853 assert_eq!(5, batch.num_columns());
1854 total_rows += batch.num_rows();
1855 }
1856 assert_eq!(3, total_rows);
1857 }
1858
1859 #[test]
1860 fn test_bulk_memtable_unsupported_operations() {
1861 let metadata = metadata_for_test();
1862 let memtable = BulkMemtable::new(
1863 111,
1864 BulkMemtableConfig::default(),
1865 metadata.clone(),
1866 None,
1867 None,
1868 false,
1869 MergeMode::LastRow,
1870 );
1871
1872 let key_values = build_key_values_with_ts_seq_values(
1873 &metadata,
1874 "test".to_string(),
1875 1,
1876 vec![1000].into_iter(),
1877 vec![Some(1.0)].into_iter(),
1878 1,
1879 );
1880
1881 let err = memtable.write(&key_values).unwrap_err();
1882 assert!(err.to_string().contains("not supported"));
1883
1884 let kv = key_values.iter().next().unwrap();
1885 let err = memtable.write_one(kv).unwrap_err();
1886 assert!(err.to_string().contains("not supported"));
1887 }
1888
1889 #[test]
1890 fn test_bulk_memtable_freeze() {
1891 let metadata = metadata_for_test();
1892 let memtable = BulkMemtable::new(
1893 222,
1894 BulkMemtableConfig::default(),
1895 metadata.clone(),
1896 None,
1897 None,
1898 false,
1899 MergeMode::LastRow,
1900 );
1901
1902 let bulk_part = create_bulk_part_with_converter(
1903 "freeze_test",
1904 10,
1905 vec![10000],
1906 vec![Some(100.0)],
1907 1000,
1908 )
1909 .unwrap();
1910
1911 memtable.write_bulk(bulk_part).unwrap();
1912 memtable.freeze().unwrap();
1913
1914 let stats_after_freeze = memtable.stats();
1915 assert_eq!(1, stats_after_freeze.num_rows);
1916 }
1917
1918 #[test]
1919 fn test_bulk_memtable_fork() {
1920 let metadata = metadata_for_test();
1921 let original_memtable = BulkMemtable::new(
1922 333,
1923 BulkMemtableConfig::default(),
1924 metadata.clone(),
1925 None,
1926 None,
1927 false,
1928 MergeMode::LastRow,
1929 );
1930
1931 let bulk_part =
1932 create_bulk_part_with_converter("fork_test", 15, vec![15000], vec![Some(150.0)], 1500)
1933 .unwrap();
1934
1935 original_memtable.write_bulk(bulk_part).unwrap();
1936
1937 let forked_memtable = original_memtable.fork(444, &metadata);
1938
1939 assert_eq!(forked_memtable.id(), 444);
1940 assert!(forked_memtable.is_empty());
1941 assert_eq!(0, forked_memtable.stats().num_rows);
1942
1943 assert!(!original_memtable.is_empty());
1944 assert_eq!(1, original_memtable.stats().num_rows);
1945 }
1946
1947 #[test]
1948 fn test_bulk_memtable_ranges_multiple_parts() {
1949 let metadata = metadata_for_test();
1950 let memtable = BulkMemtable::new(
1951 777,
1952 BulkMemtableConfig::default(),
1953 metadata.clone(),
1954 None,
1955 None,
1956 false,
1957 MergeMode::LastRow,
1958 );
1959 memtable.set_unordered_part_threshold(0);
1961
1962 let parts_data = vec![
1963 (
1964 "part1",
1965 1u32,
1966 vec![1000i64, 1100i64],
1967 vec![Some(10.0), Some(11.0)],
1968 100u64,
1969 ),
1970 (
1971 "part2",
1972 2u32,
1973 vec![2000i64, 2100i64],
1974 vec![Some(20.0), Some(21.0)],
1975 200u64,
1976 ),
1977 ("part3", 3u32, vec![3000i64], vec![Some(30.0)], 300u64),
1978 ];
1979
1980 for (k0, k1, timestamps, values, seq) in parts_data {
1981 let part = create_bulk_part_with_converter(k0, k1, timestamps, values, seq).unwrap();
1982 memtable.write_bulk(part).unwrap();
1983 }
1984
1985 let predicate_group = PredicateGroup::new(&metadata, &[]).unwrap();
1986 let ranges = memtable
1987 .ranges(
1988 None,
1989 RangesOptions::default().with_predicate(predicate_group),
1990 )
1991 .unwrap();
1992
1993 assert_eq!(3, ranges.ranges.len());
1994 let total_rows: usize = ranges.ranges.values().map(|r| r.stats().num_rows()).sum();
1995 assert_eq!(5, total_rows);
1996 assert_eq!(3, ranges.ranges.len());
1997
1998 for (range_id, range) in ranges.ranges.iter() {
1999 assert!(*range_id < 3);
2000 assert!(range.num_rows() > 0);
2001 assert!(range.is_record_batch());
2002 }
2003 }
2004
2005 #[test]
2006 fn test_bulk_memtable_ranges_with_sequence_filter() {
2007 let metadata = metadata_for_test();
2008 let memtable = BulkMemtable::new(
2009 888,
2010 BulkMemtableConfig::default(),
2011 metadata.clone(),
2012 None,
2013 None,
2014 false,
2015 MergeMode::LastRow,
2016 );
2017
2018 let part = create_bulk_part_with_converter(
2019 "seq_test",
2020 1,
2021 vec![1000, 2000, 3000],
2022 vec![Some(10.0), Some(20.0), Some(30.0)],
2023 500,
2024 )
2025 .unwrap();
2026
2027 memtable.write_bulk(part).unwrap();
2028
2029 let predicate_group = PredicateGroup::new(&metadata, &[]).unwrap();
2030 let sequence_filter = Some(SequenceRange::LtEq { max: 400 }); let ranges = memtable
2032 .ranges(
2033 None,
2034 RangesOptions::default()
2035 .with_predicate(predicate_group)
2036 .with_sequence(sequence_filter),
2037 )
2038 .unwrap();
2039
2040 assert_eq!(1, ranges.ranges.len());
2041 let range = ranges.ranges.get(&0).unwrap();
2042
2043 let mut record_batch_iter = range.build_record_batch_iter(None, None).unwrap();
2044 assert!(record_batch_iter.next().is_none());
2045 }
2046
2047 #[test]
2048 fn test_bulk_memtable_ranges_with_encoded_parts() {
2049 let metadata = metadata_for_test();
2050 let config = BulkMemtableConfig {
2051 merge_threshold: 8,
2052 ..Default::default()
2053 };
2054 let memtable = BulkMemtable::new(
2055 999,
2056 config,
2057 metadata.clone(),
2058 None,
2059 None,
2060 false,
2061 MergeMode::LastRow,
2062 );
2063 memtable.set_unordered_part_threshold(0);
2065
2066 for i in 0..10 {
2068 let part = create_bulk_part_with_converter(
2069 &format!("key_{}", i),
2070 i,
2071 vec![1000 + i as i64 * 100],
2072 vec![Some(i as f64 * 10.0)],
2073 100 + i as u64,
2074 )
2075 .unwrap();
2076 memtable.write_bulk(part).unwrap();
2077 }
2078
2079 memtable.compact(false).unwrap();
2080
2081 let predicate_group = PredicateGroup::new(&metadata, &[]).unwrap();
2082 let ranges = memtable
2083 .ranges(
2084 None,
2085 RangesOptions::default().with_predicate(predicate_group),
2086 )
2087 .unwrap();
2088
2089 assert_eq!(3, ranges.ranges.len());
2091 let total_rows: usize = ranges.ranges.values().map(|r| r.stats().num_rows()).sum();
2092 assert_eq!(10, total_rows);
2093
2094 for (_range_id, range) in ranges.ranges.iter() {
2095 assert!(range.num_rows() > 0);
2096 assert!(range.is_record_batch());
2097
2098 let record_batch_iter = range.build_record_batch_iter(None, None).unwrap();
2099 let mut total_rows = 0;
2100 for batch_result in record_batch_iter {
2101 let batch = batch_result.unwrap();
2102 total_rows += batch.num_rows();
2103 assert!(batch.num_rows() > 0);
2104 }
2105 assert_eq!(total_rows, range.num_rows());
2106 }
2107 }
2108
2109 #[test]
2110 fn test_bulk_memtable_unordered_part() {
2111 let metadata = metadata_for_test();
2112 let memtable = BulkMemtable::new(
2113 1001,
2114 BulkMemtableConfig::default(),
2115 metadata.clone(),
2116 None,
2117 None,
2118 false,
2119 MergeMode::LastRow,
2120 );
2121
2122 memtable.set_unordered_part_threshold(5);
2125 memtable.set_unordered_part_compact_threshold(10);
2127
2128 for i in 0..3 {
2130 let part = create_bulk_part_with_converter(
2131 &format!("key_{}", i),
2132 i,
2133 vec![1000 + i as i64 * 100, 1100 + i as i64 * 100],
2134 vec![Some(i as f64 * 10.0), Some(i as f64 * 10.0 + 1.0)],
2135 100 + i as u64,
2136 )
2137 .unwrap();
2138 assert_eq!(2, part.num_rows());
2139 memtable.write_bulk(part).unwrap();
2140 }
2141
2142 let stats = memtable.stats();
2144 assert_eq!(6, stats.num_rows);
2145
2146 for i in 3..5 {
2149 let part = create_bulk_part_with_converter(
2150 &format!("key_{}", i),
2151 i,
2152 vec![1000 + i as i64 * 100, 1100 + i as i64 * 100],
2153 vec![Some(i as f64 * 10.0), Some(i as f64 * 10.0 + 1.0)],
2154 100 + i as u64,
2155 )
2156 .unwrap();
2157 memtable.write_bulk(part).unwrap();
2158 }
2159
2160 let stats = memtable.stats();
2162 assert_eq!(10, stats.num_rows);
2163
2164 let predicate_group = PredicateGroup::new(&metadata, &[]).unwrap();
2166 let ranges = memtable
2167 .ranges(
2168 None,
2169 RangesOptions::default().with_predicate(predicate_group),
2170 )
2171 .unwrap();
2172
2173 assert!(!ranges.ranges.is_empty());
2175 let total_rows: usize = ranges.ranges.values().map(|r| r.stats().num_rows()).sum();
2176 assert_eq!(10, total_rows);
2177
2178 let mut total_rows_read = 0;
2180 for (_range_id, range) in ranges.ranges.iter() {
2181 assert!(range.is_record_batch());
2182 let record_batch_iter = range.build_record_batch_iter(None, None).unwrap();
2183
2184 for batch_result in record_batch_iter {
2185 let batch = batch_result.unwrap();
2186 total_rows_read += batch.num_rows();
2187 }
2188 }
2189 assert_eq!(10, total_rows_read);
2190 }
2191
2192 #[test]
2193 fn test_bulk_memtable_unordered_part_mixed_sizes() {
2194 let metadata = metadata_for_test();
2195 let memtable = BulkMemtable::new(
2196 1002,
2197 BulkMemtableConfig::default(),
2198 metadata.clone(),
2199 None,
2200 None,
2201 false,
2202 MergeMode::LastRow,
2203 );
2204
2205 memtable.set_unordered_part_threshold(4);
2207 memtable.set_unordered_part_compact_threshold(8);
2208
2209 for i in 0..2 {
2211 let part = create_bulk_part_with_converter(
2212 &format!("small_{}", i),
2213 i,
2214 vec![1000 + i as i64, 2000 + i as i64, 3000 + i as i64],
2215 vec![Some(i as f64), Some(i as f64 + 1.0), Some(i as f64 + 2.0)],
2216 10 + i as u64,
2217 )
2218 .unwrap();
2219 assert_eq!(3, part.num_rows());
2220 memtable.write_bulk(part).unwrap();
2221 }
2222
2223 let large_part = create_bulk_part_with_converter(
2225 "large_key",
2226 100,
2227 vec![5000, 6000, 7000, 8000, 9000],
2228 vec![
2229 Some(100.0),
2230 Some(101.0),
2231 Some(102.0),
2232 Some(103.0),
2233 Some(104.0),
2234 ],
2235 50,
2236 )
2237 .unwrap();
2238 assert_eq!(5, large_part.num_rows());
2239 memtable.write_bulk(large_part).unwrap();
2240
2241 let part = create_bulk_part_with_converter(
2243 "small_2",
2244 2,
2245 vec![4000, 4100],
2246 vec![Some(20.0), Some(21.0)],
2247 30,
2248 )
2249 .unwrap();
2250 memtable.write_bulk(part).unwrap();
2251
2252 let stats = memtable.stats();
2253 assert_eq!(13, stats.num_rows); let predicate_group = PredicateGroup::new(&metadata, &[]).unwrap();
2257 let ranges = memtable
2258 .ranges(
2259 None,
2260 RangesOptions::default().with_predicate(predicate_group),
2261 )
2262 .unwrap();
2263
2264 let total_rows: usize = ranges.ranges.values().map(|r| r.stats().num_rows()).sum();
2265 assert_eq!(13, total_rows);
2266
2267 let mut total_rows_read = 0;
2268 for (_range_id, range) in ranges.ranges.iter() {
2269 let record_batch_iter = range.build_record_batch_iter(None, None).unwrap();
2270 for batch_result in record_batch_iter {
2271 let batch = batch_result.unwrap();
2272 total_rows_read += batch.num_rows();
2273 }
2274 }
2275 assert_eq!(13, total_rows_read);
2276 }
2277
2278 #[test]
2279 fn test_bulk_memtable_unordered_part_byte_threshold() {
2280 let metadata = metadata_for_test();
2281 let first =
2282 create_bulk_part_with_converter("first", 1, vec![1000], vec![Some(1.0)], 1).unwrap();
2283 let second =
2284 create_bulk_part_with_converter("second", 2, vec![2000], vec![Some(2.0)], 2).unwrap();
2285 let threshold = first.estimated_size().max(second.estimated_size());
2286 let memtable = BulkMemtable::new(
2287 1004,
2288 BulkMemtableConfig {
2289 encode_bytes_threshold: threshold,
2290 ..Default::default()
2291 },
2292 metadata.clone(),
2293 None,
2294 None,
2295 false,
2296 MergeMode::LastRow,
2297 );
2298 memtable.set_unordered_part_compact_threshold(usize::MAX);
2299
2300 memtable.write_bulk(first).unwrap();
2302 assert_eq!(1, memtable.parts.read().unwrap().unordered_part.num_parts());
2303
2304 memtable.write_bulk(second).unwrap();
2306 let parts = memtable.parts.read().unwrap();
2307 assert!(parts.unordered_part.is_empty());
2308 assert_eq!(1, parts.parts.len());
2309 drop(parts);
2310
2311 let oversized =
2312 create_bulk_part_with_converter("oversized", 3, vec![3000], vec![Some(3.0)], 3)
2313 .unwrap();
2314 let oversized_memtable = BulkMemtable::new(
2315 1005,
2316 BulkMemtableConfig {
2317 encode_bytes_threshold: oversized.estimated_size() - 1,
2318 ..Default::default()
2319 },
2320 metadata,
2321 None,
2322 None,
2323 false,
2324 MergeMode::LastRow,
2325 );
2326 oversized_memtable.write_bulk(oversized).unwrap();
2327 let parts = oversized_memtable.parts.read().unwrap();
2328 assert!(parts.unordered_part.is_empty());
2329 assert_eq!(1, parts.parts.len());
2330 }
2331
2332 #[test]
2333 fn test_bulk_memtable_unordered_part_with_ranges() {
2334 let metadata = metadata_for_test();
2335 let memtable = BulkMemtable::new(
2336 1003,
2337 BulkMemtableConfig::default(),
2338 metadata.clone(),
2339 None,
2340 None,
2341 false,
2342 MergeMode::LastRow,
2343 );
2344
2345 memtable.set_unordered_part_threshold(3);
2347 memtable.set_unordered_part_compact_threshold(100); for i in 0..3 {
2351 let part = create_bulk_part_with_converter(
2352 &format!("key_{}", i),
2353 i,
2354 vec![1000 + i as i64 * 100],
2355 vec![Some(i as f64 * 10.0)],
2356 100 + i as u64,
2357 )
2358 .unwrap();
2359 assert_eq!(1, part.num_rows());
2360 memtable.write_bulk(part).unwrap();
2361 }
2362
2363 let stats = memtable.stats();
2364 assert_eq!(3, stats.num_rows);
2365
2366 let predicate_group = PredicateGroup::new(&metadata, &[]).unwrap();
2368 let ranges = memtable
2369 .ranges(
2370 None,
2371 RangesOptions::default().with_predicate(predicate_group),
2372 )
2373 .unwrap();
2374
2375 assert_eq!(1, ranges.ranges.len());
2377 let total_rows: usize = ranges.ranges.values().map(|r| r.stats().num_rows()).sum();
2378 assert_eq!(3, total_rows);
2379
2380 let range = ranges.ranges.get(&0).unwrap();
2382 let record_batch_iter = range.build_record_batch_iter(None, None).unwrap();
2383
2384 let mut total_rows = 0;
2385 for batch_result in record_batch_iter {
2386 let batch = batch_result.unwrap();
2387 total_rows += batch.num_rows();
2388 assert!(batch.num_rows() > 0);
2390 }
2391 assert_eq!(3, total_rows);
2392 }
2393
2394 fn create_bulk_part_wrapper(part: BulkPart) -> BulkPartWrapper {
2396 BulkPartWrapper {
2397 part: PartToMerge::Bulk {
2398 part,
2399 file_id: FileId::random(),
2400 },
2401 merging: false,
2402 }
2403 }
2404
2405 #[test]
2406 fn test_should_merge_parts_below_threshold() {
2407 let mut bulk_parts = BulkParts::default();
2408
2409 for i in 0..DEFAULT_MERGE_THRESHOLD - 1 {
2411 let part = create_bulk_part_with_converter(
2412 &format!("key_{}", i),
2413 i as u32,
2414 vec![1000 + i as i64 * 100],
2415 vec![Some(i as f64 * 10.0)],
2416 100 + i as u64,
2417 )
2418 .unwrap();
2419 bulk_parts.parts.push(create_bulk_part_wrapper(part));
2420 }
2421
2422 assert!(!bulk_parts.should_merge_parts(DEFAULT_MERGE_THRESHOLD, usize::MAX));
2424 }
2425
2426 #[test]
2427 fn test_should_merge_parts_at_threshold() {
2428 let mut bulk_parts = BulkParts::default();
2429 let merge_threshold = 8;
2430
2431 for i in 0..merge_threshold {
2433 let part = create_bulk_part_with_converter(
2434 &format!("key_{}", i),
2435 i as u32,
2436 vec![1000 + i as i64 * 100],
2437 vec![Some(i as f64 * 10.0)],
2438 100 + i as u64,
2439 )
2440 .unwrap();
2441 bulk_parts.parts.push(create_bulk_part_wrapper(part));
2442 }
2443
2444 assert!(bulk_parts.should_merge_parts(merge_threshold, usize::MAX));
2446 }
2447
2448 #[test]
2449 fn test_raw_parts_do_not_use_encoded_merge_limit() {
2450 let mut bulk_parts = BulkParts::default();
2451 let merge_threshold = 4;
2452 for i in 0..merge_threshold {
2453 let part = create_bulk_part_with_converter(
2454 &format!("key_{}", i),
2455 i as u32,
2456 vec![1000 + i as i64],
2457 vec![Some(i as f64)],
2458 100 + i as u64,
2459 )
2460 .unwrap();
2461 bulk_parts.parts.push(create_bulk_part_wrapper(part));
2462 }
2463
2464 let max_size = bulk_parts.parts[0].part.estimated_size();
2465 assert!(bulk_parts.should_merge_parts(merge_threshold, max_size));
2466 assert_eq!(
2467 1,
2468 bulk_parts
2469 .collect_parts_to_merge(merge_threshold, 1, max_size)
2470 .groups
2471 .len()
2472 );
2473 }
2474
2475 #[test]
2476 fn test_should_merge_parts_with_merging_flag() {
2477 let mut bulk_parts = BulkParts::default();
2478 let merge_threshold = 8;
2479
2480 for i in 0..10 {
2482 let part = create_bulk_part_with_converter(
2483 &format!("key_{}", i),
2484 i as u32,
2485 vec![1000 + i as i64 * 100],
2486 vec![Some(i as f64 * 10.0)],
2487 100 + i as u64,
2488 )
2489 .unwrap();
2490 bulk_parts.parts.push(create_bulk_part_wrapper(part));
2491 }
2492
2493 assert!(bulk_parts.should_merge_parts(merge_threshold, usize::MAX));
2495
2496 for wrapper in bulk_parts.parts.iter_mut().take(3) {
2498 wrapper.merging = true;
2499 }
2500
2501 assert!(!bulk_parts.should_merge_parts(merge_threshold, usize::MAX));
2503 }
2504
2505 #[test]
2506 fn test_collect_parts_to_merge_grouping() {
2507 let mut bulk_parts = BulkParts::default();
2508
2509 for i in 0..16 {
2511 let num_rows = (i % 4) + 1; let timestamps: Vec<i64> = (0..num_rows)
2513 .map(|j| 1000 + i as i64 * 100 + j as i64)
2514 .collect();
2515 let values: Vec<Option<f64>> =
2516 (0..num_rows).map(|j| Some((i * 10 + j) as f64)).collect();
2517 let part = create_bulk_part_with_converter(
2518 &format!("key_{}", i),
2519 i as u32,
2520 timestamps,
2521 values,
2522 100 + i as u64,
2523 )
2524 .unwrap();
2525 bulk_parts.parts.push(create_bulk_part_wrapper(part));
2526 }
2527
2528 assert!(bulk_parts.should_merge_parts(DEFAULT_MERGE_THRESHOLD, usize::MAX));
2530
2531 let collected = bulk_parts.collect_parts_to_merge(
2533 DEFAULT_MERGE_THRESHOLD,
2534 DEFAULT_MAX_MERGE_GROUPS,
2535 usize::MAX,
2536 );
2537
2538 assert!(!collected.groups.is_empty());
2540
2541 for group in &collected.groups {
2543 assert!(!group.is_empty());
2544 }
2545
2546 let total_parts: usize = collected.groups.iter().map(|g| g.len()).sum();
2548 assert_eq!(16, total_parts);
2549 }
2550
2551 #[test]
2552 fn test_encoded_merge_candidate_size_limit() {
2553 assert!(BulkParts::is_merge_candidate_by_size(false, usize::MAX, 8));
2554 assert!(BulkParts::is_merge_candidate_by_size(true, 8, 8));
2555 assert!(!BulkParts::is_merge_candidate_by_size(true, 9, 8));
2556 }
2557
2558 #[test]
2559 fn test_encoded_part_batch_size_uses_largest_uncompressed_row_group() {
2560 const NUM_ROWS: usize = 14;
2561 const ROW_GROUP_SIZE: usize = 4;
2562
2563 let metadata = metadata_for_test();
2564 let timestamps = (0..NUM_ROWS as i64).collect::<Vec<_>>();
2565 let field_values = (0..NUM_ROWS)
2566 .map(|value| Some(value as f64))
2567 .collect::<Vec<_>>();
2568 let bulk_part =
2569 create_bulk_part_with_converter("key", 0, timestamps, field_values, 0).unwrap();
2570 let encoder = BulkPartEncoder::new(metadata, ROW_GROUP_SIZE).unwrap();
2571 let encoded_part = encoder.encode_part(&bulk_part).unwrap().unwrap();
2572 let max_row_group = encoded_part
2573 .metadata()
2574 .parquet_metadata
2575 .row_groups()
2576 .iter()
2577 .max_by_key(|row_group| {
2578 row_group
2579 .columns()
2580 .iter()
2581 .map(|column| column.uncompressed_size() as u64)
2582 .sum::<u64>()
2583 })
2584 .unwrap();
2585 let max_uncompressed_size = max_row_group
2586 .columns()
2587 .iter()
2588 .map(|column| column.uncompressed_size() as u64)
2589 .sum();
2590 let expected = (max_row_group.num_rows() as u64, max_uncompressed_size);
2591 let part = PartToMerge::Encoded {
2592 part: encoded_part,
2593 file_id: FileId::random(),
2594 };
2595
2596 assert_eq!(Some(expected), part.batch_size_statistic());
2597 }
2598
2599 #[test]
2600 fn test_bulk_memtable_ranges_with_multi_bulk_part() {
2601 let metadata = metadata_for_test();
2602 let merge_threshold = 8;
2603 let config = BulkMemtableConfig {
2604 merge_threshold,
2605 ..Default::default()
2606 };
2607 let memtable = BulkMemtable::new(
2608 2005,
2609 config,
2610 metadata.clone(),
2611 None,
2612 None,
2613 false,
2614 MergeMode::LastRow,
2615 );
2616 memtable.set_unordered_part_threshold(0);
2618
2619 for i in 0..merge_threshold {
2623 let part = create_bulk_part_with_converter(
2624 &format!("key_{}", i),
2625 i as u32,
2626 vec![1000 + i as i64 * 100, 2000 + i as i64 * 100],
2627 vec![Some(i as f64 * 10.0), Some(i as f64 * 10.0 + 1.0)],
2628 100 + i as u64,
2629 )
2630 .unwrap();
2631 memtable.write_bulk(part).unwrap();
2632 }
2633
2634 memtable.compact(false).unwrap();
2636
2637 let predicate_group = PredicateGroup::new(&metadata, &[]).unwrap();
2639 let ranges = memtable
2640 .ranges(
2641 None,
2642 RangesOptions::default().with_predicate(predicate_group),
2643 )
2644 .unwrap();
2645
2646 assert_eq!(1, ranges.ranges.len());
2647 let expected_rows = merge_threshold * 2; let total_rows: usize = ranges.ranges.values().map(|r| r.stats().num_rows()).sum();
2649 assert_eq!(expected_rows, total_rows);
2650
2651 let mut total_rows_read = 0;
2653 for (_range_id, range) in ranges.ranges.iter() {
2654 assert!(range.is_record_batch());
2655 let record_batch_iter = range.build_record_batch_iter(None, None).unwrap();
2656
2657 for batch_result in record_batch_iter {
2658 let batch = batch_result.unwrap();
2659 total_rows_read += batch.num_rows();
2660 }
2661 }
2662 assert_eq!(expected_rows, total_rows_read);
2663 }
2664
2665 #[test]
2666 fn test_multi_bulk_range_iter_builder_all_pruned() {
2667 let metadata = metadata_for_test();
2668 let merge_threshold = 8;
2669 let config = BulkMemtableConfig {
2670 merge_threshold,
2671 ..Default::default()
2672 };
2673 let memtable = BulkMemtable::new(
2674 2006,
2675 config,
2676 metadata.clone(),
2677 None,
2678 None,
2679 false,
2680 MergeMode::LastRow,
2681 );
2682 memtable.set_unordered_part_threshold(0);
2683
2684 for i in 0..merge_threshold {
2686 let part = create_bulk_part_with_converter(
2687 &format!("key_{}", i),
2688 i as u32,
2689 vec![1000 + i as i64 * 100, 2000 + i as i64 * 100],
2690 vec![Some(i as f64 * 10.0), Some(i as f64 * 10.0 + 1.0)],
2691 100 + i as u64,
2692 )
2693 .unwrap();
2694 memtable.write_bulk(part).unwrap();
2695 }
2696 memtable.compact(false).unwrap();
2697
2698 let filter = datafusion_expr::col("k0").eq(datafusion_expr::lit("nonexistent"));
2700 let predicate_group = PredicateGroup::new(&metadata, &[filter]).unwrap();
2701 let ranges = memtable
2702 .ranges(
2703 None,
2704 RangesOptions::default().with_predicate(predicate_group),
2705 )
2706 .unwrap();
2707
2708 for (_range_id, range) in ranges.ranges.iter() {
2711 assert!(range.is_record_batch());
2712 let record_batch_iter = range.build_record_batch_iter(None, None).unwrap();
2713 let total_rows: usize = record_batch_iter.map(|r| r.unwrap().num_rows()).sum();
2714 assert_eq!(0, total_rows);
2715 }
2716 }
2717}