Skip to main content

mito2/memtable/
bulk.rs

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