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