Skip to main content

mito2/memtable/bulk/
part.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//! Bulk part encoder/decoder.
16
17use std::collections::{HashMap, HashSet};
18use std::sync::Arc;
19use std::time::{Duration, Instant};
20
21use api::helper::{ColumnDataTypeWrapper, to_grpc_value};
22use api::v1::bulk_wal_entry::Body;
23use api::v1::{ArrowIpc, BulkWalEntry, Mutation, OpType};
24use bytes::Bytes;
25use common_grpc::flight::{FlightDecoder, FlightEncoder, FlightMessage};
26use common_recordbatch::DfRecordBatch as RecordBatch;
27use common_time::Timestamp;
28use datafusion_common::Column;
29use datafusion_common::pruning::PruningStatistics;
30use datafusion_expr::utils::expr_to_columns;
31use datatypes::arrow;
32use datatypes::arrow::array::{
33    Array, ArrayRef, BinaryArray, BooleanArray, DictionaryArray, StringDictionaryBuilder,
34    TimestampMicrosecondArray, TimestampMillisecondArray, TimestampNanosecondArray,
35    TimestampSecondArray, UInt8Array, UInt32Array, UInt64Array,
36};
37use datatypes::arrow::compute::{SortColumn, SortOptions, concat_batches};
38use datatypes::arrow::datatypes::{
39    DataType as ArrowDataType, Field, Schema, SchemaRef, TimeUnit, UInt32Type,
40};
41use datatypes::data_type::DataType;
42use datatypes::extension::json::is_json2_extension_type;
43use datatypes::prelude::{MutableVector, Vector};
44use datatypes::value::ValueRef;
45use datatypes::vectors::Helper;
46use mito_codec::key_values::{KeyValue, KeyValues};
47use mito_codec::row_converter::{PrimaryKeyCodec, SortField, build_primary_key_codec_with_fields};
48use parquet::arrow::ArrowWriter;
49use parquet::basic::{Compression, ZstdLevel};
50use parquet::file::metadata::ParquetMetaData;
51use parquet::file::properties::WriterProperties;
52use smallvec::SmallVec;
53use snafu::{OptionExt, ResultExt};
54use store_api::codec::PrimaryKeyEncoding;
55use store_api::metadata::{RegionMetadata, RegionMetadataRef};
56use store_api::storage::consts::PRIMARY_KEY_COLUMN_NAME;
57use store_api::storage::{ColumnId, FileId, SequenceNumber, SequenceRange};
58
59use crate::error::{
60    self, ColumnNotFoundSnafu, ComputeArrowSnafu, CreateDefaultSnafu, DataTypeMismatchSnafu,
61    EncodeMemtableSnafu, EncodeSnafu, InvalidMetadataSnafu, InvalidRequestSnafu,
62    NewRecordBatchSnafu, Result,
63};
64use crate::memtable::bulk::context::{BulkIterContext, BulkIterContextRef};
65use crate::memtable::bulk::json_align::Json2Aligner;
66use crate::memtable::bulk::part_reader::EncodedBulkPartIter;
67use crate::memtable::time_series::{ValueBuilder, Values};
68use crate::memtable::{BoxedRecordBatchIterator, MemScanMetrics, MemtableStats};
69use crate::sst::SeriesEstimator;
70use crate::sst::index::IndexOutput;
71use crate::sst::parquet::flat_format::primary_key_column_index;
72use crate::sst::parquet::format::{PrimaryKeyArray, PrimaryKeyArrayBuilder};
73use crate::sst::parquet::{PARQUET_METADATA_KEY, SstInfo};
74
75const INIT_DICT_VALUE_CAPACITY: usize = 8;
76
77/// A raw bulk part in the memtable.
78#[derive(Clone)]
79pub struct BulkPart {
80    pub batch: RecordBatch,
81    pub max_timestamp: i64,
82    pub min_timestamp: i64,
83    pub sequence: u64,
84    pub timestamp_index: usize,
85    pub raw_data: Option<ArrowIpc>,
86}
87
88impl TryFrom<BulkWalEntry> for BulkPart {
89    type Error = error::Error;
90
91    fn try_from(value: BulkWalEntry) -> std::result::Result<Self, Self::Error> {
92        match value.body.expect("Entry payload should be present") {
93            Body::ArrowIpc(ipc) => {
94                let mut decoder = FlightDecoder::try_from_schema_bytes(&ipc.schema)
95                    .context(error::ConvertBulkWalEntrySnafu)?;
96                let batch = decoder
97                    .try_decode_record_batch(&ipc.data_header, &ipc.payload)
98                    .context(error::ConvertBulkWalEntrySnafu)?;
99                Ok(Self {
100                    batch,
101                    max_timestamp: value.max_ts,
102                    min_timestamp: value.min_ts,
103                    sequence: value.sequence,
104                    timestamp_index: value.timestamp_index as usize,
105                    raw_data: Some(ipc),
106                })
107            }
108        }
109    }
110}
111
112impl TryFrom<&BulkPart> for BulkWalEntry {
113    type Error = error::Error;
114
115    fn try_from(value: &BulkPart) -> Result<Self> {
116        if let Some(ipc) = &value.raw_data {
117            Ok(BulkWalEntry {
118                sequence: value.sequence,
119                max_ts: value.max_timestamp,
120                min_ts: value.min_timestamp,
121                timestamp_index: value.timestamp_index as u32,
122                body: Some(Body::ArrowIpc(ipc.clone())),
123            })
124        } else {
125            let mut encoder = FlightEncoder::default();
126            let schema_bytes = encoder
127                .encode_schema(value.batch.schema().as_ref())
128                .data_header;
129            let [rb_data] = encoder
130                .encode(FlightMessage::RecordBatch(value.batch.clone()))
131                .try_into()
132                .map_err(|_| {
133                    error::UnsupportedOperationSnafu {
134                        err_msg: "create BulkWalEntry from RecordBatch with dictionary arrays",
135                    }
136                    .build()
137                })?;
138            Ok(BulkWalEntry {
139                sequence: value.sequence,
140                max_ts: value.max_timestamp,
141                min_ts: value.min_timestamp,
142                timestamp_index: value.timestamp_index as u32,
143                body: Some(Body::ArrowIpc(ArrowIpc {
144                    schema: schema_bytes,
145                    data_header: rb_data.data_header,
146                    payload: rb_data.data_body,
147                })),
148            })
149        }
150    }
151}
152
153impl BulkPart {
154    pub(crate) fn schema(&self) -> SchemaRef {
155        self.batch.schema()
156    }
157
158    pub(crate) fn estimated_size(&self) -> usize {
159        record_batch_estimated_size(&self.batch)
160    }
161
162    /// Returns the estimated series count in this BulkPart.
163    /// This is calculated from the dictionary values count of the PrimaryKeyArray.
164    pub fn estimated_series_count(&self) -> usize {
165        let pk_column_idx = primary_key_column_index(self.batch.num_columns());
166        let pk_column = self.batch.column(pk_column_idx);
167        if let Some(dict_array) = pk_column.as_any().downcast_ref::<PrimaryKeyArray>() {
168            dict_array.values().len()
169        } else {
170            0
171        }
172    }
173
174    /// Creates MemtableStats from this BulkPart.
175    pub fn to_memtable_stats(&self, region_metadata: &RegionMetadataRef) -> MemtableStats {
176        let ts_type = region_metadata.time_index_type();
177        let min_ts = ts_type.create_timestamp(self.min_timestamp);
178        let max_ts = ts_type.create_timestamp(self.max_timestamp);
179
180        MemtableStats {
181            estimated_bytes: self.estimated_size(),
182            time_range: Some((min_ts, max_ts)),
183            num_rows: self.num_rows(),
184            num_ranges: 1,
185            max_sequence: self.sequence,
186            series_count: self.estimated_series_count(),
187        }
188    }
189
190    /// Fills missing columns in the BulkPart batch with default values.
191    ///
192    /// This function checks if the batch schema matches the region metadata schema,
193    /// and if there are missing columns, it fills them with default values (or null
194    /// for nullable columns).
195    ///
196    /// # Arguments
197    ///
198    /// * `region_metadata` - The region metadata containing the expected schema
199    pub fn fill_missing_columns(&mut self, region_metadata: &RegionMetadata) -> Result<()> {
200        // Builds a map of existing columns in the batch
201        let batch_schema = self.batch.schema();
202        let batch_columns: HashSet<_> = batch_schema
203            .fields()
204            .iter()
205            .map(|f| f.name().as_str())
206            .collect();
207
208        // Finds columns that need to be filled
209        let mut columns_to_fill = Vec::new();
210        for column_meta in &region_metadata.column_metadatas {
211            // TODO(yingwen): Returns error if it is impure default after we support filling
212            // bulk insert request in the frontend
213            if !batch_columns.contains(column_meta.column_schema.name.as_str()) {
214                columns_to_fill.push(column_meta);
215            }
216        }
217
218        if columns_to_fill.is_empty() {
219            return Ok(());
220        }
221
222        let num_rows = self.batch.num_rows();
223
224        let mut new_columns = Vec::new();
225        let mut new_fields = Vec::new();
226
227        // First, adds all existing columns
228        new_fields.extend(batch_schema.fields().iter().cloned());
229        new_columns.extend_from_slice(self.batch.columns());
230
231        let region_id = region_metadata.region_id;
232        // Then adds the missing columns with default values
233        for column_meta in columns_to_fill {
234            let default_vector = column_meta
235                .column_schema
236                .create_default_vector(num_rows)
237                .context(CreateDefaultSnafu {
238                    region_id,
239                    column: &column_meta.column_schema.name,
240                })?
241                .with_context(|| InvalidRequestSnafu {
242                    region_id,
243                    reason: format!(
244                        "column {} does not have default value",
245                        column_meta.column_schema.name
246                    ),
247                })?;
248            let arrow_array = default_vector.to_arrow_array();
249            column_meta.column_schema.data_type.as_arrow_type();
250
251            new_fields.push(Arc::new(Field::new(
252                column_meta.column_schema.name.clone(),
253                column_meta.column_schema.data_type.as_arrow_type(),
254                column_meta.column_schema.is_nullable(),
255            )));
256            new_columns.push(arrow_array);
257        }
258
259        // Create a new schema and batch with the filled columns
260        let new_schema = Arc::new(Schema::new(new_fields));
261        let new_batch =
262            RecordBatch::try_new(new_schema, new_columns).context(NewRecordBatchSnafu)?;
263
264        // Update the batch
265        self.batch = new_batch;
266
267        // The raw data no longer matches the batch after filling missing
268        // columns. Clears it so the WAL entry is re-encoded from the filled
269        // batch, otherwise replaying the entry restores a batch without the
270        // filled columns and the memtable rejects it.
271        self.raw_data = None;
272
273        Ok(())
274    }
275
276    /// Converts [BulkPart] to [Mutation] for fallback `write_bulk` implementation.
277    pub(crate) fn to_mutation(&self, region_metadata: &RegionMetadataRef) -> Result<Mutation> {
278        let vectors = region_metadata
279            .schema
280            .column_schemas()
281            .iter()
282            .map(|col| match self.batch.column_by_name(&col.name) {
283                None => Ok(None),
284                Some(col) => Helper::try_into_vector(col).map(Some),
285            })
286            .collect::<datatypes::error::Result<Vec<_>>>()
287            .context(error::ComputeVectorSnafu)?;
288
289        let rows = (0..self.num_rows())
290            .map(|row_idx| {
291                let values = (0..self.batch.num_columns())
292                    .map(|col_idx| {
293                        if let Some(v) = &vectors[col_idx] {
294                            to_grpc_value(v.get(row_idx))
295                        } else {
296                            api::v1::Value { value_data: None }
297                        }
298                    })
299                    .collect::<Vec<_>>();
300                api::v1::Row { values }
301            })
302            .collect::<Vec<_>>();
303
304        let schema = region_metadata
305            .column_metadatas
306            .iter()
307            .map(|c| {
308                let data_type_wrapper =
309                    ColumnDataTypeWrapper::try_from(c.column_schema.data_type.clone())?;
310                Ok(api::v1::ColumnSchema {
311                    column_name: c.column_schema.name.clone(),
312                    datatype: data_type_wrapper.datatype() as i32,
313                    semantic_type: c.semantic_type as i32,
314                    ..Default::default()
315                })
316            })
317            .collect::<api::error::Result<Vec<_>>>()
318            .context(error::ConvertColumnDataTypeSnafu {
319                reason: "failed to convert region metadata to column schema",
320            })?;
321
322        let rows = api::v1::Rows { schema, rows };
323
324        Ok(Mutation {
325            op_type: OpType::Put as i32,
326            sequence: self.sequence,
327            rows: Some(rows),
328            write_hint: None,
329        })
330    }
331
332    pub fn timestamps(&self) -> &ArrayRef {
333        self.batch.column(self.timestamp_index)
334    }
335
336    pub fn num_rows(&self) -> usize {
337        self.batch.num_rows()
338    }
339}
340
341/// A collection of small unordered bulk parts.
342/// Used to batch small parts together before merging them into a sorted part.
343pub struct UnorderedPart {
344    /// Small bulk parts that haven't been sorted yet.
345    parts: Vec<BulkPart>,
346    /// Total number of rows across all parts.
347    total_rows: usize,
348    /// Total estimated uncompressed bytes across all parts.
349    total_bytes: usize,
350    /// Minimum timestamp across all parts.
351    min_timestamp: i64,
352    /// Maximum timestamp across all parts.
353    max_timestamp: i64,
354    /// Maximum sequence number across all parts.
355    max_sequence: u64,
356    /// Row count threshold for accepting parts (default: 1024).
357    threshold: usize,
358    /// Row count threshold for compacting (default: 4096).
359    compact_threshold: usize,
360}
361
362impl Default for UnorderedPart {
363    fn default() -> Self {
364        Self::new()
365    }
366}
367
368impl UnorderedPart {
369    /// Creates a new empty UnorderedPart.
370    pub fn new() -> Self {
371        Self {
372            parts: Vec::new(),
373            total_rows: 0,
374            total_bytes: 0,
375            min_timestamp: i64::MAX,
376            max_timestamp: i64::MIN,
377            max_sequence: 0,
378            threshold: 1024,
379            compact_threshold: 4096,
380        }
381    }
382
383    /// Sets the threshold for accepting parts into unordered_part.
384    pub fn set_threshold(&mut self, threshold: usize) {
385        self.threshold = threshold;
386    }
387
388    /// Sets the threshold for compacting unordered_part.
389    pub fn set_compact_threshold(&mut self, compact_threshold: usize) {
390        self.compact_threshold = compact_threshold;
391    }
392
393    /// Returns the threshold for accepting parts.
394    pub fn threshold(&self) -> usize {
395        self.threshold
396    }
397
398    /// Returns the compact threshold.
399    pub fn compact_threshold(&self) -> usize {
400        self.compact_threshold
401    }
402
403    /// Returns true if this part should accept the given row count.
404    pub fn should_accept(&self, num_rows: usize) -> bool {
405        num_rows < self.threshold
406    }
407
408    /// Returns true if this part should be compacted by row count.
409    pub fn should_compact(&self) -> bool {
410        self.total_rows >= self.compact_threshold
411    }
412
413    /// Returns the total estimated uncompressed bytes across all parts.
414    pub(super) fn estimated_bytes(&self) -> usize {
415        self.total_bytes
416    }
417
418    /// Adds a BulkPart to this unordered collection.
419    pub fn push(&mut self, part: BulkPart) {
420        self.total_rows += part.num_rows();
421        self.total_bytes = self.total_bytes.saturating_add(part.estimated_size());
422        self.min_timestamp = self.min_timestamp.min(part.min_timestamp);
423        self.max_timestamp = self.max_timestamp.max(part.max_timestamp);
424        self.max_sequence = self.max_sequence.max(part.sequence);
425        self.parts.push(part);
426    }
427
428    /// Returns the total number of rows across all parts.
429    pub fn num_rows(&self) -> usize {
430        self.total_rows
431    }
432
433    /// Returns true if there are no parts.
434    pub fn is_empty(&self) -> bool {
435        self.parts.is_empty()
436    }
437
438    /// Returns the number of parts in this collection.
439    pub fn num_parts(&self) -> usize {
440        self.parts.len()
441    }
442
443    /// Concatenates and sorts all parts into a single RecordBatch.
444    /// Returns None if the collection is empty.
445    pub fn concat_and_sort(&self) -> Result<Option<RecordBatch>> {
446        if self.parts.is_empty() {
447            return Ok(None);
448        }
449
450        if self.parts.len() == 1 {
451            // If there's only one part, return its batch directly
452            return Ok(Some(self.parts[0].batch.clone()));
453        }
454
455        // Get the schema from the first part
456        let schema = self.parts[0].batch.schema();
457        let concatenated = if schema.fields().iter().any(is_json2_extension_type) {
458            let aligner = Json2Aligner::try_new(self.parts.iter().map(|part| part.batch.schema()))?;
459            let aligned_batches =
460                aligner.align_batches(self.parts.iter().map(|part| part.batch.clone()))?;
461            concat_batches(aligner.schema(), &aligned_batches).context(ComputeArrowSnafu)?
462        } else {
463            concat_batches(&schema, self.parts.iter().map(|x| &x.batch))
464                .context(ComputeArrowSnafu)?
465        };
466
467        // Sort the concatenated batch
468        let sorted_batch = sort_primary_key_record_batch(&concatenated)?;
469
470        Ok(Some(sorted_batch))
471    }
472
473    /// Converts all parts into a single sorted BulkPart.
474    /// Returns None if the collection is empty.
475    pub fn to_bulk_part(&self) -> Result<Option<BulkPart>> {
476        let Some(sorted_batch) = self.concat_and_sort()? else {
477            return Ok(None);
478        };
479
480        let timestamp_index = self.parts[0].timestamp_index;
481
482        Ok(Some(BulkPart {
483            batch: sorted_batch,
484            max_timestamp: self.max_timestamp,
485            min_timestamp: self.min_timestamp,
486            sequence: self.max_sequence,
487            timestamp_index,
488            raw_data: None,
489        }))
490    }
491
492    /// Clears all parts from this collection.
493    pub fn clear(&mut self) {
494        self.parts.clear();
495        self.total_rows = 0;
496        self.total_bytes = 0;
497        self.min_timestamp = i64::MAX;
498        self.max_timestamp = i64::MIN;
499        self.max_sequence = 0;
500    }
501}
502
503/// More accurate estimation of the size of a record batch.
504pub fn record_batch_estimated_size(batch: &RecordBatch) -> usize {
505    batch
506        .columns()
507        .iter()
508        // If can not get slice memory size, assume 0 here.
509        .map(|c| c.to_data().get_slice_memory_size().unwrap_or(0))
510        .sum()
511}
512
513/// Primary key column builder for handling strings specially.
514enum PrimaryKeyColumnBuilder {
515    /// String dictionary builder for string types.
516    StringDict(StringDictionaryBuilder<UInt32Type>),
517    /// Generic mutable vector for other types.
518    Vector(Box<dyn MutableVector>),
519}
520
521impl PrimaryKeyColumnBuilder {
522    /// Appends a value to the builder.
523    fn push_value_ref(&mut self, value: ValueRef) -> Result<()> {
524        match self {
525            PrimaryKeyColumnBuilder::StringDict(builder) => {
526                if let Some(s) = value.try_into_string().context(DataTypeMismatchSnafu)? {
527                    // We know the value is a string.
528                    builder.append_value(s);
529                } else {
530                    builder.append_null();
531                }
532            }
533            PrimaryKeyColumnBuilder::Vector(builder) => {
534                builder.push_value_ref(&value);
535            }
536        }
537        Ok(())
538    }
539
540    /// Converts the builder to an ArrayRef.
541    fn into_arrow_array(self) -> ArrayRef {
542        match self {
543            PrimaryKeyColumnBuilder::StringDict(mut builder) => Arc::new(builder.finish()),
544            PrimaryKeyColumnBuilder::Vector(mut builder) => builder.to_vector().to_arrow_array(),
545        }
546    }
547}
548
549/// Converter that converts structs into [BulkPart].
550pub struct BulkPartConverter {
551    /// Schema of the converted batch.
552    schema: SchemaRef,
553    /// Primary key codec for encoding keys
554    primary_key_codec: Arc<dyn PrimaryKeyCodec>,
555    /// Buffer for encoding primary key.
556    key_buf: Vec<u8>,
557    /// Primary key array builder.
558    key_array_builder: PrimaryKeyArrayBuilder,
559    /// Builders for non-primary key columns.
560    value_builder: ValueBuilder,
561    /// Builders for individual primary key columns.
562    /// The order of builders is the same as the order of primary key columns in the region metadata.
563    primary_key_column_builders: Vec<PrimaryKeyColumnBuilder>,
564
565    /// Max timestamp value.
566    max_ts: i64,
567    /// Min timestamp value.
568    min_ts: i64,
569    /// Max sequence number.
570    max_sequence: SequenceNumber,
571}
572
573impl BulkPartConverter {
574    /// Creates a new converter.
575    ///
576    /// If `store_primary_key_columns` is true and the encoding is not sparse encoding, it
577    /// stores primary key columns in arrays additionally.
578    pub fn new(
579        region_metadata: &RegionMetadataRef,
580        schema: SchemaRef,
581        capacity: usize,
582        primary_key_codec: Arc<dyn PrimaryKeyCodec>,
583        store_primary_key_columns: bool,
584    ) -> Self {
585        debug_assert_eq!(
586            region_metadata.primary_key_encoding,
587            primary_key_codec.encoding()
588        );
589
590        let primary_key_column_builders = if store_primary_key_columns
591            && region_metadata.primary_key_encoding != PrimaryKeyEncoding::Sparse
592        {
593            new_primary_key_column_builders(region_metadata, capacity)
594        } else {
595            Vec::new()
596        };
597
598        Self {
599            schema,
600            primary_key_codec,
601            key_buf: Vec::new(),
602            key_array_builder: PrimaryKeyArrayBuilder::new(),
603            value_builder: ValueBuilder::new(region_metadata, capacity),
604            primary_key_column_builders,
605            min_ts: i64::MAX,
606            max_ts: i64::MIN,
607            max_sequence: SequenceNumber::MIN,
608        }
609    }
610
611    /// Appends a [KeyValues] into the converter.
612    pub fn append_key_values(&mut self, key_values: &KeyValues) -> Result<()> {
613        for kv in key_values.iter() {
614            self.append_key_value(&kv)?;
615        }
616
617        Ok(())
618    }
619
620    /// Appends a [KeyValue] to builders.
621    ///
622    /// If the primary key uses sparse encoding, callers must encoded the primary key in the [KeyValue].
623    fn append_key_value(&mut self, kv: &KeyValue) -> Result<()> {
624        // Handles primary key based on encoding type
625        if self.primary_key_codec.encoding() == PrimaryKeyEncoding::Sparse {
626            // For sparse encoding, the primary key is already encoded in the KeyValue
627            // Gets the first (and only) primary key value which contains the encoded key
628            let mut primary_keys = kv.primary_keys();
629            if let Some(encoded) = primary_keys
630                .next()
631                .context(ColumnNotFoundSnafu {
632                    column: PRIMARY_KEY_COLUMN_NAME,
633                })?
634                .try_into_binary()
635                .context(DataTypeMismatchSnafu)?
636            {
637                self.key_array_builder
638                    .append(encoded)
639                    .context(ComputeArrowSnafu)?;
640            } else {
641                self.key_array_builder
642                    .append("")
643                    .context(ComputeArrowSnafu)?;
644            }
645        } else {
646            // For dense encoding, we need to encode the primary key columns
647            self.key_buf.clear();
648            self.primary_key_codec
649                .encode_key_value(kv, &mut self.key_buf)
650                .context(EncodeSnafu)?;
651            self.key_array_builder
652                .append(&self.key_buf)
653                .context(ComputeArrowSnafu)?;
654        };
655
656        // If storing primary key columns, append values to individual builders
657        if !self.primary_key_column_builders.is_empty() {
658            for (builder, pk_value) in self
659                .primary_key_column_builders
660                .iter_mut()
661                .zip(kv.primary_keys())
662            {
663                builder.push_value_ref(pk_value)?;
664            }
665        }
666
667        // Pushes other columns.
668        self.value_builder.push(
669            kv.timestamp(),
670            kv.sequence(),
671            kv.op_type() as u8,
672            kv.fields(),
673        );
674
675        // Updates statistics
676        // Safety: timestamp of kv must be both present and a valid timestamp value.
677        let ts = kv
678            .timestamp()
679            .try_into_timestamp()
680            .unwrap()
681            .unwrap()
682            .value();
683        self.min_ts = self.min_ts.min(ts);
684        self.max_ts = self.max_ts.max(ts);
685        self.max_sequence = self.max_sequence.max(kv.sequence());
686
687        Ok(())
688    }
689
690    /// Converts buffered content into a [BulkPart].
691    ///
692    /// It sorts the record batch by (primary key, timestamp, sequence desc).
693    pub fn convert(mut self) -> Result<BulkPart> {
694        let values = Values::from(self.value_builder);
695        let mut columns =
696            Vec::with_capacity(4 + values.fields.len() + self.primary_key_column_builders.len());
697
698        // Build primary key column arrays if enabled.
699        for builder in self.primary_key_column_builders {
700            columns.push(builder.into_arrow_array());
701        }
702        // Then fields columns.
703        columns.extend(values.fields.iter().map(|field| field.to_arrow_array()));
704        // Time index.
705        let timestamp_index = columns.len();
706        columns.push(values.timestamp.to_arrow_array());
707        // Primary key.
708        let pk_array = self.key_array_builder.finish();
709        columns.push(Arc::new(pk_array));
710        // Sequence and op type.
711        columns.push(values.sequence.to_arrow_array());
712        columns.push(values.op_type.to_arrow_array());
713
714        // The actual datatype of JSON array is data oriented, not to be derived from the Region
715        // metadata, which is static. So here we have to align the schema.
716        let schema = align_schema_with_json_array(self.schema, &columns);
717        let batch = RecordBatch::try_new(schema, columns).context(NewRecordBatchSnafu)?;
718        // Sorts the record batch.
719        let batch = sort_primary_key_record_batch(&batch)?;
720
721        Ok(BulkPart {
722            batch,
723            max_timestamp: self.max_ts,
724            min_timestamp: self.min_ts,
725            sequence: self.max_sequence,
726            timestamp_index,
727            raw_data: None,
728        })
729    }
730}
731
732fn align_schema_with_json_array(schema: SchemaRef, columns: &[ArrayRef]) -> SchemaRef {
733    if schema.fields().iter().all(|f| !is_json2_extension_type(f)) {
734        return schema;
735    }
736
737    let mut fields = Vec::with_capacity(schema.fields().len());
738    for (field, array) in schema.fields().iter().zip(columns) {
739        if !is_json2_extension_type(field) {
740            fields.push(field.clone());
741            continue;
742        }
743
744        let mut field = field.as_ref().clone();
745        field.set_data_type(array.data_type().clone());
746        fields.push(Arc::new(field));
747    }
748
749    Arc::new(Schema::new_with_metadata(fields, schema.metadata().clone()))
750}
751
752fn new_primary_key_column_builders(
753    metadata: &RegionMetadata,
754    capacity: usize,
755) -> Vec<PrimaryKeyColumnBuilder> {
756    metadata
757        .primary_key_columns()
758        .map(|col| {
759            if col.column_schema.data_type.is_string() {
760                PrimaryKeyColumnBuilder::StringDict(StringDictionaryBuilder::with_capacity(
761                    capacity,
762                    INIT_DICT_VALUE_CAPACITY,
763                    capacity,
764                ))
765            } else {
766                PrimaryKeyColumnBuilder::Vector(
767                    col.column_schema.data_type.create_mutable_vector(capacity),
768                )
769            }
770        })
771        .collect()
772}
773
774/// Sorts the record batch with primary key format.
775pub fn sort_primary_key_record_batch(batch: &RecordBatch) -> Result<RecordBatch> {
776    let indices = if let Some(indices) = try_sort_primary_key_indices(batch)? {
777        indices
778    } else {
779        lexsort_primary_key_indices(batch)?
780    };
781
782    datatypes::arrow::compute::take_record_batch(batch, &indices).context(ComputeArrowSnafu)
783}
784
785/// Sorts the known, non-null primary-key layout without comparing encoded keys for every row.
786///
787/// Returns `None` if the batch does not have the exact layout supported by this fast path.
788fn try_sort_primary_key_indices(batch: &RecordBatch) -> Result<Option<UInt32Array>> {
789    let total_columns = batch.num_columns();
790    let timestamp = batch.column(total_columns - 4);
791    let primary_key = batch.column(total_columns - 3);
792    let sequence = batch.column(total_columns - 2);
793
794    let Some(primary_key) = primary_key
795        .as_any()
796        .downcast_ref::<DictionaryArray<UInt32Type>>()
797    else {
798        return Ok(None);
799    };
800    let Some(sequence) = sequence.as_any().downcast_ref::<UInt64Array>() else {
801        return Ok(None);
802    };
803
804    if batch.num_rows() > u32::MAX as usize
805        || primary_key.null_count() != 0
806        || primary_key.values().data_type() != &ArrowDataType::Binary
807        || primary_key.values().null_count() != 0
808        || primary_key.values().len() > batch.num_rows()
809        || timestamp.null_count() != 0
810        || sequence.null_count() != 0
811    {
812        return Ok(None);
813    }
814
815    let indices = match timestamp.data_type() {
816        ArrowDataType::Timestamp(TimeUnit::Second, _) => timestamp
817            .as_any()
818            .downcast_ref::<TimestampSecondArray>()
819            .map(|timestamp| sort_primary_key_indices(primary_key, timestamp.values(), sequence)),
820        ArrowDataType::Timestamp(TimeUnit::Millisecond, _) => timestamp
821            .as_any()
822            .downcast_ref::<TimestampMillisecondArray>()
823            .map(|timestamp| sort_primary_key_indices(primary_key, timestamp.values(), sequence)),
824        ArrowDataType::Timestamp(TimeUnit::Microsecond, _) => timestamp
825            .as_any()
826            .downcast_ref::<TimestampMicrosecondArray>()
827            .map(|timestamp| sort_primary_key_indices(primary_key, timestamp.values(), sequence)),
828        ArrowDataType::Timestamp(TimeUnit::Nanosecond, _) => timestamp
829            .as_any()
830            .downcast_ref::<TimestampNanosecondArray>()
831            .map(|timestamp| sort_primary_key_indices(primary_key, timestamp.values(), sequence)),
832        _ => None,
833    };
834
835    indices.transpose()
836}
837
838/// Computes sorted row indices by ranking dictionary values, scattering rows into primary-key
839/// buckets, and sorting only buckets that are not already ordered by time and sequence.
840fn sort_primary_key_indices(
841    primary_key: &PrimaryKeyArray,
842    timestamps: &[i64],
843    sequences: &UInt64Array,
844) -> Result<UInt32Array> {
845    debug_assert_eq!(primary_key.len(), timestamps.len());
846    debug_assert_eq!(primary_key.len(), sequences.len());
847    debug_assert_eq!(primary_key.null_count(), 0);
848    debug_assert_eq!(primary_key.values().null_count(), 0);
849    debug_assert_eq!(sequences.null_count(), 0);
850
851    // Ranks compare in the same order as the dictionary values. Equal dictionary values receive
852    // the same rank, which is required because Arrow dictionary concatenation only deduplicates
853    // values on a best-effort basis.
854    let value_ranks = datatypes::arrow::compute::rank(
855        primary_key.values().as_ref(),
856        Some(SortOptions {
857            descending: false,
858            nulls_first: true,
859        }),
860    )
861    .context(ComputeArrowSnafu)?;
862
863    // A rank is at most the number of dictionary values. It is not necessarily dense when the
864    // dictionary contains duplicate values, so keep one extra slot and allow empty buckets.
865    let mut counts = vec![0usize; value_ranks.len() + 1];
866    for &key in primary_key.keys().values() {
867        counts[value_ranks[key as usize] as usize] += 1;
868    }
869
870    let mut bucket_offsets = vec![0usize; counts.len()];
871    let mut offset = 0;
872    for (bucket_offset, &count) in bucket_offsets.iter_mut().zip(&counts) {
873        *bucket_offset = offset;
874        offset += count;
875    }
876
877    // Scatter in input order. This preserves already sorted runs within each primary key.
878    let mut next_offsets = bucket_offsets.clone();
879    let mut indices = vec![0u32; primary_key.len()];
880    for (row, &key) in primary_key.keys().values().iter().enumerate() {
881        let rank = value_ranks[key as usize] as usize;
882        indices[next_offsets[rank]] = row as u32;
883        next_offsets[rank] += 1;
884    }
885
886    let sequences = sequences.values();
887    for (&start, count) in bucket_offsets.iter().zip(counts) {
888        let end = start + count;
889        let bucket = &mut indices[start..end];
890        if !bucket.is_sorted_by(|left, right| {
891            compare_time_sequence(timestamps, sequences, *left, *right).is_le()
892        }) {
893            bucket.sort_unstable_by(|left, right| {
894                compare_time_sequence(timestamps, sequences, *left, *right)
895            });
896        }
897    }
898
899    Ok(UInt32Array::from(indices))
900}
901
902#[inline]
903fn compare_time_sequence(
904    timestamps: &[i64],
905    sequences: &[u64],
906    left: u32,
907    right: u32,
908) -> std::cmp::Ordering {
909    let left = left as usize;
910    let right = right as usize;
911    timestamps[left]
912        .cmp(&timestamps[right])
913        .then_with(|| sequences[right].cmp(&sequences[left]))
914}
915
916fn lexsort_primary_key_indices(batch: &RecordBatch) -> Result<UInt32Array> {
917    let total_columns = batch.num_columns();
918    let sort_columns = vec![
919        // Primary key column (ascending)
920        SortColumn {
921            values: batch.column(total_columns - 3).clone(),
922            options: Some(SortOptions {
923                descending: false,
924                nulls_first: true,
925            }),
926        },
927        // Time index column (ascending)
928        SortColumn {
929            values: batch.column(total_columns - 4).clone(),
930            options: Some(SortOptions {
931                descending: false,
932                nulls_first: true,
933            }),
934        },
935        // Sequence column (descending)
936        SortColumn {
937            values: batch.column(total_columns - 2).clone(),
938            options: Some(SortOptions {
939                descending: true,
940                nulls_first: true,
941            }),
942        },
943    ];
944
945    datatypes::arrow::compute::lexsort_to_indices(&sort_columns, None).context(ComputeArrowSnafu)
946}
947
948/// Converts a `BulkPart` that is unordered and without encoded primary keys into a `BulkPart`
949/// with the same format as produced by [BulkPartConverter].
950///
951/// This function takes a `BulkPart` where:
952/// - For dense encoding: Primary key columns may be stored as individual columns
953/// - For sparse encoding: The `__primary_key` column should already be present with encoded keys
954/// - The batch may not be sorted
955///
956/// And produces a `BulkPart` where:
957/// - Primary key columns are optionally stored (depending on `store_primary_key_columns` and encoding)
958/// - An encoded `__primary_key` dictionary column is present
959/// - The batch is sorted by (primary_key, timestamp, sequence desc)
960///
961/// # Arguments
962///
963/// * `part` - The input `BulkPart` to convert
964/// * `region_metadata` - Region metadata containing schema information
965/// * `primary_key_codec` - Codec for encoding primary keys
966/// * `schema` - Target schema for the output batch
967/// * `store_primary_key_columns` - If true and encoding is not sparse, stores individual primary key columns
968///
969/// # Returns
970///
971/// Returns `None` if the input part has no rows, otherwise returns a new `BulkPart` with
972/// encoded primary keys and sorted data.
973pub fn convert_bulk_part(
974    part: BulkPart,
975    region_metadata: &RegionMetadataRef,
976    primary_key_codec: Arc<dyn PrimaryKeyCodec>,
977    schema: SchemaRef,
978    store_primary_key_columns: bool,
979) -> Result<Option<BulkPart>> {
980    if part.num_rows() == 0 {
981        return Ok(None);
982    }
983
984    let num_rows = part.num_rows();
985    let is_sparse = region_metadata.primary_key_encoding == PrimaryKeyEncoding::Sparse;
986
987    // Builds a column name-to-index map for efficient lookups
988    let input_schema = part.batch.schema();
989    let column_indices: HashMap<&str, usize> = input_schema
990        .fields()
991        .iter()
992        .enumerate()
993        .map(|(idx, field)| (field.name().as_str(), idx))
994        .collect();
995
996    // Determines the structure of the input batch by looking up columns by name
997    let mut output_columns = Vec::new();
998
999    // Extracts primary key columns if we need to encode them (dense encoding)
1000    let pk_array = if is_sparse {
1001        // For sparse encoding, the input should already have the __primary_key column
1002        // We need to find it in the input batch
1003        None
1004    } else {
1005        // For dense encoding, extract and encode primary key columns by name
1006        let pk_vectors: Result<Vec<_>> = region_metadata
1007            .primary_key_columns()
1008            .map(|col_meta| {
1009                let col_idx = column_indices
1010                    .get(col_meta.column_schema.name.as_str())
1011                    .context(ColumnNotFoundSnafu {
1012                        column: &col_meta.column_schema.name,
1013                    })?;
1014                let col = part.batch.column(*col_idx);
1015                Helper::try_into_vector(col).context(error::ComputeVectorSnafu)
1016            })
1017            .collect();
1018        let pk_vectors = pk_vectors?;
1019
1020        let mut key_array_builder = PrimaryKeyArrayBuilder::new();
1021        let mut encode_buf = Vec::new();
1022
1023        for row_idx in 0..num_rows {
1024            encode_buf.clear();
1025
1026            // Collects primary key values with column IDs for this row
1027            let pk_values_with_ids: Vec<_> = region_metadata
1028                .primary_key
1029                .iter()
1030                .zip(pk_vectors.iter())
1031                .map(|(col_id, vector)| (*col_id, vector.get_ref(row_idx)))
1032                .collect();
1033
1034            // Encodes the primary key
1035            primary_key_codec
1036                .encode_value_refs(&pk_values_with_ids, &mut encode_buf)
1037                .context(EncodeSnafu)?;
1038
1039            key_array_builder
1040                .append(&encode_buf)
1041                .context(ComputeArrowSnafu)?;
1042        }
1043
1044        Some(key_array_builder.finish())
1045    };
1046
1047    // Adds primary key columns if storing them (only for dense encoding)
1048    if store_primary_key_columns && !is_sparse {
1049        for col_meta in region_metadata.primary_key_columns() {
1050            let col_idx = column_indices
1051                .get(col_meta.column_schema.name.as_str())
1052                .context(ColumnNotFoundSnafu {
1053                    column: &col_meta.column_schema.name,
1054                })?;
1055            let col = part.batch.column(*col_idx);
1056
1057            // Converts to dictionary if needed for string types
1058            let col = if col_meta.column_schema.data_type.is_string() {
1059                let target_type = ArrowDataType::Dictionary(
1060                    Box::new(ArrowDataType::UInt32),
1061                    Box::new(ArrowDataType::Utf8),
1062                );
1063                arrow::compute::cast(col, &target_type).context(ComputeArrowSnafu)?
1064            } else {
1065                col.clone()
1066            };
1067            output_columns.push(col);
1068        }
1069    }
1070
1071    // Adds field columns
1072    for col_meta in region_metadata.field_columns() {
1073        let col_idx = column_indices
1074            .get(col_meta.column_schema.name.as_str())
1075            .context(ColumnNotFoundSnafu {
1076                column: &col_meta.column_schema.name,
1077            })?;
1078        output_columns.push(part.batch.column(*col_idx).clone());
1079    }
1080
1081    // Adds timestamp column
1082    let new_timestamp_index = output_columns.len();
1083    let ts_col_idx = column_indices
1084        .get(
1085            region_metadata
1086                .time_index_column()
1087                .column_schema
1088                .name
1089                .as_str(),
1090        )
1091        .context(ColumnNotFoundSnafu {
1092            column: &region_metadata.time_index_column().column_schema.name,
1093        })?;
1094    output_columns.push(part.batch.column(*ts_col_idx).clone());
1095
1096    // Adds encoded primary key dictionary column
1097    let pk_dictionary = if let Some(pk_dict_array) = pk_array {
1098        Arc::new(pk_dict_array) as ArrayRef
1099    } else {
1100        let pk_col_idx =
1101            column_indices
1102                .get(PRIMARY_KEY_COLUMN_NAME)
1103                .context(ColumnNotFoundSnafu {
1104                    column: PRIMARY_KEY_COLUMN_NAME,
1105                })?;
1106        let col = part.batch.column(*pk_col_idx);
1107
1108        // Casts to dictionary type if needed
1109        let target_type = ArrowDataType::Dictionary(
1110            Box::new(ArrowDataType::UInt32),
1111            Box::new(ArrowDataType::Binary),
1112        );
1113        arrow::compute::cast(col, &target_type).context(ComputeArrowSnafu)?
1114    };
1115    output_columns.push(pk_dictionary);
1116
1117    let sequence_array = UInt64Array::from(vec![part.sequence; num_rows]);
1118    output_columns.push(Arc::new(sequence_array) as ArrayRef);
1119
1120    let op_type_array = UInt8Array::from(vec![OpType::Put as u8; num_rows]);
1121    output_columns.push(Arc::new(op_type_array) as ArrayRef);
1122
1123    let batch = RecordBatch::try_new(schema, output_columns).context(NewRecordBatchSnafu)?;
1124
1125    // Sorts the batch by (primary_key, timestamp, sequence desc)
1126    let sorted_batch = sort_primary_key_record_batch(&batch)?;
1127
1128    Ok(Some(BulkPart {
1129        batch: sorted_batch,
1130        max_timestamp: part.max_timestamp,
1131        min_timestamp: part.min_timestamp,
1132        sequence: part.sequence,
1133        timestamp_index: new_timestamp_index,
1134        raw_data: None,
1135    }))
1136}
1137
1138#[derive(Debug, Clone)]
1139pub struct EncodedBulkPart {
1140    data: Bytes,
1141    metadata: BulkPartMeta,
1142    /// Cached Arrow schema to avoid rebuilding it from parquet metadata.
1143    schema: SchemaRef,
1144}
1145
1146impl EncodedBulkPart {
1147    pub fn new(data: Bytes, metadata: BulkPartMeta, schema: SchemaRef) -> Self {
1148        Self {
1149            data,
1150            metadata,
1151            schema,
1152        }
1153    }
1154
1155    pub fn metadata(&self) -> &BulkPartMeta {
1156        &self.metadata
1157    }
1158
1159    pub(crate) fn schema(&self) -> SchemaRef {
1160        self.schema.clone()
1161    }
1162
1163    /// Returns the size of the encoded data in bytes
1164    pub(crate) fn size_bytes(&self) -> usize {
1165        self.data.len()
1166    }
1167
1168    /// Returns the encoded data.
1169    pub fn data(&self) -> &Bytes {
1170        &self.data
1171    }
1172
1173    /// Creates MemtableStats from this EncodedBulkPart.
1174    pub fn to_memtable_stats(&self) -> MemtableStats {
1175        let meta = &self.metadata;
1176        let ts_type = meta.region_metadata.time_index_type();
1177        let min_ts = ts_type.create_timestamp(meta.min_timestamp);
1178        let max_ts = ts_type.create_timestamp(meta.max_timestamp);
1179
1180        MemtableStats {
1181            estimated_bytes: self.size_bytes(),
1182            time_range: Some((min_ts, max_ts)),
1183            num_rows: meta.num_rows,
1184            num_ranges: 1,
1185            max_sequence: meta.max_sequence,
1186            series_count: meta.num_series as usize,
1187        }
1188    }
1189
1190    /// Converts this `EncodedBulkPart` to `SstInfo`.
1191    ///
1192    /// # Arguments
1193    /// * `file_id` - The SST file ID to assign to this part
1194    ///
1195    /// # Returns
1196    /// Returns a `SstInfo` instance with information derived from this bulk part's metadata
1197    pub(crate) fn to_sst_info(&self, file_id: FileId) -> SstInfo {
1198        let unit = self.metadata.region_metadata.time_index_type().unit();
1199        let max_row_group_uncompressed_size: u64 = self
1200            .metadata
1201            .parquet_metadata
1202            .row_groups()
1203            .iter()
1204            .map(|rg| {
1205                rg.columns()
1206                    .iter()
1207                    .map(|c| c.uncompressed_size() as u64)
1208                    .sum::<u64>()
1209            })
1210            .max()
1211            .unwrap_or(0);
1212        SstInfo {
1213            file_id,
1214            time_range: (
1215                Timestamp::new(self.metadata.min_timestamp, unit),
1216                Timestamp::new(self.metadata.max_timestamp, unit),
1217            ),
1218            file_size: self.data.len() as u64,
1219            max_row_group_uncompressed_size,
1220            num_rows: self.metadata.num_rows,
1221            num_row_groups: self.metadata.parquet_metadata.num_row_groups() as u64,
1222            file_metadata: Some(self.metadata.parquet_metadata.clone()),
1223            index_metadata: IndexOutput::default(),
1224            num_series: self.metadata.num_series,
1225        }
1226    }
1227
1228    pub(crate) fn read(
1229        &self,
1230        context: BulkIterContextRef,
1231        sequence: Option<SequenceRange>,
1232        mem_scan_metrics: Option<MemScanMetrics>,
1233    ) -> Result<Option<BoxedRecordBatchIterator>> {
1234        // Compute skip_fields for row group pruning from the configured pre-filter mode.
1235        let skip_fields_for_pruning = context.pre_filter_mode().skip_fields();
1236
1237        // use predicate to find row groups to read.
1238        let row_groups_to_read =
1239            context.row_groups_to_read(&self.metadata.parquet_metadata, skip_fields_for_pruning);
1240
1241        if row_groups_to_read.is_empty() {
1242            // All row groups are filtered.
1243            return Ok(None);
1244        }
1245
1246        let iter = EncodedBulkPartIter::try_new(
1247            self,
1248            context,
1249            row_groups_to_read,
1250            sequence,
1251            mem_scan_metrics,
1252        )?;
1253        Ok(Some(Box::new(iter) as BoxedRecordBatchIterator))
1254    }
1255}
1256
1257// TODO(yingwen): max_sequence
1258#[derive(Debug, Clone)]
1259pub struct BulkPartMeta {
1260    /// Total rows in part.
1261    pub num_rows: usize,
1262    /// Max timestamp in part.
1263    pub max_timestamp: i64,
1264    /// Min timestamp in part.
1265    pub min_timestamp: i64,
1266    /// Part file metadata.
1267    pub parquet_metadata: Arc<ParquetMetaData>,
1268    /// Part region schema.
1269    pub region_metadata: RegionMetadataRef,
1270    /// Number of series.
1271    pub num_series: u64,
1272    /// Maximum sequence number in part.
1273    pub max_sequence: u64,
1274}
1275
1276/// Metrics for encoding a part.
1277#[derive(Default, Debug)]
1278pub struct BulkPartEncodeMetrics {
1279    /// Cost of iterating over the data.
1280    pub iter_cost: Duration,
1281    /// Cost of writing the data.
1282    pub write_cost: Duration,
1283    /// Size of data before encoding.
1284    pub raw_size: usize,
1285    /// Size of data after encoding.
1286    pub encoded_size: usize,
1287    /// Number of rows in part.
1288    pub num_rows: usize,
1289}
1290
1291pub struct BulkPartEncoder {
1292    metadata: RegionMetadataRef,
1293    writer_props: Option<WriterProperties>,
1294}
1295
1296impl BulkPartEncoder {
1297    pub fn new(metadata: RegionMetadataRef, row_group_size: usize) -> Result<BulkPartEncoder> {
1298        // TODO(yingwen): Skip arrow schema if needed.
1299        let json = metadata.to_json().context(InvalidMetadataSnafu)?;
1300        let key_value_meta =
1301            parquet::file::metadata::KeyValue::new(PARQUET_METADATA_KEY.to_string(), json);
1302
1303        // TODO(yingwen): Do we need compression?
1304        let writer_props = Some(
1305            WriterProperties::builder()
1306                .set_key_value_metadata(Some(vec![key_value_meta]))
1307                .set_write_batch_size(row_group_size)
1308                .set_max_row_group_row_count(Some(row_group_size))
1309                .set_compression(Compression::ZSTD(ZstdLevel::default()))
1310                .set_column_index_truncate_length(None)
1311                .set_statistics_truncate_length(None)
1312                .build(),
1313        );
1314
1315        Ok(Self {
1316            metadata,
1317            writer_props,
1318        })
1319    }
1320}
1321
1322impl BulkPartEncoder {
1323    /// Encodes [BoxedRecordBatchIterator] into [EncodedBulkPart] with min/max timestamps.
1324    pub fn encode_record_batch_iter(
1325        &self,
1326        iter: BoxedRecordBatchIterator,
1327        arrow_schema: SchemaRef,
1328        min_timestamp: i64,
1329        max_timestamp: i64,
1330        max_sequence: u64,
1331        metrics: &mut BulkPartEncodeMetrics,
1332    ) -> Result<Option<EncodedBulkPart>> {
1333        let mut buf = Vec::with_capacity(4096);
1334        let mut writer =
1335            ArrowWriter::try_new(&mut buf, arrow_schema.clone(), self.writer_props.clone())
1336                .context(EncodeMemtableSnafu)?;
1337        let mut total_rows = 0;
1338        let mut series_estimator = SeriesEstimator::default();
1339
1340        // Process each batch from the iterator
1341        let mut iter_start = Instant::now();
1342        for batch_result in iter {
1343            metrics.iter_cost += iter_start.elapsed();
1344            let batch = batch_result?;
1345            if batch.num_rows() == 0 {
1346                continue;
1347            }
1348
1349            series_estimator.update_flat(&batch);
1350            metrics.raw_size += record_batch_estimated_size(&batch);
1351            let write_start = Instant::now();
1352            writer.write(&batch).context(EncodeMemtableSnafu)?;
1353            metrics.write_cost += write_start.elapsed();
1354            total_rows += batch.num_rows();
1355            iter_start = Instant::now();
1356        }
1357        metrics.iter_cost += iter_start.elapsed();
1358
1359        if total_rows == 0 {
1360            return Ok(None);
1361        }
1362
1363        let close_start = Instant::now();
1364        let file_metadata = writer.close().context(EncodeMemtableSnafu)?;
1365        metrics.write_cost += close_start.elapsed();
1366        metrics.encoded_size += buf.len();
1367        metrics.num_rows += total_rows;
1368
1369        let buf = Bytes::from(buf);
1370        let parquet_metadata = Arc::new(file_metadata);
1371        let num_series = series_estimator.finish();
1372
1373        Ok(Some(EncodedBulkPart {
1374            data: buf,
1375            metadata: BulkPartMeta {
1376                num_rows: total_rows,
1377                max_timestamp,
1378                min_timestamp,
1379                parquet_metadata,
1380                region_metadata: self.metadata.clone(),
1381                num_series,
1382                max_sequence,
1383            },
1384            schema: arrow_schema,
1385        }))
1386    }
1387
1388    /// Encodes bulk part to a [EncodedBulkPart], returns the encoded data.
1389    pub fn encode_part(&self, part: &BulkPart) -> Result<Option<EncodedBulkPart>> {
1390        if part.batch.num_rows() == 0 {
1391            return Ok(None);
1392        }
1393
1394        let mut buf = Vec::with_capacity(4096);
1395        let arrow_schema = part.batch.schema();
1396
1397        let file_metadata = {
1398            let mut writer =
1399                ArrowWriter::try_new(&mut buf, arrow_schema.clone(), self.writer_props.clone())
1400                    .context(EncodeMemtableSnafu)?;
1401            writer.write(&part.batch).context(EncodeMemtableSnafu)?;
1402            writer.finish().context(EncodeMemtableSnafu)?
1403        };
1404
1405        let buf = Bytes::from(buf);
1406        let parquet_metadata = Arc::new(file_metadata);
1407
1408        Ok(Some(EncodedBulkPart {
1409            data: buf,
1410            metadata: BulkPartMeta {
1411                num_rows: part.batch.num_rows(),
1412                max_timestamp: part.max_timestamp,
1413                min_timestamp: part.min_timestamp,
1414                parquet_metadata,
1415                region_metadata: self.metadata.clone(),
1416                num_series: part.estimated_series_count() as u64,
1417                max_sequence: part.sequence,
1418            },
1419            schema: arrow_schema,
1420        }))
1421    }
1422}
1423
1424/// Per-batch min/max statistics for the first tag column in a `MultiBulkPart`.
1425///
1426/// Since batches are sorted by primary key, we can extract the min/max of the first tag
1427/// from the first/last row's encoded primary key in each batch. These statistics enable
1428/// batch-level pruning using predicates, analogous to row-group pruning in parquet.
1429#[derive(Debug, Clone)]
1430struct BatchStats {
1431    /// Number of batches.
1432    num_batches: usize,
1433    /// Column id of the first tag.
1434    first_tag_id: ColumnId,
1435    /// Min values of the first tag, one element per batch.
1436    min_values: ArrayRef,
1437    /// Max values of the first tag, one element per batch.
1438    max_values: ArrayRef,
1439}
1440
1441impl BatchStats {
1442    /// Computes batch statistics from a slice of record batches.
1443    ///
1444    /// Returns `None` if there is no primary key (no first tag to collect stats for)
1445    /// or if extracting statistics fails.
1446    fn compute(batches: &[RecordBatch], metadata: &RegionMetadata) -> Option<Self> {
1447        // `primary_key.first()` is correct for both dense and sparse encodings.
1448        // For dense, values follow the order of `metadata.primary_key`.
1449        // For sparse, `decode_leftmost` decodes the first value which also
1450        // corresponds to `primary_key.first()`. See `SparsePrimaryKeyCodec` for format details.
1451        let first_tag_id = *metadata.primary_key.first()?;
1452        let first_tag_column = metadata.column_by_id(first_tag_id)?;
1453        let data_type = &first_tag_column.column_schema.data_type;
1454
1455        let converter = build_primary_key_codec_with_fields(
1456            metadata.primary_key_encoding,
1457            [(first_tag_id, SortField::new(data_type.clone()))].into_iter(),
1458        );
1459        let pk_index = primary_key_column_index(batches.first()?.num_columns());
1460
1461        let mut min_builder = data_type.create_mutable_vector(batches.len());
1462        let mut max_builder = data_type.create_mutable_vector(batches.len());
1463
1464        for batch in batches {
1465            match Self::extract_first_tag_bounds(batch, pk_index, &*converter) {
1466                Some((min_val, max_val)) => {
1467                    min_builder.push_value_ref(&min_val.as_value_ref());
1468                    max_builder.push_value_ref(&max_val.as_value_ref());
1469                }
1470                None => {
1471                    min_builder.push_null();
1472                    max_builder.push_null();
1473                }
1474            }
1475        }
1476
1477        Some(Self {
1478            num_batches: batches.len(),
1479            first_tag_id,
1480            min_values: min_builder.to_vector().to_arrow_array(),
1481            max_values: max_builder.to_vector().to_arrow_array(),
1482        })
1483    }
1484
1485    /// Extracts the first tag value from the first and last rows of a batch.
1486    fn extract_first_tag_bounds(
1487        batch: &RecordBatch,
1488        pk_index: usize,
1489        converter: &dyn PrimaryKeyCodec,
1490    ) -> Option<(datatypes::value::Value, datatypes::value::Value)> {
1491        if batch.num_rows() == 0 {
1492            return None;
1493        }
1494
1495        let pk_dict = batch
1496            .column(pk_index)
1497            .as_any()
1498            .downcast_ref::<PrimaryKeyArray>()?;
1499        let pk_values = pk_dict.values().as_any().downcast_ref::<BinaryArray>()?;
1500
1501        let keys = pk_dict.keys();
1502        let min_key = keys.value(0);
1503        let max_key = keys.value(batch.num_rows() - 1);
1504        let min_bytes = pk_values.value(min_key as usize);
1505        let max_bytes = pk_values.value(max_key as usize);
1506
1507        Some((
1508            converter.decode_leftmost(min_bytes).ok()??,
1509            converter.decode_leftmost(max_bytes).ok()??,
1510        ))
1511    }
1512}
1513
1514/// Adapter implementing `PruningStatistics` for `BatchStats`.
1515///
1516/// Used with `Predicate::prune_with_stats()` to skip batches whose first-tag
1517/// min/max range does not match the query predicate.
1518struct BatchPruningStats<'a> {
1519    stats: &'a BatchStats,
1520    metadata: &'a RegionMetadataRef,
1521}
1522
1523impl PruningStatistics for BatchPruningStats<'_> {
1524    fn min_values(&self, column: &Column) -> Option<ArrayRef> {
1525        let col = self.metadata.column_by_name(&column.name)?;
1526        if col.column_id == self.stats.first_tag_id {
1527            Some(self.stats.min_values.clone())
1528        } else {
1529            None
1530        }
1531    }
1532
1533    fn max_values(&self, column: &Column) -> Option<ArrayRef> {
1534        let col = self.metadata.column_by_name(&column.name)?;
1535        if col.column_id == self.stats.first_tag_id {
1536            Some(self.stats.max_values.clone())
1537        } else {
1538            None
1539        }
1540    }
1541
1542    fn num_containers(&self) -> usize {
1543        self.stats.num_batches
1544    }
1545
1546    fn null_counts(&self, _column: &Column) -> Option<ArrayRef> {
1547        None
1548    }
1549
1550    fn row_counts(&self, _column: &Column) -> Option<ArrayRef> {
1551        None
1552    }
1553
1554    fn contained(
1555        &self,
1556        _column: &Column,
1557        _values: &std::collections::HashSet<datafusion_common::ScalarValue>,
1558    ) -> Option<BooleanArray> {
1559        None
1560    }
1561}
1562
1563/// Returns true if the predicate references the given column name.
1564fn predicate_references_column(predicate: &table::predicate::Predicate, column_name: &str) -> bool {
1565    let mut columns = HashSet::new();
1566    for expr in predicate.exprs() {
1567        let _ = expr_to_columns(expr, &mut columns);
1568    }
1569    columns.iter().any(|col| col.name == column_name)
1570}
1571
1572/// Returns true if the batch should be pruned (skipped) based on the first-tag min/max
1573/// statistics and the predicate in the context. Returns false if no pruning is possible
1574/// (no primary key, no predicate, or the batch matches the predicate).
1575pub(crate) fn should_prune_bulk_part(
1576    batch: &RecordBatch,
1577    context: &BulkIterContext,
1578    metadata: &RegionMetadata,
1579) -> bool {
1580    let predicate = match &context.predicate {
1581        Some(p) => p,
1582        None => return false,
1583    };
1584    // Check if the predicate references the first tag column to avoid computing
1585    // expensive batch statistics when they won't help with pruning.
1586    let first_tag_id = match metadata.primary_key.first() {
1587        Some(id) => *id,
1588        None => return false,
1589    };
1590    // Safety: `first_tag_id` comes from `metadata.primary_key` so the column always exists.
1591    let first_tag_name = &metadata
1592        .column_by_id(first_tag_id)
1593        .unwrap()
1594        .column_schema
1595        .name;
1596    if !predicate_references_column(predicate, first_tag_name) {
1597        return false;
1598    }
1599    let stats = match BatchStats::compute(std::slice::from_ref(batch), metadata) {
1600        Some(s) => s,
1601        None => return false,
1602    };
1603    let region_meta = context.read_format().metadata();
1604    let pruning_stats = BatchPruningStats {
1605        stats: &stats,
1606        metadata: region_meta,
1607    };
1608    let mask = predicate.prune_with_stats(&pruning_stats, region_meta.schema.arrow_schema());
1609    !mask.first().copied().unwrap_or(true)
1610}
1611
1612/// A collection of ordered RecordBatches representing a bulk part without parquet encoding.
1613///
1614/// Similar to `EncodedBulkPart` but stores raw RecordBatches instead of encoded parquet data.
1615/// The RecordBatches must be ordered by (primary key, timestamp, sequence desc).
1616/// Uses SmallVec to optimize for the common case of few batches while avoiding heap allocation.
1617#[derive(Debug, Clone)]
1618pub struct MultiBulkPart {
1619    /// Ordered record batches. SmallVec optimized for up to 4 batches inline.
1620    batches: SmallVec<[RecordBatch; 4]>,
1621    /// Total rows across all batches.
1622    total_rows: usize,
1623    /// Max timestamp in part.
1624    max_timestamp: i64,
1625    /// Min timestamp in part.
1626    min_timestamp: i64,
1627    /// Max sequence number in part.
1628    max_sequence: SequenceNumber,
1629    /// Number of series.
1630    series_count: usize,
1631    /// Pre-computed per-batch statistics for the first tag column.
1632    /// `None` if there is no primary key.
1633    batch_stats: Option<BatchStats>,
1634}
1635
1636impl MultiBulkPart {
1637    /// Creates a new MultiBulkPart from a single BulkPart.
1638    pub fn from_bulk_part(part: BulkPart, metadata: &RegionMetadata) -> Self {
1639        let num_rows = part.num_rows();
1640        let series_count = part.estimated_series_count();
1641        let batch_stats = BatchStats::compute(std::slice::from_ref(&part.batch), metadata);
1642        let mut batches = SmallVec::new();
1643        batches.push(part.batch);
1644
1645        Self {
1646            batches,
1647            total_rows: num_rows,
1648            max_timestamp: part.max_timestamp,
1649            min_timestamp: part.min_timestamp,
1650            max_sequence: part.sequence,
1651            series_count,
1652            batch_stats,
1653        }
1654    }
1655
1656    /// Creates a new MultiBulkPart from multiple ordered RecordBatches.
1657    ///
1658    /// # Arguments
1659    /// * `batches` - Ordered record batches
1660    /// * `min_timestamp` - Minimum timestamp across all batches
1661    /// * `max_timestamp` - Maximum timestamp across all batches
1662    /// * `max_sequence` - Maximum sequence number across all batches
1663    /// * `series_count` - Number of series in the batches
1664    /// * `metadata` - Region metadata for computing batch statistics
1665    ///
1666    /// # Panics
1667    /// Panics if batches is empty.
1668    pub fn new(
1669        batches: Vec<RecordBatch>,
1670        min_timestamp: i64,
1671        max_timestamp: i64,
1672        max_sequence: SequenceNumber,
1673        series_count: usize,
1674        metadata: &RegionMetadata,
1675    ) -> Self {
1676        assert!(!batches.is_empty(), "batches must not be empty");
1677
1678        let total_rows = batches.iter().map(|b| b.num_rows()).sum();
1679        let batch_stats = BatchStats::compute(&batches, metadata);
1680
1681        Self {
1682            batches: SmallVec::from_vec(batches),
1683            total_rows,
1684            max_timestamp,
1685            min_timestamp,
1686            max_sequence,
1687            series_count,
1688            batch_stats,
1689        }
1690    }
1691
1692    /// Returns the total number of rows across all batches.
1693    pub fn num_rows(&self) -> usize {
1694        self.total_rows
1695    }
1696
1697    pub(crate) fn schemas(&self) -> impl Iterator<Item = SchemaRef> + '_ {
1698        self.batches.iter().map(|batch| batch.schema())
1699    }
1700
1701    /// Returns the minimum timestamp.
1702    pub fn min_timestamp(&self) -> i64 {
1703        self.min_timestamp
1704    }
1705
1706    /// Returns the maximum timestamp.
1707    pub fn max_timestamp(&self) -> i64 {
1708        self.max_timestamp
1709    }
1710
1711    /// Returns the maximum sequence number.
1712    pub fn max_sequence(&self) -> SequenceNumber {
1713        self.max_sequence
1714    }
1715
1716    /// Returns the number of series.
1717    pub fn series_count(&self) -> usize {
1718        self.series_count
1719    }
1720
1721    /// Returns the number of record batches in this part.
1722    pub fn num_batches(&self) -> usize {
1723        self.batches.len()
1724    }
1725
1726    /// Returns the estimated memory size of all batches.
1727    pub(crate) fn estimated_size(&self) -> usize {
1728        self.batches.iter().map(record_batch_estimated_size).sum()
1729    }
1730
1731    /// Reads data from this part with the given context and filters.
1732    ///
1733    /// If batch-level statistics are available and a predicate is set, prunes
1734    /// batches whose first-tag min/max range doesn't match the predicate before
1735    /// creating the iterator.
1736    pub(crate) fn read(
1737        &self,
1738        context: BulkIterContextRef,
1739        sequence: Option<SequenceRange>,
1740        mem_scan_metrics: Option<MemScanMetrics>,
1741    ) -> Result<Option<BoxedRecordBatchIterator>> {
1742        if self.batches.is_empty() {
1743            return Ok(None);
1744        }
1745
1746        let batches_to_read = self.prune_batches(&context);
1747
1748        if batches_to_read.is_empty() {
1749            return Ok(None);
1750        }
1751
1752        let iter = crate::memtable::bulk::part_reader::BulkPartBatchIter::new(
1753            batches_to_read,
1754            context,
1755            sequence,
1756            self.series_count,
1757            mem_scan_metrics,
1758        );
1759        Ok(Some(Box::new(iter) as BoxedRecordBatchIterator))
1760    }
1761
1762    /// Prunes batches using the first-tag min/max statistics and the predicate.
1763    /// Returns all batches if no stats or no predicate is available.
1764    fn prune_batches(&self, context: &BulkIterContextRef) -> Vec<RecordBatch> {
1765        if let Some(stats) = &self.batch_stats
1766            && let Some(predicate) = &context.predicate
1767        {
1768            let region_meta = context.read_format().metadata();
1769            let pruning_stats = BatchPruningStats {
1770                stats,
1771                metadata: region_meta,
1772            };
1773            let mask =
1774                predicate.prune_with_stats(&pruning_stats, region_meta.schema.arrow_schema());
1775            self.batches
1776                .iter()
1777                .zip(mask.iter())
1778                .filter_map(
1779                    |(batch, &selected)| {
1780                        if selected { Some(batch.clone()) } else { None }
1781                    },
1782                )
1783                .collect()
1784        } else {
1785            self.batches.iter().cloned().collect()
1786        }
1787    }
1788
1789    /// Converts this `MultiBulkPart` to `MemtableStats`.
1790    pub fn to_memtable_stats(&self, region_metadata: &RegionMetadataRef) -> MemtableStats {
1791        let ts_type = region_metadata.time_index_type();
1792        let min_ts = ts_type.create_timestamp(self.min_timestamp);
1793        let max_ts = ts_type.create_timestamp(self.max_timestamp);
1794
1795        MemtableStats {
1796            estimated_bytes: self.estimated_size(),
1797            time_range: Some((min_ts, max_ts)),
1798            num_rows: self.num_rows(),
1799            num_ranges: 1,
1800            max_sequence: self.max_sequence,
1801            series_count: self.series_count,
1802        }
1803    }
1804}
1805
1806#[cfg(test)]
1807mod tests {
1808    use api::v1::{Row, SemanticType, WriteHint};
1809    use datafusion_common::ScalarValue;
1810    use datatypes::arrow::array::{
1811        BinaryArray, DictionaryArray, Float64Array, TimestampMillisecondArray,
1812    };
1813    use datatypes::arrow::datatypes::UInt32Type;
1814    use datatypes::prelude::{ConcreteDataType, Value};
1815    use datatypes::schema::ColumnSchema;
1816    use mito_codec::row_converter::build_primary_key_codec;
1817    use rand::rngs::StdRng;
1818    use rand::{Rng, SeedableRng};
1819    use store_api::metadata::{ColumnMetadata, RegionMetadataBuilder};
1820    use store_api::storage::RegionId;
1821    use store_api::storage::consts::ReservedColumnId;
1822    use table::predicate::Predicate;
1823
1824    use super::*;
1825    use crate::memtable::bulk::context::BulkIterContext;
1826    use crate::sst::{FlatSchemaOptions, to_flat_sst_arrow_schema};
1827    use crate::test_util::memtable_util::{build_key_values_with_ts_seq_values, metadata_for_test};
1828
1829    struct MutationInput<'a> {
1830        k0: &'a str,
1831        k1: u32,
1832        timestamps: &'a [i64],
1833        v1: &'a [Option<f64>],
1834        sequence: u64,
1835    }
1836
1837    #[test]
1838    fn test_sort_primary_key_record_batch_with_duplicate_dictionary_values() {
1839        // Dictionary keys 1 and 2 both represent "series_a". This can happen after Arrow
1840        // concatenates dictionaries because dictionary deduplication is best-effort.
1841        let primary_key = DictionaryArray::try_new(
1842            UInt32Array::from(vec![0, 1, 0, 2, 2]),
1843            Arc::new(BinaryArray::from_vec(vec![
1844                b"series_b".as_slice(),
1845                b"series_a".as_slice(),
1846                b"series_a".as_slice(),
1847            ])),
1848        )
1849        .unwrap();
1850        let timestamps = TimestampMillisecondArray::from(vec![20, 30, 10, 30, 10]);
1851        let sequences = UInt64Array::from(vec![5, 6, 4, 8, 9]);
1852        let row_ids = UInt32Array::from_iter_values(0..5);
1853        let op_types = UInt8Array::from_value(OpType::Put as u8, 5);
1854        let schema = Arc::new(Schema::new(vec![
1855            Field::new("row_id", ArrowDataType::UInt32, false),
1856            Field::new(
1857                "ts",
1858                ArrowDataType::Timestamp(TimeUnit::Millisecond, None),
1859                false,
1860            ),
1861            Field::new(
1862                PRIMARY_KEY_COLUMN_NAME,
1863                ArrowDataType::Dictionary(
1864                    Box::new(ArrowDataType::UInt32),
1865                    Box::new(ArrowDataType::Binary),
1866                ),
1867                false,
1868            ),
1869            Field::new("__sequence", ArrowDataType::UInt64, false),
1870            Field::new("__op_type", ArrowDataType::UInt8, false),
1871        ]));
1872        let batch = RecordBatch::try_new(
1873            schema,
1874            vec![
1875                Arc::new(row_ids),
1876                Arc::new(timestamps),
1877                Arc::new(primary_key),
1878                Arc::new(sequences),
1879                Arc::new(op_types),
1880            ],
1881        )
1882        .unwrap();
1883
1884        let expected_indices = lexsort_primary_key_indices(&batch).unwrap();
1885        let expected =
1886            datatypes::arrow::compute::take_record_batch(&batch, &expected_indices).unwrap();
1887        let actual = sort_primary_key_record_batch(&batch).unwrap();
1888
1889        let expected_row_ids = expected
1890            .column(0)
1891            .as_any()
1892            .downcast_ref::<UInt32Array>()
1893            .unwrap();
1894        let actual_row_ids = actual
1895            .column(0)
1896            .as_any()
1897            .downcast_ref::<UInt32Array>()
1898            .unwrap();
1899        assert_eq!(expected_row_ids, actual_row_ids);
1900        assert_eq!(actual_row_ids.values(), &[4, 3, 1, 2, 0]);
1901    }
1902
1903    #[test]
1904    fn test_sort_primary_key_record_batch_against_lexsort() {
1905        let mut rng = StdRng::seed_from_u64(0x5eed);
1906        for _ in 0..100 {
1907            let num_rows = rng.random_range(0..128);
1908            let keys = UInt32Array::from_iter_values((0..num_rows).map(|_| rng.random_range(0..8)));
1909            let primary_key = DictionaryArray::try_new(
1910                keys,
1911                Arc::new(BinaryArray::from_vec(vec![
1912                    b"d".as_slice(),
1913                    b"a".as_slice(),
1914                    b"c".as_slice(),
1915                    b"b".as_slice(),
1916                    b"a".as_slice(),
1917                    b"d".as_slice(),
1918                    b"b".as_slice(),
1919                    b"c".as_slice(),
1920                ])),
1921            )
1922            .unwrap();
1923            let timestamps = TimestampMillisecondArray::from_iter_values(
1924                (0..num_rows).map(|_| rng.random_range(-1000..1000)),
1925            );
1926            // Unique sequences ensure that the oracle and fast path have no fully equal sort keys.
1927            let sequences = UInt64Array::from_iter_values(0..num_rows as u64);
1928            let row_ids = UInt32Array::from_iter_values(0..num_rows as u32);
1929            let op_types = UInt8Array::from_value(OpType::Put as u8, num_rows);
1930            let schema = Arc::new(Schema::new(vec![
1931                Field::new("row_id", ArrowDataType::UInt32, false),
1932                Field::new(
1933                    "ts",
1934                    ArrowDataType::Timestamp(TimeUnit::Millisecond, None),
1935                    false,
1936                ),
1937                Field::new(
1938                    PRIMARY_KEY_COLUMN_NAME,
1939                    ArrowDataType::Dictionary(
1940                        Box::new(ArrowDataType::UInt32),
1941                        Box::new(ArrowDataType::Binary),
1942                    ),
1943                    false,
1944                ),
1945                Field::new("__sequence", ArrowDataType::UInt64, false),
1946                Field::new("__op_type", ArrowDataType::UInt8, false),
1947            ]));
1948            let batch = RecordBatch::try_new(
1949                schema,
1950                vec![
1951                    Arc::new(row_ids),
1952                    Arc::new(timestamps),
1953                    Arc::new(primary_key),
1954                    Arc::new(sequences),
1955                    Arc::new(op_types),
1956                ],
1957            )
1958            .unwrap();
1959
1960            let expected_indices = lexsort_primary_key_indices(&batch).unwrap();
1961            let expected =
1962                datatypes::arrow::compute::take_record_batch(&batch, &expected_indices).unwrap();
1963            let actual = sort_primary_key_record_batch(&batch).unwrap();
1964            assert_eq!(expected.column(0), actual.column(0));
1965        }
1966    }
1967
1968    #[test]
1969    fn test_unordered_part_tracks_estimated_bytes() {
1970        let mut part = UnorderedPart::new();
1971        let bulk_part = BulkPart {
1972            batch: RecordBatch::new_empty(Arc::new(arrow::datatypes::Schema::empty())),
1973            max_timestamp: 0,
1974            min_timestamp: 0,
1975            sequence: 0,
1976            timestamp_index: 0,
1977            raw_data: None,
1978        };
1979        let estimated_size = bulk_part.estimated_size();
1980
1981        part.push(bulk_part);
1982        assert_eq!(estimated_size, part.estimated_bytes());
1983        part.clear();
1984        assert_eq!(0, part.estimated_bytes());
1985        assert!(part.is_empty());
1986    }
1987
1988    #[test]
1989    fn test_unordered_part_should_accept() {
1990        let mut part = UnorderedPart::new();
1991        part.set_threshold(10);
1992        assert!(part.should_accept(9));
1993        assert!(!part.should_accept(10));
1994    }
1995
1996    fn encode(input: &[MutationInput]) -> EncodedBulkPart {
1997        let metadata = metadata_for_test();
1998        let kvs = input
1999            .iter()
2000            .map(|m| {
2001                build_key_values_with_ts_seq_values(
2002                    &metadata,
2003                    m.k0.to_string(),
2004                    m.k1,
2005                    m.timestamps.iter().copied(),
2006                    m.v1.iter().copied(),
2007                    m.sequence,
2008                )
2009            })
2010            .collect::<Vec<_>>();
2011        let schema = to_flat_sst_arrow_schema(&metadata, &FlatSchemaOptions::default());
2012        let primary_key_codec = build_primary_key_codec(&metadata);
2013        let mut converter = BulkPartConverter::new(&metadata, schema, 64, primary_key_codec, true);
2014        for kv in kvs {
2015            converter.append_key_values(&kv).unwrap();
2016        }
2017        let part = converter.convert().unwrap();
2018        let encoder = BulkPartEncoder::new(metadata, 1024).unwrap();
2019        encoder.encode_part(&part).unwrap().unwrap()
2020    }
2021
2022    #[test]
2023    fn test_write_and_read_part_projection() {
2024        let part = encode(&[
2025            MutationInput {
2026                k0: "a",
2027                k1: 0,
2028                timestamps: &[1],
2029                v1: &[Some(0.1)],
2030                sequence: 0,
2031            },
2032            MutationInput {
2033                k0: "b",
2034                k1: 0,
2035                timestamps: &[1],
2036                v1: &[Some(0.0)],
2037                sequence: 0,
2038            },
2039            MutationInput {
2040                k0: "a",
2041                k1: 0,
2042                timestamps: &[2],
2043                v1: &[Some(0.2)],
2044                sequence: 1,
2045            },
2046        ]);
2047
2048        let projection = &[4u32];
2049        let reader = part
2050            .read(
2051                Arc::new(
2052                    BulkIterContext::new(
2053                        part.metadata.region_metadata.clone(),
2054                        Some(projection.as_slice()),
2055                        None,
2056                        false,
2057                        crate::sst::parquet::DEFAULT_READ_BATCH_SIZE,
2058                    )
2059                    .unwrap(),
2060                ),
2061                None,
2062                None,
2063            )
2064            .unwrap()
2065            .expect("expect at least one row group");
2066
2067        let mut total_rows_read = 0;
2068        let mut field: Vec<f64> = vec![];
2069        for res in reader {
2070            let batch = res.unwrap();
2071            assert_eq!(5, batch.num_columns());
2072            field.extend_from_slice(
2073                batch
2074                    .column(0)
2075                    .as_any()
2076                    .downcast_ref::<Float64Array>()
2077                    .unwrap()
2078                    .values(),
2079            );
2080            total_rows_read += batch.num_rows();
2081        }
2082        assert_eq!(3, total_rows_read);
2083        assert_eq!(vec![0.1, 0.2, 0.0], field);
2084    }
2085
2086    fn prepare(key_values: Vec<(&str, u32, (i64, i64), u64)>) -> EncodedBulkPart {
2087        let metadata = metadata_for_test();
2088        let kvs = key_values
2089            .into_iter()
2090            .map(|(k0, k1, (start, end), sequence)| {
2091                let ts = start..end;
2092                let v1 = (start..end).map(|_| None);
2093                build_key_values_with_ts_seq_values(&metadata, k0.to_string(), k1, ts, v1, sequence)
2094            })
2095            .collect::<Vec<_>>();
2096        let schema = to_flat_sst_arrow_schema(&metadata, &FlatSchemaOptions::default());
2097        let primary_key_codec = build_primary_key_codec(&metadata);
2098        let mut converter = BulkPartConverter::new(&metadata, schema, 64, primary_key_codec, true);
2099        for kv in kvs {
2100            converter.append_key_values(&kv).unwrap();
2101        }
2102        let part = converter.convert().unwrap();
2103        let encoder = BulkPartEncoder::new(metadata, 1024).unwrap();
2104        encoder.encode_part(&part).unwrap().unwrap()
2105    }
2106
2107    fn check_prune_row_group(
2108        part: &EncodedBulkPart,
2109        predicate: Option<Predicate>,
2110        expected_rows: usize,
2111    ) {
2112        let context = Arc::new(
2113            BulkIterContext::new(
2114                part.metadata.region_metadata.clone(),
2115                None,
2116                predicate,
2117                false,
2118                crate::sst::parquet::DEFAULT_READ_BATCH_SIZE,
2119            )
2120            .unwrap(),
2121        );
2122        let reader = part
2123            .read(context, None, None)
2124            .unwrap()
2125            .expect("expect at least one row group");
2126        let mut total_rows_read = 0;
2127        for res in reader {
2128            let batch = res.unwrap();
2129            total_rows_read += batch.num_rows();
2130        }
2131        // Should only read row group 1.
2132        assert_eq!(expected_rows, total_rows_read);
2133    }
2134
2135    #[test]
2136    fn test_prune_row_groups() {
2137        let part = prepare(vec![
2138            ("a", 0, (0, 40), 1),
2139            ("a", 1, (0, 60), 1),
2140            ("b", 0, (0, 100), 2),
2141            ("b", 1, (100, 180), 3),
2142            ("b", 1, (180, 210), 4),
2143        ]);
2144
2145        let context = Arc::new(
2146            BulkIterContext::new(
2147                part.metadata.region_metadata.clone(),
2148                None,
2149                Some(Predicate::new(vec![datafusion_expr::col("ts").eq(
2150                    datafusion_expr::lit(ScalarValue::TimestampMillisecond(Some(300), None)),
2151                )])),
2152                false,
2153                crate::sst::parquet::DEFAULT_READ_BATCH_SIZE,
2154            )
2155            .unwrap(),
2156        );
2157        assert!(part.read(context, None, None).unwrap().is_none());
2158
2159        check_prune_row_group(&part, None, 310);
2160
2161        check_prune_row_group(
2162            &part,
2163            Some(Predicate::new(vec![
2164                datafusion_expr::col("k0").eq(datafusion_expr::lit("a")),
2165                datafusion_expr::col("k1").eq(datafusion_expr::lit(0u32)),
2166            ])),
2167            40,
2168        );
2169
2170        check_prune_row_group(
2171            &part,
2172            Some(Predicate::new(vec![
2173                datafusion_expr::col("k0").eq(datafusion_expr::lit("a")),
2174                datafusion_expr::col("k1").eq(datafusion_expr::lit(1u32)),
2175            ])),
2176            60,
2177        );
2178
2179        check_prune_row_group(
2180            &part,
2181            Some(Predicate::new(vec![
2182                datafusion_expr::col("k0").eq(datafusion_expr::lit("a")),
2183            ])),
2184            100,
2185        );
2186
2187        check_prune_row_group(
2188            &part,
2189            Some(Predicate::new(vec![
2190                datafusion_expr::col("k0").eq(datafusion_expr::lit("b")),
2191                datafusion_expr::col("k1").eq(datafusion_expr::lit(0u32)),
2192            ])),
2193            100,
2194        );
2195
2196        // Predicates over field column can do precise filtering.
2197        check_prune_row_group(
2198            &part,
2199            Some(Predicate::new(vec![
2200                datafusion_expr::col("v0").eq(datafusion_expr::lit(150i64)),
2201            ])),
2202            1,
2203        );
2204    }
2205
2206    #[test]
2207    fn test_bulk_part_converter_append_and_convert() {
2208        let metadata = metadata_for_test();
2209        let capacity = 100;
2210        let primary_key_codec = build_primary_key_codec(&metadata);
2211        let schema = to_flat_sst_arrow_schema(
2212            &metadata,
2213            &FlatSchemaOptions::from_encoding(metadata.primary_key_encoding),
2214        );
2215
2216        let mut converter =
2217            BulkPartConverter::new(&metadata, schema, capacity, primary_key_codec, true);
2218
2219        let key_values1 = build_key_values_with_ts_seq_values(
2220            &metadata,
2221            "key1".to_string(),
2222            1u32,
2223            vec![1000, 2000].into_iter(),
2224            vec![Some(1.0), Some(2.0)].into_iter(),
2225            1,
2226        );
2227
2228        let key_values2 = build_key_values_with_ts_seq_values(
2229            &metadata,
2230            "key2".to_string(),
2231            2u32,
2232            vec![1500].into_iter(),
2233            vec![Some(3.0)].into_iter(),
2234            2,
2235        );
2236
2237        converter.append_key_values(&key_values1).unwrap();
2238        converter.append_key_values(&key_values2).unwrap();
2239
2240        let bulk_part = converter.convert().unwrap();
2241
2242        assert_eq!(bulk_part.num_rows(), 3);
2243        assert_eq!(bulk_part.min_timestamp, 1000);
2244        assert_eq!(bulk_part.max_timestamp, 2000);
2245        assert_eq!(bulk_part.sequence, 2);
2246        assert_eq!(bulk_part.timestamp_index, bulk_part.batch.num_columns() - 4);
2247
2248        // Validate primary key columns are stored
2249        // Schema should include primary key columns k0 and k1 at the beginning
2250        let schema = bulk_part.batch.schema();
2251        let field_names: Vec<&str> = schema.fields().iter().map(|f| f.name().as_str()).collect();
2252        assert_eq!(
2253            field_names,
2254            vec![
2255                "k0",
2256                "k1",
2257                "v0",
2258                "v1",
2259                "ts",
2260                "__primary_key",
2261                "__sequence",
2262                "__op_type"
2263            ]
2264        );
2265    }
2266
2267    #[test]
2268    fn test_bulk_part_converter_sorting() {
2269        let metadata = metadata_for_test();
2270        let capacity = 100;
2271        let primary_key_codec = build_primary_key_codec(&metadata);
2272        let schema = to_flat_sst_arrow_schema(
2273            &metadata,
2274            &FlatSchemaOptions::from_encoding(metadata.primary_key_encoding),
2275        );
2276
2277        let mut converter =
2278            BulkPartConverter::new(&metadata, schema, capacity, primary_key_codec, true);
2279
2280        let key_values1 = build_key_values_with_ts_seq_values(
2281            &metadata,
2282            "z_key".to_string(),
2283            3u32,
2284            vec![3000].into_iter(),
2285            vec![Some(3.0)].into_iter(),
2286            3,
2287        );
2288
2289        let key_values2 = build_key_values_with_ts_seq_values(
2290            &metadata,
2291            "a_key".to_string(),
2292            1u32,
2293            vec![1000].into_iter(),
2294            vec![Some(1.0)].into_iter(),
2295            1,
2296        );
2297
2298        let key_values3 = build_key_values_with_ts_seq_values(
2299            &metadata,
2300            "m_key".to_string(),
2301            2u32,
2302            vec![2000].into_iter(),
2303            vec![Some(2.0)].into_iter(),
2304            2,
2305        );
2306
2307        converter.append_key_values(&key_values1).unwrap();
2308        converter.append_key_values(&key_values2).unwrap();
2309        converter.append_key_values(&key_values3).unwrap();
2310
2311        let bulk_part = converter.convert().unwrap();
2312
2313        assert_eq!(bulk_part.num_rows(), 3);
2314
2315        let ts_column = bulk_part.batch.column(bulk_part.timestamp_index);
2316        let seq_column = bulk_part.batch.column(bulk_part.batch.num_columns() - 2);
2317
2318        let ts_array = ts_column
2319            .as_any()
2320            .downcast_ref::<TimestampMillisecondArray>()
2321            .unwrap();
2322        let seq_array = seq_column.as_any().downcast_ref::<UInt64Array>().unwrap();
2323
2324        assert_eq!(ts_array.values(), &[1000, 2000, 3000]);
2325        assert_eq!(seq_array.values(), &[1, 2, 3]);
2326
2327        // Validate primary key columns are stored
2328        let schema = bulk_part.batch.schema();
2329        let field_names: Vec<&str> = schema.fields().iter().map(|f| f.name().as_str()).collect();
2330        assert_eq!(
2331            field_names,
2332            vec![
2333                "k0",
2334                "k1",
2335                "v0",
2336                "v1",
2337                "ts",
2338                "__primary_key",
2339                "__sequence",
2340                "__op_type"
2341            ]
2342        );
2343    }
2344
2345    #[test]
2346    fn test_bulk_part_converter_empty() {
2347        let metadata = metadata_for_test();
2348        let capacity = 10;
2349        let primary_key_codec = build_primary_key_codec(&metadata);
2350        let schema = to_flat_sst_arrow_schema(
2351            &metadata,
2352            &FlatSchemaOptions::from_encoding(metadata.primary_key_encoding),
2353        );
2354
2355        let converter =
2356            BulkPartConverter::new(&metadata, schema, capacity, primary_key_codec, true);
2357
2358        let bulk_part = converter.convert().unwrap();
2359
2360        assert_eq!(bulk_part.num_rows(), 0);
2361        assert_eq!(bulk_part.min_timestamp, i64::MAX);
2362        assert_eq!(bulk_part.max_timestamp, i64::MIN);
2363        assert_eq!(bulk_part.sequence, SequenceNumber::MIN);
2364
2365        // Validate primary key columns are present in schema even for empty batch
2366        let schema = bulk_part.batch.schema();
2367        let field_names: Vec<&str> = schema.fields().iter().map(|f| f.name().as_str()).collect();
2368        assert_eq!(
2369            field_names,
2370            vec![
2371                "k0",
2372                "k1",
2373                "v0",
2374                "v1",
2375                "ts",
2376                "__primary_key",
2377                "__sequence",
2378                "__op_type"
2379            ]
2380        );
2381    }
2382
2383    #[test]
2384    fn test_bulk_part_converter_without_primary_key_columns() {
2385        let metadata = metadata_for_test();
2386        let primary_key_codec = build_primary_key_codec(&metadata);
2387        let schema = to_flat_sst_arrow_schema(
2388            &metadata,
2389            &FlatSchemaOptions {
2390                raw_pk_columns: false,
2391                string_pk_use_dict: true,
2392                ..Default::default()
2393            },
2394        );
2395
2396        let capacity = 100;
2397        let mut converter =
2398            BulkPartConverter::new(&metadata, schema, capacity, primary_key_codec, false);
2399
2400        let key_values1 = build_key_values_with_ts_seq_values(
2401            &metadata,
2402            "key1".to_string(),
2403            1u32,
2404            vec![1000, 2000].into_iter(),
2405            vec![Some(1.0), Some(2.0)].into_iter(),
2406            1,
2407        );
2408
2409        let key_values2 = build_key_values_with_ts_seq_values(
2410            &metadata,
2411            "key2".to_string(),
2412            2u32,
2413            vec![1500].into_iter(),
2414            vec![Some(3.0)].into_iter(),
2415            2,
2416        );
2417
2418        converter.append_key_values(&key_values1).unwrap();
2419        converter.append_key_values(&key_values2).unwrap();
2420
2421        let bulk_part = converter.convert().unwrap();
2422
2423        assert_eq!(bulk_part.num_rows(), 3);
2424        assert_eq!(bulk_part.min_timestamp, 1000);
2425        assert_eq!(bulk_part.max_timestamp, 2000);
2426        assert_eq!(bulk_part.sequence, 2);
2427        assert_eq!(bulk_part.timestamp_index, bulk_part.batch.num_columns() - 4);
2428
2429        // Validate primary key columns are NOT stored individually
2430        let schema = bulk_part.batch.schema();
2431        let field_names: Vec<&str> = schema.fields().iter().map(|f| f.name().as_str()).collect();
2432        assert_eq!(
2433            field_names,
2434            vec!["v0", "v1", "ts", "__primary_key", "__sequence", "__op_type"]
2435        );
2436    }
2437
2438    #[allow(clippy::too_many_arguments)]
2439    fn build_key_values_with_sparse_encoding(
2440        metadata: &RegionMetadataRef,
2441        primary_key_codec: &Arc<dyn PrimaryKeyCodec>,
2442        table_id: u32,
2443        tsid: u64,
2444        k0: String,
2445        k1: String,
2446        timestamps: impl Iterator<Item = i64>,
2447        values: impl Iterator<Item = Option<f64>>,
2448        sequence: SequenceNumber,
2449    ) -> KeyValues {
2450        // Encode the primary key (__table_id, __tsid, k0, k1) into binary format using the sparse codec
2451        let pk_values = vec![
2452            (ReservedColumnId::table_id(), Value::UInt32(table_id)),
2453            (ReservedColumnId::tsid(), Value::UInt64(tsid)),
2454            (0, Value::String(k0.clone().into())),
2455            (1, Value::String(k1.clone().into())),
2456        ];
2457        let mut encoded_key = Vec::new();
2458        primary_key_codec
2459            .encode_values(&pk_values, &mut encoded_key)
2460            .unwrap();
2461        assert!(!encoded_key.is_empty());
2462
2463        // Create schema for sparse encoding: __primary_key, ts, v0, v1
2464        let column_schema = vec![
2465            api::v1::ColumnSchema {
2466                column_name: PRIMARY_KEY_COLUMN_NAME.to_string(),
2467                datatype: api::helper::ColumnDataTypeWrapper::try_from(
2468                    ConcreteDataType::binary_datatype(),
2469                )
2470                .unwrap()
2471                .datatype() as i32,
2472                semantic_type: api::v1::SemanticType::Tag as i32,
2473                ..Default::default()
2474            },
2475            api::v1::ColumnSchema {
2476                column_name: "ts".to_string(),
2477                datatype: api::helper::ColumnDataTypeWrapper::try_from(
2478                    ConcreteDataType::timestamp_millisecond_datatype(),
2479                )
2480                .unwrap()
2481                .datatype() as i32,
2482                semantic_type: api::v1::SemanticType::Timestamp as i32,
2483                ..Default::default()
2484            },
2485            api::v1::ColumnSchema {
2486                column_name: "v0".to_string(),
2487                datatype: api::helper::ColumnDataTypeWrapper::try_from(
2488                    ConcreteDataType::int64_datatype(),
2489                )
2490                .unwrap()
2491                .datatype() as i32,
2492                semantic_type: api::v1::SemanticType::Field as i32,
2493                ..Default::default()
2494            },
2495            api::v1::ColumnSchema {
2496                column_name: "v1".to_string(),
2497                datatype: api::helper::ColumnDataTypeWrapper::try_from(
2498                    ConcreteDataType::float64_datatype(),
2499                )
2500                .unwrap()
2501                .datatype() as i32,
2502                semantic_type: api::v1::SemanticType::Field as i32,
2503                ..Default::default()
2504            },
2505        ];
2506
2507        let rows = timestamps
2508            .zip(values)
2509            .map(|(ts, v)| Row {
2510                values: vec![
2511                    api::v1::Value {
2512                        value_data: Some(api::v1::value::ValueData::BinaryValue(
2513                            encoded_key.clone(),
2514                        )),
2515                    },
2516                    api::v1::Value {
2517                        value_data: Some(api::v1::value::ValueData::TimestampMillisecondValue(ts)),
2518                    },
2519                    api::v1::Value {
2520                        value_data: Some(api::v1::value::ValueData::I64Value(ts)),
2521                    },
2522                    api::v1::Value {
2523                        value_data: v.map(api::v1::value::ValueData::F64Value),
2524                    },
2525                ],
2526            })
2527            .collect();
2528
2529        let mutation = api::v1::Mutation {
2530            op_type: 1,
2531            sequence,
2532            rows: Some(api::v1::Rows {
2533                schema: column_schema,
2534                rows,
2535            }),
2536            write_hint: Some(WriteHint {
2537                primary_key_encoding: api::v1::PrimaryKeyEncoding::Sparse.into(),
2538            }),
2539        };
2540        KeyValues::new(metadata.as_ref(), mutation).unwrap()
2541    }
2542
2543    #[test]
2544    fn test_bulk_part_converter_sparse_primary_key_encoding() {
2545        use api::v1::SemanticType;
2546        use datatypes::schema::ColumnSchema;
2547        use store_api::metadata::{ColumnMetadata, RegionMetadataBuilder};
2548        use store_api::storage::RegionId;
2549
2550        let mut builder = RegionMetadataBuilder::new(RegionId::new(123, 456));
2551        builder
2552            .push_column_metadata(ColumnMetadata {
2553                column_schema: ColumnSchema::new("k0", ConcreteDataType::string_datatype(), false),
2554                semantic_type: SemanticType::Tag,
2555                column_id: 0,
2556            })
2557            .push_column_metadata(ColumnMetadata {
2558                column_schema: ColumnSchema::new("k1", ConcreteDataType::string_datatype(), false),
2559                semantic_type: SemanticType::Tag,
2560                column_id: 1,
2561            })
2562            .push_column_metadata(ColumnMetadata {
2563                column_schema: ColumnSchema::new(
2564                    "ts",
2565                    ConcreteDataType::timestamp_millisecond_datatype(),
2566                    false,
2567                ),
2568                semantic_type: SemanticType::Timestamp,
2569                column_id: 2,
2570            })
2571            .push_column_metadata(ColumnMetadata {
2572                column_schema: ColumnSchema::new("v0", ConcreteDataType::int64_datatype(), true),
2573                semantic_type: SemanticType::Field,
2574                column_id: 3,
2575            })
2576            .push_column_metadata(ColumnMetadata {
2577                column_schema: ColumnSchema::new("v1", ConcreteDataType::float64_datatype(), true),
2578                semantic_type: SemanticType::Field,
2579                column_id: 4,
2580            })
2581            .primary_key(vec![0, 1])
2582            .primary_key_encoding(PrimaryKeyEncoding::Sparse);
2583        let metadata = Arc::new(builder.build().unwrap());
2584
2585        let primary_key_codec = build_primary_key_codec(&metadata);
2586        let schema = to_flat_sst_arrow_schema(
2587            &metadata,
2588            &FlatSchemaOptions::from_encoding(metadata.primary_key_encoding),
2589        );
2590
2591        assert_eq!(metadata.primary_key_encoding, PrimaryKeyEncoding::Sparse);
2592        assert_eq!(primary_key_codec.encoding(), PrimaryKeyEncoding::Sparse);
2593
2594        let capacity = 100;
2595        let mut converter =
2596            BulkPartConverter::new(&metadata, schema, capacity, primary_key_codec.clone(), true);
2597
2598        let key_values1 = build_key_values_with_sparse_encoding(
2599            &metadata,
2600            &primary_key_codec,
2601            2048u32, // table_id
2602            100u64,  // tsid
2603            "key11".to_string(),
2604            "key21".to_string(),
2605            vec![1000, 2000].into_iter(),
2606            vec![Some(1.0), Some(2.0)].into_iter(),
2607            1,
2608        );
2609
2610        let key_values2 = build_key_values_with_sparse_encoding(
2611            &metadata,
2612            &primary_key_codec,
2613            4096u32, // table_id
2614            200u64,  // tsid
2615            "key12".to_string(),
2616            "key22".to_string(),
2617            vec![1500].into_iter(),
2618            vec![Some(3.0)].into_iter(),
2619            2,
2620        );
2621
2622        converter.append_key_values(&key_values1).unwrap();
2623        converter.append_key_values(&key_values2).unwrap();
2624
2625        let bulk_part = converter.convert().unwrap();
2626
2627        assert_eq!(bulk_part.num_rows(), 3);
2628        assert_eq!(bulk_part.min_timestamp, 1000);
2629        assert_eq!(bulk_part.max_timestamp, 2000);
2630        assert_eq!(bulk_part.sequence, 2);
2631        assert_eq!(bulk_part.timestamp_index, bulk_part.batch.num_columns() - 4);
2632
2633        // For sparse encoding, primary key columns should NOT be stored individually
2634        // even when store_primary_key_columns is true, because sparse encoding
2635        // stores the encoded primary key in the __primary_key column
2636        let schema = bulk_part.batch.schema();
2637        let field_names: Vec<&str> = schema.fields().iter().map(|f| f.name().as_str()).collect();
2638        assert_eq!(
2639            field_names,
2640            vec!["v0", "v1", "ts", "__primary_key", "__sequence", "__op_type"]
2641        );
2642
2643        // Verify the __primary_key column contains encoded sparse keys
2644        let primary_key_column = bulk_part.batch.column_by_name("__primary_key").unwrap();
2645        let dict_array = primary_key_column
2646            .as_any()
2647            .downcast_ref::<DictionaryArray<UInt32Type>>()
2648            .unwrap();
2649
2650        // Should have non-zero entries indicating encoded primary keys
2651        assert!(!dict_array.is_empty());
2652        assert_eq!(dict_array.len(), 3); // 3 rows total
2653
2654        // Verify values are properly encoded binary data (not empty)
2655        let values = dict_array
2656            .values()
2657            .as_any()
2658            .downcast_ref::<BinaryArray>()
2659            .unwrap();
2660        for i in 0..values.len() {
2661            assert!(
2662                !values.value(i).is_empty(),
2663                "Encoded primary key should not be empty"
2664            );
2665        }
2666    }
2667
2668    #[test]
2669    fn test_convert_bulk_part_empty() {
2670        let metadata = metadata_for_test();
2671        let schema = to_flat_sst_arrow_schema(
2672            &metadata,
2673            &FlatSchemaOptions::from_encoding(metadata.primary_key_encoding),
2674        );
2675        let primary_key_codec = build_primary_key_codec(&metadata);
2676
2677        // Create empty batch
2678        let empty_batch = RecordBatch::new_empty(schema.clone());
2679        let empty_part = BulkPart {
2680            batch: empty_batch,
2681            max_timestamp: 0,
2682            min_timestamp: 0,
2683            sequence: 0,
2684            timestamp_index: 0,
2685            raw_data: None,
2686        };
2687
2688        let result =
2689            convert_bulk_part(empty_part, &metadata, primary_key_codec, schema, true).unwrap();
2690        assert!(result.is_none());
2691    }
2692
2693    #[test]
2694    fn test_convert_bulk_part_dense_with_pk_columns() {
2695        let metadata = metadata_for_test();
2696        let primary_key_codec = build_primary_key_codec(&metadata);
2697
2698        let k0_array = Arc::new(arrow::array::StringArray::from(vec![
2699            "key1", "key2", "key1",
2700        ]));
2701        let k1_array = Arc::new(arrow::array::UInt32Array::from(vec![1, 2, 1]));
2702        let v0_array = Arc::new(arrow::array::Int64Array::from(vec![100, 200, 300]));
2703        let v1_array = Arc::new(arrow::array::Float64Array::from(vec![1.0, 2.0, 3.0]));
2704        let ts_array = Arc::new(TimestampMillisecondArray::from(vec![1000, 2000, 1500]));
2705
2706        let input_schema = Arc::new(Schema::new(vec![
2707            Field::new("k0", ArrowDataType::Utf8, false),
2708            Field::new("k1", ArrowDataType::UInt32, false),
2709            Field::new("v0", ArrowDataType::Int64, true),
2710            Field::new("v1", ArrowDataType::Float64, true),
2711            Field::new(
2712                "ts",
2713                ArrowDataType::Timestamp(arrow::datatypes::TimeUnit::Millisecond, None),
2714                false,
2715            ),
2716        ]));
2717
2718        let input_batch = RecordBatch::try_new(
2719            input_schema,
2720            vec![k0_array, k1_array, v0_array, v1_array, ts_array],
2721        )
2722        .unwrap();
2723
2724        let part = BulkPart {
2725            batch: input_batch,
2726            max_timestamp: 2000,
2727            min_timestamp: 1000,
2728            sequence: 5,
2729            timestamp_index: 4,
2730            raw_data: None,
2731        };
2732
2733        let output_schema = to_flat_sst_arrow_schema(
2734            &metadata,
2735            &FlatSchemaOptions::from_encoding(metadata.primary_key_encoding),
2736        );
2737
2738        let result = convert_bulk_part(
2739            part,
2740            &metadata,
2741            primary_key_codec,
2742            output_schema,
2743            true, // store primary key columns
2744        )
2745        .unwrap();
2746
2747        let converted = result.unwrap();
2748
2749        assert_eq!(converted.num_rows(), 3);
2750        assert_eq!(converted.max_timestamp, 2000);
2751        assert_eq!(converted.min_timestamp, 1000);
2752        assert_eq!(converted.sequence, 5);
2753
2754        let schema = converted.batch.schema();
2755        let field_names: Vec<&str> = schema.fields().iter().map(|f| f.name().as_str()).collect();
2756        assert_eq!(
2757            field_names,
2758            vec![
2759                "k0",
2760                "k1",
2761                "v0",
2762                "v1",
2763                "ts",
2764                "__primary_key",
2765                "__sequence",
2766                "__op_type"
2767            ]
2768        );
2769
2770        let k0_col = converted.batch.column_by_name("k0").unwrap();
2771        assert!(matches!(
2772            k0_col.data_type(),
2773            ArrowDataType::Dictionary(_, _)
2774        ));
2775
2776        let pk_col = converted.batch.column_by_name("__primary_key").unwrap();
2777        let dict_array = pk_col
2778            .as_any()
2779            .downcast_ref::<DictionaryArray<UInt32Type>>()
2780            .unwrap();
2781        let keys = dict_array.keys();
2782
2783        assert_eq!(keys.len(), 3);
2784    }
2785
2786    #[test]
2787    fn test_convert_bulk_part_dense_without_pk_columns() {
2788        let metadata = metadata_for_test();
2789        let primary_key_codec = build_primary_key_codec(&metadata);
2790
2791        // Create input batch with primary key columns (k0, k1)
2792        let k0_array = Arc::new(arrow::array::StringArray::from(vec!["key1", "key2"]));
2793        let k1_array = Arc::new(arrow::array::UInt32Array::from(vec![1, 2]));
2794        let v0_array = Arc::new(arrow::array::Int64Array::from(vec![100, 200]));
2795        let v1_array = Arc::new(arrow::array::Float64Array::from(vec![1.0, 2.0]));
2796        let ts_array = Arc::new(TimestampMillisecondArray::from(vec![1000, 2000]));
2797
2798        let input_schema = Arc::new(Schema::new(vec![
2799            Field::new("k0", ArrowDataType::Utf8, false),
2800            Field::new("k1", ArrowDataType::UInt32, false),
2801            Field::new("v0", ArrowDataType::Int64, true),
2802            Field::new("v1", ArrowDataType::Float64, true),
2803            Field::new(
2804                "ts",
2805                ArrowDataType::Timestamp(arrow::datatypes::TimeUnit::Millisecond, None),
2806                false,
2807            ),
2808        ]));
2809
2810        let input_batch = RecordBatch::try_new(
2811            input_schema,
2812            vec![k0_array, k1_array, v0_array, v1_array, ts_array],
2813        )
2814        .unwrap();
2815
2816        let part = BulkPart {
2817            batch: input_batch,
2818            max_timestamp: 2000,
2819            min_timestamp: 1000,
2820            sequence: 3,
2821            timestamp_index: 4,
2822            raw_data: None,
2823        };
2824
2825        let output_schema = to_flat_sst_arrow_schema(
2826            &metadata,
2827            &FlatSchemaOptions {
2828                raw_pk_columns: false,
2829                string_pk_use_dict: true,
2830                ..Default::default()
2831            },
2832        );
2833
2834        let result = convert_bulk_part(
2835            part,
2836            &metadata,
2837            primary_key_codec,
2838            output_schema,
2839            false, // don't store primary key columns
2840        )
2841        .unwrap();
2842
2843        let converted = result.unwrap();
2844
2845        assert_eq!(converted.num_rows(), 2);
2846        assert_eq!(converted.max_timestamp, 2000);
2847        assert_eq!(converted.min_timestamp, 1000);
2848        assert_eq!(converted.sequence, 3);
2849
2850        // Verify schema does NOT include individual primary key columns
2851        let schema = converted.batch.schema();
2852        let field_names: Vec<&str> = schema.fields().iter().map(|f| f.name().as_str()).collect();
2853        assert_eq!(
2854            field_names,
2855            vec!["v0", "v1", "ts", "__primary_key", "__sequence", "__op_type"]
2856        );
2857
2858        // Verify __primary_key column is present and is a dictionary
2859        let pk_col = converted.batch.column_by_name("__primary_key").unwrap();
2860        assert!(matches!(
2861            pk_col.data_type(),
2862            ArrowDataType::Dictionary(_, _)
2863        ));
2864    }
2865
2866    #[test]
2867    fn test_convert_bulk_part_sparse_encoding() {
2868        let mut builder = RegionMetadataBuilder::new(RegionId::new(123, 456));
2869        builder
2870            .push_column_metadata(ColumnMetadata {
2871                column_schema: ColumnSchema::new("k0", ConcreteDataType::string_datatype(), false),
2872                semantic_type: SemanticType::Tag,
2873                column_id: 0,
2874            })
2875            .push_column_metadata(ColumnMetadata {
2876                column_schema: ColumnSchema::new("k1", ConcreteDataType::string_datatype(), false),
2877                semantic_type: SemanticType::Tag,
2878                column_id: 1,
2879            })
2880            .push_column_metadata(ColumnMetadata {
2881                column_schema: ColumnSchema::new(
2882                    "ts",
2883                    ConcreteDataType::timestamp_millisecond_datatype(),
2884                    false,
2885                ),
2886                semantic_type: SemanticType::Timestamp,
2887                column_id: 2,
2888            })
2889            .push_column_metadata(ColumnMetadata {
2890                column_schema: ColumnSchema::new("v0", ConcreteDataType::int64_datatype(), true),
2891                semantic_type: SemanticType::Field,
2892                column_id: 3,
2893            })
2894            .push_column_metadata(ColumnMetadata {
2895                column_schema: ColumnSchema::new("v1", ConcreteDataType::float64_datatype(), true),
2896                semantic_type: SemanticType::Field,
2897                column_id: 4,
2898            })
2899            .primary_key(vec![0, 1])
2900            .primary_key_encoding(PrimaryKeyEncoding::Sparse);
2901        let metadata = Arc::new(builder.build().unwrap());
2902
2903        let primary_key_codec = build_primary_key_codec(&metadata);
2904
2905        // Create input batch with __primary_key column (sparse encoding)
2906        let pk_array = Arc::new(arrow::array::BinaryArray::from(vec![
2907            b"encoded_key_1".as_slice(),
2908            b"encoded_key_2".as_slice(),
2909        ]));
2910        let v0_array = Arc::new(arrow::array::Int64Array::from(vec![100, 200]));
2911        let v1_array = Arc::new(arrow::array::Float64Array::from(vec![1.0, 2.0]));
2912        let ts_array = Arc::new(TimestampMillisecondArray::from(vec![1000, 2000]));
2913
2914        let input_schema = Arc::new(Schema::new(vec![
2915            Field::new("__primary_key", ArrowDataType::Binary, false),
2916            Field::new("v0", ArrowDataType::Int64, true),
2917            Field::new("v1", ArrowDataType::Float64, true),
2918            Field::new(
2919                "ts",
2920                ArrowDataType::Timestamp(arrow::datatypes::TimeUnit::Millisecond, None),
2921                false,
2922            ),
2923        ]));
2924
2925        let input_batch =
2926            RecordBatch::try_new(input_schema, vec![pk_array, v0_array, v1_array, ts_array])
2927                .unwrap();
2928
2929        let part = BulkPart {
2930            batch: input_batch,
2931            max_timestamp: 2000,
2932            min_timestamp: 1000,
2933            sequence: 7,
2934            timestamp_index: 3,
2935            raw_data: None,
2936        };
2937
2938        let output_schema = to_flat_sst_arrow_schema(
2939            &metadata,
2940            &FlatSchemaOptions::from_encoding(metadata.primary_key_encoding),
2941        );
2942
2943        let result = convert_bulk_part(
2944            part,
2945            &metadata,
2946            primary_key_codec,
2947            output_schema,
2948            true, // store_primary_key_columns (ignored for sparse)
2949        )
2950        .unwrap();
2951
2952        let converted = result.unwrap();
2953
2954        assert_eq!(converted.num_rows(), 2);
2955        assert_eq!(converted.max_timestamp, 2000);
2956        assert_eq!(converted.min_timestamp, 1000);
2957        assert_eq!(converted.sequence, 7);
2958
2959        // Verify schema does NOT include individual primary key columns (sparse encoding)
2960        let schema = converted.batch.schema();
2961        let field_names: Vec<&str> = schema.fields().iter().map(|f| f.name().as_str()).collect();
2962        assert_eq!(
2963            field_names,
2964            vec!["v0", "v1", "ts", "__primary_key", "__sequence", "__op_type"]
2965        );
2966
2967        // Verify __primary_key is dictionary encoded
2968        let pk_col = converted.batch.column_by_name("__primary_key").unwrap();
2969        assert!(matches!(
2970            pk_col.data_type(),
2971            ArrowDataType::Dictionary(_, _)
2972        ));
2973    }
2974
2975    #[test]
2976    fn test_convert_bulk_part_sorting_with_multiple_series() {
2977        let metadata = metadata_for_test();
2978        let primary_key_codec = build_primary_key_codec(&metadata);
2979
2980        // Create unsorted batch with multiple series and timestamps
2981        let k0_array = Arc::new(arrow::array::StringArray::from(vec![
2982            "series_b", "series_a", "series_b", "series_a",
2983        ]));
2984        let k1_array = Arc::new(arrow::array::UInt32Array::from(vec![2, 1, 2, 1]));
2985        let v0_array = Arc::new(arrow::array::Int64Array::from(vec![200, 100, 400, 300]));
2986        let v1_array = Arc::new(arrow::array::Float64Array::from(vec![2.0, 1.0, 4.0, 3.0]));
2987        let ts_array = Arc::new(TimestampMillisecondArray::from(vec![
2988            2000, 1000, 4000, 3000,
2989        ]));
2990
2991        let input_schema = Arc::new(Schema::new(vec![
2992            Field::new("k0", ArrowDataType::Utf8, false),
2993            Field::new("k1", ArrowDataType::UInt32, false),
2994            Field::new("v0", ArrowDataType::Int64, true),
2995            Field::new("v1", ArrowDataType::Float64, true),
2996            Field::new(
2997                "ts",
2998                ArrowDataType::Timestamp(arrow::datatypes::TimeUnit::Millisecond, None),
2999                false,
3000            ),
3001        ]));
3002
3003        let input_batch = RecordBatch::try_new(
3004            input_schema,
3005            vec![k0_array, k1_array, v0_array, v1_array, ts_array],
3006        )
3007        .unwrap();
3008
3009        let part = BulkPart {
3010            batch: input_batch,
3011            max_timestamp: 4000,
3012            min_timestamp: 1000,
3013            sequence: 10,
3014            timestamp_index: 4,
3015            raw_data: None,
3016        };
3017
3018        let output_schema = to_flat_sst_arrow_schema(
3019            &metadata,
3020            &FlatSchemaOptions::from_encoding(metadata.primary_key_encoding),
3021        );
3022
3023        let result =
3024            convert_bulk_part(part, &metadata, primary_key_codec, output_schema, true).unwrap();
3025
3026        let converted = result.unwrap();
3027
3028        assert_eq!(converted.num_rows(), 4);
3029
3030        // Verify data is sorted by (primary_key, timestamp, sequence desc)
3031        let ts_col = converted.batch.column(converted.timestamp_index);
3032        let ts_array = ts_col
3033            .as_any()
3034            .downcast_ref::<TimestampMillisecondArray>()
3035            .unwrap();
3036
3037        // After sorting by (pk, ts), we should have:
3038        // series_a,1: ts=1000, 3000
3039        // series_b,2: ts=2000, 4000
3040        let timestamps: Vec<i64> = ts_array.values().to_vec();
3041        assert_eq!(timestamps, vec![1000, 3000, 2000, 4000]);
3042    }
3043
3044    /// Helper to create a converted BulkPart (with __primary_key column) from MutationInputs.
3045    fn build_converted_bulk_part(inputs: &[MutationInput]) -> BulkPart {
3046        let metadata = metadata_for_test();
3047        let kvs = inputs
3048            .iter()
3049            .map(|m| {
3050                build_key_values_with_ts_seq_values(
3051                    &metadata,
3052                    m.k0.to_string(),
3053                    m.k1,
3054                    m.timestamps.iter().copied(),
3055                    m.v1.iter().copied(),
3056                    m.sequence,
3057                )
3058            })
3059            .collect::<Vec<_>>();
3060        let schema = to_flat_sst_arrow_schema(&metadata, &FlatSchemaOptions::default());
3061        let primary_key_codec = build_primary_key_codec(&metadata);
3062        let mut converter = BulkPartConverter::new(&metadata, schema, 64, primary_key_codec, true);
3063        for kv in kvs {
3064            converter.append_key_values(&kv).unwrap();
3065        }
3066        converter.convert().unwrap()
3067    }
3068
3069    /// Helper to create a MultiBulkPart where each group becomes a separate batch.
3070    fn build_multi_bulk_part(groups: &[&[MutationInput]]) -> (MultiBulkPart, RegionMetadataRef) {
3071        let metadata = metadata_for_test();
3072        let mut all_batches = Vec::new();
3073        let mut min_ts = i64::MAX;
3074        let mut max_ts = i64::MIN;
3075        let mut max_seq = 0u64;
3076
3077        for inputs in groups {
3078            let part = build_converted_bulk_part(inputs);
3079            min_ts = min_ts.min(part.min_timestamp);
3080            max_ts = max_ts.max(part.max_timestamp);
3081            max_seq = max_seq.max(part.sequence);
3082            all_batches.push(part.batch);
3083        }
3084
3085        let multi = MultiBulkPart::new(
3086            all_batches,
3087            min_ts,
3088            max_ts,
3089            max_seq,
3090            groups.len(),
3091            &metadata,
3092        );
3093        (multi, metadata)
3094    }
3095
3096    #[test]
3097    fn test_multi_bulk_part_prune_batches() {
3098        // Three batches with distinct k0 ranges: ["a"], ["m"], ["z"].
3099        let (multi, metadata) = build_multi_bulk_part(&[
3100            &[MutationInput {
3101                k0: "a",
3102                k1: 0,
3103                timestamps: &[1, 2],
3104                v1: &[Some(1.0), Some(2.0)],
3105                sequence: 0,
3106            }],
3107            &[MutationInput {
3108                k0: "m",
3109                k1: 0,
3110                timestamps: &[3, 4],
3111                v1: &[Some(3.0), Some(4.0)],
3112                sequence: 1,
3113            }],
3114            &[MutationInput {
3115                k0: "z",
3116                k1: 0,
3117                timestamps: &[5, 6],
3118                v1: &[Some(5.0), Some(6.0)],
3119                sequence: 2,
3120            }],
3121        ]);
3122        assert_eq!(multi.num_rows(), 6);
3123        assert_eq!(multi.num_batches(), 3);
3124
3125        // k0 = "m" => only middle batch (2 rows).
3126        let context = Arc::new(
3127            BulkIterContext::new(
3128                metadata.clone(),
3129                None,
3130                Some(Predicate::new(vec![
3131                    datafusion_expr::col("k0").eq(datafusion_expr::lit("m")),
3132                ])),
3133                false,
3134                crate::sst::parquet::DEFAULT_READ_BATCH_SIZE,
3135            )
3136            .unwrap(),
3137        );
3138        let reader = multi
3139            .read(context, None, None)
3140            .unwrap()
3141            .expect("should have results");
3142        let total_rows: usize = reader.map(|r| r.unwrap().num_rows()).sum();
3143        assert_eq!(total_rows, 2);
3144
3145        // k0 = "nonexistent" => all pruned, returns None.
3146        let context = Arc::new(
3147            BulkIterContext::new(
3148                metadata.clone(),
3149                None,
3150                Some(Predicate::new(vec![
3151                    datafusion_expr::col("k0").eq(datafusion_expr::lit("nonexistent")),
3152                ])),
3153                false,
3154                crate::sst::parquet::DEFAULT_READ_BATCH_SIZE,
3155            )
3156            .unwrap(),
3157        );
3158        assert!(multi.read(context, None, None).unwrap().is_none());
3159
3160        // No predicate => all 6 rows.
3161        let context = Arc::new(
3162            BulkIterContext::new(
3163                metadata.clone(),
3164                None,
3165                None,
3166                false,
3167                crate::sst::parquet::DEFAULT_READ_BATCH_SIZE,
3168            )
3169            .unwrap(),
3170        );
3171        let reader = multi
3172            .read(context, None, None)
3173            .unwrap()
3174            .expect("should have results");
3175        let total_rows: usize = reader.map(|r| r.unwrap().num_rows()).sum();
3176        assert_eq!(total_rows, 6);
3177    }
3178
3179    #[test]
3180    fn test_fill_missing_columns_resets_raw_data() {
3181        let metadata = metadata_for_test();
3182
3183        // A batch missing the `v1` field column, e.g. built against an older
3184        // schema before `v1` was added.
3185        let k0_array = Arc::new(arrow::array::StringArray::from(vec!["key1", "key2"]));
3186        let k1_array = Arc::new(arrow::array::UInt32Array::from(vec![1u32, 2]));
3187        let v0_array = Arc::new(arrow::array::Int64Array::from(vec![100i64, 200]));
3188        let ts_array = Arc::new(TimestampMillisecondArray::from(vec![1000i64, 2000]));
3189        let input_schema = Arc::new(Schema::new(vec![
3190            Field::new("k0", ArrowDataType::Utf8, false),
3191            Field::new("k1", ArrowDataType::UInt32, false),
3192            Field::new("v0", ArrowDataType::Int64, true),
3193            Field::new(
3194                "ts",
3195                ArrowDataType::Timestamp(arrow::datatypes::TimeUnit::Millisecond, None),
3196                false,
3197            ),
3198        ]));
3199        let batch =
3200            RecordBatch::try_new(input_schema, vec![k0_array, k1_array, v0_array, ts_array])
3201                .unwrap();
3202
3203        // Encodes the batch as raw data like the bulk insert request does.
3204        let mut encoder = FlightEncoder::default();
3205        let schema_bytes = encoder.encode_schema(batch.schema().as_ref()).data_header;
3206        let [flight_data] = encoder
3207            .encode(FlightMessage::RecordBatch(batch.clone()))
3208            .try_into()
3209            .unwrap();
3210        let mut part = BulkPart {
3211            batch,
3212            max_timestamp: 2000,
3213            min_timestamp: 1000,
3214            sequence: 5,
3215            timestamp_index: 3,
3216            raw_data: Some(ArrowIpc {
3217                schema: schema_bytes,
3218                data_header: flight_data.data_header,
3219                payload: flight_data.data_body,
3220            }),
3221        };
3222
3223        part.fill_missing_columns(&metadata).unwrap();
3224        assert!(part.batch.column_by_name("v1").is_some());
3225        // The raw data no longer matches the batch after filling missing
3226        // columns and must be cleared, otherwise the WAL entry built from
3227        // this part loses the filled columns.
3228        assert!(part.raw_data.is_none());
3229
3230        // The WAL entry round trip keeps the filled column.
3231        let entry = BulkWalEntry::try_from(&part).unwrap();
3232        let replayed = BulkPart::try_from(entry).unwrap();
3233        assert_eq!(2, replayed.num_rows());
3234        assert!(replayed.batch.column_by_name("v1").is_some());
3235    }
3236}