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