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