Skip to main content

mito2/sst/parquet/
format.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//! Format to store in parquet.
16//!
17//! We store three internal columns in parquet:
18//! - `__primary_key`, the primary key of the row (tags). Type: dictionary(uint32, binary)
19//! - `__sequence`, the sequence number of a row. Type: uint64
20//! - `__op_type`, the op type of the row. Type: uint8
21//!
22//! The schema of a parquet file is:
23//! ```text
24//! field 0, field 1, ..., field N, time index, primary key, sequence, op type
25//! ```
26//!
27//! We stores fields in the same order as [RegionMetadata::field_columns()](store_api::metadata::RegionMetadata::field_columns()).
28
29use std::borrow::Borrow;
30use std::collections::{HashMap, VecDeque};
31use std::sync::Arc;
32
33use api::v1::SemanticType;
34use common_time::Timestamp;
35use datafusion_common::ScalarValue;
36use datatypes::arrow::array::{
37    ArrayRef, BinaryArray, BinaryDictionaryBuilder, DictionaryArray, UInt64Array,
38};
39use datatypes::arrow::datatypes::{DataType as ArrowDataType, SchemaRef, UInt32Type};
40use datatypes::arrow::record_batch::RecordBatch;
41use datatypes::prelude::DataType;
42use datatypes::types::json_type::JsonNativeType;
43use datatypes::vectors::Helper;
44use mito_codec::row_converter::{
45    CompositeValues, PrimaryKeyCodec, SortField, build_primary_key_codec_with_fields,
46};
47use parquet::file::metadata::{ParquetMetaData, RowGroupMetaData};
48use parquet::file::statistics::Statistics;
49use snafu::{OptionExt, ResultExt, ensure};
50use store_api::metadata::{ColumnMetadata, RegionMetadataRef};
51use store_api::storage::{ColumnId, NestedPath, SequenceNumber};
52
53use crate::error::{
54    ConvertVectorSnafu, DecodeSnafu, InvalidRecordBatchSnafu, NewRecordBatchSnafu, Result,
55};
56use crate::read::read_columns::ReadColumns;
57use crate::read::{Batch, BatchBuilder, BatchColumn};
58use crate::sst::file::{FileMeta, FileTimeRange};
59use crate::sst::parquet::read_columns::{ParquetReadColumn, ParquetReadColumns};
60use crate::sst::to_sst_arrow_schema;
61
62/// Arrow array type for the primary key dictionary.
63pub(crate) type PrimaryKeyArray = DictionaryArray<UInt32Type>;
64/// Builder type for primary key dictionary array.
65pub(crate) type PrimaryKeyArrayBuilder = BinaryDictionaryBuilder<UInt32Type>;
66
67/// Number of columns that have fixed positions.
68///
69/// Contains: time index and internal columns.
70pub(crate) const FIXED_POS_COLUMN_NUM: usize = 4;
71/// Number of internal columns.
72pub(crate) const INTERNAL_COLUMN_NUM: usize = 3;
73
74/// Helper for writing the SST format with primary key.
75pub(crate) struct PrimaryKeyWriteFormat {
76    /// SST file schema.
77    arrow_schema: SchemaRef,
78    override_sequence: Option<SequenceNumber>,
79}
80
81impl PrimaryKeyWriteFormat {
82    /// Creates a new helper.
83    pub(crate) fn new(metadata: RegionMetadataRef) -> PrimaryKeyWriteFormat {
84        let arrow_schema = to_sst_arrow_schema(&metadata);
85        PrimaryKeyWriteFormat {
86            arrow_schema,
87            override_sequence: None,
88        }
89    }
90
91    /// Set override sequence.
92    pub(crate) fn with_override_sequence(
93        mut self,
94        override_sequence: Option<SequenceNumber>,
95    ) -> Self {
96        self.override_sequence = override_sequence;
97        self
98    }
99
100    /// Gets the arrow schema to store in parquet.
101    #[cfg(test)]
102    pub(crate) fn arrow_schema(&self) -> &SchemaRef {
103        &self.arrow_schema
104    }
105
106    /// Convert a flat `RecordBatch` to primary-key format, retaining only
107    /// field columns, time index, and internal columns.
108    ///
109    /// `num_fields` is the number of field columns. The method strips
110    /// leading tag columns: `num_tag_columns = batch.num_columns() - num_fields - FIXED_POS_COLUMN_NUM`.
111    pub(crate) fn convert_flat_batch(
112        &self,
113        batch: &RecordBatch,
114        num_fields: usize,
115    ) -> Result<RecordBatch> {
116        let num_tag_columns = batch.num_columns() - num_fields - FIXED_POS_COLUMN_NUM;
117        let mut columns: Vec<ArrayRef> = batch.columns()[num_tag_columns..].to_vec();
118
119        if let Some(override_sequence) = self.override_sequence {
120            let num_cols = columns.len();
121            // sequence is at num_cols - 2 (before op_type)
122            columns[num_cols - 2] =
123                Arc::new(UInt64Array::from(vec![override_sequence; batch.num_rows()]));
124        }
125
126        RecordBatch::try_new(self.arrow_schema.clone(), columns).context(NewRecordBatchSnafu)
127    }
128}
129
130/// Returns min/max values of specific columns.
131/// Returns None if the column does not have statistics.
132/// The column should not be encoded as a part of a primary key.
133pub(crate) fn column_values(
134    row_groups: &[impl Borrow<RowGroupMetaData>],
135    column: &ColumnMetadata,
136    column_index: usize,
137    is_min: bool,
138) -> Option<ArrayRef> {
139    column_values_by_type(
140        row_groups,
141        &column.column_schema.data_type.as_arrow_type(),
142        column_index,
143        is_min,
144    )
145}
146
147/// Returns min/max values of a parquet column with the given Arrow data type.
148pub(crate) fn column_values_by_type(
149    row_groups: &[impl Borrow<RowGroupMetaData>],
150    data_type: &ArrowDataType,
151    column_index: usize,
152    is_min: bool,
153) -> Option<ArrayRef> {
154    let null_scalar: ScalarValue = data_type.try_into().ok()?;
155    let scalar_values = row_groups
156        .iter()
157        .map(|meta| {
158            let stats = meta.borrow().column(column_index).statistics()?;
159            match stats {
160                Statistics::Boolean(s) => Some(ScalarValue::Boolean(Some(if is_min {
161                    *s.min_opt()?
162                } else {
163                    *s.max_opt()?
164                }))),
165                Statistics::Int32(s) => {
166                    let value = if is_min { *s.min_opt()? } else { *s.max_opt()? };
167                    if data_type == &ArrowDataType::UInt32 {
168                        Some(ScalarValue::UInt32(Some(value as u32)))
169                    } else {
170                        Some(ScalarValue::Int32(Some(value)))
171                    }
172                }
173                Statistics::Int64(s) => {
174                    let value = if is_min { *s.min_opt()? } else { *s.max_opt()? };
175                    if data_type == &ArrowDataType::UInt64 {
176                        Some(ScalarValue::UInt64(Some(value as u64)))
177                    } else {
178                        Some(ScalarValue::Int64(Some(value)))
179                    }
180                }
181                Statistics::Int96(_) => None,
182                Statistics::Float(s) => Some(ScalarValue::Float32(Some(if is_min {
183                    *s.min_opt()?
184                } else {
185                    *s.max_opt()?
186                }))),
187                Statistics::Double(s) => Some(ScalarValue::Float64(Some(if is_min {
188                    *s.min_opt()?
189                } else {
190                    *s.max_opt()?
191                }))),
192                Statistics::ByteArray(s) => {
193                    let bytes = if is_min {
194                        s.min_bytes_opt()?
195                    } else {
196                        s.max_bytes_opt()?
197                    };
198                    let s = String::from_utf8(bytes.to_vec()).ok();
199                    Some(ScalarValue::Utf8(s))
200                }
201                Statistics::FixedLenByteArray(_) => None,
202            }
203        })
204        .map(|maybe_scalar| maybe_scalar.unwrap_or_else(|| null_scalar.clone()))
205        .collect::<Vec<ScalarValue>>();
206    debug_assert_eq!(scalar_values.len(), row_groups.len());
207    ScalarValue::iter_to_array(scalar_values).ok()
208}
209
210/// Returns null counts of specific columns.
211/// The column should not be encoded as a part of a primary key.
212pub(crate) fn column_null_counts(
213    row_groups: &[impl Borrow<RowGroupMetaData>],
214    column_index: usize,
215) -> Option<ArrayRef> {
216    let values = row_groups.iter().map(|meta| {
217        let col = meta.borrow().column(column_index);
218        let stat = col.statistics()?;
219        stat.null_count_opt()
220    });
221    Some(Arc::new(UInt64Array::from_iter(values)))
222}
223
224/// Helper for reading the SST format.
225pub struct PrimaryKeyReadFormat {
226    /// The metadata stored in the SST.
227    metadata: RegionMetadataRef,
228    /// SST file schema.
229    arrow_schema: SchemaRef,
230    /// Field column id to its index in `schema` (SST schema).
231    /// In SST schema, fields are stored in the front of the schema.
232    field_id_to_index: HashMap<ColumnId, usize>,
233    /// Indices of columns to read from the SST. It contains all internal columns.
234    parquet_read_cols: ParquetReadColumns,
235    /// Field column id to their index in the projected schema (
236    /// the schema of [Batch]).
237    field_id_to_projected_index: HashMap<ColumnId, usize>,
238    /// Codec used to decode primary key values if eager decoding is enabled.
239    primary_key_codec: Option<Arc<dyn PrimaryKeyCodec>>,
240}
241
242impl PrimaryKeyReadFormat {
243    /// Creates a helper with existing `metadata` and `column_ids` to read.
244    pub fn new(metadata: RegionMetadataRef, read_cols: ReadColumns) -> PrimaryKeyReadFormat {
245        let field_id_to_index: HashMap<_, _> = metadata
246            .field_columns()
247            .enumerate()
248            .map(|(index, column)| (column.column_id, index))
249            .collect();
250        let arrow_schema = to_sst_arrow_schema(&metadata);
251
252        let format_projection = FormatProjection::compute_format_projection(
253            &metadata,
254            &field_id_to_index,
255            arrow_schema.fields.len(),
256            read_cols,
257        );
258
259        PrimaryKeyReadFormat {
260            metadata,
261            arrow_schema,
262            field_id_to_index,
263            parquet_read_cols: format_projection.parquet_read_cols,
264            field_id_to_projected_index: format_projection.column_id_to_projected_index,
265            primary_key_codec: None,
266        }
267    }
268
269    /// Gets the arrow schema of the SST file.
270    ///
271    /// This schema is computed from the region metadata but should be the same
272    /// as the arrow schema decoded from the file metadata.
273    pub(crate) fn arrow_schema(&self) -> &SchemaRef {
274        &self.arrow_schema
275    }
276
277    /// Gets the metadata of the SST.
278    pub(crate) fn metadata(&self) -> &RegionMetadataRef {
279        &self.metadata
280    }
281
282    pub(crate) fn parquet_read_columns(&self) -> &ParquetReadColumns {
283        &self.parquet_read_cols
284    }
285
286    /// Gets the field id to projected index.
287    pub(crate) fn field_id_to_projected_index(&self) -> &HashMap<ColumnId, usize> {
288        &self.field_id_to_projected_index
289    }
290
291    /// Convert a arrow record batch into `batches`.
292    ///
293    /// The length of `override_sequence_array` must be larger than the length of the record batch.
294    /// Note that the `record_batch` may only contains a subset of columns if it is projected.
295    pub fn convert_record_batch(
296        &self,
297        record_batch: &RecordBatch,
298        override_sequence_array: Option<&ArrayRef>,
299        batches: &mut VecDeque<Batch>,
300    ) -> Result<()> {
301        debug_assert!(batches.is_empty());
302
303        // The record batch must has time index and internal columns.
304        ensure!(
305            record_batch.num_columns() >= FIXED_POS_COLUMN_NUM,
306            InvalidRecordBatchSnafu {
307                reason: format!(
308                    "record batch only has {} columns",
309                    record_batch.num_columns()
310                ),
311            }
312        );
313
314        let mut fixed_pos_columns = record_batch
315            .columns()
316            .iter()
317            .rev()
318            .take(FIXED_POS_COLUMN_NUM);
319        // Safety: We have checked the column number.
320        let op_type_array = fixed_pos_columns.next().unwrap();
321        let mut sequence_array = fixed_pos_columns.next().unwrap().clone();
322        let pk_array = fixed_pos_columns.next().unwrap();
323        let ts_array = fixed_pos_columns.next().unwrap();
324        let field_batch_columns = self.get_field_batch_columns(record_batch)?;
325
326        // Override sequence array if provided.
327        if let Some(override_array) = override_sequence_array {
328            assert!(override_array.len() >= sequence_array.len());
329            // It's fine to assign the override array directly, but we slice it to make
330            // sure it matches the length of the original sequence array.
331            sequence_array = if override_array.len() > sequence_array.len() {
332                override_array.slice(0, sequence_array.len())
333            } else {
334                override_array.clone()
335            };
336        }
337
338        // Compute primary key offsets.
339        let pk_dict_array = pk_array
340            .as_any()
341            .downcast_ref::<PrimaryKeyArray>()
342            .with_context(|| InvalidRecordBatchSnafu {
343                reason: format!("primary key array should not be {:?}", pk_array.data_type()),
344            })?;
345        let offsets = primary_key_offsets(pk_dict_array)?;
346        if offsets.is_empty() {
347            return Ok(());
348        }
349
350        // Split record batch according to pk offsets.
351        let keys = pk_dict_array.keys();
352        let pk_values = pk_dict_array
353            .values()
354            .as_any()
355            .downcast_ref::<BinaryArray>()
356            .with_context(|| InvalidRecordBatchSnafu {
357                reason: format!(
358                    "values of primary key array should not be {:?}",
359                    pk_dict_array.values().data_type()
360                ),
361            })?;
362        for (i, start) in offsets[..offsets.len() - 1].iter().enumerate() {
363            let end = offsets[i + 1];
364            let rows_in_batch = end - start;
365            let dict_key = keys.value(*start);
366            let primary_key = pk_values.value(dict_key as usize).to_vec();
367
368            let mut builder = BatchBuilder::new(primary_key);
369            builder
370                .timestamps_array(ts_array.slice(*start, rows_in_batch))?
371                .sequences_array(sequence_array.slice(*start, rows_in_batch))?
372                .op_types_array(op_type_array.slice(*start, rows_in_batch))?;
373            // Push all fields
374            for batch_column in &field_batch_columns {
375                builder.push_field(BatchColumn {
376                    column_id: batch_column.column_id,
377                    data: batch_column.data.slice(*start, rows_in_batch),
378                });
379            }
380
381            let mut batch = builder.build()?;
382            if let Some(codec) = &self.primary_key_codec {
383                let pk_values: CompositeValues =
384                    codec.decode(batch.primary_key()).context(DecodeSnafu)?;
385                batch.set_pk_values(pk_values);
386            }
387            batches.push_back(batch);
388        }
389
390        Ok(())
391    }
392
393    /// Returns min values of specific column in row groups.
394    pub fn min_values(
395        &self,
396        row_groups: &[impl Borrow<RowGroupMetaData>],
397        column_id: ColumnId,
398    ) -> StatValues {
399        let Some(column) = self.metadata.column_by_id(column_id) else {
400            // No such column in the SST.
401            return StatValues::NoColumn;
402        };
403        match column.semantic_type {
404            SemanticType::Tag => self.tag_values(row_groups, column, true),
405            SemanticType::Field => {
406                // Safety: `field_id_to_index` is initialized by the semantic type.
407                let index = self.field_id_to_index.get(&column_id).unwrap();
408                let stats = column_values(row_groups, column, *index, true);
409                StatValues::from_stats_opt(stats)
410            }
411            SemanticType::Timestamp => {
412                let index = self.time_index_position();
413                let stats = column_values(row_groups, column, index, true);
414                StatValues::from_stats_opt(stats)
415            }
416        }
417    }
418
419    /// Returns max values of specific column in row groups.
420    pub fn max_values(
421        &self,
422        row_groups: &[impl Borrow<RowGroupMetaData>],
423        column_id: ColumnId,
424    ) -> StatValues {
425        let Some(column) = self.metadata.column_by_id(column_id) else {
426            // No such column in the SST.
427            return StatValues::NoColumn;
428        };
429        match column.semantic_type {
430            SemanticType::Tag => self.tag_values(row_groups, column, false),
431            SemanticType::Field => {
432                // Safety: `field_id_to_index` is initialized by the semantic type.
433                let index = self.field_id_to_index.get(&column_id).unwrap();
434                let stats = column_values(row_groups, column, *index, false);
435                StatValues::from_stats_opt(stats)
436            }
437            SemanticType::Timestamp => {
438                let index = self.time_index_position();
439                let stats = column_values(row_groups, column, index, false);
440                StatValues::from_stats_opt(stats)
441            }
442        }
443    }
444
445    /// Returns null counts of specific column in row groups.
446    pub fn null_counts(
447        &self,
448        row_groups: &[impl Borrow<RowGroupMetaData>],
449        column_id: ColumnId,
450    ) -> StatValues {
451        let Some(column) = self.metadata.column_by_id(column_id) else {
452            // No such column in the SST.
453            return StatValues::NoColumn;
454        };
455        match column.semantic_type {
456            SemanticType::Tag => StatValues::NoStats,
457            SemanticType::Field => {
458                // Safety: `field_id_to_index` is initialized by the semantic type.
459                let index = self.field_id_to_index.get(&column_id).unwrap();
460                let stats = column_null_counts(row_groups, *index);
461                StatValues::from_stats_opt(stats)
462            }
463            SemanticType::Timestamp => {
464                let index = self.time_index_position();
465                let stats = column_null_counts(row_groups, index);
466                StatValues::from_stats_opt(stats)
467            }
468        }
469    }
470
471    /// Get fields from `record_batch`.
472    fn get_field_batch_columns(&self, record_batch: &RecordBatch) -> Result<Vec<BatchColumn>> {
473        record_batch
474            .columns()
475            .iter()
476            .zip(record_batch.schema().fields())
477            .take(record_batch.num_columns() - FIXED_POS_COLUMN_NUM) // Take all field columns.
478            .map(|(array, field)| {
479                let vector = Helper::try_into_vector(array.clone()).context(ConvertVectorSnafu)?;
480                let column = self
481                    .metadata
482                    .column_by_name(field.name())
483                    .with_context(|| InvalidRecordBatchSnafu {
484                        reason: format!("column {} not found in metadata", field.name()),
485                    })?;
486
487                Ok(BatchColumn {
488                    column_id: column.column_id,
489                    data: vector,
490                })
491            })
492            .collect()
493    }
494
495    /// Returns min/max values of specific tag.
496    fn tag_values(
497        &self,
498        row_groups: &[impl Borrow<RowGroupMetaData>],
499        column: &ColumnMetadata,
500        is_min: bool,
501    ) -> StatValues {
502        let is_first_tag = self
503            .metadata
504            .primary_key
505            .first()
506            .map(|id| *id == column.column_id)
507            .unwrap_or(false);
508        if !is_first_tag {
509            // Only the min-max of the first tag is available in the primary key.
510            return StatValues::NoStats;
511        }
512
513        StatValues::from_stats_opt(self.first_tag_values(row_groups, column, is_min))
514    }
515
516    /// Returns min/max values of the first tag.
517    /// Returns None if the tag does not have statistics.
518    fn first_tag_values(
519        &self,
520        row_groups: &[impl Borrow<RowGroupMetaData>],
521        column: &ColumnMetadata,
522        is_min: bool,
523    ) -> Option<ArrayRef> {
524        debug_assert!(
525            self.metadata
526                .primary_key
527                .first()
528                .map(|id| *id == column.column_id)
529                .unwrap_or(false)
530        );
531
532        let primary_key_encoding = self.metadata.primary_key_encoding;
533        let converter = build_primary_key_codec_with_fields(
534            primary_key_encoding,
535            [(
536                column.column_id,
537                SortField::new(column.column_schema.data_type.clone()),
538            )]
539            .into_iter(),
540        );
541
542        let values = row_groups.iter().map(|meta| {
543            let stats = meta
544                .borrow()
545                .column(self.primary_key_position())
546                .statistics()?;
547            match stats {
548                Statistics::Boolean(_) => None,
549                Statistics::Int32(_) => None,
550                Statistics::Int64(_) => None,
551                Statistics::Int96(_) => None,
552                Statistics::Float(_) => None,
553                Statistics::Double(_) => None,
554                Statistics::ByteArray(s) => {
555                    let bytes = if is_min {
556                        s.min_bytes_opt()?
557                    } else {
558                        s.max_bytes_opt()?
559                    };
560                    converter.decode_leftmost(bytes).ok()?
561                }
562                Statistics::FixedLenByteArray(_) => None,
563            }
564        });
565        let mut builder = column
566            .column_schema
567            .data_type
568            .create_mutable_vector(row_groups.len());
569        for value_opt in values {
570            match value_opt {
571                // Safety: We use the same data type to create the converter.
572                Some(v) => builder.push_value_ref(&v.as_value_ref()),
573                None => builder.push_null(),
574            }
575        }
576        let vector = builder.to_vector();
577
578        Some(vector.to_arrow_array())
579    }
580
581    /// Index in SST of the primary key.
582    fn primary_key_position(&self) -> usize {
583        self.arrow_schema.fields.len() - 3
584    }
585
586    /// Index in SST of the time index.
587    fn time_index_position(&self) -> usize {
588        self.arrow_schema.fields.len() - FIXED_POS_COLUMN_NUM
589    }
590
591    /// Index of a field column by its column id.
592    pub fn field_index_by_id(&self, column_id: ColumnId) -> Option<usize> {
593        self.field_id_to_projected_index.get(&column_id).copied()
594    }
595}
596
597/// Helper to compute the projection for the SST.
598pub(crate) struct FormatProjection {
599    /// The columns to read from the SST. It contains all internal columns.
600    pub(crate) parquet_read_cols: ParquetReadColumns,
601    /// Column id to their index in the projected schema (
602    /// the schema after projection).
603    ///
604    /// It doesn't contain time index column if it is not present in the projection.
605    pub(crate) column_id_to_projected_index: HashMap<ColumnId, usize>,
606}
607
608impl FormatProjection {
609    /// Computes the projection.
610    ///
611    /// `id_to_index` is a mapping from column id to the index of the column in the SST.
612    pub(crate) fn compute_format_projection(
613        metadata: &RegionMetadataRef,
614        id_to_index: &HashMap<ColumnId, usize>,
615        sst_column_num: usize,
616        cols: ReadColumns,
617    ) -> Self {
618        let mut projected_columns: Vec<_> = cols
619            .col_ids
620            .iter()
621            .copied()
622            .filter_map(|col_id| {
623                id_to_index.get(&col_id).copied().map(|index_of_sst| {
624                    let nested_paths = json_target_nested_paths(metadata, &cols, col_id);
625                    (col_id, index_of_sst, nested_paths)
626                })
627            })
628            .collect();
629        // Sorts columns by their indices in the SST. SST uses a bitmap for projection.
630        // This ensures the schema of `projected_columns` is the same as the batch returned from the SST.
631        projected_columns.sort_unstable_by_key(|(_, index, _)| *index);
632
633        let mut parquet_read_cols: Vec<ParquetReadColumn> =
634            Vec::with_capacity(projected_columns.len() + FIXED_POS_COLUMN_NUM);
635        // Creates a map from column id to the index of that column in the projected record batch.
636        let mut column_id_to_projected_index = HashMap::with_capacity(projected_columns.len());
637
638        for (col_id, index_of_sst, nested_paths) in projected_columns {
639            Self::merge_or_push_parquet_column(&mut parquet_read_cols, index_of_sst, nested_paths);
640
641            column_id_to_projected_index
642                .entry(col_id)
643                .or_insert_with(|| parquet_read_cols.len() - 1);
644        }
645
646        // In SST schema, fixed-position columns are always in the tail:
647        // `time index, __primary_key, __sequence, __op_type`.
648        Self::append_time_index_if_needed(&mut parquet_read_cols, sst_column_num);
649        Self::append_fixed_internal_columns(&mut parquet_read_cols, sst_column_num);
650
651        Self {
652            parquet_read_cols: ParquetReadColumns::from_deduped(parquet_read_cols),
653            column_id_to_projected_index,
654        }
655    }
656
657    fn merge_or_push_parquet_column(
658        parquet_read_cols: &mut Vec<ParquetReadColumn>,
659        index_of_sst: usize,
660        nested_paths: Vec<Vec<String>>,
661    ) {
662        // `projected_columns` is sorted by parquet root index, so repeated reads
663        // for the same root column are always adjacent.
664        if let Some(last_col) = parquet_read_cols.last_mut()
665            && last_col.root_index() == index_of_sst
666        {
667            last_col.merge_nested_paths(nested_paths);
668            return;
669        }
670
671        let parquet_col = ParquetReadColumn::new(index_of_sst).with_nested_paths(nested_paths);
672        parquet_read_cols.push(parquet_col);
673    }
674
675    fn append_time_index_if_needed(
676        parquet_read_cols: &mut Vec<ParquetReadColumn>,
677        sst_column_num: usize,
678    ) {
679        let time_index = sst_column_num - FIXED_POS_COLUMN_NUM;
680        // Existing projected roots are already sorted by SST root index, and may
681        // already include the time index, so we compare against the last root to
682        // decide whether we still need to append `time index`.
683        let needs_time_index = parquet_read_cols
684            .last()
685            .map(|col| col.root_index() != time_index)
686            .unwrap_or(true);
687        if needs_time_index {
688            parquet_read_cols.push(ParquetReadColumn::new(time_index));
689        }
690    }
691
692    // Append internal columns in fixed order: `__primary_key`, `__sequence`,
693    // `__op_type`.
694    fn append_fixed_internal_columns(
695        parquet_read_cols: &mut Vec<ParquetReadColumn>,
696        sst_column_num: usize,
697    ) {
698        for index in sst_column_num - INTERNAL_COLUMN_NUM..sst_column_num {
699            parquet_read_cols.push(ParquetReadColumn::new(index));
700        }
701    }
702}
703
704fn json_target_nested_paths(
705    metadata: &RegionMetadataRef,
706    read_columns: &ReadColumns,
707    column_id: ColumnId,
708) -> Vec<NestedPath> {
709    let Some(target_type) = read_columns.json_target_type(column_id) else {
710        return Vec::new();
711    };
712    let Some(column) = metadata.column_by_id(column_id) else {
713        return Vec::new();
714    };
715
716    json_nested_paths(&column.column_schema.name, target_type)
717}
718
719fn json_nested_paths(column_name: &str, json_type: &JsonNativeType) -> Vec<NestedPath> {
720    let mut paths = Vec::new();
721    let mut current = vec![column_name.to_string()];
722    collect_json_nested_paths(json_type, &mut current, &mut paths);
723    paths
724}
725
726fn collect_json_nested_paths(
727    json_type: &JsonNativeType,
728    current: &mut NestedPath,
729    paths: &mut Vec<NestedPath>,
730) {
731    match json_type {
732        JsonNativeType::Object(fields) if !fields.is_empty() => {
733            for (field, child) in fields {
734                current.push(field.clone());
735                collect_json_nested_paths(child, current, paths);
736                current.pop();
737            }
738        }
739        _ => paths.push(current.clone()),
740    }
741}
742
743/// Values of column statistics of the SST.
744///
745/// It also distinguishes the case that a column is not found and
746/// the column exists but has no statistics.
747pub enum StatValues {
748    /// Values of each row group.
749    Values(ArrayRef),
750    /// No such column.
751    NoColumn,
752    /// Column exists but has no statistics.
753    NoStats,
754}
755
756impl StatValues {
757    /// Creates a new `StatValues` instance from optional statistics.
758    pub fn from_stats_opt(stats: Option<ArrayRef>) -> Self {
759        match stats {
760            Some(stats) => StatValues::Values(stats),
761            None => StatValues::NoStats,
762        }
763    }
764}
765
766#[cfg(test)]
767impl PrimaryKeyReadFormat {
768    /// Creates a helper with existing `metadata` and all columns.
769    pub fn new_with_all_columns(metadata: RegionMetadataRef) -> PrimaryKeyReadFormat {
770        Self::new(
771            Arc::clone(&metadata),
772            ReadColumns::new(metadata.column_metadatas.iter().map(|c| c.column_id)),
773        )
774    }
775}
776
777/// Compute offsets of different primary keys in the array.
778pub(crate) fn primary_key_offsets(pk_dict_array: &PrimaryKeyArray) -> Result<Vec<usize>> {
779    if pk_dict_array.is_empty() {
780        return Ok(Vec::new());
781    }
782
783    // Init offsets.
784    let mut offsets = vec![0];
785    let keys = pk_dict_array.keys();
786    // We know that primary keys are always not null so we iterate `keys.values()` directly.
787    let pk_indices = keys.values();
788    for (i, key) in pk_indices.iter().take(keys.len() - 1).enumerate() {
789        // Compare each key with next key
790        if *key != pk_indices[i + 1] {
791            // We meet a new key, push the next index as end of the offset.
792            offsets.push(i + 1);
793        }
794    }
795    offsets.push(keys.len());
796
797    Ok(offsets)
798}
799
800/// Gets the min/max time index of the row group from the parquet meta.
801/// It assumes the parquet is created by the mito engine.
802pub(crate) fn parquet_row_group_time_range(
803    file_meta: &FileMeta,
804    parquet_meta: &ParquetMetaData,
805    row_group_idx: usize,
806) -> Option<FileTimeRange> {
807    let row_group_meta = parquet_meta.row_group(row_group_idx);
808    let num_columns = parquet_meta.file_metadata().schema_descr().num_columns();
809    assert!(
810        num_columns >= FIXED_POS_COLUMN_NUM,
811        "file only has {} columns",
812        num_columns
813    );
814    let time_index_pos = num_columns - FIXED_POS_COLUMN_NUM;
815
816    let stats = row_group_meta.column(time_index_pos).statistics()?;
817    // The physical type for the timestamp should be i64.
818    let (min, max) = match stats {
819        Statistics::Int64(value_stats) => (*value_stats.min_opt()?, *value_stats.max_opt()?),
820        Statistics::Int32(_)
821        | Statistics::Boolean(_)
822        | Statistics::Int96(_)
823        | Statistics::Float(_)
824        | Statistics::Double(_)
825        | Statistics::ByteArray(_)
826        | Statistics::FixedLenByteArray(_) => {
827            common_telemetry::warn!(
828                "Invalid statistics {:?} for time index in parquet in {}",
829                stats,
830                file_meta.file_id
831            );
832            return None;
833        }
834    };
835
836    debug_assert!(min >= file_meta.time_range.0.value() && min <= file_meta.time_range.1.value());
837    debug_assert!(max >= file_meta.time_range.0.value() && max <= file_meta.time_range.1.value());
838    let unit = file_meta.time_range.0.unit();
839
840    Some((Timestamp::new(min, unit), Timestamp::new(max, unit)))
841}
842
843/// Checks if sequence override is needed based on all row groups' statistics.
844/// Returns true if ALL row groups have sequence min-max values of 0.
845pub(crate) fn need_override_sequence(parquet_meta: &ParquetMetaData) -> bool {
846    let num_columns = parquet_meta.file_metadata().schema_descr().num_columns();
847    if num_columns < FIXED_POS_COLUMN_NUM {
848        return false;
849    }
850
851    // The sequence column is the second-to-last column (before op_type)
852    let sequence_pos = num_columns - 2;
853
854    // Check all row groups - all must have sequence min-max of 0
855    for row_group in parquet_meta.row_groups() {
856        if let Some(Statistics::Int64(value_stats)) = row_group.column(sequence_pos).statistics() {
857            if let (Some(min_val), Some(max_val)) = (value_stats.min_opt(), value_stats.max_opt()) {
858                // If any row group doesn't have min=0 and max=0, return false
859                if *min_val != 0 || *max_val != 0 {
860                    return false;
861                }
862            } else {
863                // If any row group doesn't have statistics, return false
864                return false;
865            }
866        } else {
867            // If any row group doesn't have Int64 statistics, return false
868            return false;
869        }
870    }
871
872    // All row groups have sequence min-max of 0, or there are no row groups
873    !parquet_meta.row_groups().is_empty()
874}
875
876#[cfg(test)]
877mod tests {
878    use std::collections::BTreeMap;
879    use std::sync::Arc;
880
881    use api::v1::OpType;
882    use datatypes::arrow::array::{
883        Int64Array, StringArray, TimestampMillisecondArray, UInt8Array, UInt32Array, UInt64Array,
884    };
885    use datatypes::arrow::datatypes::{DataType as ArrowDataType, Field, Schema, TimeUnit};
886    use datatypes::prelude::ConcreteDataType;
887    use datatypes::schema::ColumnSchema;
888    use datatypes::types::json_type::{JsonNativeType, JsonObjectType};
889    use datatypes::value::ValueRef;
890    use datatypes::vectors::{Int64Vector, TimestampMillisecondVector, UInt8Vector, UInt64Vector};
891    use mito_codec::row_converter::{
892        DensePrimaryKeyCodec, PrimaryKeyCodec, PrimaryKeyCodecExt, SparsePrimaryKeyCodec,
893    };
894    use store_api::codec::PrimaryKeyEncoding;
895    use store_api::metadata::{ColumnMetadata, RegionMetadataBuilder};
896    use store_api::storage::RegionId;
897    use store_api::storage::consts::ReservedColumnId;
898
899    use super::*;
900    use crate::error::InvalidMetadataSnafu;
901    use crate::sst::parquet::flat_format::{
902        FlatReadFormat, FlatWriteFormat, sequence_column_index, sst_column_id_indices,
903    };
904    use crate::sst::{
905        FlatSchemaOptions, OP_TYPE_PARQUET_FIELD_ID, PRIMARY_KEY_PARQUET_FIELD_ID,
906        SEQUENCE_PARQUET_FIELD_ID, to_flat_sst_arrow_schema, with_field_id,
907    };
908
909    const TEST_SEQUENCE: u64 = 1;
910    const TEST_OP_TYPE: u8 = OpType::Put as u8;
911
912    fn build_test_region_metadata() -> RegionMetadataRef {
913        let mut builder = RegionMetadataBuilder::new(RegionId::new(1, 1));
914        builder
915            .push_column_metadata(ColumnMetadata {
916                column_schema: ColumnSchema::new("tag0", ConcreteDataType::int64_datatype(), true),
917                semantic_type: SemanticType::Tag,
918                column_id: 1,
919            })
920            .push_column_metadata(ColumnMetadata {
921                column_schema: ColumnSchema::new(
922                    "field1",
923                    ConcreteDataType::int64_datatype(),
924                    true,
925                ),
926                semantic_type: SemanticType::Field,
927                column_id: 4, // We change the order of fields columns.
928            })
929            .push_column_metadata(ColumnMetadata {
930                column_schema: ColumnSchema::new("tag1", ConcreteDataType::int64_datatype(), true),
931                semantic_type: SemanticType::Tag,
932                column_id: 3,
933            })
934            .push_column_metadata(ColumnMetadata {
935                column_schema: ColumnSchema::new(
936                    "field0",
937                    ConcreteDataType::int64_datatype(),
938                    true,
939                ),
940                semantic_type: SemanticType::Field,
941                column_id: 2,
942            })
943            .push_column_metadata(ColumnMetadata {
944                column_schema: ColumnSchema::new(
945                    "ts",
946                    ConcreteDataType::timestamp_millisecond_datatype(),
947                    false,
948                ),
949                semantic_type: SemanticType::Timestamp,
950                column_id: 5,
951            })
952            .primary_key(vec![1, 3]);
953        Arc::new(builder.build().unwrap())
954    }
955
956    fn build_test_arrow_schema() -> SchemaRef {
957        let fields = vec![
958            make_field("field1", ArrowDataType::Int64, true, Some(4)),
959            make_field("field0", ArrowDataType::Int64, true, Some(2)),
960            make_field(
961                "ts",
962                ArrowDataType::Timestamp(TimeUnit::Millisecond, None),
963                false,
964                Some(5),
965            ),
966            make_field(
967                "__primary_key",
968                ArrowDataType::Dictionary(
969                    Box::new(ArrowDataType::UInt32),
970                    Box::new(ArrowDataType::Binary),
971                ),
972                false,
973                Some(PRIMARY_KEY_PARQUET_FIELD_ID),
974            ),
975            make_field(
976                "__sequence",
977                ArrowDataType::UInt64,
978                false,
979                Some(SEQUENCE_PARQUET_FIELD_ID),
980            ),
981            make_field(
982                "__op_type",
983                ArrowDataType::UInt8,
984                false,
985                Some(OP_TYPE_PARQUET_FIELD_ID),
986            ),
987        ];
988        Arc::new(Schema::new(fields))
989    }
990
991    fn new_batch(primary_key: &[u8], start_ts: i64, start_field: i64, num_rows: usize) -> Batch {
992        new_batch_with_sequence(primary_key, start_ts, start_field, num_rows, TEST_SEQUENCE)
993    }
994
995    fn new_batch_with_sequence(
996        primary_key: &[u8],
997        start_ts: i64,
998        start_field: i64,
999        num_rows: usize,
1000        sequence: u64,
1001    ) -> Batch {
1002        let ts_values = (0..num_rows).map(|i| start_ts + i as i64);
1003        let timestamps = Arc::new(TimestampMillisecondVector::from_values(ts_values));
1004        let sequences = Arc::new(UInt64Vector::from_vec(vec![sequence; num_rows]));
1005        let op_types = Arc::new(UInt8Vector::from_vec(vec![TEST_OP_TYPE; num_rows]));
1006        let fields = vec![
1007            BatchColumn {
1008                column_id: 4,
1009                data: Arc::new(Int64Vector::from_vec(vec![start_field; num_rows])),
1010            }, // field1
1011            BatchColumn {
1012                column_id: 2,
1013                data: Arc::new(Int64Vector::from_vec(vec![start_field + 1; num_rows])),
1014            }, // field0
1015        ];
1016
1017        BatchBuilder::with_required_columns(primary_key.to_vec(), timestamps, sequences, op_types)
1018            .with_fields(fields)
1019            .build()
1020            .unwrap()
1021    }
1022
1023    #[test]
1024    fn test_to_sst_arrow_schema() {
1025        let metadata = build_test_region_metadata();
1026        let write_format = PrimaryKeyWriteFormat::new(metadata);
1027        assert_eq!(&build_test_arrow_schema(), write_format.arrow_schema());
1028    }
1029
1030    fn build_test_pk_array(pk_row_nums: &[(Vec<u8>, usize)]) -> Arc<PrimaryKeyArray> {
1031        let values = Arc::new(BinaryArray::from_iter_values(
1032            pk_row_nums.iter().map(|v| &v.0),
1033        ));
1034        let mut keys = vec![];
1035        for (index, num_rows) in pk_row_nums.iter().map(|v| v.1).enumerate() {
1036            keys.extend(std::iter::repeat_n(index as u32, num_rows));
1037        }
1038        let keys = UInt32Array::from(keys);
1039        Arc::new(DictionaryArray::new(keys, values))
1040    }
1041
1042    #[test]
1043    fn test_projection_indices() {
1044        let metadata = build_test_region_metadata();
1045        // Only read tag1
1046        let read_format = PrimaryKeyReadFormat::new(metadata.clone(), ReadColumns::new([3]));
1047        assert_eq!(
1048            &[2, 3, 4, 5],
1049            read_format.parquet_read_columns().root_indices()
1050        );
1051        // Only read field1
1052        let read_format = PrimaryKeyReadFormat::new(metadata.clone(), ReadColumns::new([4]));
1053        assert_eq!(
1054            &[0, 2, 3, 4, 5],
1055            read_format.parquet_read_columns().root_indices()
1056        );
1057        // Only read ts
1058        let read_format = PrimaryKeyReadFormat::new(metadata.clone(), ReadColumns::new([5]));
1059        assert_eq!(
1060            &[2, 3, 4, 5],
1061            read_format.parquet_read_columns().root_indices()
1062        );
1063        // Read field0, tag0, ts
1064        let read_format = PrimaryKeyReadFormat::new(metadata, ReadColumns::new([2, 1, 5]));
1065        assert_eq!(
1066            &[1, 2, 3, 4, 5],
1067            read_format.parquet_read_columns().root_indices()
1068        );
1069    }
1070
1071    #[test]
1072    fn test_format_projection_preserves_nested_paths() -> Result<()> {
1073        let mut builder = RegionMetadataBuilder::new(RegionId::new(1, 1));
1074        builder
1075            .push_column_metadata(ColumnMetadata {
1076                column_schema: ColumnSchema::new("tag0", ConcreteDataType::string_datatype(), true),
1077                semantic_type: SemanticType::Tag,
1078                column_id: 1,
1079            })
1080            .push_column_metadata(ColumnMetadata {
1081                column_schema: ColumnSchema::new(
1082                    "j",
1083                    ConcreteDataType::json2(JsonNativeType::Object(JsonObjectType::from([
1084                        ("a".to_string(), JsonNativeType::i64()),
1085                        ("b".to_string(), JsonNativeType::String),
1086                    ]))),
1087                    true,
1088                ),
1089                semantic_type: SemanticType::Field,
1090                column_id: 4,
1091            })
1092            .push_column_metadata(ColumnMetadata {
1093                column_schema: ColumnSchema::new(
1094                    "ts",
1095                    ConcreteDataType::timestamp_millisecond_datatype(),
1096                    false,
1097                ),
1098                semantic_type: SemanticType::Timestamp,
1099                column_id: 5,
1100            })
1101            .primary_key(vec![1]);
1102        let metadata = Arc::new(builder.build().context(InvalidMetadataSnafu)?);
1103        let column_id_to_parquet_index = sst_column_id_indices(&metadata);
1104        let projection = FormatProjection::compute_format_projection(
1105            &metadata,
1106            &column_id_to_parquet_index,
1107            metadata.column_metadatas.len() + FIXED_POS_COLUMN_NUM,
1108            ReadColumns::new([4]).with_json_target_types(BTreeMap::from([(
1109                4,
1110                JsonNativeType::Object(JsonObjectType::from([(
1111                    "a".to_string(),
1112                    JsonNativeType::i64(),
1113                )])),
1114            )])),
1115        );
1116
1117        let columns = projection.parquet_read_cols.columns();
1118        assert_eq!(1, columns[0].root_index());
1119        assert_eq!(
1120            &[vec!["j".to_string(), "a".to_string()]],
1121            columns[0].nested_paths()
1122        );
1123        Ok(())
1124    }
1125
1126    #[test]
1127    fn test_empty_primary_key_offsets() {
1128        let array = build_test_pk_array(&[]);
1129        assert!(primary_key_offsets(&array).unwrap().is_empty());
1130    }
1131
1132    #[test]
1133    fn test_primary_key_offsets_one_series() {
1134        let array = build_test_pk_array(&[(b"one".to_vec(), 1)]);
1135        assert_eq!(vec![0, 1], primary_key_offsets(&array).unwrap());
1136
1137        let array = build_test_pk_array(&[(b"one".to_vec(), 1), (b"two".to_vec(), 1)]);
1138        assert_eq!(vec![0, 1, 2], primary_key_offsets(&array).unwrap());
1139
1140        let array = build_test_pk_array(&[
1141            (b"one".to_vec(), 1),
1142            (b"two".to_vec(), 1),
1143            (b"three".to_vec(), 1),
1144        ]);
1145        assert_eq!(vec![0, 1, 2, 3], primary_key_offsets(&array).unwrap());
1146    }
1147
1148    #[test]
1149    fn test_primary_key_offsets_multi_series() {
1150        let array = build_test_pk_array(&[(b"one".to_vec(), 1), (b"two".to_vec(), 3)]);
1151        assert_eq!(vec![0, 1, 4], primary_key_offsets(&array).unwrap());
1152
1153        let array = build_test_pk_array(&[(b"one".to_vec(), 3), (b"two".to_vec(), 1)]);
1154        assert_eq!(vec![0, 3, 4], primary_key_offsets(&array).unwrap());
1155
1156        let array = build_test_pk_array(&[(b"one".to_vec(), 3), (b"two".to_vec(), 3)]);
1157        assert_eq!(vec![0, 3, 6], primary_key_offsets(&array).unwrap());
1158    }
1159
1160    #[test]
1161    fn test_convert_empty_record_batch() {
1162        let metadata = build_test_region_metadata();
1163        let arrow_schema = build_test_arrow_schema();
1164        let column_ids: Vec<_> = metadata
1165            .column_metadatas
1166            .iter()
1167            .map(|col| col.column_id)
1168            .collect();
1169        let read_format = PrimaryKeyReadFormat::new(metadata, ReadColumns::new(column_ids));
1170        assert_eq!(arrow_schema, *read_format.arrow_schema());
1171
1172        let record_batch = RecordBatch::new_empty(arrow_schema);
1173        let mut batches = VecDeque::new();
1174        read_format
1175            .convert_record_batch(&record_batch, None, &mut batches)
1176            .unwrap();
1177        assert!(batches.is_empty());
1178    }
1179
1180    #[test]
1181    fn test_convert_record_batch() {
1182        let metadata = build_test_region_metadata();
1183        let column_ids: Vec<_> = metadata
1184            .column_metadatas
1185            .iter()
1186            .map(|col| col.column_id)
1187            .collect();
1188        let read_format = PrimaryKeyReadFormat::new(metadata, ReadColumns::new(column_ids));
1189
1190        let columns: Vec<ArrayRef> = vec![
1191            Arc::new(Int64Array::from(vec![1, 1, 10, 10])), // field1
1192            Arc::new(Int64Array::from(vec![2, 2, 11, 11])), // field0
1193            Arc::new(TimestampMillisecondArray::from(vec![1, 2, 11, 12])), // ts
1194            build_test_pk_array(&[(b"one".to_vec(), 2), (b"two".to_vec(), 2)]), // primary key
1195            Arc::new(UInt64Array::from(vec![TEST_SEQUENCE; 4])), // sequence
1196            Arc::new(UInt8Array::from(vec![TEST_OP_TYPE; 4])), // op type
1197        ];
1198        let arrow_schema = build_test_arrow_schema();
1199        let record_batch = RecordBatch::try_new(arrow_schema, columns).unwrap();
1200        let mut batches = VecDeque::new();
1201        read_format
1202            .convert_record_batch(&record_batch, None, &mut batches)
1203            .unwrap();
1204
1205        assert_eq!(
1206            vec![new_batch(b"one", 1, 1, 2), new_batch(b"two", 11, 10, 2)],
1207            batches.into_iter().collect::<Vec<_>>(),
1208        );
1209    }
1210
1211    #[test]
1212    fn test_convert_record_batch_with_override_sequence() {
1213        let metadata = build_test_region_metadata();
1214        let read_format = PrimaryKeyReadFormat::new(
1215            metadata.clone(),
1216            ReadColumns::new(metadata.column_metadatas.iter().map(|c| c.column_id)),
1217        );
1218
1219        let columns: Vec<ArrayRef> = vec![
1220            Arc::new(Int64Array::from(vec![1, 1, 10, 10])), // field1
1221            Arc::new(Int64Array::from(vec![2, 2, 11, 11])), // field0
1222            Arc::new(TimestampMillisecondArray::from(vec![1, 2, 11, 12])), // ts
1223            build_test_pk_array(&[(b"one".to_vec(), 2), (b"two".to_vec(), 2)]), // primary key
1224            Arc::new(UInt64Array::from(vec![TEST_SEQUENCE; 4])), // sequence
1225            Arc::new(UInt8Array::from(vec![TEST_OP_TYPE; 4])), // op type
1226        ];
1227        let arrow_schema = build_test_arrow_schema();
1228        let record_batch = RecordBatch::try_new(arrow_schema, columns).unwrap();
1229
1230        // Create override sequence array with custom values
1231        let override_sequence: u64 = 12345;
1232        let override_sequence_array: ArrayRef =
1233            Arc::new(UInt64Array::from_value(override_sequence, 4));
1234
1235        let mut batches = VecDeque::new();
1236        read_format
1237            .convert_record_batch(&record_batch, Some(&override_sequence_array), &mut batches)
1238            .unwrap();
1239
1240        // Create expected batches with override sequence
1241        let expected_batch1 = new_batch_with_sequence(b"one", 1, 1, 2, override_sequence);
1242        let expected_batch2 = new_batch_with_sequence(b"two", 11, 10, 2, override_sequence);
1243
1244        assert_eq!(
1245            vec![expected_batch1, expected_batch2],
1246            batches.into_iter().collect::<Vec<_>>(),
1247        );
1248    }
1249
1250    fn make_field(name: &str, dt: ArrowDataType, nullable: bool, field_id: Option<u32>) -> Field {
1251        let mut field = Field::new(name, dt, nullable);
1252        if let Some(id) = field_id {
1253            field = with_field_id(field, id);
1254        }
1255        field
1256    }
1257
1258    fn build_test_flat_sst_schema() -> SchemaRef {
1259        let fields = vec![
1260            Field::new("tag0", ArrowDataType::Int64, true), // primary key columns first
1261            Field::new("tag1", ArrowDataType::Int64, true),
1262            Field::new("field1", ArrowDataType::Int64, true), // then field columns
1263            Field::new("field0", ArrowDataType::Int64, true),
1264            Field::new(
1265                "ts",
1266                ArrowDataType::Timestamp(TimeUnit::Millisecond, None),
1267                false,
1268            ),
1269            Field::new(
1270                "__primary_key",
1271                ArrowDataType::Dictionary(
1272                    Box::new(ArrowDataType::UInt32),
1273                    Box::new(ArrowDataType::Binary),
1274                ),
1275                false,
1276            ),
1277            Field::new("__sequence", ArrowDataType::UInt64, false),
1278            Field::new("__op_type", ArrowDataType::UInt8, false),
1279        ];
1280        Arc::new(Schema::new(fields))
1281    }
1282
1283    fn build_test_flat_sst_schema_with_field_ids() -> SchemaRef {
1284        let ids = [
1285            Some(1u32),
1286            Some(3),
1287            Some(4),
1288            Some(2),
1289            Some(5),
1290            Some(PRIMARY_KEY_PARQUET_FIELD_ID),
1291            Some(SEQUENCE_PARQUET_FIELD_ID),
1292            Some(OP_TYPE_PARQUET_FIELD_ID),
1293        ];
1294        let fields: Vec<_> = build_test_flat_sst_schema()
1295            .fields()
1296            .iter()
1297            .zip(ids)
1298            .map(|(f, id)| match id {
1299                Some(id) => Arc::new(with_field_id((**f).clone(), id)) as _,
1300                None => f.clone(),
1301            })
1302            .collect();
1303        Arc::new(Schema::new(fields))
1304    }
1305
1306    #[test]
1307    fn test_flat_to_sst_arrow_schema() {
1308        let metadata = build_test_region_metadata();
1309        let format = FlatWriteFormat::new(metadata, &FlatSchemaOptions::default());
1310        assert_eq!(
1311            &build_test_flat_sst_schema_with_field_ids(),
1312            format.arrow_schema()
1313        );
1314    }
1315
1316    fn input_columns_for_flat_batch(num_rows: usize) -> Vec<ArrayRef> {
1317        vec![
1318            Arc::new(Int64Array::from(vec![1; num_rows])), // tag0
1319            Arc::new(Int64Array::from(vec![1; num_rows])), // tag1
1320            Arc::new(Int64Array::from(vec![2; num_rows])), // field1
1321            Arc::new(Int64Array::from(vec![3; num_rows])), // field0
1322            Arc::new(TimestampMillisecondArray::from(vec![1, 2, 3, 4])), // ts
1323            build_test_pk_array(&[(b"test".to_vec(), num_rows)]), // __primary_key
1324            Arc::new(UInt64Array::from(vec![TEST_SEQUENCE; num_rows])), // sequence
1325            Arc::new(UInt8Array::from(vec![TEST_OP_TYPE; num_rows])), // op type
1326        ]
1327    }
1328
1329    #[test]
1330    fn test_flat_convert_batch() {
1331        let metadata = build_test_region_metadata();
1332        let format = FlatWriteFormat::new(metadata, &FlatSchemaOptions::default());
1333
1334        let num_rows = 4;
1335        let columns: Vec<ArrayRef> = input_columns_for_flat_batch(num_rows);
1336        let batch =
1337            RecordBatch::try_new(build_test_flat_sst_schema_with_field_ids(), columns.clone())
1338                .unwrap();
1339        let expect_record =
1340            RecordBatch::try_new(build_test_flat_sst_schema_with_field_ids(), columns).unwrap();
1341
1342        let actual = format.convert_batch(&batch).unwrap();
1343        assert_eq!(expect_record, actual);
1344    }
1345
1346    #[test]
1347    fn test_flat_convert_with_override_sequence() {
1348        let metadata = build_test_region_metadata();
1349        let format = FlatWriteFormat::new(metadata, &FlatSchemaOptions::default())
1350            .with_override_sequence(Some(415411));
1351
1352        let num_rows = 4;
1353        let columns: Vec<ArrayRef> = input_columns_for_flat_batch(num_rows);
1354        let batch =
1355            RecordBatch::try_new(build_test_flat_sst_schema_with_field_ids(), columns).unwrap();
1356
1357        let expected_columns: Vec<ArrayRef> = vec![
1358            Arc::new(Int64Array::from(vec![1; num_rows])), // tag0
1359            Arc::new(Int64Array::from(vec![1; num_rows])), // tag1
1360            Arc::new(Int64Array::from(vec![2; num_rows])), // field1
1361            Arc::new(Int64Array::from(vec![3; num_rows])), // field0
1362            Arc::new(TimestampMillisecondArray::from(vec![1, 2, 3, 4])), // ts
1363            build_test_pk_array(&[(b"test".to_vec(), num_rows)]), // __primary_key
1364            Arc::new(UInt64Array::from(vec![415411; num_rows])), // overridden sequence
1365            Arc::new(UInt8Array::from(vec![TEST_OP_TYPE; num_rows])), // op type
1366        ];
1367        let expected_record = RecordBatch::try_new(
1368            build_test_flat_sst_schema_with_field_ids(),
1369            expected_columns,
1370        )
1371        .unwrap();
1372
1373        let actual = format.convert_batch(&batch).unwrap();
1374        assert_eq!(expected_record, actual);
1375    }
1376
1377    #[test]
1378    fn test_flat_projection_indices() {
1379        let metadata = build_test_region_metadata();
1380        // Based on flat format: tag0(0), tag1(1), field1(2), field0(3), ts(4), __primary_key(5), __sequence(6), __op_type(7)
1381        // The projection includes all "fixed position" columns: ts(4), __primary_key(5), __sequence(6), __op_type(7)
1382
1383        // Only read tag1 (column_id=3, index=1) + fixed columns
1384        let read_format =
1385            FlatReadFormat::new(metadata.clone(), ReadColumns::new([3]), None, "test", false)
1386                .unwrap();
1387        assert_eq!(
1388            &[1, 4, 5, 6, 7],
1389            read_format.parquet_read_columns().root_indices()
1390        );
1391
1392        // Only read field1 (column_id=4, index=2) + fixed columns
1393        let read_format =
1394            FlatReadFormat::new(metadata.clone(), ReadColumns::new([4]), None, "test", false)
1395                .unwrap();
1396        assert_eq!(
1397            &[2, 4, 5, 6, 7],
1398            read_format.parquet_read_columns().root_indices()
1399        );
1400
1401        // Only read ts (column_id=5, index=4) + fixed columns (ts is already included in fixed)
1402        let read_format =
1403            FlatReadFormat::new(metadata.clone(), ReadColumns::new([5]), None, "test", false)
1404                .unwrap();
1405        assert_eq!(
1406            &[4, 5, 6, 7],
1407            read_format.parquet_read_columns().root_indices()
1408        );
1409
1410        // Read field0(column_id=2, index=3), tag0(column_id=1, index=0), ts(column_id=5, index=4) + fixed columns
1411        let read_format =
1412            FlatReadFormat::new(metadata, ReadColumns::new([2, 1, 5]), None, "test", false)
1413                .unwrap();
1414        assert_eq!(
1415            &[0, 3, 4, 5, 6, 7],
1416            read_format.parquet_read_columns().root_indices()
1417        );
1418    }
1419
1420    #[test]
1421    fn test_flat_read_format_convert_batch() {
1422        let metadata = build_test_region_metadata();
1423        let mut format = FlatReadFormat::new(
1424            metadata,
1425            ReadColumns::new(std::iter::once(1)), // Just read tag0
1426            Some(build_test_flat_sst_schema()),
1427            "test",
1428            false,
1429        )
1430        .unwrap();
1431
1432        let num_rows = 4;
1433        let original_sequence = 100u64;
1434        let override_sequence = 200u64;
1435
1436        // Create a test record batch
1437        let columns: Vec<ArrayRef> = input_columns_for_flat_batch(num_rows);
1438        let mut test_columns = columns.clone();
1439        // Replace sequence column with original sequence values
1440        test_columns[6] = Arc::new(UInt64Array::from(vec![original_sequence; num_rows]));
1441        let record_batch =
1442            RecordBatch::try_new(format.arrow_schema().clone(), test_columns).unwrap();
1443
1444        // Test without override sequence - should return clone
1445        let result = format.convert_batch(record_batch.clone(), None).unwrap();
1446        let sequence_column = result.column(sequence_column_index(result.num_columns()));
1447        let sequence_array = sequence_column
1448            .as_any()
1449            .downcast_ref::<UInt64Array>()
1450            .unwrap();
1451
1452        let expected_original = UInt64Array::from(vec![original_sequence; num_rows]);
1453        assert_eq!(sequence_array, &expected_original);
1454
1455        // Set override sequence and test with new_override_sequence_array
1456        format.set_override_sequence(Some(override_sequence));
1457        let override_sequence_array = format.new_override_sequence_array(num_rows).unwrap();
1458        let result = format
1459            .convert_batch(record_batch, Some(&override_sequence_array))
1460            .unwrap();
1461        let sequence_column = result.column(sequence_column_index(result.num_columns()));
1462        let sequence_array = sequence_column
1463            .as_any()
1464            .downcast_ref::<UInt64Array>()
1465            .unwrap();
1466
1467        let expected_override = UInt64Array::from(vec![override_sequence; num_rows]);
1468        assert_eq!(sequence_array, &expected_override);
1469    }
1470
1471    #[test]
1472    fn test_need_convert_to_flat() {
1473        let metadata = build_test_region_metadata();
1474
1475        // Test case 1: Same number of columns, no conversion needed
1476        // For flat format: all columns (5) + internal columns (3)
1477        let expected_columns = metadata.column_metadatas.len() + 3;
1478        let result =
1479            FlatReadFormat::is_legacy_format(&metadata, expected_columns, "test.parquet").unwrap();
1480        assert!(
1481            !result,
1482            "Should not need conversion when column counts match"
1483        );
1484
1485        // Test case 2: Different number of columns, need conversion
1486        // Missing primary key columns (2 primary keys in test metadata)
1487        let num_columns_without_pk = expected_columns - metadata.primary_key.len();
1488        let result =
1489            FlatReadFormat::is_legacy_format(&metadata, num_columns_without_pk, "test.parquet")
1490                .unwrap();
1491        assert!(
1492            result,
1493            "Should need conversion when primary key columns are missing"
1494        );
1495
1496        // Test case 3: Invalid case - actual columns more than expected
1497        let too_many_columns = expected_columns + 1;
1498        let err = FlatReadFormat::is_legacy_format(&metadata, too_many_columns, "test.parquet")
1499            .unwrap_err();
1500        assert!(err.to_string().contains("Expected columns"), "{err:?}");
1501
1502        // Test case 4: Invalid case - column difference doesn't match primary key count
1503        let wrong_diff_columns = expected_columns - 1; // Difference of 1, but we have 2 primary keys
1504        let err = FlatReadFormat::is_legacy_format(&metadata, wrong_diff_columns, "test.parquet")
1505            .unwrap_err();
1506        assert!(
1507            err.to_string().contains("Column number difference"),
1508            "{err:?}"
1509        );
1510    }
1511
1512    fn build_test_dense_pk_array(
1513        codec: &DensePrimaryKeyCodec,
1514        pk_values_per_row: &[&[Option<i64>]],
1515    ) -> Arc<PrimaryKeyArray> {
1516        let mut builder = PrimaryKeyArrayBuilder::with_capacity(pk_values_per_row.len(), 1024, 0);
1517
1518        for pk_values_row in pk_values_per_row {
1519            let values: Vec<ValueRef> = pk_values_row
1520                .iter()
1521                .map(|opt| match opt {
1522                    Some(val) => ValueRef::Int64(*val),
1523                    None => ValueRef::Null,
1524                })
1525                .collect();
1526
1527            let encoded = codec.encode(values.into_iter()).unwrap();
1528            builder.append_value(&encoded);
1529        }
1530
1531        Arc::new(builder.finish())
1532    }
1533
1534    fn build_test_sparse_region_metadata() -> RegionMetadataRef {
1535        let mut builder = RegionMetadataBuilder::new(RegionId::new(1, 1));
1536        builder
1537            .push_column_metadata(ColumnMetadata {
1538                column_schema: ColumnSchema::new(
1539                    "__table_id",
1540                    ConcreteDataType::uint32_datatype(),
1541                    false,
1542                ),
1543                semantic_type: SemanticType::Tag,
1544                column_id: ReservedColumnId::table_id(),
1545            })
1546            .push_column_metadata(ColumnMetadata {
1547                column_schema: ColumnSchema::new(
1548                    "__tsid",
1549                    ConcreteDataType::uint64_datatype(),
1550                    false,
1551                ),
1552                semantic_type: SemanticType::Tag,
1553                column_id: ReservedColumnId::tsid(),
1554            })
1555            .push_column_metadata(ColumnMetadata {
1556                column_schema: ColumnSchema::new("tag0", ConcreteDataType::string_datatype(), true),
1557                semantic_type: SemanticType::Tag,
1558                column_id: 1,
1559            })
1560            .push_column_metadata(ColumnMetadata {
1561                column_schema: ColumnSchema::new("tag1", ConcreteDataType::string_datatype(), true),
1562                semantic_type: SemanticType::Tag,
1563                column_id: 3,
1564            })
1565            .push_column_metadata(ColumnMetadata {
1566                column_schema: ColumnSchema::new(
1567                    "field1",
1568                    ConcreteDataType::int64_datatype(),
1569                    true,
1570                ),
1571                semantic_type: SemanticType::Field,
1572                column_id: 4,
1573            })
1574            .push_column_metadata(ColumnMetadata {
1575                column_schema: ColumnSchema::new(
1576                    "field0",
1577                    ConcreteDataType::int64_datatype(),
1578                    true,
1579                ),
1580                semantic_type: SemanticType::Field,
1581                column_id: 2,
1582            })
1583            .push_column_metadata(ColumnMetadata {
1584                column_schema: ColumnSchema::new(
1585                    "ts",
1586                    ConcreteDataType::timestamp_millisecond_datatype(),
1587                    false,
1588                ),
1589                semantic_type: SemanticType::Timestamp,
1590                column_id: 5,
1591            })
1592            .primary_key(vec![
1593                ReservedColumnId::table_id(),
1594                ReservedColumnId::tsid(),
1595                1,
1596                3,
1597            ])
1598            .primary_key_encoding(PrimaryKeyEncoding::Sparse);
1599        Arc::new(builder.build().unwrap())
1600    }
1601
1602    fn build_test_sparse_pk_array(
1603        codec: &SparsePrimaryKeyCodec,
1604        pk_values_per_row: &[SparseTestRow],
1605    ) -> Arc<PrimaryKeyArray> {
1606        let mut builder = PrimaryKeyArrayBuilder::with_capacity(pk_values_per_row.len(), 1024, 0);
1607        for row in pk_values_per_row {
1608            let values = vec![
1609                (ReservedColumnId::table_id(), ValueRef::UInt32(row.table_id)),
1610                (ReservedColumnId::tsid(), ValueRef::UInt64(row.tsid)),
1611                (1, ValueRef::String(&row.tag0)),
1612                (3, ValueRef::String(&row.tag1)),
1613            ];
1614
1615            let mut buffer = Vec::new();
1616            codec.encode_value_refs(&values, &mut buffer).unwrap();
1617            builder.append_value(&buffer);
1618        }
1619
1620        Arc::new(builder.finish())
1621    }
1622
1623    #[derive(Clone)]
1624    struct SparseTestRow {
1625        table_id: u32,
1626        tsid: u64,
1627        tag0: String,
1628        tag1: String,
1629    }
1630
1631    #[test]
1632    fn test_flat_read_format_convert_format_with_dense_encoding() {
1633        let metadata = build_test_region_metadata();
1634
1635        let column_ids: Vec<_> = metadata
1636            .column_metadatas
1637            .iter()
1638            .map(|c| c.column_id)
1639            .collect();
1640        let format = FlatReadFormat::new(
1641            metadata.clone(),
1642            ReadColumns::new(column_ids),
1643            Some(build_test_arrow_schema()),
1644            "test",
1645            false,
1646        )
1647        .unwrap();
1648
1649        let num_rows = 4;
1650        let original_sequence = 100u64;
1651
1652        // Create primary key values for each row: tag0=1, tag1=1 for all rows
1653        let pk_values_per_row = vec![
1654                &[Some(1i64), Some(1i64)][..]; num_rows  // All rows have same primary key values
1655            ];
1656
1657        // Create a test record batch in old format using dense encoding
1658        let codec = DensePrimaryKeyCodec::new(&metadata);
1659        let dense_pk_array = build_test_dense_pk_array(&codec, &pk_values_per_row);
1660        let columns: Vec<ArrayRef> = vec![
1661            Arc::new(Int64Array::from(vec![2; num_rows])), // field1
1662            Arc::new(Int64Array::from(vec![3; num_rows])), // field0
1663            Arc::new(TimestampMillisecondArray::from(vec![1, 2, 3, 4])), // ts
1664            dense_pk_array.clone(),                        // __primary_key (dense encoding)
1665            Arc::new(UInt64Array::from(vec![original_sequence; num_rows])), // sequence
1666            Arc::new(UInt8Array::from(vec![TEST_OP_TYPE; num_rows])), // op type
1667        ];
1668
1669        // Create schema for old format (without primary key columns)
1670        let old_schema = build_test_arrow_schema();
1671        let record_batch = RecordBatch::try_new(old_schema, columns).unwrap();
1672
1673        // Test conversion with dense encoding
1674        let result = format.convert_batch(record_batch, None).unwrap();
1675
1676        // Construct expected RecordBatch in flat format with decoded primary key columns
1677        let expected_columns: Vec<ArrayRef> = vec![
1678            Arc::new(Int64Array::from(vec![1; num_rows])), // tag0 (decoded from primary key)
1679            Arc::new(Int64Array::from(vec![1; num_rows])), // tag1 (decoded from primary key)
1680            Arc::new(Int64Array::from(vec![2; num_rows])), // field1
1681            Arc::new(Int64Array::from(vec![3; num_rows])), // field0
1682            Arc::new(TimestampMillisecondArray::from(vec![1, 2, 3, 4])), // ts
1683            dense_pk_array,                                // __primary_key (preserved)
1684            Arc::new(UInt64Array::from(vec![original_sequence; num_rows])), // sequence
1685            Arc::new(UInt8Array::from(vec![TEST_OP_TYPE; num_rows])), // op type
1686        ];
1687        let expected_record_batch = RecordBatch::try_new(
1688            build_test_flat_sst_schema_with_field_ids(),
1689            expected_columns,
1690        )
1691        .unwrap();
1692
1693        // Compare the actual result with the expected record batch
1694        assert_eq!(expected_record_batch, result);
1695    }
1696
1697    #[test]
1698    fn test_flat_read_format_convert_format_with_sparse_encoding() {
1699        let metadata = build_test_sparse_region_metadata();
1700
1701        let column_ids: Vec<_> = metadata
1702            .column_metadatas
1703            .iter()
1704            .map(|c| c.column_id)
1705            .collect();
1706        let format = FlatReadFormat::new(
1707            metadata.clone(),
1708            ReadColumns::new(column_ids.clone()),
1709            None,
1710            "test",
1711            false,
1712        )
1713        .unwrap();
1714
1715        let num_rows = 4;
1716        let original_sequence = 100u64;
1717
1718        // Create sparse test data with table_id, tsid and string tags
1719        let pk_test_rows = vec![
1720            SparseTestRow {
1721                table_id: 1,
1722                tsid: 123,
1723                tag0: "frontend".to_string(),
1724                tag1: "pod1".to_string(),
1725            };
1726            num_rows
1727        ];
1728
1729        let codec = SparsePrimaryKeyCodec::new(&metadata);
1730        let sparse_pk_array = build_test_sparse_pk_array(&codec, &pk_test_rows);
1731        // Create a test record batch in old format using sparse encoding
1732        let columns: Vec<ArrayRef> = vec![
1733            Arc::new(Int64Array::from(vec![2; num_rows])), // field1
1734            Arc::new(Int64Array::from(vec![3; num_rows])), // field0
1735            Arc::new(TimestampMillisecondArray::from(vec![1, 2, 3, 4])), // ts
1736            sparse_pk_array.clone(),                       // __primary_key (sparse encoding)
1737            Arc::new(UInt64Array::from(vec![original_sequence; num_rows])), // sequence
1738            Arc::new(UInt8Array::from(vec![TEST_OP_TYPE; num_rows])), // op type
1739        ];
1740
1741        // Create schema for old format (without primary key columns)
1742        let old_schema = build_test_arrow_schema();
1743        let record_batch = RecordBatch::try_new(old_schema, columns).unwrap();
1744
1745        // Test conversion with sparse encoding
1746        let result = format.convert_batch(record_batch.clone(), None).unwrap();
1747
1748        // Construct expected RecordBatch in flat format with decoded primary key columns
1749        let tag0_array = Arc::new(DictionaryArray::new(
1750            UInt32Array::from(vec![0; num_rows]),
1751            Arc::new(StringArray::from(vec!["frontend"])),
1752        ));
1753        let tag1_array = Arc::new(DictionaryArray::new(
1754            UInt32Array::from(vec![0; num_rows]),
1755            Arc::new(StringArray::from(vec!["pod1"])),
1756        ));
1757        let expected_columns: Vec<ArrayRef> = vec![
1758            Arc::new(UInt32Array::from(vec![1; num_rows])), // __table_id (decoded from primary key)
1759            Arc::new(UInt64Array::from(vec![123; num_rows])), // __tsid (decoded from primary key)
1760            tag0_array,                                     // tag0 (decoded from primary key)
1761            tag1_array,                                     // tag1 (decoded from primary key)
1762            Arc::new(Int64Array::from(vec![2; num_rows])),  // field1
1763            Arc::new(Int64Array::from(vec![3; num_rows])),  // field0
1764            Arc::new(TimestampMillisecondArray::from(vec![1, 2, 3, 4])), // ts
1765            sparse_pk_array,                                // __primary_key (preserved)
1766            Arc::new(UInt64Array::from(vec![original_sequence; num_rows])), // sequence
1767            Arc::new(UInt8Array::from(vec![TEST_OP_TYPE; num_rows])), // op type
1768        ];
1769        let expected_schema = to_flat_sst_arrow_schema(&metadata, &FlatSchemaOptions::default());
1770        let expected_record_batch =
1771            RecordBatch::try_new(expected_schema, expected_columns).unwrap();
1772
1773        // Compare the actual result with the expected record batch
1774        assert_eq!(expected_record_batch, result);
1775
1776        let format = FlatReadFormat::new(
1777            metadata.clone(),
1778            ReadColumns::new(column_ids),
1779            None,
1780            "test",
1781            true,
1782        )
1783        .unwrap();
1784        // Test conversion with sparse encoding and skip convert.
1785        let result = format.convert_batch(record_batch.clone(), None).unwrap();
1786        assert_eq!(record_batch, result);
1787    }
1788
1789    #[test]
1790    fn test_convert_flat_batch() {
1791        let metadata = build_test_region_metadata();
1792        let write_format = PrimaryKeyWriteFormat::new(metadata);
1793
1794        let num_rows = 4;
1795        // Build a flat record batch: tag0, tag1, field1, field0, ts, __primary_key, __sequence, __op_type
1796        let flat_columns: Vec<ArrayRef> = input_columns_for_flat_batch(num_rows);
1797        let flat_batch = RecordBatch::try_new(build_test_flat_sst_schema(), flat_columns).unwrap();
1798
1799        // num_fields = 2 (field1, field0)
1800        let result = write_format.convert_flat_batch(&flat_batch, 2).unwrap();
1801
1802        // Expected: tag columns stripped, only field1, field0, ts, __primary_key, __sequence, __op_type
1803        let expected_columns: Vec<ArrayRef> = vec![
1804            Arc::new(Int64Array::from(vec![2; num_rows])), // field1
1805            Arc::new(Int64Array::from(vec![3; num_rows])), // field0
1806            Arc::new(TimestampMillisecondArray::from(vec![1, 2, 3, 4])), // ts
1807            build_test_pk_array(&[(b"test".to_vec(), num_rows)]), // __primary_key
1808            Arc::new(UInt64Array::from(vec![TEST_SEQUENCE; num_rows])), // __sequence
1809            Arc::new(UInt8Array::from(vec![TEST_OP_TYPE; num_rows])), // __op_type
1810        ];
1811        let expected = RecordBatch::try_new(build_test_arrow_schema(), expected_columns).unwrap();
1812
1813        assert_eq!(expected, result);
1814    }
1815
1816    #[test]
1817    fn test_convert_flat_batch_with_override_sequence() {
1818        let metadata = build_test_region_metadata();
1819        let write_format = PrimaryKeyWriteFormat::new(metadata).with_override_sequence(Some(999));
1820
1821        let num_rows = 4;
1822        let flat_columns: Vec<ArrayRef> = input_columns_for_flat_batch(num_rows);
1823        let flat_batch = RecordBatch::try_new(build_test_flat_sst_schema(), flat_columns).unwrap();
1824
1825        let result = write_format.convert_flat_batch(&flat_batch, 2).unwrap();
1826
1827        let expected_columns: Vec<ArrayRef> = vec![
1828            Arc::new(Int64Array::from(vec![2; num_rows])), // field1
1829            Arc::new(Int64Array::from(vec![3; num_rows])), // field0
1830            Arc::new(TimestampMillisecondArray::from(vec![1, 2, 3, 4])), // ts
1831            build_test_pk_array(&[(b"test".to_vec(), num_rows)]), // __primary_key
1832            Arc::new(UInt64Array::from(vec![999; num_rows])), // overridden __sequence
1833            Arc::new(UInt8Array::from(vec![TEST_OP_TYPE; num_rows])), // __op_type
1834        ];
1835        let expected = RecordBatch::try_new(build_test_arrow_schema(), expected_columns).unwrap();
1836
1837        assert_eq!(expected, result);
1838    }
1839
1840    #[test]
1841    fn test_convert_flat_batch_no_tags() {
1842        // Test with a region that has no primary key columns (no tags to strip).
1843        let mut builder = RegionMetadataBuilder::new(RegionId::new(1, 1));
1844        builder
1845            .push_column_metadata(ColumnMetadata {
1846                column_schema: ColumnSchema::new(
1847                    "field0",
1848                    ConcreteDataType::int64_datatype(),
1849                    true,
1850                ),
1851                semantic_type: SemanticType::Field,
1852                column_id: 1,
1853            })
1854            .push_column_metadata(ColumnMetadata {
1855                column_schema: ColumnSchema::new(
1856                    "ts",
1857                    ConcreteDataType::timestamp_millisecond_datatype(),
1858                    false,
1859                ),
1860                semantic_type: SemanticType::Timestamp,
1861                column_id: 2,
1862            });
1863        let metadata = Arc::new(builder.build().unwrap());
1864        let write_format = PrimaryKeyWriteFormat::new(metadata);
1865
1866        let num_rows = 3;
1867        // No tag columns, so flat batch is: field0, ts, __primary_key, __sequence, __op_type
1868        let sst_schema = write_format.arrow_schema().clone();
1869        let columns: Vec<ArrayRef> = vec![
1870            Arc::new(Int64Array::from(vec![10; num_rows])), // field0
1871            Arc::new(TimestampMillisecondArray::from(vec![1, 2, 3])), // ts
1872            build_test_pk_array(&[(b"".to_vec(), num_rows)]), // __primary_key
1873            Arc::new(UInt64Array::from(vec![TEST_SEQUENCE; num_rows])), // __sequence
1874            Arc::new(UInt8Array::from(vec![TEST_OP_TYPE; num_rows])), // __op_type
1875        ];
1876        let flat_batch = RecordBatch::try_new(sst_schema.clone(), columns.clone()).unwrap();
1877
1878        // num_fields = 1, num_tag_columns = 5 - 1 - 4 = 0, so nothing is stripped
1879        let result = write_format.convert_flat_batch(&flat_batch, 1).unwrap();
1880        let expected = RecordBatch::try_new(sst_schema, columns).unwrap();
1881
1882        assert_eq!(expected, result);
1883    }
1884}