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