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