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