Skip to main content

mito2/read/
compat.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//! Utilities to adapt readers with different schema.
16
17use std::collections::HashMap;
18use std::sync::Arc;
19
20use api::v1::SemanticType;
21use datatypes::arrow::array::{
22    Array, ArrayRef, BinaryArray, BinaryBuilder, DictionaryArray, UInt32Array,
23};
24use datatypes::arrow::compute::{TakeOptions, take};
25use datatypes::arrow::datatypes::{FieldRef, Schema, SchemaRef};
26use datatypes::arrow::record_batch::RecordBatch;
27use datatypes::data_type::ConcreteDataType;
28use datatypes::extension::json::is_json2_extension_type;
29use datatypes::prelude::DataType;
30use datatypes::value::Value;
31use datatypes::vectors::VectorRef;
32use datatypes::vectors::json::array::JsonArray;
33use mito_codec::row_converter::{
34    CompositeValues, PrimaryKeyCodec, SortField, build_primary_key_codec,
35    build_primary_key_codec_with_fields,
36};
37use snafu::{OptionExt, ResultExt, ensure};
38use store_api::codec::PrimaryKeyEncoding;
39use store_api::metadata::{RegionMetadata, RegionMetadataRef};
40use store_api::storage::ColumnId;
41
42use crate::error::{
43    CompatReaderSnafu, ComputeArrowSnafu, ConvertValueSnafu, CreateDefaultSnafu, DecodeSnafu,
44    EncodeSnafu, NewRecordBatchSnafu, Result, UnexpectedSnafu, UnsupportedOperationSnafu,
45};
46use crate::read::flat_projection::{FlatProjectionMapper, flat_projected_columns};
47use crate::sst::parquet::flat_format::{FlatReadFormat, primary_key_column_index};
48use crate::sst::parquet::format::{INTERNAL_COLUMN_NUM, PrimaryKeyArray};
49use crate::sst::{internal_fields, tag_maybe_to_dictionary_field, with_field_id};
50
51/// Returns true if the columns in the `projection_mapper` and `read_format` have same data types
52/// and primary key encodings.
53pub(crate) fn has_same_columns_and_pk_encoding(
54    projection_mapper: &FlatProjectionMapper,
55    read_format: &FlatReadFormat,
56    compaction: bool,
57) -> bool {
58    let left = projection_mapper.metadata();
59    let right = read_format.metadata();
60    if left.primary_key_encoding != right.primary_key_encoding {
61        return false;
62    }
63
64    if left.column_metadatas.len() != right.column_metadatas.len() {
65        return false;
66    }
67
68    for (left_col, right_col) in left.column_metadatas.iter().zip(&right.column_metadatas) {
69        if left_col.column_id != right_col.column_id {
70            return false;
71        }
72        debug_assert_eq!(left_col.semantic_type, right_col.semantic_type);
73    }
74
75    &projection_mapper.input_arrow_schema(compaction) == read_format.arrow_schema()
76}
77
78/// A helper struct to adapt schema of the batch to an expected schema.
79pub(crate) struct FlatCompatBatch {
80    /// Indices to convert actual fields to expect fields.
81    index_or_defaults: Vec<IndexOrDefault>,
82    /// Expected arrow schema.
83    arrow_schema: SchemaRef,
84    /// Primary key adapter.
85    compat_pk: FlatCompatPrimaryKey,
86}
87
88impl FlatCompatBatch {
89    /// Creates a [FlatCompatBatch].
90    ///
91    /// - `mapper` is built from the metadata users expect to see.
92    /// - `read_format` is the [FlatReadFormat] of the input parquet.
93    /// - `compaction` indicates whether the reader is for compaction.
94    pub(crate) fn try_new(
95        mapper: &FlatProjectionMapper,
96        read_format: &FlatReadFormat,
97        compaction: bool,
98    ) -> Result<Option<Self>> {
99        let actual = read_format.metadata();
100        let format_projection = read_format.format_projection();
101        let mut actual_schema = flat_projected_columns(actual, format_projection);
102        if read_format
103            .arrow_schema()
104            .fields()
105            .iter()
106            .any(is_json2_extension_type)
107        {
108            for field in read_format.arrow_schema().fields() {
109                if is_json2_extension_type(field)
110                    && let Some(column_id) =
111                        actual.column_by_name(field.name()).map(|x| x.column_id)
112                    && let Some(i) = actual_schema.iter().position(|x| x.0 == column_id)
113                {
114                    actual_schema[i].1 = ConcreteDataType::from_arrow_type(field.data_type());
115                }
116            }
117        }
118
119        let expect_schema = mapper.batch_schema();
120        if expect_schema == actual_schema
121            && actual.primary_key == mapper.metadata().primary_key
122            && actual.primary_key_encoding == mapper.metadata().primary_key_encoding
123        {
124            // Although the SST has a different schema, but the schema after projection is the same
125            // as expected schema.
126            return Ok(None);
127        }
128
129        if actual.primary_key_encoding == PrimaryKeyEncoding::Sparse && compaction {
130            // Special handling for sparse encoding in compaction.
131            return FlatCompatBatch::try_new_compact_sparse(mapper, actual);
132        }
133
134        let (index_or_defaults, fields) =
135            Self::compute_index_and_fields(&actual_schema, expect_schema, mapper.metadata())?;
136
137        let compat_pk = FlatCompatPrimaryKey::new(mapper.metadata(), actual)?;
138
139        Ok(Some(Self {
140            index_or_defaults,
141            arrow_schema: Arc::new(Schema::new(fields)),
142            compat_pk,
143        }))
144    }
145
146    fn compute_index_and_fields(
147        actual_schema: &[(ColumnId, ConcreteDataType)],
148        expect_schema: &[(ColumnId, ConcreteDataType)],
149        expect_metadata: &RegionMetadata,
150    ) -> Result<(Vec<IndexOrDefault>, Vec<FieldRef>)> {
151        // Maps column id to the index and data type in the actual schema.
152        let actual_schema_index: HashMap<_, _> = actual_schema
153            .iter()
154            .enumerate()
155            .map(|(idx, (column_id, data_type))| (*column_id, (idx, data_type)))
156            .collect();
157
158        let mut index_or_defaults = Vec::with_capacity(expect_schema.len());
159        let mut fields = Vec::with_capacity(expect_schema.len());
160        for (column_id, expect_data_type) in expect_schema {
161            // Safety: expect_schema comes from the same mapper.
162            let column_index = expect_metadata.column_index_by_id(*column_id).unwrap();
163            let expect_column = &expect_metadata.column_metadatas[column_index];
164            let column_field = &expect_metadata.schema.arrow_schema().fields()[column_index];
165            // For tag columns, we need to create a dictionary field.
166            if expect_column.semantic_type == SemanticType::Tag {
167                let field = tag_maybe_to_dictionary_field(
168                    &expect_column.column_schema.data_type,
169                    column_field,
170                );
171                fields.push(Arc::new(with_field_id(
172                    (*field).clone(),
173                    expect_column.column_id,
174                )));
175            } else {
176                let field = with_field_id(
177                    Arc::unwrap_or_clone(column_field.clone()),
178                    expect_column.column_id,
179                )
180                .with_data_type(expect_data_type.as_arrow_type());
181                fields.push(Arc::new(field))
182            };
183
184            if let Some((index, actual_data_type)) = actual_schema_index.get(column_id) {
185                let mut cast_type = None;
186
187                // Same column different type.
188                if expect_data_type != *actual_data_type {
189                    cast_type = Some(expect_data_type.clone())
190                }
191                // Source has this column.
192                index_or_defaults.push(IndexOrDefault::Index {
193                    pos: *index,
194                    cast_type,
195                });
196            } else {
197                // Create a default vector with 1 element for that column.
198                let default_vector = expect_column
199                    .column_schema
200                    .create_default_vector(1)
201                    .context(CreateDefaultSnafu {
202                        region_id: expect_metadata.region_id,
203                        column: &expect_column.column_schema.name,
204                    })?
205                    .with_context(|| CompatReaderSnafu {
206                        region_id: expect_metadata.region_id,
207                        reason: format!(
208                            "column {} does not have a default value to read",
209                            expect_column.column_schema.name
210                        ),
211                    })?;
212                index_or_defaults.push(IndexOrDefault::DefaultValue {
213                    default_vector,
214                    semantic_type: expect_column.semantic_type,
215                });
216            };
217        }
218        fields.extend_from_slice(&internal_fields());
219
220        Ok((index_or_defaults, fields))
221    }
222
223    fn try_new_compact_sparse(
224        mapper: &FlatProjectionMapper,
225        actual: &RegionMetadataRef,
226    ) -> Result<Option<Self>> {
227        // Currently, we don't support converting sparse encoding back to dense encoding in
228        // flat format.
229        ensure!(
230            mapper.metadata().primary_key_encoding == PrimaryKeyEncoding::Sparse,
231            UnsupportedOperationSnafu {
232                err_msg: "Flat format doesn't support converting sparse encoding back to dense encoding"
233            }
234        );
235
236        // For sparse encoding, we don't need to check the primary keys.
237        // Since this is for compaction, we always read all columns.
238        let actual_schema: Vec<_> = actual
239            .field_columns()
240            .chain([actual.time_index_column()])
241            .map(|col| (col.column_id, col.column_schema.data_type.clone()))
242            .collect();
243        let expect_schema: Vec<_> = mapper
244            .metadata()
245            .field_columns()
246            .chain([mapper.metadata().time_index_column()])
247            .map(|col| (col.column_id, col.column_schema.data_type.clone()))
248            .collect();
249
250        let (index_or_defaults, fields) =
251            Self::compute_index_and_fields(&actual_schema, &expect_schema, mapper.metadata())?;
252
253        let compat_pk = FlatCompatPrimaryKey::default();
254
255        Ok(Some(Self {
256            index_or_defaults,
257            arrow_schema: Arc::new(Schema::new(fields)),
258            compat_pk,
259        }))
260    }
261
262    /// Make columns of the `batch` compatible.
263    pub(crate) fn compat(&self, batch: RecordBatch) -> Result<RecordBatch> {
264        let len = batch.num_rows();
265        let columns = self
266            .index_or_defaults
267            .iter()
268            .map(|index_or_default| match index_or_default {
269                IndexOrDefault::Index { pos, cast_type } => {
270                    let old_column = batch.column(*pos);
271
272                    if let Some(ty) = cast_type {
273                        let casted = if let Some(json_type) = ty.as_json()
274                            && json_type.is_json2()
275                        {
276                            JsonArray::from(old_column)
277                                .project_to(&json_type.as_arrow_type())
278                                .context(ConvertValueSnafu)?
279                        } else {
280                            datatypes::arrow::compute::cast(old_column, &ty.as_arrow_type())
281                                .context(ComputeArrowSnafu)?
282                        };
283                        Ok(casted)
284                    } else {
285                        Ok(old_column.clone())
286                    }
287                }
288                IndexOrDefault::DefaultValue {
289                    default_vector,
290                    semantic_type,
291                } => repeat_vector(default_vector, len, *semantic_type == SemanticType::Tag),
292            })
293            .chain(
294                // Adds internal columns.
295                batch.columns()[batch.num_columns() - INTERNAL_COLUMN_NUM..]
296                    .iter()
297                    .map(|col| Ok(col.clone())),
298            )
299            .collect::<Result<Vec<_>>>()?;
300
301        let mut columns = columns;
302        let primary_key_index = primary_key_column_index(columns.len());
303        columns[primary_key_index] = self.compat_primary_key(&columns[primary_key_index])?;
304
305        RecordBatch::try_new(self.arrow_schema.clone(), columns).context(NewRecordBatchSnafu)
306    }
307
308    /// Makes an encoded primary-key array compatible with the expected metadata.
309    pub(crate) fn compat_primary_key(&self, primary_key: &ArrayRef) -> Result<ArrayRef> {
310        self.compat_pk.compat(primary_key)
311    }
312}
313
314/// Repeats the vector value `to_len` times.
315fn repeat_vector(vector: &VectorRef, to_len: usize, is_tag: bool) -> Result<ArrayRef> {
316    assert_eq!(1, vector.len());
317    let data_type = vector.data_type();
318    if is_tag && data_type.is_string() {
319        let values = vector.to_arrow_array();
320        if values.is_null(0) {
321            // Creates a dictionary array with `to_len` null keys.
322            let keys = UInt32Array::new_null(to_len);
323            Ok(Arc::new(DictionaryArray::new(keys, values.slice(0, 0))))
324        } else {
325            let keys = UInt32Array::from_value(0, to_len);
326            Ok(Arc::new(DictionaryArray::new(keys, values)))
327        }
328    } else {
329        let keys = UInt32Array::from_value(0, to_len);
330        take(
331            &vector.to_arrow_array(),
332            &keys,
333            Some(TakeOptions {
334                check_bounds: false,
335            }),
336        )
337        .context(ComputeArrowSnafu)
338    }
339}
340
341/// Returns true if the actual primary keys is the same as expected.
342fn is_primary_key_same(expect: &RegionMetadata, actual: &RegionMetadata) -> Result<bool> {
343    ensure!(
344        actual.primary_key.len() <= expect.primary_key.len(),
345        CompatReaderSnafu {
346            region_id: expect.region_id,
347            reason: format!(
348                "primary key has more columns {} than expect {}",
349                actual.primary_key.len(),
350                expect.primary_key.len()
351            ),
352        }
353    );
354    ensure!(
355        actual.primary_key == expect.primary_key[..actual.primary_key.len()],
356        CompatReaderSnafu {
357            region_id: expect.region_id,
358            reason: format!(
359                "primary key has different prefix, expect: {:?}, actual: {:?}",
360                expect.primary_key, actual.primary_key
361            ),
362        }
363    );
364
365    Ok(actual.primary_key.len() == expect.primary_key.len())
366}
367
368/// Index in source batch or a default value to fill a column.
369#[derive(Debug)]
370enum IndexOrDefault {
371    /// Index of the column in source batch.
372    Index {
373        pos: usize,
374        cast_type: Option<ConcreteDataType>,
375    },
376    /// Default value for the column.
377    DefaultValue {
378        /// Default value. The vector has only 1 element.
379        default_vector: VectorRef,
380        /// Semantic type of the column.
381        semantic_type: SemanticType,
382    },
383}
384
385/// Helper to rewrite primary key to another encoding for flat format.
386struct FlatRewritePrimaryKey {
387    /// New primary key encoder.
388    codec: Arc<dyn PrimaryKeyCodec>,
389    /// Metadata of the expected region.
390    metadata: RegionMetadataRef,
391    /// Original primary key codec.
392    /// If we need to rewrite the primary key.
393    old_codec: Arc<dyn PrimaryKeyCodec>,
394}
395
396impl FlatRewritePrimaryKey {
397    fn new(
398        expect: &RegionMetadataRef,
399        actual: &RegionMetadataRef,
400    ) -> Option<FlatRewritePrimaryKey> {
401        if expect.primary_key_encoding == actual.primary_key_encoding {
402            return None;
403        }
404        let codec = build_primary_key_codec(expect);
405        let old_codec = build_primary_key_codec(actual);
406
407        Some(FlatRewritePrimaryKey {
408            codec,
409            metadata: expect.clone(),
410            old_codec,
411        })
412    }
413
414    /// Rewrites the primary key of the `batch`.
415    /// It also appends the values to the primary key.
416    fn rewrite_key(
417        &self,
418        append_values: &[(ColumnId, Value)],
419        primary_key: &ArrayRef,
420    ) -> Result<ArrayRef> {
421        if let Some(old_pk_dict_array) = primary_key.as_any().downcast_ref::<PrimaryKeyArray>() {
422            let old_pk_values_array = old_pk_dict_array
423                .values()
424                .as_any()
425                .downcast_ref::<BinaryArray>()
426                .context(UnexpectedSnafu {
427                    reason: "Primary-key dictionary values are not binary",
428                })?;
429            let new_pk_values_array =
430                Arc::new(self.rewrite_values(append_values, old_pk_values_array)?);
431            return Ok(Arc::new(PrimaryKeyArray::new(
432                old_pk_dict_array.keys().clone(),
433                new_pk_values_array,
434            )));
435        }
436
437        let old_pk_values_array =
438            primary_key
439                .as_any()
440                .downcast_ref::<BinaryArray>()
441                .context(UnexpectedSnafu {
442                    reason: format!(
443                        "Primary-key column is neither binary nor dictionary, got {:?}",
444                        primary_key.data_type()
445                    ),
446                })?;
447        Ok(Arc::new(
448            self.rewrite_values(append_values, old_pk_values_array)?,
449        ))
450    }
451
452    fn rewrite_values(
453        &self,
454        append_values: &[(ColumnId, Value)],
455        old_pk_values_array: &BinaryArray,
456    ) -> Result<BinaryArray> {
457        let mut builder = BinaryBuilder::with_capacity(
458            old_pk_values_array.len(),
459            old_pk_values_array.value_data().len(),
460        );
461
462        // Binary buffer for the primary key.
463        let mut buffer = Vec::with_capacity(
464            old_pk_values_array.value_data().len() / old_pk_values_array.len().max(1),
465        );
466        let mut column_id_values = Vec::new();
467        // Iterates the binary array and rewrites the primary key.
468        for value in old_pk_values_array.iter() {
469            let Some(old_pk) = value else {
470                builder.append_null();
471                continue;
472            };
473            // Decodes the old primary key.
474            let mut pk_values = self.old_codec.decode(old_pk).context(DecodeSnafu)?;
475            pk_values.extend(append_values);
476
477            buffer.clear();
478            column_id_values.clear();
479            // Encodes the new primary key.
480            match pk_values {
481                CompositeValues::Dense(dense_values) => {
482                    self.codec
483                        .encode_values(dense_values.as_slice(), &mut buffer)
484                        .context(EncodeSnafu)?;
485                }
486                CompositeValues::Sparse(sparse_values) => {
487                    for id in &self.metadata.primary_key {
488                        let value = sparse_values.get_or_null(*id);
489                        column_id_values.push((*id, value.clone()));
490                    }
491                    self.codec
492                        .encode_values(&column_id_values, &mut buffer)
493                        .context(EncodeSnafu)?;
494                }
495            }
496            builder.append_value(&buffer);
497        }
498        Ok(builder.finish())
499    }
500}
501
502/// Helper to make primary key compatible for flat format.
503#[derive(Default)]
504struct FlatCompatPrimaryKey {
505    /// Primary key rewriter.
506    rewriter: Option<FlatRewritePrimaryKey>,
507    /// Converter to append values to primary keys.
508    converter: Option<Arc<dyn PrimaryKeyCodec>>,
509    /// Default values to append.
510    values: Vec<(ColumnId, Value)>,
511}
512
513impl FlatCompatPrimaryKey {
514    fn new(expect: &RegionMetadataRef, actual: &RegionMetadataRef) -> Result<Self> {
515        let rewriter = FlatRewritePrimaryKey::new(expect, actual);
516
517        if is_primary_key_same(expect, actual)? {
518            return Ok(Self {
519                rewriter,
520                converter: None,
521                values: Vec::new(),
522            });
523        }
524
525        // We need to append default values to the primary key.
526        let to_add = &expect.primary_key[actual.primary_key.len()..];
527        let mut values = Vec::with_capacity(to_add.len());
528        let mut fields = Vec::with_capacity(to_add.len());
529        for column_id in to_add {
530            // Safety: The id comes from expect region metadata.
531            let column = expect.column_by_id(*column_id).unwrap();
532            fields.push((
533                *column_id,
534                SortField::new(column.column_schema.data_type.clone()),
535            ));
536            let default_value = column
537                .column_schema
538                .create_default()
539                .context(CreateDefaultSnafu {
540                    region_id: expect.region_id,
541                    column: &column.column_schema.name,
542                })?
543                .with_context(|| CompatReaderSnafu {
544                    region_id: expect.region_id,
545                    reason: format!(
546                        "key column {} does not have a default value to read",
547                        column.column_schema.name
548                    ),
549                })?;
550            values.push((*column_id, default_value));
551        }
552        // is_primary_key_same() is false so we have different number of primary key columns.
553        debug_assert!(!fields.is_empty());
554
555        // Create converter to append values.
556        let converter = Some(build_primary_key_codec_with_fields(
557            expect.primary_key_encoding,
558            fields.into_iter(),
559        ));
560
561        Ok(Self {
562            rewriter,
563            converter,
564            values,
565        })
566    }
567
568    /// Makes an encoded primary-key array compatible.
569    fn compat(&self, primary_key: &ArrayRef) -> Result<ArrayRef> {
570        if let Some(rewriter) = &self.rewriter {
571            // If we have different encoding, rewrite the whole primary key.
572            return rewriter.rewrite_key(&self.values, primary_key);
573        }
574
575        self.append_key(primary_key)
576    }
577
578    /// Appends values to the primary key array.
579    fn append_key(&self, primary_key: &ArrayRef) -> Result<ArrayRef> {
580        let Some(converter) = &self.converter else {
581            return Ok(primary_key.clone());
582        };
583
584        if let Some(old_pk_dict_array) = primary_key.as_any().downcast_ref::<PrimaryKeyArray>() {
585            let old_pk_values_array = old_pk_dict_array
586                .values()
587                .as_any()
588                .downcast_ref::<BinaryArray>()
589                .context(UnexpectedSnafu {
590                    reason: "Primary-key dictionary values are not binary",
591                })?;
592            let new_pk_values_array =
593                Arc::new(self.append_values(old_pk_values_array, converter.as_ref())?);
594            return Ok(Arc::new(PrimaryKeyArray::new(
595                old_pk_dict_array.keys().clone(),
596                new_pk_values_array,
597            )));
598        }
599
600        let old_pk_values_array =
601            primary_key
602                .as_any()
603                .downcast_ref::<BinaryArray>()
604                .context(UnexpectedSnafu {
605                    reason: format!(
606                        "Primary-key column is neither binary nor dictionary, got {:?}",
607                        primary_key.data_type()
608                    ),
609                })?;
610        Ok(Arc::new(
611            self.append_values(old_pk_values_array, converter.as_ref())?,
612        ))
613    }
614
615    fn append_values(
616        &self,
617        old_pk_values_array: &BinaryArray,
618        converter: &dyn PrimaryKeyCodec,
619    ) -> Result<BinaryArray> {
620        let mut builder = BinaryBuilder::with_capacity(
621            old_pk_values_array.len(),
622            old_pk_values_array.value_data().len()
623                + converter.estimated_size().unwrap_or_default() * old_pk_values_array.len(),
624        );
625
626        // Binary buffer for the primary key.
627        let mut buffer = Vec::with_capacity(
628            old_pk_values_array.value_data().len() / old_pk_values_array.len().max(1)
629                + converter.estimated_size().unwrap_or_default(),
630        );
631
632        // Iterates the binary array and appends values to the primary key.
633        for value in old_pk_values_array.iter() {
634            let Some(old_pk) = value else {
635                builder.append_null();
636                continue;
637            };
638
639            buffer.clear();
640            buffer.extend_from_slice(old_pk);
641            converter
642                .encode_values(&self.values, &mut buffer)
643                .context(EncodeSnafu)?;
644
645            builder.append_value(&buffer);
646        }
647
648        Ok(builder.finish())
649    }
650}
651
652#[cfg(test)]
653mod tests {
654    use std::sync::Arc;
655
656    use api::v1::{OpType, SemanticType};
657    use datatypes::arrow::array::{
658        ArrayRef, BinaryArray, BinaryDictionaryBuilder, Int64Array, StringDictionaryBuilder,
659        TimestampMillisecondArray, UInt8Array, UInt64Array,
660    };
661    use datatypes::arrow::datatypes::UInt32Type;
662    use datatypes::arrow::record_batch::RecordBatch;
663    use datatypes::prelude::ConcreteDataType;
664    use datatypes::schema::ColumnSchema;
665    use datatypes::value::ValueRef;
666    use mito_codec::row_converter::{
667        DensePrimaryKeyCodec, PrimaryKeyCodecExt, SparsePrimaryKeyCodec,
668    };
669    use store_api::codec::PrimaryKeyEncoding;
670    use store_api::metadata::{ColumnMetadata, RegionMetadataBuilder};
671    use store_api::storage::RegionId;
672
673    use super::*;
674    use crate::read::flat_projection::FlatProjectionMapper;
675    use crate::read::read_columns::ReadColumns;
676    use crate::sst::parquet::flat_format::FlatReadFormat;
677    use crate::sst::{FlatSchemaOptions, to_flat_sst_arrow_schema};
678
679    /// Creates a new [RegionMetadata].
680    fn new_metadata(
681        semantic_types: &[(ColumnId, SemanticType, ConcreteDataType)],
682        primary_key: &[ColumnId],
683    ) -> RegionMetadata {
684        let mut builder = RegionMetadataBuilder::new(RegionId::new(1, 1));
685        for (id, semantic_type, data_type) in semantic_types {
686            let column_schema = match semantic_type {
687                SemanticType::Tag => {
688                    ColumnSchema::new(format!("tag_{id}"), data_type.clone(), true)
689                }
690                SemanticType::Field => {
691                    ColumnSchema::new(format!("field_{id}"), data_type.clone(), true)
692                }
693                SemanticType::Timestamp => ColumnSchema::new("ts", data_type.clone(), false),
694            };
695
696            builder.push_column_metadata(ColumnMetadata {
697                column_schema,
698                semantic_type: *semantic_type,
699                column_id: *id,
700            });
701        }
702        builder.primary_key(primary_key.to_vec());
703        builder.build().unwrap()
704    }
705
706    /// Encode primary key.
707    fn encode_key(keys: &[Option<&str>]) -> Vec<u8> {
708        let fields = (0..keys.len())
709            .map(|_| (0, SortField::new(ConcreteDataType::string_datatype())))
710            .collect();
711        let converter = DensePrimaryKeyCodec::with_fields(fields);
712        let row = keys.iter().map(|str_opt| match str_opt {
713            Some(v) => ValueRef::String(v),
714            None => ValueRef::Null,
715        });
716
717        converter.encode(row).unwrap()
718    }
719
720    /// Encode sparse primary key.
721    fn encode_sparse_key(keys: &[(ColumnId, Option<&str>)]) -> Vec<u8> {
722        let fields = (0..keys.len())
723            .map(|_| (1, SortField::new(ConcreteDataType::string_datatype())))
724            .collect();
725        let converter = SparsePrimaryKeyCodec::with_fields(fields);
726        let row = keys
727            .iter()
728            .map(|(id, str_opt)| match str_opt {
729                Some(v) => (*id, ValueRef::String(v)),
730                None => (*id, ValueRef::Null),
731            })
732            .collect::<Vec<_>>();
733        let mut buffer = vec![];
734        converter.encode_value_refs(&row, &mut buffer).unwrap();
735        buffer
736    }
737
738    /// Creates a primary key array for flat format testing.
739    fn build_flat_test_pk_array(primary_keys: &[&[u8]]) -> ArrayRef {
740        let mut builder = BinaryDictionaryBuilder::<UInt32Type>::new();
741        for &pk in primary_keys {
742            builder.append(pk).unwrap();
743        }
744        Arc::new(builder.finish())
745    }
746
747    #[test]
748    fn test_flat_compat_batch_with_missing_columns() {
749        let actual_metadata = Arc::new(new_metadata(
750            &[
751                (
752                    0,
753                    SemanticType::Timestamp,
754                    ConcreteDataType::timestamp_millisecond_datatype(),
755                ),
756                (1, SemanticType::Tag, ConcreteDataType::string_datatype()),
757                (2, SemanticType::Field, ConcreteDataType::int64_datatype()),
758            ],
759            &[1],
760        ));
761
762        let expected_metadata = Arc::new(new_metadata(
763            &[
764                (
765                    0,
766                    SemanticType::Timestamp,
767                    ConcreteDataType::timestamp_millisecond_datatype(),
768                ),
769                (1, SemanticType::Tag, ConcreteDataType::string_datatype()),
770                (2, SemanticType::Field, ConcreteDataType::int64_datatype()),
771                // Adds a new field.
772                (3, SemanticType::Field, ConcreteDataType::int64_datatype()),
773            ],
774            &[1],
775        ));
776
777        let mapper = FlatProjectionMapper::all(&expected_metadata).unwrap();
778        let read_format = FlatReadFormat::new(
779            actual_metadata.clone(),
780            ReadColumns::from_deduped_column_ids([0, 1, 2, 3]),
781            None,
782            "test",
783            false,
784        )
785        .unwrap();
786
787        let compat_batch = FlatCompatBatch::try_new(&mapper, &read_format, false)
788            .unwrap()
789            .unwrap();
790
791        let mut tag_builder = StringDictionaryBuilder::<UInt32Type>::new();
792        tag_builder.append_value("tag1");
793        tag_builder.append_value("tag1");
794        let tag_dict_array = Arc::new(tag_builder.finish());
795
796        let k1 = encode_key(&[Some("tag1")]);
797        let input_columns: Vec<ArrayRef> = vec![
798            tag_dict_array.clone(),
799            Arc::new(Int64Array::from(vec![100, 200])),
800            Arc::new(TimestampMillisecondArray::from_iter_values([1000, 2000])),
801            build_flat_test_pk_array(&[&k1, &k1]),
802            Arc::new(UInt64Array::from_iter_values([1, 2])),
803            Arc::new(UInt8Array::from_iter_values([
804                OpType::Put as u8,
805                OpType::Put as u8,
806            ])),
807        ];
808        let input_schema =
809            to_flat_sst_arrow_schema(&actual_metadata, &FlatSchemaOptions::default());
810        let input_batch = RecordBatch::try_new(input_schema, input_columns).unwrap();
811
812        let result = compat_batch.compat(input_batch).unwrap();
813
814        let expected_schema =
815            to_flat_sst_arrow_schema(&expected_metadata, &FlatSchemaOptions::default());
816
817        let expected_columns: Vec<ArrayRef> = vec![
818            tag_dict_array.clone(),
819            Arc::new(Int64Array::from(vec![100, 200])),
820            Arc::new(Int64Array::from(vec![None::<i64>, None::<i64>])),
821            Arc::new(TimestampMillisecondArray::from_iter_values([1000, 2000])),
822            build_flat_test_pk_array(&[&k1, &k1]),
823            Arc::new(UInt64Array::from_iter_values([1, 2])),
824            Arc::new(UInt8Array::from_iter_values([
825                OpType::Put as u8,
826                OpType::Put as u8,
827            ])),
828        ];
829        let expected_batch = RecordBatch::try_new(expected_schema, expected_columns).unwrap();
830
831        assert_eq!(expected_batch, result);
832    }
833
834    #[test]
835    fn test_flat_compat_batch_with_read_projection_superset() {
836        let actual_metadata = Arc::new(new_metadata(
837            &[
838                (
839                    0,
840                    SemanticType::Timestamp,
841                    ConcreteDataType::timestamp_millisecond_datatype(),
842                ),
843                (1, SemanticType::Tag, ConcreteDataType::string_datatype()),
844                (2, SemanticType::Field, ConcreteDataType::int64_datatype()),
845            ],
846            &[1],
847        ));
848
849        let expected_metadata = Arc::new(new_metadata(
850            &[
851                (
852                    0,
853                    SemanticType::Timestamp,
854                    ConcreteDataType::timestamp_millisecond_datatype(),
855                ),
856                (1, SemanticType::Tag, ConcreteDataType::string_datatype()),
857                (2, SemanticType::Field, ConcreteDataType::int64_datatype()),
858                // Adds a new field.
859                (3, SemanticType::Field, ConcreteDataType::int64_datatype()),
860            ],
861            &[1],
862        ));
863
864        let mapper = FlatProjectionMapper::new_with_read_columns(
865            &expected_metadata,
866            vec![1, 2],
867            ReadColumns::from_deduped_column_ids([1, 2, 3]),
868            None,
869        )
870        .unwrap();
871        let read_format = FlatReadFormat::new(
872            actual_metadata.clone(),
873            ReadColumns::from_deduped_column_ids([1, 2, 3]),
874            None,
875            "test",
876            false,
877        )
878        .unwrap();
879
880        let compat_batch = FlatCompatBatch::try_new(&mapper, &read_format, false)
881            .unwrap()
882            .unwrap();
883
884        let mut tag_builder = StringDictionaryBuilder::<UInt32Type>::new();
885        tag_builder.append_value("tag1");
886        tag_builder.append_value("tag1");
887        let tag_dict_array = Arc::new(tag_builder.finish());
888
889        let k1 = encode_key(&[Some("tag1")]);
890        let input_columns: Vec<ArrayRef> = vec![
891            tag_dict_array.clone(),
892            Arc::new(Int64Array::from(vec![100, 200])),
893            Arc::new(TimestampMillisecondArray::from_iter_values([1000, 2000])),
894            build_flat_test_pk_array(&[&k1, &k1]),
895            Arc::new(UInt64Array::from_iter_values([1, 2])),
896            Arc::new(UInt8Array::from_iter_values([
897                OpType::Put as u8,
898                OpType::Put as u8,
899            ])),
900        ];
901        let input_schema =
902            to_flat_sst_arrow_schema(&actual_metadata, &FlatSchemaOptions::default());
903        let input_batch = RecordBatch::try_new(input_schema, input_columns).unwrap();
904
905        let result = compat_batch.compat(input_batch).unwrap();
906
907        let expected_schema =
908            to_flat_sst_arrow_schema(&expected_metadata, &FlatSchemaOptions::default());
909        let expected_columns: Vec<ArrayRef> = vec![
910            tag_dict_array.clone(),
911            Arc::new(Int64Array::from(vec![100, 200])),
912            Arc::new(Int64Array::from(vec![None::<i64>, None::<i64>])),
913            Arc::new(TimestampMillisecondArray::from_iter_values([1000, 2000])),
914            build_flat_test_pk_array(&[&k1, &k1]),
915            Arc::new(UInt64Array::from_iter_values([1, 2])),
916            Arc::new(UInt8Array::from_iter_values([
917                OpType::Put as u8,
918                OpType::Put as u8,
919            ])),
920        ];
921        let expected_batch = RecordBatch::try_new(expected_schema, expected_columns).unwrap();
922
923        assert_eq!(expected_batch, result);
924    }
925
926    #[test]
927    fn test_flat_compat_batch_with_different_pk_encoding() {
928        let mut actual_metadata = new_metadata(
929            &[
930                (
931                    0,
932                    SemanticType::Timestamp,
933                    ConcreteDataType::timestamp_millisecond_datatype(),
934                ),
935                (1, SemanticType::Tag, ConcreteDataType::string_datatype()),
936                (2, SemanticType::Field, ConcreteDataType::int64_datatype()),
937            ],
938            &[1],
939        );
940        actual_metadata.primary_key_encoding = PrimaryKeyEncoding::Dense;
941        let actual_metadata = Arc::new(actual_metadata);
942
943        let mut expected_metadata = new_metadata(
944            &[
945                (
946                    0,
947                    SemanticType::Timestamp,
948                    ConcreteDataType::timestamp_millisecond_datatype(),
949                ),
950                (1, SemanticType::Tag, ConcreteDataType::string_datatype()),
951                (2, SemanticType::Field, ConcreteDataType::int64_datatype()),
952                (3, SemanticType::Tag, ConcreteDataType::string_datatype()),
953            ],
954            &[1, 3],
955        );
956        expected_metadata.primary_key_encoding = PrimaryKeyEncoding::Sparse;
957        let expected_metadata = Arc::new(expected_metadata);
958
959        let mapper = FlatProjectionMapper::all(&expected_metadata).unwrap();
960        let read_format = FlatReadFormat::new(
961            actual_metadata.clone(),
962            ReadColumns::from_deduped_column_ids([0, 1, 2, 3]),
963            None,
964            "test",
965            false,
966        )
967        .unwrap();
968
969        let compat_batch = FlatCompatBatch::try_new(&mapper, &read_format, false)
970            .unwrap()
971            .unwrap();
972
973        // Tag array.
974        let mut tag1_builder = StringDictionaryBuilder::<UInt32Type>::new();
975        tag1_builder.append_value("tag1");
976        tag1_builder.append_value("tag1");
977        let tag1_dict_array = Arc::new(tag1_builder.finish());
978
979        let k1 = encode_key(&[Some("tag1")]);
980        let input_columns: Vec<ArrayRef> = vec![
981            tag1_dict_array.clone(),
982            Arc::new(Int64Array::from(vec![100, 200])),
983            Arc::new(TimestampMillisecondArray::from_iter_values([1000, 2000])),
984            build_flat_test_pk_array(&[&k1, &k1]),
985            Arc::new(UInt64Array::from_iter_values([1, 2])),
986            Arc::new(UInt8Array::from_iter_values([
987                OpType::Put as u8,
988                OpType::Put as u8,
989            ])),
990        ];
991        let input_schema =
992            to_flat_sst_arrow_schema(&actual_metadata, &FlatSchemaOptions::default());
993        let input_batch = RecordBatch::try_new(input_schema, input_columns).unwrap();
994
995        let result = compat_batch.compat(input_batch).unwrap();
996
997        let sparse_k1 = encode_sparse_key(&[(1, Some("tag1")), (3, None)]);
998        let mut null_tag_builder = StringDictionaryBuilder::<UInt32Type>::new();
999        null_tag_builder.append_nulls(2);
1000        let null_tag_dict_array = Arc::new(null_tag_builder.finish());
1001        let expected_columns: Vec<ArrayRef> = vec![
1002            tag1_dict_array.clone(),
1003            null_tag_dict_array,
1004            Arc::new(Int64Array::from(vec![100, 200])),
1005            Arc::new(TimestampMillisecondArray::from_iter_values([1000, 2000])),
1006            build_flat_test_pk_array(&[&sparse_k1, &sparse_k1]),
1007            Arc::new(UInt64Array::from_iter_values([1, 2])),
1008            Arc::new(UInt8Array::from_iter_values([
1009                OpType::Put as u8,
1010                OpType::Put as u8,
1011            ])),
1012        ];
1013        let output_schema =
1014            to_flat_sst_arrow_schema(&expected_metadata, &FlatSchemaOptions::default());
1015        let expected_batch = RecordBatch::try_new(output_schema, expected_columns).unwrap();
1016
1017        assert_eq!(expected_batch, result);
1018    }
1019
1020    #[test]
1021    fn test_compat_primary_key_with_different_encoding_only() {
1022        let mut actual_metadata = new_metadata(
1023            &[
1024                (
1025                    0,
1026                    SemanticType::Timestamp,
1027                    ConcreteDataType::timestamp_millisecond_datatype(),
1028                ),
1029                (1, SemanticType::Tag, ConcreteDataType::string_datatype()),
1030                (2, SemanticType::Field, ConcreteDataType::int64_datatype()),
1031            ],
1032            &[1],
1033        );
1034        actual_metadata.primary_key_encoding = PrimaryKeyEncoding::Dense;
1035        let actual_metadata = Arc::new(actual_metadata);
1036
1037        let mut expected_metadata = (*actual_metadata).clone();
1038        expected_metadata.primary_key_encoding = PrimaryKeyEncoding::Sparse;
1039        let expected_metadata = Arc::new(expected_metadata);
1040
1041        let mapper = FlatProjectionMapper::all(&expected_metadata).unwrap();
1042        let read_format = FlatReadFormat::new(
1043            actual_metadata,
1044            ReadColumns::from_deduped_column_ids([0, 1, 2]),
1045            None,
1046            "test",
1047            false,
1048        )
1049        .unwrap();
1050        let compat = FlatCompatBatch::try_new(&mapper, &read_format, false)
1051            .unwrap()
1052            .unwrap();
1053
1054        let dense_key = encode_key(&[Some("tag1")]);
1055        let sparse_key = encode_sparse_key(&[(1, Some("tag1"))]);
1056
1057        let dictionary_key = build_flat_test_pk_array(&[&dense_key, &dense_key]);
1058        let result = compat.compat_primary_key(&dictionary_key).unwrap();
1059        let result = result.as_any().downcast_ref::<PrimaryKeyArray>().unwrap();
1060        let values = result
1061            .values()
1062            .as_any()
1063            .downcast_ref::<BinaryArray>()
1064            .unwrap();
1065        assert_eq!(values.value(result.keys().value(0) as usize), sparse_key);
1066        assert_eq!(values.value(result.keys().value(1) as usize), sparse_key);
1067
1068        let binary_key: ArrayRef = Arc::new(BinaryArray::from(vec![
1069            Some(dense_key.as_slice()),
1070            Some(dense_key.as_slice()),
1071        ]));
1072        let result = compat.compat_primary_key(&binary_key).unwrap();
1073        let result = result.as_any().downcast_ref::<BinaryArray>().unwrap();
1074        assert_eq!(result.value(0), sparse_key);
1075        assert_eq!(result.value(1), sparse_key);
1076    }
1077
1078    #[test]
1079    fn test_flat_compat_batch_compact_sparse() {
1080        let mut actual_metadata = new_metadata(
1081            &[
1082                (
1083                    0,
1084                    SemanticType::Timestamp,
1085                    ConcreteDataType::timestamp_millisecond_datatype(),
1086                ),
1087                (2, SemanticType::Field, ConcreteDataType::int64_datatype()),
1088            ],
1089            &[],
1090        );
1091        actual_metadata.primary_key_encoding = PrimaryKeyEncoding::Sparse;
1092        let actual_metadata = Arc::new(actual_metadata);
1093
1094        let mut expected_metadata = new_metadata(
1095            &[
1096                (
1097                    0,
1098                    SemanticType::Timestamp,
1099                    ConcreteDataType::timestamp_millisecond_datatype(),
1100                ),
1101                (2, SemanticType::Field, ConcreteDataType::int64_datatype()),
1102                (3, SemanticType::Field, ConcreteDataType::int64_datatype()),
1103            ],
1104            &[],
1105        );
1106        expected_metadata.primary_key_encoding = PrimaryKeyEncoding::Sparse;
1107        let expected_metadata = Arc::new(expected_metadata);
1108
1109        let mapper = FlatProjectionMapper::all(&expected_metadata).unwrap();
1110        let read_format = FlatReadFormat::new(
1111            actual_metadata.clone(),
1112            ReadColumns::from_deduped_column_ids([0, 2, 3]),
1113            None,
1114            "test",
1115            true,
1116        )
1117        .unwrap();
1118
1119        let compat_batch = FlatCompatBatch::try_new(&mapper, &read_format, true)
1120            .unwrap()
1121            .unwrap();
1122
1123        let sparse_k1 = encode_sparse_key(&[]);
1124        let input_columns: Vec<ArrayRef> = vec![
1125            Arc::new(Int64Array::from(vec![100, 200])),
1126            Arc::new(TimestampMillisecondArray::from_iter_values([1000, 2000])),
1127            build_flat_test_pk_array(&[&sparse_k1, &sparse_k1]),
1128            Arc::new(UInt64Array::from_iter_values([1, 2])),
1129            Arc::new(UInt8Array::from_iter_values([
1130                OpType::Put as u8,
1131                OpType::Put as u8,
1132            ])),
1133        ];
1134        let input_schema =
1135            to_flat_sst_arrow_schema(&actual_metadata, &FlatSchemaOptions::default());
1136        let input_batch = RecordBatch::try_new(input_schema, input_columns).unwrap();
1137
1138        let result = compat_batch.compat(input_batch).unwrap();
1139
1140        let expected_columns: Vec<ArrayRef> = vec![
1141            Arc::new(Int64Array::from(vec![100, 200])),
1142            Arc::new(Int64Array::from(vec![None::<i64>, None::<i64>])),
1143            Arc::new(TimestampMillisecondArray::from_iter_values([1000, 2000])),
1144            build_flat_test_pk_array(&[&sparse_k1, &sparse_k1]),
1145            Arc::new(UInt64Array::from_iter_values([1, 2])),
1146            Arc::new(UInt8Array::from_iter_values([
1147                OpType::Put as u8,
1148                OpType::Put as u8,
1149            ])),
1150        ];
1151        let output_schema =
1152            to_flat_sst_arrow_schema(&expected_metadata, &FlatSchemaOptions::default());
1153        let expected_batch = RecordBatch::try_new(output_schema, expected_columns).unwrap();
1154
1155        assert_eq!(expected_batch, result);
1156    }
1157}