1use 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
72pub(crate) struct FlatWriteFormat {
74 arrow_schema: SchemaRef,
76 override_sequence: Option<SequenceNumber>,
77}
78
79impl FlatWriteFormat {
80 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 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 #[cfg(test)]
100 pub(crate) fn arrow_schema(&self) -> &SchemaRef {
101 &self.arrow_schema
102 }
103
104 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
120pub(crate) fn sequence_column_index(num_columns: usize) -> usize {
122 num_columns - 2
123}
124
125pub(crate) fn time_index_column_index(num_columns: usize) -> usize {
127 num_columns - 4
128}
129
130pub(crate) fn primary_key_column_index(num_columns: usize) -> usize {
132 num_columns - 3
133}
134
135pub(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
161pub(crate) fn op_type_column_index(num_columns: usize) -> usize {
163 num_columns - 1
164}
165
166pub(crate) fn field_column_start(metadata: &RegionMetadata, num_columns: usize) -> usize {
175 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
181pub struct FlatReadFormat {
188 override_sequence: Option<SequenceNumber>,
190 read_cols: ReadColumns,
192 parquet_adapter: ParquetAdapter,
194 pk_dict_wrap_schema: Option<SchemaRef>,
196}
197
198impl FlatReadFormat {
199 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 if metadata.primary_key_encoding == PrimaryKeyEncoding::Sparse {
218 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 pub(crate) fn set_override_sequence(&mut self, sequence: Option<SequenceNumber>) {
247 self.override_sequence = sequence;
248 }
249
250 pub(crate) fn set_pk_as_binary(&mut self, output_schema: SchemaRef) {
252 self.pk_dict_wrap_schema = Some(output_schema);
253 }
254
255 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 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 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 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 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 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 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 ParquetAdapter::PrimaryKeyToFlat(p) => p
342 .format
343 .field_id_to_projected_index()
344 .get(&column_id)
345 .copied(),
346 }
347 }
348
349 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 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 pub(crate) fn json_target_types(&self) -> &JsonTargetTypes {
367 self.read_cols.json_target_types()
368 }
369
370 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 pub(crate) fn batch_has_raw_pk_columns(&self) -> bool {
385 matches!(&self.parquet_adapter, ParquetAdapter::Flat(_))
386 }
387
388 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 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 let mut batch = match &self.parquet_adapter {
411 ParquetAdapter::Flat(_) => record_batch,
412 ParquetAdapter::PrimaryKeyToFlat(p) => p.convert_batch(record_batch)?,
413 };
414
415 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 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 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 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 let expected_columns = metadata.column_metadatas.len() + INTERNAL_COLUMN_NUM;
483
484 if expected_columns == num_columns {
485 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 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
519enum ParquetAdapter {
521 Flat(ParquetFlat),
522 PrimaryKeyToFlat(ParquetPrimaryKeyToFlat),
523}
524
525struct ParquetPrimaryKeyToFlat {
527 format: PrimaryKeyReadFormat,
529 convert_format: Option<FlatConvertFormat>,
531 format_projection: FormatProjection,
533}
534
535impl ParquetPrimaryKeyToFlat {
536 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 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 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
593struct ParquetFlat {
595 metadata: RegionMetadataRef,
597 arrow_schema: SchemaRef,
599 format_projection: FormatProjection,
601 column_id_to_sst_index: HashMap<ColumnId, usize>,
604}
605
606impl ParquetFlat {
607 fn new(
609 metadata: RegionMetadataRef,
610 read_cols: ReadColumns,
611 arrow_schema: SchemaRef,
612 ) -> ParquetFlat {
613 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 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 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 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 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 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
683pub(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 for pk_id in &metadata.primary_key {
691 id_to_index.insert(*pk_id, column_index);
692 column_index += 1;
693 }
694 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 id_to_index.insert(metadata.time_index_column().column_id, column_index);
703
704 id_to_index
705}
706
707pub 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 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 ¤t_key in pk_indices.iter().take(keys.len()) {
745 if let Some(prev) = prev_key
747 && prev == current_key
748 {
749 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
785enum DecodedKeysInner {
787 Dense {
789 codec: DensePrimaryKeyCodec,
790 values: BinaryArray,
791 distinct_keys: Vec<u32>,
792 offsets: Vec<Vec<usize>>,
794 },
795 Sparse {
798 values: BinaryArray,
800 distinct_keys: Vec<u32>,
802 },
803}
804
805pub struct DecodedPrimaryKeys {
807 inner: DecodedKeysInner,
808 keys_array: UInt32Array,
810 pk_offsets: SparseOffsetsCache,
812 value_buf: Vec<u8>,
814}
815
816fn 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
828fn 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 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 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 if column_type.is_string() {
911 let dict_array = DictionaryArray::new(keys_array.clone(), values_array);
914 Ok(Arc::new(dict_array))
915 } else {
916 let taken_array = take(&values_array, keys_array, None).context(ComputeArrowSnafu)?;
918 Ok(taken_array)
919 }
920 }
921
922 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 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 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 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
1026pub(crate) struct FlatConvertFormat {
1029 metadata: RegionMetadataRef,
1031 codec: Arc<dyn PrimaryKeyCodec>,
1033 projected_primary_keys: Vec<(ColumnId, usize, usize)>,
1035}
1036
1037impl FlatConvertFormat {
1038 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 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 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 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 let mut decoded_columns = Vec::new();
1085 if self.codec.encoding() == PrimaryKeyEncoding::Sparse {
1086 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 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 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 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 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 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 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 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 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 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 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 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 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 #[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 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 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 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 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 #[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 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 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 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 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 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}