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