Skip to main content

mito2/sst/parquet/
flat_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//! It can store both encoded primary key and raw key columns.
18//!
19//! We store two additional internal columns at last:
20//! - `__primary_key`, the encoded primary key of the row (tags). Type: dictionary(uint32, binary)
21//! - `__sequence`, the sequence number of a row. Type: uint64
22//! - `__op_type`, the op type of the row. Type: uint8
23//!
24//! The format is
25//! ```text
26//! primary key columns, field columns, time index, encoded primary key, __sequence, __op_type.
27//!
28//! It stores field columns in the same order as [RegionMetadata::field_columns()](store_api::metadata::RegionMetadata::field_columns())
29//! and stores primary key columns in the same order as [RegionMetadata::primary_key].
30
31use std::borrow::Borrow;
32use std::collections::HashMap;
33use std::sync::Arc;
34
35use api::v1::SemanticType;
36use datatypes::arrow::array::{
37    Array, ArrayRef, BinaryArray, DictionaryArray, UInt32Array, UInt64Array,
38};
39use datatypes::arrow::compute::kernels::take::take;
40use datatypes::arrow::datatypes::{DataType as ArrowDataType, Schema, SchemaRef};
41use datatypes::arrow::record_batch::RecordBatch;
42use datatypes::prelude::{ConcreteDataType, DataType};
43use datatypes::value::ValueRef;
44use datatypes::vectors::MutableVector;
45use mito_codec::row_converter::sparse::{
46    RESERVED_COLUMN_ID_TABLE_ID, RESERVED_COLUMN_ID_TSID, SparsePrimaryKeyView,
47};
48use mito_codec::row_converter::{
49    DensePrimaryKeyCodec, PrimaryKeyCodec, SparseOffsetsCache, build_primary_key_codec,
50};
51use parquet::file::metadata::RowGroupMetaData;
52use snafu::{OptionExt, ResultExt, ensure};
53use store_api::codec::PrimaryKeyEncoding;
54use store_api::metadata::{RegionMetadata, RegionMetadataRef};
55use store_api::storage::{ColumnId, SequenceNumber};
56
57use crate::error::{
58    ComputeArrowSnafu, DecodeSnafu, InvalidParquetSnafu, InvalidRecordBatchSnafu,
59    NewRecordBatchSnafu, Result,
60};
61use crate::read::read_columns::{JsonTargetTypes, ReadColumns};
62use crate::sst::parquet::format::{
63    FIXED_POS_COLUMN_NUM, FormatProjection, INTERNAL_COLUMN_NUM, PrimaryKeyArray,
64    PrimaryKeyReadFormat, StatValues, column_null_counts, column_values,
65};
66use crate::sst::parquet::read_columns::ParquetReadColumns;
67use crate::sst::{
68    FlatSchemaOptions, flat_sst_arrow_schema_column_num, tag_maybe_to_dictionary_field,
69    to_flat_sst_arrow_schema, with_field_id,
70};
71
72/// Helper for writing the SST format.
73pub(crate) struct FlatWriteFormat {
74    /// SST file schema.
75    arrow_schema: SchemaRef,
76    override_sequence: Option<SequenceNumber>,
77}
78
79impl FlatWriteFormat {
80    /// Creates a new helper.
81    pub(crate) fn new(metadata: RegionMetadataRef, options: &FlatSchemaOptions) -> FlatWriteFormat {
82        let arrow_schema = to_flat_sst_arrow_schema(&metadata, options);
83        FlatWriteFormat {
84            arrow_schema,
85            override_sequence: None,
86        }
87    }
88
89    /// Set override sequence.
90    pub(crate) fn with_override_sequence(
91        mut self,
92        override_sequence: Option<SequenceNumber>,
93    ) -> Self {
94        self.override_sequence = override_sequence;
95        self
96    }
97
98    /// Gets the arrow schema to store in parquet.
99    #[cfg(test)]
100    pub(crate) fn arrow_schema(&self) -> &SchemaRef {
101        &self.arrow_schema
102    }
103
104    /// Convert `batch` to a arrow record batch to store in parquet.
105    pub(crate) fn convert_batch(&self, batch: &RecordBatch) -> Result<RecordBatch> {
106        debug_assert_eq!(batch.num_columns(), self.arrow_schema.fields().len());
107
108        let Some(override_sequence) = self.override_sequence else {
109            return Ok(batch.clone());
110        };
111
112        let mut columns = batch.columns().to_vec();
113        let sequence_array = Arc::new(UInt64Array::from(vec![override_sequence; batch.num_rows()]));
114        columns[sequence_column_index(batch.num_columns())] = sequence_array;
115
116        RecordBatch::try_new(batch.schema(), columns).context(NewRecordBatchSnafu)
117    }
118}
119
120/// Returns the position of the sequence column.
121pub(crate) fn sequence_column_index(num_columns: usize) -> usize {
122    num_columns - 2
123}
124
125/// Returns the position of the time index column.
126pub(crate) fn time_index_column_index(num_columns: usize) -> usize {
127    num_columns - 4
128}
129
130/// Returns the position of the primary key column.
131pub(crate) fn primary_key_column_index(num_columns: usize) -> usize {
132    num_columns - 3
133}
134
135/// Wraps the `__primary_key` `BinaryArray` back into a `DictionaryArray<UInt32, Binary>` with identity keys.
136pub(crate) fn wrap_pk_binary_to_dict(
137    record_batch: RecordBatch,
138    dict_schema: &SchemaRef,
139) -> Result<RecordBatch> {
140    let pk_idx = primary_key_column_index(record_batch.num_columns());
141    let pk_column = record_batch.column(pk_idx);
142    let binary_array = pk_column
143        .as_any()
144        .downcast_ref::<BinaryArray>()
145        .with_context(|| InvalidRecordBatchSnafu {
146            reason: format!(
147                "expected BinaryArray for __primary_key, got {:?}",
148                pk_column.data_type()
149            ),
150        })?;
151    let n = binary_array.len();
152    let keys = UInt32Array::from_iter_values(0..n as u32);
153    let dict_array: ArrayRef = Arc::new(DictionaryArray::new(keys, pk_column.clone()));
154
155    let mut columns = record_batch.columns().to_vec();
156    columns[pk_idx] = dict_array;
157
158    RecordBatch::try_new(dict_schema.clone(), columns).context(NewRecordBatchSnafu)
159}
160
161/// Returns the position of the op type key column.
162pub(crate) fn op_type_column_index(num_columns: usize) -> usize {
163    num_columns - 1
164}
165
166/// Returns the start index of field columns in a flat batch.
167///
168/// `num_columns` is the total number of columns in the flat batch schema,
169/// including tag columns (if present), field columns, and fixed position columns
170/// (time index, primary key, sequence, op type).
171///
172/// For Dense encoding (raw PK columns included): field_column_start = primary_key.len()
173/// For Sparse encoding (no raw PK columns): field_column_start = 0
174pub(crate) fn field_column_start(metadata: &RegionMetadata, num_columns: usize) -> usize {
175    // Calculates field column start: total columns - fixed columns - field columns
176    // Field column count = total metadata columns - time index column - primary key columns
177    let field_column_count = metadata.column_metadatas.len() - 1 - metadata.primary_key.len();
178    num_columns - FIXED_POS_COLUMN_NUM - field_column_count
179}
180
181// TODO(yingwen): Add an option to skip reading internal columns if the region is
182// append only and doesn't use sparse encoding (We need to check the table id under
183// sparse encoding).
184/// Helper for reading the flat SST format with projection.
185///
186/// It only supports flat format that stores primary keys additionally.
187pub struct FlatReadFormat {
188    /// Sequence number to override the sequence read from the SST.
189    override_sequence: Option<SequenceNumber>,
190    /// Logical columns requested by this read.
191    read_cols: ReadColumns,
192    /// Parquet format adapter.
193    parquet_adapter: ParquetAdapter,
194    /// Output schema to wrap binary `__primary_key` back to a dictionary; `None` disables wrapping.
195    pk_dict_wrap_schema: Option<SchemaRef>,
196}
197
198impl FlatReadFormat {
199    /// Creates a helper with existing `metadata` and `column_ids` to read.
200    ///
201    /// If `skip_auto_convert` is true, skips auto conversion of format when the encoding is sparse encoding.
202    pub fn new(
203        metadata: RegionMetadataRef,
204        read_cols: ReadColumns,
205        file_schema: Option<SchemaRef>,
206        file_path: &str,
207        skip_auto_convert: bool,
208    ) -> Result<FlatReadFormat> {
209        let num_columns = file_schema.as_ref().map(|x| x.fields().len());
210        let is_legacy = match num_columns {
211            Some(num) => Self::is_legacy_format(&metadata, num, file_path)?,
212            None => metadata.primary_key_encoding == PrimaryKeyEncoding::Sparse,
213        };
214
215        let parquet_adapter = if is_legacy {
216            // Safety: is_legacy_format() ensures primary_key is not empty.
217            if metadata.primary_key_encoding == PrimaryKeyEncoding::Sparse {
218                // Only skip auto convert when the primary key encoding is sparse.
219                ParquetAdapter::PrimaryKeyToFlat(ParquetPrimaryKeyToFlat::new(
220                    metadata,
221                    read_cols.clone(),
222                    skip_auto_convert,
223                ))
224            } else {
225                ParquetAdapter::PrimaryKeyToFlat(ParquetPrimaryKeyToFlat::new(
226                    metadata,
227                    read_cols.clone(),
228                    false,
229                ))
230            }
231        } else {
232            let file_schema = file_schema
233                .unwrap_or_else(|| to_flat_sst_arrow_schema(&metadata, &Default::default()));
234            ParquetAdapter::Flat(ParquetFlat::new(metadata, read_cols.clone(), file_schema))
235        };
236
237        Ok(FlatReadFormat {
238            override_sequence: None,
239            read_cols,
240            parquet_adapter,
241            pk_dict_wrap_schema: None,
242        })
243    }
244
245    /// Sets the sequence number to override.
246    pub(crate) fn set_override_sequence(&mut self, sequence: Option<SequenceNumber>) {
247        self.override_sequence = sequence;
248    }
249
250    /// Enables wrapping binary `__primary_key` batches back to a dictionary in [`Self::convert_batch`].
251    pub(crate) fn set_pk_as_binary(&mut self, output_schema: SchemaRef) {
252        self.pk_dict_wrap_schema = Some(output_schema);
253    }
254
255    /// Index of a column in the projected batch by its column id.
256    pub fn projected_index_by_id(&self, column_id: ColumnId) -> Option<usize> {
257        self.format_projection()
258            .column_id_to_projected_index
259            .get(&column_id)
260            .copied()
261    }
262
263    /// Returns min values of specific column in row groups.
264    pub fn min_values(
265        &self,
266        row_groups: &[impl Borrow<RowGroupMetaData>],
267        column_id: ColumnId,
268    ) -> StatValues {
269        match &self.parquet_adapter {
270            ParquetAdapter::Flat(p) => p.min_values(row_groups, column_id),
271            ParquetAdapter::PrimaryKeyToFlat(p) => p.format.min_values(row_groups, column_id),
272        }
273    }
274
275    /// Returns max values of specific column in row groups.
276    pub fn max_values(
277        &self,
278        row_groups: &[impl Borrow<RowGroupMetaData>],
279        column_id: ColumnId,
280    ) -> StatValues {
281        match &self.parquet_adapter {
282            ParquetAdapter::Flat(p) => p.max_values(row_groups, column_id),
283            ParquetAdapter::PrimaryKeyToFlat(p) => p.format.max_values(row_groups, column_id),
284        }
285    }
286
287    /// Returns null counts of specific column in row groups.
288    pub fn null_counts(
289        &self,
290        row_groups: &[impl Borrow<RowGroupMetaData>],
291        column_id: ColumnId,
292    ) -> StatValues {
293        match &self.parquet_adapter {
294            ParquetAdapter::Flat(p) => p.null_counts(row_groups, column_id),
295            ParquetAdapter::PrimaryKeyToFlat(p) => p.format.null_counts(row_groups, column_id),
296        }
297    }
298
299    /// Gets the arrow schema of the SST file.
300    pub(crate) fn arrow_schema(&self) -> &SchemaRef {
301        match &self.parquet_adapter {
302            ParquetAdapter::Flat(p) => &p.arrow_schema,
303            ParquetAdapter::PrimaryKeyToFlat(p) => p.format.arrow_schema(),
304        }
305    }
306
307    /// Gets the projected output schema expected by the scan.
308    pub(crate) fn output_arrow_schema(&self) -> Result<SchemaRef> {
309        let projection = self.parquet_read_columns().root_indices();
310        let mut schema = self
311            .arrow_schema()
312            .project(projection)
313            .context(ComputeArrowSnafu)?;
314        let mut fields = schema.fields().iter().cloned().collect::<Vec<_>>();
315        for (column_id, target) in self.json_target_types().iter() {
316            let Some(index) = self.parquet_projected_index_by_id(*column_id) else {
317                continue;
318            };
319            let Some(field) = schema.fields().get(index) else {
320                continue;
321            };
322            let mut field = field.as_ref().clone();
323            field.set_data_type(ConcreteDataType::json2(target.clone()).as_arrow_type());
324            fields[index] = Arc::new(field);
325        }
326        schema.fields = fields.into();
327        Ok(Arc::new(schema))
328    }
329
330    /// Index of a column in the projected schema produced directly by parquet
331    /// reading, before any primary-key-to-flat conversion.
332    pub(crate) fn parquet_projected_index_by_id(&self, column_id: ColumnId) -> Option<usize> {
333        match &self.parquet_adapter {
334            ParquetAdapter::Flat(p) => p
335                .format_projection
336                .column_id_to_projected_index
337                .get(&column_id)
338                .copied(),
339            // `format_projection` addresses the post-conversion flat batch here.
340            // This helper needs the raw primary-key projection used by parquet reading.
341            ParquetAdapter::PrimaryKeyToFlat(p) => p
342                .format
343                .field_id_to_projected_index()
344                .get(&column_id)
345                .copied(),
346        }
347    }
348
349    /// Gets the metadata of the SST.
350    pub(crate) fn metadata(&self) -> &RegionMetadataRef {
351        match &self.parquet_adapter {
352            ParquetAdapter::Flat(p) => &p.metadata,
353            ParquetAdapter::PrimaryKeyToFlat(p) => p.format.metadata(),
354        }
355    }
356
357    /// Get the sorted read columns to read from the sst file.
358    pub(crate) fn parquet_read_columns(&self) -> &ParquetReadColumns {
359        match &self.parquet_adapter {
360            ParquetAdapter::Flat(p) => &p.format_projection.parquet_read_cols,
361            ParquetAdapter::PrimaryKeyToFlat(p) => p.format.parquet_read_columns(),
362        }
363    }
364
365    /// Gets JSON2 read targets.
366    pub(crate) fn json_target_types(&self) -> &JsonTargetTypes {
367        self.read_cols.json_target_types()
368    }
369
370    /// Gets the projection in the flat format.
371    ///
372    /// When `skip_auto_convert` is enabled (primary-key format read), this returns the
373    /// primary-key format projection so filter/prune can resolve projected indices.
374    pub(crate) fn format_projection(&self) -> &FormatProjection {
375        match &self.parquet_adapter {
376            ParquetAdapter::Flat(p) => &p.format_projection,
377            ParquetAdapter::PrimaryKeyToFlat(p) => &p.format_projection,
378        }
379    }
380
381    /// Returns `true` if raw batches from parquet use the flat layout and
382    /// stores primary key columns as raw columns.
383    /// Returns `false` for the legacy primary-key-to-flat conversion path.
384    pub(crate) fn batch_has_raw_pk_columns(&self) -> bool {
385        matches!(&self.parquet_adapter, ParquetAdapter::Flat(_))
386    }
387
388    /// Creates a sequence array to override.
389    pub(crate) fn new_override_sequence_array(&self, length: usize) -> Option<ArrayRef> {
390        self.override_sequence
391            .map(|seq| Arc::new(UInt64Array::from_value(seq, length)) as ArrayRef)
392    }
393
394    /// Convert a record batch to apply flat format conversion and override sequence array.
395    ///
396    /// Returns a new RecordBatch with flat format conversion applied first (if enabled),
397    /// then the sequence column replaced by the override sequence array.
398    pub(crate) fn convert_batch(
399        &self,
400        record_batch: RecordBatch,
401        override_sequence_array: Option<&ArrayRef>,
402    ) -> Result<RecordBatch> {
403        let record_batch = if let Some(dict_schema) = &self.pk_dict_wrap_schema {
404            wrap_pk_binary_to_dict(record_batch, dict_schema)?
405        } else {
406            record_batch
407        };
408
409        // First, apply flat format conversion.
410        let mut batch = match &self.parquet_adapter {
411            ParquetAdapter::Flat(_) => record_batch,
412            ParquetAdapter::PrimaryKeyToFlat(p) => p.convert_batch(record_batch)?,
413        };
414
415        // Normalize nested field names and metadata to the SST's region metadata
416        // before schema compatibility and merging with memtables. This removes
417        // Parquet-added field IDs; equals_datatype also permits nested name differences.
418        for index in 0..batch.num_columns() {
419            let array = batch.column(index);
420            if !matches!(array.data_type(), ArrowDataType::Struct(_)) {
421                continue;
422            }
423            let field = batch.schema_ref().field(index);
424            let Some(column) = self.metadata().column_by_name(field.name()) else {
425                continue;
426            };
427            let target = column.column_schema.data_type.as_arrow_type();
428            if array.data_type() != &target && array.data_type().equals_datatype(&target) {
429                let array =
430                    datatypes::arrow::compute::cast(array, &target).context(ComputeArrowSnafu)?;
431                let mut fields = batch.schema().fields().to_vec();
432                fields[index] = Arc::new(field.clone().with_data_type(target));
433                let mut columns = batch.columns().to_vec();
434                columns[index] = array;
435                batch = RecordBatch::try_new(
436                    Arc::new(Schema::new_with_metadata(
437                        fields,
438                        batch.schema().metadata().clone(),
439                    )),
440                    columns,
441                )
442                .context(NewRecordBatchSnafu)?;
443            }
444        }
445
446        // Then apply sequence override if provided
447        let Some(override_array) = override_sequence_array else {
448            return Ok(batch);
449        };
450
451        let mut columns = batch.columns().to_vec();
452        let sequence_column_idx = sequence_column_index(batch.num_columns());
453
454        // Use the provided override sequence array, slicing if necessary to match batch length
455        let sequence_array = if override_array.len() > batch.num_rows() {
456            override_array.slice(0, batch.num_rows())
457        } else {
458            override_array.clone()
459        };
460
461        columns[sequence_column_idx] = sequence_array;
462
463        RecordBatch::try_new(batch.schema(), columns).context(NewRecordBatchSnafu)
464    }
465
466    /// Checks whether the batch from the parquet file needs to be converted to match the flat format.
467    ///
468    /// * `metadata` is the region metadata (always assumes flat format).
469    /// * `num_columns` is the number of columns in the parquet file.
470    /// * `file_path` is the path to the parquet file, for error message.
471    pub(crate) fn is_legacy_format(
472        metadata: &RegionMetadata,
473        num_columns: usize,
474        file_path: &str,
475    ) -> Result<bool> {
476        if metadata.primary_key.is_empty() {
477            return Ok(false);
478        }
479
480        // For flat format, compute expected column number:
481        // all columns + internal columns (pk, sequence, op_type)
482        let expected_columns = metadata.column_metadatas.len() + INTERNAL_COLUMN_NUM;
483
484        if expected_columns == num_columns {
485            // Same number of columns, no conversion needed
486            Ok(false)
487        } else {
488            ensure!(
489                expected_columns >= num_columns,
490                InvalidParquetSnafu {
491                    file: file_path,
492                    reason: format!(
493                        "Expected columns {} should be >= actual columns {}",
494                        expected_columns, num_columns
495                    )
496                }
497            );
498
499            // Different number of columns, check if the difference matches primary key count
500            let column_diff = expected_columns - num_columns;
501
502            ensure!(
503                column_diff == metadata.primary_key.len(),
504                InvalidParquetSnafu {
505                    file: file_path,
506                    reason: format!(
507                        "Column number difference {} does not match primary key count {}",
508                        column_diff,
509                        metadata.primary_key.len()
510                    )
511                }
512            );
513
514            Ok(true)
515        }
516    }
517}
518
519/// Wraps the parquet helper for different formats.
520enum ParquetAdapter {
521    Flat(ParquetFlat),
522    PrimaryKeyToFlat(ParquetPrimaryKeyToFlat),
523}
524
525/// Helper to reads the parquet from primary key format into the flat format.
526struct ParquetPrimaryKeyToFlat {
527    /// The primary key format to read the parquet.
528    format: PrimaryKeyReadFormat,
529    /// Format converter for handling flat format conversion.
530    convert_format: Option<FlatConvertFormat>,
531    /// Projection computed for the flat format.
532    format_projection: FormatProjection,
533}
534
535impl ParquetPrimaryKeyToFlat {
536    /// Creates a helper with existing `metadata` and `column_ids` to read.
537    fn new(
538        metadata: RegionMetadataRef,
539        read_cols: ReadColumns,
540        skip_auto_convert: bool,
541    ) -> ParquetPrimaryKeyToFlat {
542        assert!(if skip_auto_convert {
543            metadata.primary_key_encoding == PrimaryKeyEncoding::Sparse
544        } else {
545            true
546        });
547
548        // Creates a map to lookup index based on the new format.
549        let id_to_index = sst_column_id_indices(&metadata);
550        let sst_column_num =
551            flat_sst_arrow_schema_column_num(&metadata, &FlatSchemaOptions::default());
552
553        let codec = build_primary_key_codec(&metadata);
554        let format = PrimaryKeyReadFormat::new(metadata.clone(), read_cols.clone());
555        let (convert_format, format_projection) = if skip_auto_convert {
556            (
557                None,
558                FormatProjection {
559                    parquet_read_cols: format.parquet_read_columns().clone(),
560                    column_id_to_projected_index: format.field_id_to_projected_index().clone(),
561                },
562            )
563        } else {
564            // Computes the format projection for the new format.
565            let format_projection = FormatProjection::compute_format_projection(
566                &metadata,
567                &id_to_index,
568                sst_column_num,
569                read_cols.clone(),
570            );
571            (
572                FlatConvertFormat::new(Arc::clone(&metadata), &format_projection, codec),
573                format_projection,
574            )
575        };
576
577        Self {
578            format,
579            convert_format,
580            format_projection,
581        }
582    }
583
584    fn convert_batch(&self, record_batch: RecordBatch) -> Result<RecordBatch> {
585        if let Some(convert_format) = &self.convert_format {
586            convert_format.convert(record_batch)
587        } else {
588            Ok(record_batch)
589        }
590    }
591}
592
593/// Helper to reads the parquet in flat format directly.
594struct ParquetFlat {
595    /// The metadata stored in the SST.
596    metadata: RegionMetadataRef,
597    /// SST file schema.
598    arrow_schema: SchemaRef,
599    /// Projection computed for the flat format.
600    format_projection: FormatProjection,
601    /// Column id to top-level SST index. Shared statistics helpers resolve
602    /// physical leaves from each file's actual Parquet schema.
603    column_id_to_sst_index: HashMap<ColumnId, usize>,
604}
605
606impl ParquetFlat {
607    /// Creates a helper with existing `metadata` and `column_ids` to read.
608    fn new(
609        metadata: RegionMetadataRef,
610        read_cols: ReadColumns,
611        arrow_schema: SchemaRef,
612    ) -> ParquetFlat {
613        // Creates a map to lookup index.
614        let id_to_index = sst_column_id_indices(&metadata);
615        let sst_column_num =
616            flat_sst_arrow_schema_column_num(&metadata, &FlatSchemaOptions::default());
617        let format_projection = FormatProjection::compute_format_projection(
618            &metadata,
619            &id_to_index,
620            sst_column_num,
621            read_cols,
622        );
623
624        Self {
625            metadata,
626            arrow_schema,
627            format_projection,
628            column_id_to_sst_index: id_to_index,
629        }
630    }
631
632    /// Returns min values of specific column in row groups.
633    fn min_values(
634        &self,
635        row_groups: &[impl Borrow<RowGroupMetaData>],
636        column_id: ColumnId,
637    ) -> StatValues {
638        self.get_stat_values(row_groups, column_id, true)
639    }
640
641    /// Returns max values of specific column in row groups.
642    fn max_values(
643        &self,
644        row_groups: &[impl Borrow<RowGroupMetaData>],
645        column_id: ColumnId,
646    ) -> StatValues {
647        self.get_stat_values(row_groups, column_id, false)
648    }
649
650    /// Returns null counts of specific column in row groups.
651    fn null_counts(
652        &self,
653        row_groups: &[impl Borrow<RowGroupMetaData>],
654        column_id: ColumnId,
655    ) -> StatValues {
656        let Some(index) = self.column_id_to_sst_index.get(&column_id) else {
657            // No such column in the SST.
658            return StatValues::NoColumn;
659        };
660        let stats = column_null_counts(row_groups, *index);
661        StatValues::from_stats_opt(stats)
662    }
663
664    fn get_stat_values(
665        &self,
666        row_groups: &[impl Borrow<RowGroupMetaData>],
667        column_id: ColumnId,
668        is_min: bool,
669    ) -> StatValues {
670        let Some(column) = self.metadata.column_by_id(column_id) else {
671            // No such column in the SST.
672            return StatValues::NoColumn;
673        };
674        let Some(index) = self.column_id_to_sst_index.get(&column_id) else {
675            return StatValues::NoStats;
676        };
677
678        let stats = column_values(row_groups, column, *index, is_min);
679        StatValues::from_stats_opt(stats)
680    }
681}
682
683/// Returns a map that the key is the column id and the value is the column position
684/// in the SST.
685/// It only supports SSTs with raw primary key columns.
686pub(crate) fn sst_column_id_indices(metadata: &RegionMetadata) -> HashMap<ColumnId, usize> {
687    let mut id_to_index = HashMap::with_capacity(metadata.column_metadatas.len());
688    let mut column_index = 0;
689    // keys
690    for pk_id in &metadata.primary_key {
691        id_to_index.insert(*pk_id, column_index);
692        column_index += 1;
693    }
694    // fields
695    for column in &metadata.column_metadatas {
696        if column.semantic_type == SemanticType::Field {
697            id_to_index.insert(column.column_id, column_index);
698            column_index += 1;
699        }
700    }
701    // time index
702    id_to_index.insert(metadata.time_index_column().column_id, column_index);
703
704    id_to_index
705}
706
707/// Decodes primary keys from a batch and returns decoded primary key information.
708///
709/// The batch must contain a primary key column at the expected index.
710/// Primary keys stay encoded: tag values are extracted lazily per column
711/// in [`DecodedPrimaryKeys::get_tag_column`] without decoding every label.
712/// The codec must describe the source key's field order and types.
713pub fn decode_primary_keys(
714    codec: &dyn PrimaryKeyCodec,
715    batch: &RecordBatch,
716) -> Result<DecodedPrimaryKeys> {
717    let primary_key_index = primary_key_column_index(batch.num_columns());
718    let pk_dict_array = batch
719        .column(primary_key_index)
720        .as_any()
721        .downcast_ref::<PrimaryKeyArray>()
722        .with_context(|| InvalidRecordBatchSnafu {
723            reason: "Primary key column is not a dictionary array".to_string(),
724        })?;
725    let pk_values_array = pk_dict_array
726        .values()
727        .as_any()
728        .downcast_ref::<BinaryArray>()
729        .with_context(|| InvalidRecordBatchSnafu {
730            reason: "Primary key values are not binary array".to_string(),
731        })?;
732
733    let keys = pk_dict_array.keys();
734
735    // Collects consecutive runs of identical dictionary keys, preserving their
736    // row order. Maps original key index -> new decoded value index.
737    // The parquet reader may read the whole dictionary page into the dictionary
738    // values, so we may decode many primary keys not in this batch if we decode
739    // the values array directly.
740    let pk_indices = keys.values();
741    let mut key_to_decoded_index = Vec::with_capacity(keys.len());
742    let mut distinct_keys: Vec<u32> = Vec::new();
743    let mut prev_key: Option<u32> = None;
744    for &current_key in pk_indices.iter().take(keys.len()) {
745        // Check if current key is the same as previous key
746        if let Some(prev) = prev_key
747            && prev == current_key
748        {
749            // Reuse the last decoded index
750            key_to_decoded_index.push((distinct_keys.len() - 1) as u32);
751            continue;
752        }
753
754        distinct_keys.push(current_key);
755        key_to_decoded_index.push((distinct_keys.len() - 1) as u32);
756        prev_key = Some(current_key);
757    }
758
759    let inner = match codec.encoding() {
760        PrimaryKeyEncoding::Sparse => DecodedKeysInner::Sparse {
761            values: pk_values_array.clone(),
762            distinct_keys,
763        },
764        PrimaryKeyEncoding::Dense => DecodedKeysInner::Dense {
765            codec: codec
766                .as_dense()
767                .context(InvalidRecordBatchSnafu {
768                    reason: "expected dense primary key codec",
769                })?
770                .clone(),
771            offsets: Vec::new(),
772            values: pk_values_array.clone(),
773            distinct_keys,
774        },
775    };
776
777    Ok(DecodedPrimaryKeys {
778        inner,
779        keys_array: UInt32Array::from(key_to_decoded_index),
780        pk_offsets: SparseOffsetsCache::new(),
781        value_buf: Vec::new(),
782    })
783}
784
785/// Encoded primary keys and encoding-specific lookup state.
786enum DecodedKeysInner {
787    /// Dense keys share positional offsets across projected tag columns.
788    Dense {
789        codec: DensePrimaryKeyCodec,
790        values: BinaryArray,
791        distinct_keys: Vec<u32>,
792        /// Initialized by positional access; full extraction needs no offsets.
793        offsets: Vec<Vec<usize>>,
794    },
795    /// Sparse primary keys stay encoded; each tag column extracts only its own
796    /// values from the raw keys.
797    Sparse {
798        /// Dictionary values of the primary key array.
799        values: BinaryArray,
800        /// Dictionary key of each distinct consecutive run, in first-appearance order.
801        distinct_keys: Vec<u32>,
802    },
803}
804
805/// Holds decoded primary key values for unique keys and their indices.
806pub struct DecodedPrimaryKeys {
807    inner: DecodedKeysInner,
808    /// Prebuilt keys array for creating dictionary arrays.
809    keys_array: UInt32Array,
810    /// Scratch offsets shared when extracting sparse tag values.
811    pk_offsets: SparseOffsetsCache,
812    /// Reusable buffer for extracting sparse tag values.
813    value_buf: Vec<u8>,
814}
815
816/// Pushes the value of `column_id` in `pk` into `builder`.
817fn push_sparse_tag_value(
818    pk: &[u8],
819    column_id: ColumnId,
820    builder: &mut dyn MutableVector,
821    pk_offsets: &mut SparseOffsetsCache,
822    value_buf: &mut Vec<u8>,
823) -> Result<()> {
824    let mut view = SparsePrimaryKeyView::new(pk, pk_offsets).context(DecodeSnafu)?;
825    push_sparse_tag_value_in_view(&mut view, column_id, builder, value_buf)
826}
827
828/// Pushes the value of `column_id` into `builder` using an existing view, so
829/// multiple columns of the same key share offset discovery.
830fn push_sparse_tag_value_in_view(
831    view: &mut SparsePrimaryKeyView,
832    column_id: ColumnId,
833    builder: &mut dyn MutableVector,
834    value_buf: &mut Vec<u8>,
835) -> Result<()> {
836    match column_id {
837        RESERVED_COLUMN_ID_TABLE_ID => builder.push_value_ref(&ValueRef::UInt32(view.table_id())),
838        RESERVED_COLUMN_ID_TSID => builder.push_value_ref(&ValueRef::UInt64(view.tsid())),
839        _ => {
840            let value = view.label(column_id, value_buf).context(DecodeSnafu)?;
841            match value {
842                None => builder.push_null(),
843                Some(value) => builder.push_value_ref(&ValueRef::String(value)),
844            }
845        }
846    }
847    Ok(())
848}
849
850impl DecodedPrimaryKeys {
851    /// Gets a tag column array by column id and data type.
852    ///
853    /// For sparse encoding, extracts the column lazily from the encoded keys.
854    /// For dense encoding, uses pk_index to decode only the requested field.
855    pub fn get_tag_column(
856        &mut self,
857        column_id: ColumnId,
858        pk_index: Option<usize>,
859        column_type: &ConcreteDataType,
860    ) -> Result<ArrayRef> {
861        let Self {
862            inner,
863            keys_array,
864            pk_offsets,
865            value_buf,
866        } = self;
867
868        // Gets values from the primary key.
869        let values_vector = match inner {
870            DecodedKeysInner::Dense {
871                codec,
872                values,
873                distinct_keys,
874                offsets,
875            } => {
876                let pk_idx = pk_index.context(InvalidRecordBatchSnafu {
877                    reason: "pk_index required for dense encoding",
878                })?;
879                let mut builder = column_type.create_mutable_vector(distinct_keys.len());
880                if offsets.is_empty() {
881                    *offsets = vec![Vec::new(); distinct_keys.len()];
882                }
883                for (&key, offsets) in distinct_keys.iter().zip(offsets) {
884                    if pk_idx < codec.num_fields() {
885                        let value = codec
886                            .decode_value_at(values.value(key as usize), pk_idx, offsets)
887                            .context(DecodeSnafu)?;
888                        builder.push_value_ref(&value.as_value_ref());
889                    } else {
890                        builder.push_null();
891                    }
892                }
893                builder.to_vector()
894            }
895            DecodedKeysInner::Sparse {
896                values,
897                distinct_keys,
898            } => {
899                let mut builder = column_type.create_mutable_vector(distinct_keys.len());
900                for &key in distinct_keys.iter() {
901                    let pk = values.value(key as usize);
902                    push_sparse_tag_value(pk, column_id, &mut *builder, pk_offsets, value_buf)?;
903                }
904                builder.to_vector()
905            }
906        };
907        let values_array = values_vector.to_arrow_array();
908
909        // Only creates dictionary array for string types, otherwise take values by keys
910        if column_type.is_string() {
911            // Creates dictionary array using the same keys for string types
912            // Note that the dictionary values may have nulls.
913            let dict_array = DictionaryArray::new(keys_array.clone(), values_array);
914            Ok(Arc::new(dict_array))
915        } else {
916            // For non-string types, takes values by keys indices to create a regular array
917            let taken_array = take(&values_array, keys_array, None).context(ComputeArrowSnafu)?;
918            Ok(taken_array)
919        }
920    }
921
922    /// Materializes all Dense tags in source-schema order with one sequential
923    /// traversal per key. Values are consumed immediately by their Arrow builders.
924    pub fn get_dense_tag_columns(&self) -> Result<Vec<ArrayRef>> {
925        let DecodedKeysInner::Dense {
926            codec,
927            values,
928            distinct_keys,
929            ..
930        } = &self.inner
931        else {
932            return InvalidRecordBatchSnafu {
933                reason: "expected dense primary key values",
934            }
935            .fail();
936        };
937        let mut builders: Vec<_> = codec
938            .fields()
939            .iter()
940            .map(|(_, field)| field.data_type().create_mutable_vector(distinct_keys.len()))
941            .collect();
942        let mut value_buf = Vec::new();
943        for &key in distinct_keys {
944            codec
945                .decode_dense_with(values.value(key as usize), &mut value_buf, |pos, value| {
946                    builders[pos].push_value_ref(&value);
947                })
948                .context(DecodeSnafu)?;
949        }
950        codec
951            .fields()
952            .iter()
953            .zip(builders)
954            .map(|((_, field), mut builder)| {
955                let values = builder.to_vector().to_arrow_array();
956                if field.data_type().is_string() {
957                    Ok(Arc::new(DictionaryArray::new(self.keys_array.clone(), values)) as ArrayRef)
958                } else {
959                    take(&values, &self.keys_array, None).context(ComputeArrowSnafu)
960                }
961            })
962            .collect()
963    }
964
965    /// Gets multiple sparse tag column arrays in one pass over the distinct keys,
966    /// sharing each key's offset discovery between the columns.
967    ///
968    /// Must only be called for sparse primary keys. Columns hold their values in
969    /// the order of `columns`, whose entries are tag column ids and data types.
970    pub fn get_sparse_tag_columns(
971        &mut self,
972        columns: &[(ColumnId, ConcreteDataType)],
973    ) -> Result<Vec<ArrayRef>> {
974        let Self {
975            inner,
976            keys_array,
977            pk_offsets,
978            value_buf,
979        } = self;
980        let DecodedKeysInner::Sparse {
981            values,
982            distinct_keys,
983        } = inner
984        else {
985            return InvalidRecordBatchSnafu {
986                reason: "expected sparse primary key values",
987            }
988            .fail();
989        };
990
991        let mut builders: Vec<_> = columns
992            .iter()
993            .map(|(_, column_type)| column_type.create_mutable_vector(distinct_keys.len()))
994            .collect();
995        for &key in distinct_keys.iter() {
996            let pk = values.value(key as usize);
997            let mut view = SparsePrimaryKeyView::new(pk, pk_offsets).context(DecodeSnafu)?;
998            // Visit all columns before moving to the next key so offset
999            // discovery is shared between columns.
1000            for ((column_id, _), builder) in columns.iter().zip(&mut builders) {
1001                push_sparse_tag_value_in_view(&mut view, *column_id, &mut **builder, value_buf)?;
1002            }
1003        }
1004
1005        columns
1006            .iter()
1007            .zip(builders)
1008            .map(|((_, column_type), mut builder)| {
1009                let values_array = builder.to_vector().to_arrow_array();
1010                if column_type.is_string() {
1011                    // Note that the dictionary values may have nulls.
1012                    Ok(
1013                        Arc::new(DictionaryArray::new(keys_array.clone(), values_array))
1014                            as ArrayRef,
1015                    )
1016                } else {
1017                    take(&values_array, keys_array, None)
1018                        .context(ComputeArrowSnafu)
1019                        .map(|array| array as ArrayRef)
1020                }
1021            })
1022            .collect()
1023    }
1024}
1025
1026/// Converts a batch that doesn't have decoded primary key columns into a batch that has decoded
1027/// primary key columns in flat format.
1028pub(crate) struct FlatConvertFormat {
1029    /// Metadata of the region.
1030    metadata: RegionMetadataRef,
1031    /// Primary key codec to decode primary keys.
1032    codec: Arc<dyn PrimaryKeyCodec>,
1033    /// Projected primary key column information: (column_id, pk_index, column_index in metadata).
1034    projected_primary_keys: Vec<(ColumnId, usize, usize)>,
1035}
1036
1037impl FlatConvertFormat {
1038    /// Creates a new `FlatConvertFormat`.
1039    ///
1040    /// The `format_projection` is the projection computed in the [FlatReadFormat] with the `metadata`.
1041    /// The `codec` is the primary key codec of the `metadata`.
1042    ///
1043    /// Returns `None` if there is no primary key.
1044    pub(crate) fn new(
1045        metadata: RegionMetadataRef,
1046        format_projection: &FormatProjection,
1047        codec: Arc<dyn PrimaryKeyCodec>,
1048    ) -> Option<Self> {
1049        if metadata.primary_key.is_empty() {
1050            return None;
1051        }
1052
1053        // Builds projected primary keys list maintaining the order of RegionMetadata::primary_key
1054        let mut projected_primary_keys = Vec::new();
1055        for (pk_index, &column_id) in metadata.primary_key.iter().enumerate() {
1056            if format_projection
1057                .column_id_to_projected_index
1058                .contains_key(&column_id)
1059            {
1060                // We expect the format_projection is built from the metadata.
1061                let column_index = metadata.column_index_by_id(column_id).unwrap();
1062                projected_primary_keys.push((column_id, pk_index, column_index));
1063            }
1064        }
1065
1066        Some(Self {
1067            metadata,
1068            codec,
1069            projected_primary_keys,
1070        })
1071    }
1072
1073    /// Converts a batch to have decoded primary key columns in flat format.
1074    ///
1075    /// The primary key array in the batch is a dictionary array.
1076    pub(crate) fn convert(&self, batch: RecordBatch) -> Result<RecordBatch> {
1077        if self.projected_primary_keys.is_empty() {
1078            return Ok(batch);
1079        }
1080
1081        let mut decoded_pks = decode_primary_keys(self.codec.as_ref(), &batch)?;
1082
1083        // Builds decoded tag column arrays.
1084        let mut decoded_columns = Vec::new();
1085        if self.codec.encoding() == PrimaryKeyEncoding::Sparse {
1086            // Extract all projected tag columns in one pass over the distinct
1087            // keys, sharing offset discovery between the columns.
1088            let columns: Vec<_> = self
1089                .projected_primary_keys
1090                .iter()
1091                .map(|(column_id, _, column_index)| {
1092                    let column_metadata = &self.metadata.column_metadatas[*column_index];
1093                    (*column_id, column_metadata.column_schema.data_type.clone())
1094                })
1095                .collect();
1096            decoded_columns.extend(decoded_pks.get_sparse_tag_columns(&columns)?);
1097        } else if self.projected_primary_keys.len() == self.metadata.primary_key.len() {
1098            decoded_columns.extend(decoded_pks.get_dense_tag_columns()?);
1099        } else {
1100            for (column_id, pk_index, column_index) in &self.projected_primary_keys {
1101                let column_metadata = &self.metadata.column_metadatas[*column_index];
1102                let tag_column = decoded_pks.get_tag_column(
1103                    *column_id,
1104                    Some(*pk_index),
1105                    &column_metadata.column_schema.data_type,
1106                )?;
1107                decoded_columns.push(tag_column);
1108            }
1109        }
1110
1111        // Builds new columns: decoded tag columns first, then original columns
1112        let mut new_columns = Vec::with_capacity(batch.num_columns() + decoded_columns.len());
1113        new_columns.extend(decoded_columns);
1114        new_columns.extend_from_slice(batch.columns());
1115
1116        // Builds new schema
1117        let mut new_fields =
1118            Vec::with_capacity(batch.schema().fields().len() + self.projected_primary_keys.len());
1119        for (column_id, _, column_index) in &self.projected_primary_keys {
1120            let column_metadata = &self.metadata.column_metadatas[*column_index];
1121            let old_field = &self.metadata.schema.arrow_schema().fields()[*column_index];
1122            let field =
1123                tag_maybe_to_dictionary_field(&column_metadata.column_schema.data_type, old_field);
1124            new_fields.push(Arc::new(with_field_id((*field).clone(), *column_id)));
1125        }
1126        new_fields.extend(batch.schema().fields().iter().cloned());
1127
1128        let new_schema = Arc::new(Schema::new(new_fields));
1129        RecordBatch::try_new(new_schema, new_columns).context(NewRecordBatchSnafu)
1130    }
1131}
1132
1133#[cfg(test)]
1134impl FlatReadFormat {
1135    /// Creates a helper with existing `metadata` and all columns.
1136    pub fn new_with_all_columns(metadata: RegionMetadataRef) -> FlatReadFormat {
1137        Self::new(
1138            Arc::clone(&metadata),
1139            ReadColumns::new(metadata.column_metadatas.iter().map(|c| c.column_id)),
1140            None,
1141            "test",
1142            false,
1143        )
1144        .unwrap()
1145    }
1146}
1147
1148#[cfg(test)]
1149mod tests {
1150    use std::sync::Arc;
1151
1152    use api::v1::SemanticType;
1153    use datatypes::arrow::array::{
1154        ArrayRef, BinaryArray, Int64Array, TimestampMillisecondArray, UInt8Array, UInt32Array,
1155        UInt64Array,
1156    };
1157    use datatypes::arrow::datatypes::{DataType as ArrowDataType, Field, TimeUnit};
1158    use datatypes::arrow::record_batch::RecordBatch;
1159    use datatypes::prelude::ConcreteDataType;
1160    use datatypes::schema::ColumnSchema;
1161    use datatypes::types::json_type::{JsonNativeType, JsonObjectType};
1162    use parquet::arrow::ArrowSchemaConverter;
1163    use parquet::basic::{Repetition, Type as PhysicalType};
1164    use parquet::file::metadata::{ColumnChunkMetaData, RowGroupMetaData};
1165    use parquet::file::statistics::Statistics;
1166    use parquet::schema::types::{SchemaDescriptor, Type};
1167    use store_api::codec::PrimaryKeyEncoding;
1168    use store_api::metadata::{ColumnMetadata, RegionMetadata, RegionMetadataBuilder};
1169    use store_api::storage::RegionId;
1170    use store_api::storage::consts::{
1171        OP_TYPE_COLUMN_NAME, PRIMARY_KEY_COLUMN_NAME, SEQUENCE_COLUMN_NAME,
1172    };
1173
1174    use super::*;
1175    use crate::read::read_columns::ReadColumns;
1176    use crate::sst::{
1177        FlatSchemaOptions, PARQUET_FIELD_ID_KEY, PRIMARY_KEY_PARQUET_FIELD_ID,
1178        flat_sst_arrow_schema_column_num, override_pk_field_to_binary, to_flat_sst_arrow_schema,
1179    };
1180
1181    #[test]
1182    fn dense_tag_columns_match_eager_decoding() {
1183        use datatypes::arrow::datatypes::UInt32Type;
1184        use datatypes::value::Value;
1185        use datatypes::vectors::Helper;
1186        use mito_codec::row_converter::{DensePrimaryKeyCodec, PrimaryKeyCodecExt, SortField};
1187
1188        let columns = [
1189            (17, ConcreteDataType::string_datatype()),
1190            (9, ConcreteDataType::int64_datatype()),
1191            (2, ConcreteDataType::binary_datatype()),
1192        ];
1193        let codec = DensePrimaryKeyCodec::with_fields(
1194            columns
1195                .iter()
1196                .map(|(id, ty)| (*id, SortField::new(ty.clone())))
1197                .collect(),
1198        );
1199        let rows = [
1200            vec![
1201                Value::from("中文\0abcdefgh"),
1202                Value::Int64(-42),
1203                Value::Binary(vec![0, 1, 255].into()),
1204            ],
1205            vec![Value::from(""), Value::Null, Value::Binary(vec![].into())],
1206            vec![Value::Null, Value::Int64(i64::MAX), Value::Null],
1207        ];
1208        let mut encoded: Vec<_> = rows
1209            .iter()
1210            .map(|row| codec.encode(row.iter().map(Value::as_value_ref)).unwrap())
1211            .collect();
1212        // Unreferenced dictionary entries must not be inspected.
1213        encoded.push(vec![255]);
1214        let row_keys = [2, 2, 0, 1, 0, 0];
1215        let pk = DictionaryArray::<UInt32Type>::new(
1216            UInt32Array::from(row_keys.to_vec()),
1217            Arc::new(BinaryArray::from_iter_values(&encoded)),
1218        );
1219        let batch = RecordBatch::try_from_iter([
1220            (
1221                "ts",
1222                Arc::new(TimestampMillisecondArray::from_iter_values(0..6)) as ArrayRef,
1223            ),
1224            ("__primary_key", Arc::new(pk) as ArrayRef),
1225            (
1226                "__sequence",
1227                Arc::new(UInt64Array::from(vec![1; 6])) as ArrayRef,
1228            ),
1229            (
1230                "__op_type",
1231                Arc::new(UInt8Array::from(vec![1; 6])) as ArrayRef,
1232            ),
1233        ])
1234        .unwrap();
1235        let mut metadata = RegionMetadataBuilder::new(RegionId::new(1, 1));
1236        for (id, ty) in columns.iter().rev() {
1237            metadata.push_column_metadata(ColumnMetadata {
1238                column_id: *id,
1239                column_schema: ColumnSchema::new(format!("tag_{id}"), ty.clone(), true),
1240                semantic_type: SemanticType::Tag,
1241            });
1242        }
1243        metadata
1244            .push_column_metadata(ColumnMetadata {
1245                column_id: 23,
1246                column_schema: ColumnSchema::new(
1247                    "ts",
1248                    ConcreteDataType::timestamp_millisecond_datatype(),
1249                    false,
1250                ),
1251                semantic_type: SemanticType::Timestamp,
1252            })
1253            .primary_key(vec![17, 9, 2]);
1254        let metadata = Arc::new(metadata.build().unwrap());
1255        let format = FlatReadFormat::new(
1256            metadata.clone(),
1257            ReadColumns::new([17, 9, 2, 23]),
1258            Some(batch.schema()),
1259            "test",
1260            false,
1261        )
1262        .unwrap();
1263        for (offset, len) in [(0, 6), (1, 4), (3, 0)] {
1264            let mut decoded = decode_primary_keys(&codec, &batch.slice(offset, len)).unwrap();
1265            let all = decoded.get_dense_tag_columns().unwrap();
1266            let converted = format
1267                .convert_batch(batch.slice(offset, len), None)
1268                .unwrap();
1269            // Out-of-order and repeated projections share the same key offsets.
1270            for pos in [2, 0, 1, 0] {
1271                let array = decoded
1272                    .get_tag_column(columns[pos].0, Some(pos), &columns[pos].1)
1273                    .unwrap();
1274                assert_eq!(array.to_data(), all[pos].to_data());
1275                assert_eq!(converted.column(pos).to_data(), all[pos].to_data());
1276                let actual = Helper::try_into_vector(array).unwrap();
1277                for (row, &key) in row_keys[offset..offset + len].iter().enumerate() {
1278                    let eager = codec
1279                        .decode_dense_without_column_id(&encoded[key as usize])
1280                        .unwrap();
1281                    assert_eq!(actual.get(row), eager[pos]);
1282                }
1283            }
1284            // A position beyond the source schema remains NULL, never decoded
1285            // with a newer schema's field order or type.
1286            let missing = decoded
1287                .get_tag_column(99, Some(3), &ConcreteDataType::int64_datatype())
1288                .unwrap();
1289            assert_eq!(missing.null_count(), len);
1290            assert!(decoded.get_tag_column(17, None, &columns[0].1).is_err());
1291        }
1292        let mut malformed_columns = batch.columns().to_vec();
1293        malformed_columns[1] = Arc::new(DictionaryArray::<UInt32Type>::new(
1294            UInt32Array::from(vec![0; 6]),
1295            Arc::new(BinaryArray::from(vec![&[0, 1][..]])),
1296        ));
1297        let malformed = RecordBatch::try_new(batch.schema(), malformed_columns).unwrap();
1298        assert!(format.convert_batch(malformed.clone(), None).is_err());
1299        // A partial projection must still avoid decoding an unneeded broken suffix.
1300        let partial = FlatReadFormat::new(
1301            metadata,
1302            ReadColumns::new([17, 23]),
1303            Some(batch.schema()),
1304            "test",
1305            false,
1306        )
1307        .unwrap();
1308        let partial = partial.convert_batch(malformed, None).unwrap();
1309        let tag = Helper::try_into_vector(partial.column(0).clone()).unwrap();
1310        assert!((0..6).all(|row| tag.get(row).is_null()));
1311    }
1312
1313    /// Builds a `RegionMetadata` with the given number of tags and fields.
1314    fn build_metadata(
1315        num_tags: usize,
1316        num_fields: usize,
1317        encoding: PrimaryKeyEncoding,
1318    ) -> RegionMetadata {
1319        let mut builder = RegionMetadataBuilder::new(RegionId::new(0, 0));
1320        let mut col_id = 0u32;
1321
1322        for i in 0..num_tags {
1323            builder.push_column_metadata(ColumnMetadata {
1324                column_schema: ColumnSchema::new(
1325                    format!("tag_{i}"),
1326                    ConcreteDataType::string_datatype(),
1327                    true,
1328                ),
1329                semantic_type: SemanticType::Tag,
1330                column_id: col_id,
1331            });
1332            col_id += 1;
1333        }
1334
1335        for i in 0..num_fields {
1336            builder.push_column_metadata(ColumnMetadata {
1337                column_schema: ColumnSchema::new(
1338                    format!("field_{i}"),
1339                    ConcreteDataType::uint64_datatype(),
1340                    true,
1341                ),
1342                semantic_type: SemanticType::Field,
1343                column_id: col_id,
1344            });
1345            col_id += 1;
1346        }
1347
1348        builder.push_column_metadata(ColumnMetadata {
1349            column_schema: ColumnSchema::new(
1350                "ts".to_string(),
1351                ConcreteDataType::timestamp_millisecond_datatype(),
1352                false,
1353            ),
1354            semantic_type: SemanticType::Timestamp,
1355            column_id: col_id,
1356        });
1357
1358        let primary_key: Vec<u32> = (0..num_tags as u32).collect();
1359        builder.primary_key(primary_key);
1360        builder.primary_key_encoding(encoding);
1361        builder.build().unwrap()
1362    }
1363
1364    /// Builds the metadata of a table with a JSON2 struct field column:
1365    /// `[tag_0, field_0, payload, nullable_after_payload, ts]` with primary key `tag_0`.
1366    fn metadata_with_struct_field() -> RegionMetadata {
1367        let mut builder = RegionMetadataBuilder::new(RegionId::new(0, 0));
1368        builder
1369            .push_column_metadata(ColumnMetadata {
1370                column_schema: ColumnSchema::new(
1371                    "tag_0".to_string(),
1372                    ConcreteDataType::string_datatype(),
1373                    true,
1374                ),
1375                semantic_type: SemanticType::Tag,
1376                column_id: 0,
1377            })
1378            .push_column_metadata(ColumnMetadata {
1379                column_schema: ColumnSchema::new(
1380                    "field_0".to_string(),
1381                    ConcreteDataType::int64_datatype(),
1382                    true,
1383                ),
1384                semantic_type: SemanticType::Field,
1385                column_id: 1,
1386            })
1387            .push_column_metadata(ColumnMetadata {
1388                column_schema: ColumnSchema::new(
1389                    "payload".to_string(),
1390                    ConcreteDataType::json2(JsonNativeType::Object(JsonObjectType::new())),
1391                    false,
1392                ),
1393                semantic_type: SemanticType::Field,
1394                column_id: 2,
1395            })
1396            .push_column_metadata(ColumnMetadata {
1397                column_schema: ColumnSchema::new(
1398                    "nullable_after_payload".to_string(),
1399                    ConcreteDataType::int64_datatype(),
1400                    true,
1401                ),
1402                semantic_type: SemanticType::Field,
1403                column_id: 4,
1404            })
1405            .push_column_metadata(ColumnMetadata {
1406                column_schema: ColumnSchema::new(
1407                    "ts".to_string(),
1408                    ConcreteDataType::timestamp_nanosecond_datatype(),
1409                    false,
1410                ),
1411                semantic_type: SemanticType::Timestamp,
1412                column_id: 3,
1413            });
1414        builder.primary_key(vec![0]);
1415        builder.primary_key_encoding(PrimaryKeyEncoding::Dense);
1416        builder.build().unwrap()
1417    }
1418
1419    /// Builds a file schema and a row group in which the `payload` struct
1420    /// column expands to three leaf columns. The `ns_edge` leaf carries small
1421    /// Int64 statistics that must not be mistaken for the statistics of `ts`.
1422    fn struct_column_file_and_row_group(raw_pk_columns: bool) -> (SchemaRef, RowGroupMetaData) {
1423        let mut arrow_fields = vec![
1424            Field::new("tag_0", ArrowDataType::Utf8, true),
1425            Field::new("field_0", ArrowDataType::Int64, true),
1426            Field::new(
1427                "payload",
1428                ArrowDataType::Struct(
1429                    vec![
1430                        Field::new("metadata", ArrowDataType::Binary, false),
1431                        Field::new("value", ArrowDataType::Binary, false),
1432                        Field::new("ns_edge", ArrowDataType::Int64, true),
1433                    ]
1434                    .into(),
1435                ),
1436                false,
1437            ),
1438            Field::new("nullable_after_payload", ArrowDataType::Int64, true),
1439            Field::new(
1440                "ts",
1441                ArrowDataType::Timestamp(TimeUnit::Nanosecond, None),
1442                false,
1443            ),
1444            Field::new(PRIMARY_KEY_COLUMN_NAME, ArrowDataType::Binary, false),
1445            Field::new(SEQUENCE_COLUMN_NAME, ArrowDataType::UInt64, false),
1446            Field::new(OP_TYPE_COLUMN_NAME, ArrowDataType::UInt8, false),
1447        ];
1448        if !raw_pk_columns {
1449            arrow_fields.remove(0);
1450        }
1451        let file_schema = Arc::new(Schema::new(arrow_fields));
1452
1453        let leaf = |name: &str, physical: PhysicalType| {
1454            Arc::new(
1455                Type::primitive_type_builder(name, physical)
1456                    .with_repetition(Repetition::OPTIONAL)
1457                    .build()
1458                    .unwrap(),
1459            )
1460        };
1461        let payload = Arc::new(
1462            Type::group_type_builder("payload")
1463                .with_repetition(Repetition::OPTIONAL)
1464                .with_fields(vec![
1465                    leaf("metadata", PhysicalType::BYTE_ARRAY),
1466                    leaf("value", PhysicalType::BYTE_ARRAY),
1467                    leaf("ns_edge", PhysicalType::INT64),
1468                ])
1469                .build()
1470                .unwrap(),
1471        );
1472        let mut parquet_fields = vec![
1473            leaf("tag_0", PhysicalType::BYTE_ARRAY),
1474            leaf("field_0", PhysicalType::INT64),
1475            payload,
1476            leaf("nullable_after_payload", PhysicalType::INT64),
1477            leaf("ts", PhysicalType::INT64),
1478            leaf(PRIMARY_KEY_COLUMN_NAME, PhysicalType::BYTE_ARRAY),
1479            leaf(SEQUENCE_COLUMN_NAME, PhysicalType::INT64),
1480            leaf(OP_TYPE_COLUMN_NAME, PhysicalType::INT32),
1481        ];
1482        if !raw_pk_columns {
1483            parquet_fields.remove(0);
1484        }
1485        let schema_descr = Arc::new(SchemaDescriptor::new(Arc::new(
1486            Type::group_type_builder("schema")
1487                .with_fields(parquet_fields)
1488                .build()
1489                .unwrap(),
1490        )));
1491
1492        // Omitted raw tags shift all subsequent leaves in primary-key SSTs.
1493        let ns_edge_leaf = if raw_pk_columns { 4 } else { 3 };
1494        let nullable_leaf = ns_edge_leaf + 1;
1495        let ts_leaf = nullable_leaf + 1;
1496        let chunks: Vec<_> = (0..schema_descr.num_columns())
1497            .map(|i| {
1498                let mut builder = ColumnChunkMetaData::builder(schema_descr.column(i));
1499                if i == ns_edge_leaf {
1500                    // Small values from the JSON payload, not timestamps.
1501                    builder = builder.set_statistics(Statistics::int64(
1502                        Some(0),
1503                        Some(86_400_000_000_000),
1504                        None,
1505                        Some(65),
1506                        true,
1507                    ));
1508                } else if i == nullable_leaf {
1509                    builder = builder.set_statistics(Statistics::int64(
1510                        Some(100),
1511                        Some(200),
1512                        None,
1513                        Some(7),
1514                        true,
1515                    ));
1516                } else if i == ts_leaf {
1517                    builder = builder.set_statistics(Statistics::int64(
1518                        Some(1_788_998_400_000_000_000),
1519                        Some(1_789_084_800_000_000_000),
1520                        None,
1521                        Some(0),
1522                        true,
1523                    ));
1524                }
1525                builder.build().unwrap()
1526            })
1527            .collect();
1528        let row_group = RowGroupMetaData::builder(schema_descr)
1529            .set_num_rows(69)
1530            .set_total_byte_size(0)
1531            .set_column_metadata(chunks)
1532            .build()
1533            .unwrap();
1534
1535        (file_schema, row_group)
1536    }
1537
1538    /// Regression test: row group statistics must be looked up by parquet leaf
1539    /// column index. A struct field column (e.g. JSON2) expands to multiple
1540    /// leaf columns, so statistics of columns after it must not be read from
1541    /// the struct's leaves. Otherwise min-max pruning can drop a whole row
1542    /// group by mistake (e.g. pruning `ts` with the small `ns_edge` stats),
1543    /// which caused data loss during SWCS compaction.
1544    #[test]
1545    fn test_stats_with_struct_field_column() {
1546        for (encoding, raw_pk_columns) in [
1547            (PrimaryKeyEncoding::Dense, true),
1548            (PrimaryKeyEncoding::Dense, false),
1549            (PrimaryKeyEncoding::Sparse, false),
1550        ] {
1551            let mut metadata = metadata_with_struct_field();
1552            metadata.primary_key_encoding = encoding;
1553            let metadata = Arc::new(metadata);
1554            let (file_schema, row_group) = struct_column_file_and_row_group(raw_pk_columns);
1555            let read_format = FlatReadFormat::new(
1556                metadata,
1557                ReadColumns::new([0, 1, 2, 3, 4]),
1558                Some(file_schema),
1559                "test",
1560                false,
1561            )
1562            .unwrap();
1563            let row_groups = [&row_group];
1564
1565            // Statistics of `ts` come from the `ts` leaf column, not the leaves of
1566            // the payload struct.
1567            let StatValues::Values(min) = read_format.min_values(&row_groups, 3) else {
1568                panic!("expected ts min values")
1569            };
1570            let min = min.as_any().downcast_ref::<Int64Array>().unwrap();
1571            assert_eq!(1_788_998_400_000_000_000, min.value(0));
1572            let StatValues::Values(max) = read_format.max_values(&row_groups, 3) else {
1573                panic!("expected ts max values")
1574            };
1575            let max = max.as_any().downcast_ref::<Int64Array>().unwrap();
1576            assert_eq!(1_789_084_800_000_000_000, max.value(0));
1577
1578            let stats = crate::sst::parquet::stats::RowGroupPruningStats::new(
1579                &row_groups,
1580                &read_format,
1581                None,
1582                false,
1583            );
1584            for (start, end, keep) in [
1585                (1_788_998_400_000_000_000, 1_789_084_800_000_000_000, true),
1586                (1_789_084_800_000_000_001, 1_789_171_200_000_000_000, false),
1587            ] {
1588                let predicate = table::predicate::Predicate::new(vec![
1589                    datafusion_expr::col("ts").gt_eq(datafusion_expr::lit(
1590                        datafusion_common::ScalarValue::TimestampNanosecond(Some(start), None),
1591                    )),
1592                    datafusion_expr::col("ts").lt(datafusion_expr::lit(
1593                        datafusion_common::ScalarValue::TimestampNanosecond(Some(end), None),
1594                    )),
1595                ]);
1596                assert_eq!(
1597                    vec![keep],
1598                    predicate
1599                        .prune_with_stats(&stats, read_format.metadata().schema.arrow_schema(),)
1600                );
1601            }
1602
1603            // Null counts of `ts` also read the correct leaf column.
1604            let StatValues::Values(nulls) = read_format.null_counts(&row_groups, 3) else {
1605                panic!("expected ts null counts")
1606            };
1607            let nulls = nulls.as_any().downcast_ref::<UInt64Array>().unwrap();
1608            assert!(nulls.is_valid(0));
1609            assert_eq!(0, nulls.value(0));
1610
1611            // A null slot may contain an underlying zero. Check validity and
1612            // a nonzero count to distinguish unknown or wrong-leaf statistics.
1613            let StatValues::Values(nulls) = read_format.null_counts(&row_groups, 4) else {
1614                panic!("expected nullable field null counts")
1615            };
1616            let nulls = nulls.as_any().downcast_ref::<UInt64Array>().unwrap();
1617            assert!(nulls.is_valid(0));
1618            assert_eq!(7, nulls.value(0));
1619
1620            // A column that expands to multiple leaf columns has no single column
1621            // statistics.
1622            assert!(matches!(
1623                read_format.min_values(&row_groups, 2),
1624                StatValues::NoStats
1625            ));
1626            assert!(matches!(
1627                read_format.max_values(&row_groups, 2),
1628                StatValues::NoStats
1629            ));
1630            assert!(matches!(
1631                read_format.null_counts(&row_groups, 2),
1632                StatValues::NoStats
1633            ));
1634        }
1635    }
1636
1637    /// Even one leaf cannot supply the null count of a nested parent:
1638    /// {"a": null} and [null] are non-null parents with null children.
1639    #[test]
1640    fn test_single_leaf_nested_columns_have_no_root_stats() {
1641        let child = Arc::new(Field::new("a", ArrowDataType::Int64, true));
1642        for nested_type in [
1643            ArrowDataType::Struct(vec![child.clone()].into()),
1644            ArrowDataType::List(child.clone()),
1645            ArrowDataType::LargeList(child.clone()),
1646            ArrowDataType::FixedSizeList(child, 1),
1647        ] {
1648            let metadata = Arc::new(metadata_with_struct_field());
1649            let (schema, _) = struct_column_file_and_row_group(true);
1650            let mut fields = schema.fields().to_vec();
1651            fields[2] = Arc::new(Field::new("payload", nested_type.clone(), true));
1652            let schema = Arc::new(Schema::new(fields));
1653            let descriptor = Arc::new(ArrowSchemaConverter::new().convert(&schema).unwrap());
1654            let chunks = descriptor
1655                .columns()
1656                .iter()
1657                .map(|column| {
1658                    ColumnChunkMetaData::builder(column.clone())
1659                        .build()
1660                        .unwrap()
1661                })
1662                .collect();
1663            let row_group = RowGroupMetaData::builder(descriptor)
1664                .set_num_rows(1)
1665                .set_total_byte_size(0)
1666                .set_column_metadata(chunks)
1667                .build()
1668                .unwrap();
1669            let format = ParquetFlat::new(metadata, ReadColumns::new([0, 1, 2, 3]), schema);
1670
1671            // The nested root has no usable statistics, independently of the
1672            // values stored in any row group. Its leaf still occupies a slot.
1673            let groups = &[row_group];
1674            assert!(
1675                matches!(format.min_values(groups, 2), StatValues::NoStats),
1676                "{nested_type:?}"
1677            );
1678            assert!(
1679                matches!(format.max_values(groups, 2), StatValues::NoStats),
1680                "{nested_type:?}"
1681            );
1682            assert!(
1683                matches!(format.null_counts(groups, 2), StatValues::NoStats),
1684                "{nested_type:?}"
1685            );
1686        }
1687    }
1688
1689    #[test]
1690    fn test_field_column_start() {
1691        // (num_tags, num_fields, encoding, expected)
1692        let cases = [
1693            (1, 1, PrimaryKeyEncoding::Dense, 1),
1694            (2, 2, PrimaryKeyEncoding::Dense, 2),
1695            (0, 2, PrimaryKeyEncoding::Dense, 0),
1696            (2, 2, PrimaryKeyEncoding::Sparse, 0),
1697        ];
1698
1699        for (num_tags, num_fields, encoding, expected) in cases {
1700            let metadata = build_metadata(num_tags, num_fields, encoding);
1701            let options = FlatSchemaOptions::from_encoding(encoding);
1702            let num_columns = flat_sst_arrow_schema_column_num(&metadata, &options);
1703            let result = field_column_start(&metadata, num_columns);
1704            assert_eq!(
1705                result, expected,
1706                "num_tags={num_tags}, num_fields={num_fields}, encoding={encoding:?}"
1707            );
1708        }
1709    }
1710
1711    #[test]
1712    fn test_convert_batch_wraps_binary_pk_to_dict() {
1713        use datatypes::arrow::array::{Array, DictionaryArray, StringArray};
1714        use datatypes::arrow::datatypes::UInt32Type;
1715
1716        // build_metadata(1, 1, Dense) projects to:
1717        // [tag_0: Dict<UInt32, Utf8>, field_0: UInt64, ts: Timestamp(ms),
1718        //  __primary_key: Dict<UInt32, Binary>, __sequence: UInt64, __op_type: UInt8]
1719        let metadata = Arc::new(build_metadata(1, 1, PrimaryKeyEncoding::Dense));
1720        let column_ids: Vec<u32> = metadata
1721            .column_metadatas
1722            .iter()
1723            .map(|c| c.column_id)
1724            .collect();
1725        let mut read_format = FlatReadFormat::new(
1726            metadata.clone(),
1727            ReadColumns::new(column_ids),
1728            None,
1729            "test",
1730            false,
1731        )
1732        .unwrap();
1733        let output_schema = read_format.output_arrow_schema().unwrap();
1734        read_format.set_pk_as_binary(output_schema.clone());
1735        let binary_schema = override_pk_field_to_binary(&output_schema);
1736
1737        // The __primary_key field must preserve its field_id metadata after
1738        // being converted from dictionary to plain binary.
1739        let pk_field = binary_schema
1740            .field_with_name(PRIMARY_KEY_COLUMN_NAME)
1741            .unwrap();
1742        assert_eq!(
1743            pk_field.metadata().get(PARQUET_FIELD_ID_KEY),
1744            Some(&PRIMARY_KEY_PARQUET_FIELD_ID.to_string()),
1745            "__primary_key field must retain its PARQUET:field_id after override_pk_field_to_binary"
1746        );
1747
1748        // Repeat the second pk to verify identity keys (no dedup).
1749        let tag_keys = UInt32Array::from(vec![0u32, 1, 1]);
1750        let tag_values = Arc::new(StringArray::from(vec!["t0", "t1"]));
1751        let tag_array: ArrayRef =
1752            Arc::new(DictionaryArray::<UInt32Type>::new(tag_keys, tag_values));
1753        let field_array: ArrayRef = Arc::new(UInt64Array::from(vec![10u64, 11, 12]));
1754        let ts_array: ArrayRef = Arc::new(TimestampMillisecondArray::from(vec![1i64, 2, 3]));
1755        let pk_array: ArrayRef = Arc::new(BinaryArray::from_iter_values(
1756            [b"alpha".as_ref(), b"beta", b"beta"].iter().copied(),
1757        ));
1758        let seq_array: ArrayRef = Arc::new(UInt64Array::from(vec![100u64, 101, 102]));
1759        let op_array: ArrayRef = Arc::new(UInt8Array::from(vec![1u8, 1, 1]));
1760
1761        let batch = RecordBatch::try_new(
1762            binary_schema,
1763            vec![
1764                tag_array,
1765                field_array,
1766                ts_array,
1767                pk_array,
1768                seq_array,
1769                op_array,
1770            ],
1771        )
1772        .unwrap();
1773
1774        let wrapped = read_format.convert_batch(batch, None).unwrap();
1775        assert_eq!(wrapped.schema(), output_schema);
1776
1777        let pk_idx = primary_key_column_index(wrapped.num_columns());
1778        let pk_col = wrapped.column(pk_idx);
1779        assert_eq!(
1780            pk_col.data_type(),
1781            &ArrowDataType::Dictionary(
1782                Box::new(ArrowDataType::UInt32),
1783                Box::new(ArrowDataType::Binary)
1784            )
1785        );
1786        let dict = pk_col
1787            .as_any()
1788            .downcast_ref::<DictionaryArray<UInt32Type>>()
1789            .unwrap();
1790        assert_eq!(dict.keys().values(), &[0, 1, 2]);
1791        let values = dict
1792            .values()
1793            .as_any()
1794            .downcast_ref::<BinaryArray>()
1795            .unwrap();
1796        assert_eq!(values.value(0), b"alpha");
1797        assert_eq!(values.value(1), b"beta");
1798        assert_eq!(values.value(2), b"beta");
1799    }
1800
1801    #[test]
1802    fn test_output_arrow_schema_uses_projection() {
1803        let metadata = Arc::new(build_metadata(1, 2, PrimaryKeyEncoding::Dense));
1804        let read_format = FlatReadFormat::new(
1805            metadata.clone(),
1806            ReadColumns::new([0_u32, 2_u32]),
1807            None,
1808            "test",
1809            false,
1810        )
1811        .unwrap();
1812
1813        let output_schema = read_format.output_arrow_schema().unwrap();
1814        let projection = read_format.parquet_read_columns().root_indices();
1815        let expected = Arc::new(
1816            to_flat_sst_arrow_schema(&metadata, &FlatSchemaOptions::default())
1817                .project(projection)
1818                .unwrap(),
1819        );
1820
1821        assert_eq!(expected, output_schema);
1822    }
1823}