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 snafu::{OptionExt, ResultExt, ensure};
50use store_api::metadata::{ColumnMetadata, RegionMetadataRef};
51use store_api::storage::{ColumnId, NestedPath, SequenceNumber};
52
53use crate::error::{
54 ConvertVectorSnafu, DecodeSnafu, InvalidRecordBatchSnafu, NewRecordBatchSnafu, Result,
55};
56use crate::read::read_columns::ReadColumns;
57use crate::read::{Batch, BatchBuilder, BatchColumn};
58use crate::sst::file::{FileMeta, FileTimeRange};
59use crate::sst::parquet::read_columns::{ParquetReadColumn, ParquetReadColumns};
60use crate::sst::to_sst_arrow_schema;
61
62pub(crate) type PrimaryKeyArray = DictionaryArray<UInt32Type>;
64pub(crate) type PrimaryKeyArrayBuilder = BinaryDictionaryBuilder<UInt32Type>;
66
67pub(crate) const FIXED_POS_COLUMN_NUM: usize = 4;
71pub(crate) const INTERNAL_COLUMN_NUM: usize = 3;
73
74pub(crate) struct PrimaryKeyWriteFormat {
76 arrow_schema: SchemaRef,
78 override_sequence: Option<SequenceNumber>,
79}
80
81impl PrimaryKeyWriteFormat {
82 pub(crate) fn new(metadata: RegionMetadataRef) -> PrimaryKeyWriteFormat {
84 let arrow_schema = to_sst_arrow_schema(&metadata);
85 PrimaryKeyWriteFormat {
86 arrow_schema,
87 override_sequence: None,
88 }
89 }
90
91 pub(crate) fn with_override_sequence(
93 mut self,
94 override_sequence: Option<SequenceNumber>,
95 ) -> Self {
96 self.override_sequence = override_sequence;
97 self
98 }
99
100 #[cfg(test)]
102 pub(crate) fn arrow_schema(&self) -> &SchemaRef {
103 &self.arrow_schema
104 }
105
106 pub(crate) fn convert_flat_batch(
112 &self,
113 batch: &RecordBatch,
114 num_fields: usize,
115 ) -> Result<RecordBatch> {
116 let num_tag_columns = batch.num_columns() - num_fields - FIXED_POS_COLUMN_NUM;
117 let mut columns: Vec<ArrayRef> = batch.columns()[num_tag_columns..].to_vec();
118
119 if let Some(override_sequence) = self.override_sequence {
120 let num_cols = columns.len();
121 columns[num_cols - 2] =
123 Arc::new(UInt64Array::from(vec![override_sequence; batch.num_rows()]));
124 }
125
126 RecordBatch::try_new(self.arrow_schema.clone(), columns).context(NewRecordBatchSnafu)
127 }
128}
129
130pub(crate) fn column_values(
134 row_groups: &[impl Borrow<RowGroupMetaData>],
135 column: &ColumnMetadata,
136 column_index: usize,
137 is_min: bool,
138) -> Option<ArrayRef> {
139 column_values_by_type(
140 row_groups,
141 &column.column_schema.data_type.as_arrow_type(),
142 column_index,
143 is_min,
144 )
145}
146
147pub(crate) fn column_values_by_type(
149 row_groups: &[impl Borrow<RowGroupMetaData>],
150 data_type: &ArrowDataType,
151 column_index: usize,
152 is_min: bool,
153) -> Option<ArrayRef> {
154 let null_scalar: ScalarValue = data_type.try_into().ok()?;
155 let scalar_values = row_groups
156 .iter()
157 .map(|meta| {
158 let stats = meta.borrow().column(column_index).statistics()?;
159 match stats {
160 Statistics::Boolean(s) => Some(ScalarValue::Boolean(Some(if is_min {
161 *s.min_opt()?
162 } else {
163 *s.max_opt()?
164 }))),
165 Statistics::Int32(s) => {
166 let value = if is_min { *s.min_opt()? } else { *s.max_opt()? };
167 if data_type == &ArrowDataType::UInt32 {
168 Some(ScalarValue::UInt32(Some(value as u32)))
169 } else {
170 Some(ScalarValue::Int32(Some(value)))
171 }
172 }
173 Statistics::Int64(s) => {
174 let value = if is_min { *s.min_opt()? } else { *s.max_opt()? };
175 if data_type == &ArrowDataType::UInt64 {
176 Some(ScalarValue::UInt64(Some(value as u64)))
177 } else {
178 Some(ScalarValue::Int64(Some(value)))
179 }
180 }
181 Statistics::Int96(_) => None,
182 Statistics::Float(s) => Some(ScalarValue::Float32(Some(if is_min {
183 *s.min_opt()?
184 } else {
185 *s.max_opt()?
186 }))),
187 Statistics::Double(s) => Some(ScalarValue::Float64(Some(if is_min {
188 *s.min_opt()?
189 } else {
190 *s.max_opt()?
191 }))),
192 Statistics::ByteArray(s) => {
193 let bytes = if is_min {
194 s.min_bytes_opt()?
195 } else {
196 s.max_bytes_opt()?
197 };
198 let s = String::from_utf8(bytes.to_vec()).ok();
199 Some(ScalarValue::Utf8(s))
200 }
201 Statistics::FixedLenByteArray(_) => None,
202 }
203 })
204 .map(|maybe_scalar| maybe_scalar.unwrap_or_else(|| null_scalar.clone()))
205 .collect::<Vec<ScalarValue>>();
206 debug_assert_eq!(scalar_values.len(), row_groups.len());
207 ScalarValue::iter_to_array(scalar_values).ok()
208}
209
210pub(crate) fn column_null_counts(
213 row_groups: &[impl Borrow<RowGroupMetaData>],
214 column_index: usize,
215) -> Option<ArrayRef> {
216 let values = row_groups.iter().map(|meta| {
217 let col = meta.borrow().column(column_index);
218 let stat = col.statistics()?;
219 stat.null_count_opt()
220 });
221 Some(Arc::new(UInt64Array::from_iter(values)))
222}
223
224pub struct PrimaryKeyReadFormat {
226 metadata: RegionMetadataRef,
228 arrow_schema: SchemaRef,
230 field_id_to_index: HashMap<ColumnId, usize>,
233 parquet_read_cols: ParquetReadColumns,
235 field_id_to_projected_index: HashMap<ColumnId, usize>,
238 primary_key_codec: Option<Arc<dyn PrimaryKeyCodec>>,
240}
241
242impl PrimaryKeyReadFormat {
243 pub fn new(metadata: RegionMetadataRef, read_cols: ReadColumns) -> PrimaryKeyReadFormat {
245 let field_id_to_index: HashMap<_, _> = metadata
246 .field_columns()
247 .enumerate()
248 .map(|(index, column)| (column.column_id, index))
249 .collect();
250 let arrow_schema = to_sst_arrow_schema(&metadata);
251
252 let format_projection = FormatProjection::compute_format_projection(
253 &metadata,
254 &field_id_to_index,
255 arrow_schema.fields.len(),
256 read_cols,
257 );
258
259 PrimaryKeyReadFormat {
260 metadata,
261 arrow_schema,
262 field_id_to_index,
263 parquet_read_cols: format_projection.parquet_read_cols,
264 field_id_to_projected_index: format_projection.column_id_to_projected_index,
265 primary_key_codec: None,
266 }
267 }
268
269 pub(crate) fn arrow_schema(&self) -> &SchemaRef {
274 &self.arrow_schema
275 }
276
277 pub(crate) fn metadata(&self) -> &RegionMetadataRef {
279 &self.metadata
280 }
281
282 pub(crate) fn parquet_read_columns(&self) -> &ParquetReadColumns {
283 &self.parquet_read_cols
284 }
285
286 pub(crate) fn field_id_to_projected_index(&self) -> &HashMap<ColumnId, usize> {
288 &self.field_id_to_projected_index
289 }
290
291 pub fn convert_record_batch(
296 &self,
297 record_batch: &RecordBatch,
298 override_sequence_array: Option<&ArrayRef>,
299 batches: &mut VecDeque<Batch>,
300 ) -> Result<()> {
301 debug_assert!(batches.is_empty());
302
303 ensure!(
305 record_batch.num_columns() >= FIXED_POS_COLUMN_NUM,
306 InvalidRecordBatchSnafu {
307 reason: format!(
308 "record batch only has {} columns",
309 record_batch.num_columns()
310 ),
311 }
312 );
313
314 let mut fixed_pos_columns = record_batch
315 .columns()
316 .iter()
317 .rev()
318 .take(FIXED_POS_COLUMN_NUM);
319 let op_type_array = fixed_pos_columns.next().unwrap();
321 let mut sequence_array = fixed_pos_columns.next().unwrap().clone();
322 let pk_array = fixed_pos_columns.next().unwrap();
323 let ts_array = fixed_pos_columns.next().unwrap();
324 let field_batch_columns = self.get_field_batch_columns(record_batch)?;
325
326 if let Some(override_array) = override_sequence_array {
328 assert!(override_array.len() >= sequence_array.len());
329 sequence_array = if override_array.len() > sequence_array.len() {
332 override_array.slice(0, sequence_array.len())
333 } else {
334 override_array.clone()
335 };
336 }
337
338 let pk_dict_array = pk_array
340 .as_any()
341 .downcast_ref::<PrimaryKeyArray>()
342 .with_context(|| InvalidRecordBatchSnafu {
343 reason: format!("primary key array should not be {:?}", pk_array.data_type()),
344 })?;
345 let offsets = primary_key_offsets(pk_dict_array)?;
346 if offsets.is_empty() {
347 return Ok(());
348 }
349
350 let keys = pk_dict_array.keys();
352 let pk_values = pk_dict_array
353 .values()
354 .as_any()
355 .downcast_ref::<BinaryArray>()
356 .with_context(|| InvalidRecordBatchSnafu {
357 reason: format!(
358 "values of primary key array should not be {:?}",
359 pk_dict_array.values().data_type()
360 ),
361 })?;
362 for (i, start) in offsets[..offsets.len() - 1].iter().enumerate() {
363 let end = offsets[i + 1];
364 let rows_in_batch = end - start;
365 let dict_key = keys.value(*start);
366 let primary_key = pk_values.value(dict_key as usize).to_vec();
367
368 let mut builder = BatchBuilder::new(primary_key);
369 builder
370 .timestamps_array(ts_array.slice(*start, rows_in_batch))?
371 .sequences_array(sequence_array.slice(*start, rows_in_batch))?
372 .op_types_array(op_type_array.slice(*start, rows_in_batch))?;
373 for batch_column in &field_batch_columns {
375 builder.push_field(BatchColumn {
376 column_id: batch_column.column_id,
377 data: batch_column.data.slice(*start, rows_in_batch),
378 });
379 }
380
381 let mut batch = builder.build()?;
382 if let Some(codec) = &self.primary_key_codec {
383 let pk_values: CompositeValues =
384 codec.decode(batch.primary_key()).context(DecodeSnafu)?;
385 batch.set_pk_values(pk_values);
386 }
387 batches.push_back(batch);
388 }
389
390 Ok(())
391 }
392
393 pub fn min_values(
395 &self,
396 row_groups: &[impl Borrow<RowGroupMetaData>],
397 column_id: ColumnId,
398 ) -> StatValues {
399 let Some(column) = self.metadata.column_by_id(column_id) else {
400 return StatValues::NoColumn;
402 };
403 match column.semantic_type {
404 SemanticType::Tag => self.tag_values(row_groups, column, true),
405 SemanticType::Field => {
406 let index = self.field_id_to_index.get(&column_id).unwrap();
408 let stats = column_values(row_groups, column, *index, true);
409 StatValues::from_stats_opt(stats)
410 }
411 SemanticType::Timestamp => {
412 let index = self.time_index_position();
413 let stats = column_values(row_groups, column, index, true);
414 StatValues::from_stats_opt(stats)
415 }
416 }
417 }
418
419 pub fn max_values(
421 &self,
422 row_groups: &[impl Borrow<RowGroupMetaData>],
423 column_id: ColumnId,
424 ) -> StatValues {
425 let Some(column) = self.metadata.column_by_id(column_id) else {
426 return StatValues::NoColumn;
428 };
429 match column.semantic_type {
430 SemanticType::Tag => self.tag_values(row_groups, column, false),
431 SemanticType::Field => {
432 let index = self.field_id_to_index.get(&column_id).unwrap();
434 let stats = column_values(row_groups, column, *index, false);
435 StatValues::from_stats_opt(stats)
436 }
437 SemanticType::Timestamp => {
438 let index = self.time_index_position();
439 let stats = column_values(row_groups, column, index, false);
440 StatValues::from_stats_opt(stats)
441 }
442 }
443 }
444
445 pub fn null_counts(
447 &self,
448 row_groups: &[impl Borrow<RowGroupMetaData>],
449 column_id: ColumnId,
450 ) -> StatValues {
451 let Some(column) = self.metadata.column_by_id(column_id) else {
452 return StatValues::NoColumn;
454 };
455 match column.semantic_type {
456 SemanticType::Tag => StatValues::NoStats,
457 SemanticType::Field => {
458 let index = self.field_id_to_index.get(&column_id).unwrap();
460 let stats = column_null_counts(row_groups, *index);
461 StatValues::from_stats_opt(stats)
462 }
463 SemanticType::Timestamp => {
464 let index = self.time_index_position();
465 let stats = column_null_counts(row_groups, index);
466 StatValues::from_stats_opt(stats)
467 }
468 }
469 }
470
471 fn get_field_batch_columns(&self, record_batch: &RecordBatch) -> Result<Vec<BatchColumn>> {
473 record_batch
474 .columns()
475 .iter()
476 .zip(record_batch.schema().fields())
477 .take(record_batch.num_columns() - FIXED_POS_COLUMN_NUM) .map(|(array, field)| {
479 let vector = Helper::try_into_vector(array.clone()).context(ConvertVectorSnafu)?;
480 let column = self
481 .metadata
482 .column_by_name(field.name())
483 .with_context(|| InvalidRecordBatchSnafu {
484 reason: format!("column {} not found in metadata", field.name()),
485 })?;
486
487 Ok(BatchColumn {
488 column_id: column.column_id,
489 data: vector,
490 })
491 })
492 .collect()
493 }
494
495 fn tag_values(
497 &self,
498 row_groups: &[impl Borrow<RowGroupMetaData>],
499 column: &ColumnMetadata,
500 is_min: bool,
501 ) -> StatValues {
502 let is_first_tag = self
503 .metadata
504 .primary_key
505 .first()
506 .map(|id| *id == column.column_id)
507 .unwrap_or(false);
508 if !is_first_tag {
509 return StatValues::NoStats;
511 }
512
513 StatValues::from_stats_opt(self.first_tag_values(row_groups, column, is_min))
514 }
515
516 fn first_tag_values(
519 &self,
520 row_groups: &[impl Borrow<RowGroupMetaData>],
521 column: &ColumnMetadata,
522 is_min: bool,
523 ) -> Option<ArrayRef> {
524 debug_assert!(
525 self.metadata
526 .primary_key
527 .first()
528 .map(|id| *id == column.column_id)
529 .unwrap_or(false)
530 );
531
532 let primary_key_encoding = self.metadata.primary_key_encoding;
533 let converter = build_primary_key_codec_with_fields(
534 primary_key_encoding,
535 [(
536 column.column_id,
537 SortField::new(column.column_schema.data_type.clone()),
538 )]
539 .into_iter(),
540 );
541
542 let values = row_groups.iter().map(|meta| {
543 let stats = meta
544 .borrow()
545 .column(self.primary_key_position())
546 .statistics()?;
547 match stats {
548 Statistics::Boolean(_) => None,
549 Statistics::Int32(_) => None,
550 Statistics::Int64(_) => None,
551 Statistics::Int96(_) => None,
552 Statistics::Float(_) => None,
553 Statistics::Double(_) => None,
554 Statistics::ByteArray(s) => {
555 let bytes = if is_min {
556 s.min_bytes_opt()?
557 } else {
558 s.max_bytes_opt()?
559 };
560 converter.decode_leftmost(bytes).ok()?
561 }
562 Statistics::FixedLenByteArray(_) => None,
563 }
564 });
565 let mut builder = column
566 .column_schema
567 .data_type
568 .create_mutable_vector(row_groups.len());
569 for value_opt in values {
570 match value_opt {
571 Some(v) => builder.push_value_ref(&v.as_value_ref()),
573 None => builder.push_null(),
574 }
575 }
576 let vector = builder.to_vector();
577
578 Some(vector.to_arrow_array())
579 }
580
581 fn primary_key_position(&self) -> usize {
583 self.arrow_schema.fields.len() - 3
584 }
585
586 fn time_index_position(&self) -> usize {
588 self.arrow_schema.fields.len() - FIXED_POS_COLUMN_NUM
589 }
590
591 pub fn field_index_by_id(&self, column_id: ColumnId) -> Option<usize> {
593 self.field_id_to_projected_index.get(&column_id).copied()
594 }
595}
596
597pub(crate) struct FormatProjection {
599 pub(crate) parquet_read_cols: ParquetReadColumns,
601 pub(crate) column_id_to_projected_index: HashMap<ColumnId, usize>,
606}
607
608impl FormatProjection {
609 pub(crate) fn compute_format_projection(
613 metadata: &RegionMetadataRef,
614 id_to_index: &HashMap<ColumnId, usize>,
615 sst_column_num: usize,
616 cols: ReadColumns,
617 ) -> Self {
618 let mut projected_columns: Vec<_> = cols
619 .col_ids
620 .iter()
621 .copied()
622 .filter_map(|col_id| {
623 id_to_index.get(&col_id).copied().map(|index_of_sst| {
624 let nested_paths = json_target_nested_paths(metadata, &cols, col_id);
625 (col_id, index_of_sst, nested_paths)
626 })
627 })
628 .collect();
629 projected_columns.sort_unstable_by_key(|(_, index, _)| *index);
632
633 let mut parquet_read_cols: Vec<ParquetReadColumn> =
634 Vec::with_capacity(projected_columns.len() + FIXED_POS_COLUMN_NUM);
635 let mut column_id_to_projected_index = HashMap::with_capacity(projected_columns.len());
637
638 for (col_id, index_of_sst, nested_paths) in projected_columns {
639 Self::merge_or_push_parquet_column(&mut parquet_read_cols, index_of_sst, nested_paths);
640
641 column_id_to_projected_index
642 .entry(col_id)
643 .or_insert_with(|| parquet_read_cols.len() - 1);
644 }
645
646 Self::append_time_index_if_needed(&mut parquet_read_cols, sst_column_num);
649 Self::append_fixed_internal_columns(&mut parquet_read_cols, sst_column_num);
650
651 Self {
652 parquet_read_cols: ParquetReadColumns::from_deduped(parquet_read_cols),
653 column_id_to_projected_index,
654 }
655 }
656
657 fn merge_or_push_parquet_column(
658 parquet_read_cols: &mut Vec<ParquetReadColumn>,
659 index_of_sst: usize,
660 nested_paths: Vec<Vec<String>>,
661 ) {
662 if let Some(last_col) = parquet_read_cols.last_mut()
665 && last_col.root_index() == index_of_sst
666 {
667 last_col.merge_nested_paths(nested_paths);
668 return;
669 }
670
671 let parquet_col = ParquetReadColumn::new(index_of_sst).with_nested_paths(nested_paths);
672 parquet_read_cols.push(parquet_col);
673 }
674
675 fn append_time_index_if_needed(
676 parquet_read_cols: &mut Vec<ParquetReadColumn>,
677 sst_column_num: usize,
678 ) {
679 let time_index = sst_column_num - FIXED_POS_COLUMN_NUM;
680 let needs_time_index = parquet_read_cols
684 .last()
685 .map(|col| col.root_index() != time_index)
686 .unwrap_or(true);
687 if needs_time_index {
688 parquet_read_cols.push(ParquetReadColumn::new(time_index));
689 }
690 }
691
692 fn append_fixed_internal_columns(
695 parquet_read_cols: &mut Vec<ParquetReadColumn>,
696 sst_column_num: usize,
697 ) {
698 for index in sst_column_num - INTERNAL_COLUMN_NUM..sst_column_num {
699 parquet_read_cols.push(ParquetReadColumn::new(index));
700 }
701 }
702}
703
704fn json_target_nested_paths(
705 metadata: &RegionMetadataRef,
706 read_columns: &ReadColumns,
707 column_id: ColumnId,
708) -> Vec<NestedPath> {
709 let Some(target_type) = read_columns.json_target_type(column_id) else {
710 return Vec::new();
711 };
712 let Some(column) = metadata.column_by_id(column_id) else {
713 return Vec::new();
714 };
715
716 json_nested_paths(&column.column_schema.name, target_type)
717}
718
719fn json_nested_paths(column_name: &str, json_type: &JsonNativeType) -> Vec<NestedPath> {
720 let mut paths = Vec::new();
721 let mut current = vec![column_name.to_string()];
722 collect_json_nested_paths(json_type, &mut current, &mut paths);
723 paths
724}
725
726fn collect_json_nested_paths(
727 json_type: &JsonNativeType,
728 current: &mut NestedPath,
729 paths: &mut Vec<NestedPath>,
730) {
731 match json_type {
732 JsonNativeType::Object(fields) if !fields.is_empty() => {
733 for (field, child) in fields {
734 current.push(field.clone());
735 collect_json_nested_paths(child, current, paths);
736 current.pop();
737 }
738 }
739 _ => paths.push(current.clone()),
740 }
741}
742
743pub enum StatValues {
748 Values(ArrayRef),
750 NoColumn,
752 NoStats,
754}
755
756impl StatValues {
757 pub fn from_stats_opt(stats: Option<ArrayRef>) -> Self {
759 match stats {
760 Some(stats) => StatValues::Values(stats),
761 None => StatValues::NoStats,
762 }
763 }
764}
765
766#[cfg(test)]
767impl PrimaryKeyReadFormat {
768 pub fn new_with_all_columns(metadata: RegionMetadataRef) -> PrimaryKeyReadFormat {
770 Self::new(
771 Arc::clone(&metadata),
772 ReadColumns::new(metadata.column_metadatas.iter().map(|c| c.column_id)),
773 )
774 }
775}
776
777pub(crate) fn primary_key_offsets(pk_dict_array: &PrimaryKeyArray) -> Result<Vec<usize>> {
779 if pk_dict_array.is_empty() {
780 return Ok(Vec::new());
781 }
782
783 let mut offsets = vec![0];
785 let keys = pk_dict_array.keys();
786 let pk_indices = keys.values();
788 for (i, key) in pk_indices.iter().take(keys.len() - 1).enumerate() {
789 if *key != pk_indices[i + 1] {
791 offsets.push(i + 1);
793 }
794 }
795 offsets.push(keys.len());
796
797 Ok(offsets)
798}
799
800pub(crate) fn parquet_row_group_time_range(
803 file_meta: &FileMeta,
804 parquet_meta: &ParquetMetaData,
805 row_group_idx: usize,
806) -> Option<FileTimeRange> {
807 let row_group_meta = parquet_meta.row_group(row_group_idx);
808 let num_columns = parquet_meta.file_metadata().schema_descr().num_columns();
809 assert!(
810 num_columns >= FIXED_POS_COLUMN_NUM,
811 "file only has {} columns",
812 num_columns
813 );
814 let time_index_pos = num_columns - FIXED_POS_COLUMN_NUM;
815
816 let stats = row_group_meta.column(time_index_pos).statistics()?;
817 let (min, max) = match stats {
819 Statistics::Int64(value_stats) => (*value_stats.min_opt()?, *value_stats.max_opt()?),
820 Statistics::Int32(_)
821 | Statistics::Boolean(_)
822 | Statistics::Int96(_)
823 | Statistics::Float(_)
824 | Statistics::Double(_)
825 | Statistics::ByteArray(_)
826 | Statistics::FixedLenByteArray(_) => {
827 common_telemetry::warn!(
828 "Invalid statistics {:?} for time index in parquet in {}",
829 stats,
830 file_meta.file_id
831 );
832 return None;
833 }
834 };
835
836 debug_assert!(min >= file_meta.time_range.0.value() && min <= file_meta.time_range.1.value());
837 debug_assert!(max >= file_meta.time_range.0.value() && max <= file_meta.time_range.1.value());
838 let unit = file_meta.time_range.0.unit();
839
840 Some((Timestamp::new(min, unit), Timestamp::new(max, unit)))
841}
842
843pub(crate) fn need_override_sequence(parquet_meta: &ParquetMetaData) -> bool {
846 let num_columns = parquet_meta.file_metadata().schema_descr().num_columns();
847 if num_columns < FIXED_POS_COLUMN_NUM {
848 return false;
849 }
850
851 let sequence_pos = num_columns - 2;
853
854 for row_group in parquet_meta.row_groups() {
856 if let Some(Statistics::Int64(value_stats)) = row_group.column(sequence_pos).statistics() {
857 if let (Some(min_val), Some(max_val)) = (value_stats.min_opt(), value_stats.max_opt()) {
858 if *min_val != 0 || *max_val != 0 {
860 return false;
861 }
862 } else {
863 return false;
865 }
866 } else {
867 return false;
869 }
870 }
871
872 !parquet_meta.row_groups().is_empty()
874}
875
876#[cfg(test)]
877mod tests {
878 use std::collections::BTreeMap;
879 use std::sync::Arc;
880
881 use api::v1::OpType;
882 use datatypes::arrow::array::{
883 Int64Array, StringArray, TimestampMillisecondArray, UInt8Array, UInt32Array, UInt64Array,
884 };
885 use datatypes::arrow::datatypes::{DataType as ArrowDataType, Field, Schema, TimeUnit};
886 use datatypes::prelude::ConcreteDataType;
887 use datatypes::schema::ColumnSchema;
888 use datatypes::types::json_type::{JsonNativeType, JsonObjectType};
889 use datatypes::value::ValueRef;
890 use datatypes::vectors::{Int64Vector, TimestampMillisecondVector, UInt8Vector, UInt64Vector};
891 use mito_codec::row_converter::{
892 DensePrimaryKeyCodec, PrimaryKeyCodec, PrimaryKeyCodecExt, SparsePrimaryKeyCodec,
893 };
894 use store_api::codec::PrimaryKeyEncoding;
895 use store_api::metadata::{ColumnMetadata, RegionMetadataBuilder};
896 use store_api::storage::RegionId;
897 use store_api::storage::consts::ReservedColumnId;
898
899 use super::*;
900 use crate::error::InvalidMetadataSnafu;
901 use crate::sst::parquet::flat_format::{
902 FlatReadFormat, FlatWriteFormat, sequence_column_index, sst_column_id_indices,
903 };
904 use crate::sst::{
905 FlatSchemaOptions, OP_TYPE_PARQUET_FIELD_ID, PRIMARY_KEY_PARQUET_FIELD_ID,
906 SEQUENCE_PARQUET_FIELD_ID, to_flat_sst_arrow_schema, with_field_id,
907 };
908
909 const TEST_SEQUENCE: u64 = 1;
910 const TEST_OP_TYPE: u8 = OpType::Put as u8;
911
912 fn build_test_region_metadata() -> RegionMetadataRef {
913 let mut builder = RegionMetadataBuilder::new(RegionId::new(1, 1));
914 builder
915 .push_column_metadata(ColumnMetadata {
916 column_schema: ColumnSchema::new("tag0", ConcreteDataType::int64_datatype(), true),
917 semantic_type: SemanticType::Tag,
918 column_id: 1,
919 })
920 .push_column_metadata(ColumnMetadata {
921 column_schema: ColumnSchema::new(
922 "field1",
923 ConcreteDataType::int64_datatype(),
924 true,
925 ),
926 semantic_type: SemanticType::Field,
927 column_id: 4, })
929 .push_column_metadata(ColumnMetadata {
930 column_schema: ColumnSchema::new("tag1", ConcreteDataType::int64_datatype(), true),
931 semantic_type: SemanticType::Tag,
932 column_id: 3,
933 })
934 .push_column_metadata(ColumnMetadata {
935 column_schema: ColumnSchema::new(
936 "field0",
937 ConcreteDataType::int64_datatype(),
938 true,
939 ),
940 semantic_type: SemanticType::Field,
941 column_id: 2,
942 })
943 .push_column_metadata(ColumnMetadata {
944 column_schema: ColumnSchema::new(
945 "ts",
946 ConcreteDataType::timestamp_millisecond_datatype(),
947 false,
948 ),
949 semantic_type: SemanticType::Timestamp,
950 column_id: 5,
951 })
952 .primary_key(vec![1, 3]);
953 Arc::new(builder.build().unwrap())
954 }
955
956 fn build_test_arrow_schema() -> SchemaRef {
957 let fields = vec![
958 make_field("field1", ArrowDataType::Int64, true, Some(4)),
959 make_field("field0", ArrowDataType::Int64, true, Some(2)),
960 make_field(
961 "ts",
962 ArrowDataType::Timestamp(TimeUnit::Millisecond, None),
963 false,
964 Some(5),
965 ),
966 make_field(
967 "__primary_key",
968 ArrowDataType::Dictionary(
969 Box::new(ArrowDataType::UInt32),
970 Box::new(ArrowDataType::Binary),
971 ),
972 false,
973 Some(PRIMARY_KEY_PARQUET_FIELD_ID),
974 ),
975 make_field(
976 "__sequence",
977 ArrowDataType::UInt64,
978 false,
979 Some(SEQUENCE_PARQUET_FIELD_ID),
980 ),
981 make_field(
982 "__op_type",
983 ArrowDataType::UInt8,
984 false,
985 Some(OP_TYPE_PARQUET_FIELD_ID),
986 ),
987 ];
988 Arc::new(Schema::new(fields))
989 }
990
991 fn new_batch(primary_key: &[u8], start_ts: i64, start_field: i64, num_rows: usize) -> Batch {
992 new_batch_with_sequence(primary_key, start_ts, start_field, num_rows, TEST_SEQUENCE)
993 }
994
995 fn new_batch_with_sequence(
996 primary_key: &[u8],
997 start_ts: i64,
998 start_field: i64,
999 num_rows: usize,
1000 sequence: u64,
1001 ) -> Batch {
1002 let ts_values = (0..num_rows).map(|i| start_ts + i as i64);
1003 let timestamps = Arc::new(TimestampMillisecondVector::from_values(ts_values));
1004 let sequences = Arc::new(UInt64Vector::from_vec(vec![sequence; num_rows]));
1005 let op_types = Arc::new(UInt8Vector::from_vec(vec![TEST_OP_TYPE; num_rows]));
1006 let fields = vec![
1007 BatchColumn {
1008 column_id: 4,
1009 data: Arc::new(Int64Vector::from_vec(vec![start_field; num_rows])),
1010 }, BatchColumn {
1012 column_id: 2,
1013 data: Arc::new(Int64Vector::from_vec(vec![start_field + 1; num_rows])),
1014 }, ];
1016
1017 BatchBuilder::with_required_columns(primary_key.to_vec(), timestamps, sequences, op_types)
1018 .with_fields(fields)
1019 .build()
1020 .unwrap()
1021 }
1022
1023 #[test]
1024 fn test_to_sst_arrow_schema() {
1025 let metadata = build_test_region_metadata();
1026 let write_format = PrimaryKeyWriteFormat::new(metadata);
1027 assert_eq!(&build_test_arrow_schema(), write_format.arrow_schema());
1028 }
1029
1030 fn build_test_pk_array(pk_row_nums: &[(Vec<u8>, usize)]) -> Arc<PrimaryKeyArray> {
1031 let values = Arc::new(BinaryArray::from_iter_values(
1032 pk_row_nums.iter().map(|v| &v.0),
1033 ));
1034 let mut keys = vec![];
1035 for (index, num_rows) in pk_row_nums.iter().map(|v| v.1).enumerate() {
1036 keys.extend(std::iter::repeat_n(index as u32, num_rows));
1037 }
1038 let keys = UInt32Array::from(keys);
1039 Arc::new(DictionaryArray::new(keys, values))
1040 }
1041
1042 #[test]
1043 fn test_projection_indices() {
1044 let metadata = build_test_region_metadata();
1045 let read_format = PrimaryKeyReadFormat::new(metadata.clone(), ReadColumns::new([3]));
1047 assert_eq!(
1048 &[2, 3, 4, 5],
1049 read_format.parquet_read_columns().root_indices()
1050 );
1051 let read_format = PrimaryKeyReadFormat::new(metadata.clone(), ReadColumns::new([4]));
1053 assert_eq!(
1054 &[0, 2, 3, 4, 5],
1055 read_format.parquet_read_columns().root_indices()
1056 );
1057 let read_format = PrimaryKeyReadFormat::new(metadata.clone(), ReadColumns::new([5]));
1059 assert_eq!(
1060 &[2, 3, 4, 5],
1061 read_format.parquet_read_columns().root_indices()
1062 );
1063 let read_format = PrimaryKeyReadFormat::new(metadata, ReadColumns::new([2, 1, 5]));
1065 assert_eq!(
1066 &[1, 2, 3, 4, 5],
1067 read_format.parquet_read_columns().root_indices()
1068 );
1069 }
1070
1071 #[test]
1072 fn test_format_projection_preserves_nested_paths() -> Result<()> {
1073 let mut builder = RegionMetadataBuilder::new(RegionId::new(1, 1));
1074 builder
1075 .push_column_metadata(ColumnMetadata {
1076 column_schema: ColumnSchema::new("tag0", ConcreteDataType::string_datatype(), true),
1077 semantic_type: SemanticType::Tag,
1078 column_id: 1,
1079 })
1080 .push_column_metadata(ColumnMetadata {
1081 column_schema: ColumnSchema::new(
1082 "j",
1083 ConcreteDataType::json2(JsonNativeType::Object(JsonObjectType::from([
1084 ("a".to_string(), JsonNativeType::i64()),
1085 ("b".to_string(), JsonNativeType::String),
1086 ]))),
1087 true,
1088 ),
1089 semantic_type: SemanticType::Field,
1090 column_id: 4,
1091 })
1092 .push_column_metadata(ColumnMetadata {
1093 column_schema: ColumnSchema::new(
1094 "ts",
1095 ConcreteDataType::timestamp_millisecond_datatype(),
1096 false,
1097 ),
1098 semantic_type: SemanticType::Timestamp,
1099 column_id: 5,
1100 })
1101 .primary_key(vec![1]);
1102 let metadata = Arc::new(builder.build().context(InvalidMetadataSnafu)?);
1103 let column_id_to_parquet_index = sst_column_id_indices(&metadata);
1104 let projection = FormatProjection::compute_format_projection(
1105 &metadata,
1106 &column_id_to_parquet_index,
1107 metadata.column_metadatas.len() + FIXED_POS_COLUMN_NUM,
1108 ReadColumns::new([4]).with_json_target_types(BTreeMap::from([(
1109 4,
1110 JsonNativeType::Object(JsonObjectType::from([(
1111 "a".to_string(),
1112 JsonNativeType::i64(),
1113 )])),
1114 )])),
1115 );
1116
1117 let columns = projection.parquet_read_cols.columns();
1118 assert_eq!(1, columns[0].root_index());
1119 assert_eq!(
1120 &[vec!["j".to_string(), "a".to_string()]],
1121 columns[0].nested_paths()
1122 );
1123 Ok(())
1124 }
1125
1126 #[test]
1127 fn test_empty_primary_key_offsets() {
1128 let array = build_test_pk_array(&[]);
1129 assert!(primary_key_offsets(&array).unwrap().is_empty());
1130 }
1131
1132 #[test]
1133 fn test_primary_key_offsets_one_series() {
1134 let array = build_test_pk_array(&[(b"one".to_vec(), 1)]);
1135 assert_eq!(vec![0, 1], primary_key_offsets(&array).unwrap());
1136
1137 let array = build_test_pk_array(&[(b"one".to_vec(), 1), (b"two".to_vec(), 1)]);
1138 assert_eq!(vec![0, 1, 2], primary_key_offsets(&array).unwrap());
1139
1140 let array = build_test_pk_array(&[
1141 (b"one".to_vec(), 1),
1142 (b"two".to_vec(), 1),
1143 (b"three".to_vec(), 1),
1144 ]);
1145 assert_eq!(vec![0, 1, 2, 3], primary_key_offsets(&array).unwrap());
1146 }
1147
1148 #[test]
1149 fn test_primary_key_offsets_multi_series() {
1150 let array = build_test_pk_array(&[(b"one".to_vec(), 1), (b"two".to_vec(), 3)]);
1151 assert_eq!(vec![0, 1, 4], primary_key_offsets(&array).unwrap());
1152
1153 let array = build_test_pk_array(&[(b"one".to_vec(), 3), (b"two".to_vec(), 1)]);
1154 assert_eq!(vec![0, 3, 4], primary_key_offsets(&array).unwrap());
1155
1156 let array = build_test_pk_array(&[(b"one".to_vec(), 3), (b"two".to_vec(), 3)]);
1157 assert_eq!(vec![0, 3, 6], primary_key_offsets(&array).unwrap());
1158 }
1159
1160 #[test]
1161 fn test_convert_empty_record_batch() {
1162 let metadata = build_test_region_metadata();
1163 let arrow_schema = build_test_arrow_schema();
1164 let column_ids: Vec<_> = metadata
1165 .column_metadatas
1166 .iter()
1167 .map(|col| col.column_id)
1168 .collect();
1169 let read_format = PrimaryKeyReadFormat::new(metadata, ReadColumns::new(column_ids));
1170 assert_eq!(arrow_schema, *read_format.arrow_schema());
1171
1172 let record_batch = RecordBatch::new_empty(arrow_schema);
1173 let mut batches = VecDeque::new();
1174 read_format
1175 .convert_record_batch(&record_batch, None, &mut batches)
1176 .unwrap();
1177 assert!(batches.is_empty());
1178 }
1179
1180 #[test]
1181 fn test_convert_record_batch() {
1182 let metadata = build_test_region_metadata();
1183 let column_ids: Vec<_> = metadata
1184 .column_metadatas
1185 .iter()
1186 .map(|col| col.column_id)
1187 .collect();
1188 let read_format = PrimaryKeyReadFormat::new(metadata, ReadColumns::new(column_ids));
1189
1190 let columns: Vec<ArrayRef> = vec![
1191 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])), ];
1198 let arrow_schema = build_test_arrow_schema();
1199 let record_batch = RecordBatch::try_new(arrow_schema, columns).unwrap();
1200 let mut batches = VecDeque::new();
1201 read_format
1202 .convert_record_batch(&record_batch, None, &mut batches)
1203 .unwrap();
1204
1205 assert_eq!(
1206 vec![new_batch(b"one", 1, 1, 2), new_batch(b"two", 11, 10, 2)],
1207 batches.into_iter().collect::<Vec<_>>(),
1208 );
1209 }
1210
1211 #[test]
1212 fn test_convert_record_batch_with_override_sequence() {
1213 let metadata = build_test_region_metadata();
1214 let read_format = PrimaryKeyReadFormat::new(
1215 metadata.clone(),
1216 ReadColumns::new(metadata.column_metadatas.iter().map(|c| c.column_id)),
1217 );
1218
1219 let columns: Vec<ArrayRef> = vec![
1220 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])), ];
1227 let arrow_schema = build_test_arrow_schema();
1228 let record_batch = RecordBatch::try_new(arrow_schema, columns).unwrap();
1229
1230 let override_sequence: u64 = 12345;
1232 let override_sequence_array: ArrayRef =
1233 Arc::new(UInt64Array::from_value(override_sequence, 4));
1234
1235 let mut batches = VecDeque::new();
1236 read_format
1237 .convert_record_batch(&record_batch, Some(&override_sequence_array), &mut batches)
1238 .unwrap();
1239
1240 let expected_batch1 = new_batch_with_sequence(b"one", 1, 1, 2, override_sequence);
1242 let expected_batch2 = new_batch_with_sequence(b"two", 11, 10, 2, override_sequence);
1243
1244 assert_eq!(
1245 vec![expected_batch1, expected_batch2],
1246 batches.into_iter().collect::<Vec<_>>(),
1247 );
1248 }
1249
1250 fn make_field(name: &str, dt: ArrowDataType, nullable: bool, field_id: Option<u32>) -> Field {
1251 let mut field = Field::new(name, dt, nullable);
1252 if let Some(id) = field_id {
1253 field = with_field_id(field, id);
1254 }
1255 field
1256 }
1257
1258 fn build_test_flat_sst_schema() -> SchemaRef {
1259 let fields = vec![
1260 Field::new("tag0", ArrowDataType::Int64, true), Field::new("tag1", ArrowDataType::Int64, true),
1262 Field::new("field1", ArrowDataType::Int64, true), Field::new("field0", ArrowDataType::Int64, true),
1264 Field::new(
1265 "ts",
1266 ArrowDataType::Timestamp(TimeUnit::Millisecond, None),
1267 false,
1268 ),
1269 Field::new(
1270 "__primary_key",
1271 ArrowDataType::Dictionary(
1272 Box::new(ArrowDataType::UInt32),
1273 Box::new(ArrowDataType::Binary),
1274 ),
1275 false,
1276 ),
1277 Field::new("__sequence", ArrowDataType::UInt64, false),
1278 Field::new("__op_type", ArrowDataType::UInt8, false),
1279 ];
1280 Arc::new(Schema::new(fields))
1281 }
1282
1283 fn build_test_flat_sst_schema_with_field_ids() -> SchemaRef {
1284 let ids = [
1285 Some(1u32),
1286 Some(3),
1287 Some(4),
1288 Some(2),
1289 Some(5),
1290 Some(PRIMARY_KEY_PARQUET_FIELD_ID),
1291 Some(SEQUENCE_PARQUET_FIELD_ID),
1292 Some(OP_TYPE_PARQUET_FIELD_ID),
1293 ];
1294 let fields: Vec<_> = build_test_flat_sst_schema()
1295 .fields()
1296 .iter()
1297 .zip(ids)
1298 .map(|(f, id)| match id {
1299 Some(id) => Arc::new(with_field_id((**f).clone(), id)) as _,
1300 None => f.clone(),
1301 })
1302 .collect();
1303 Arc::new(Schema::new(fields))
1304 }
1305
1306 #[test]
1307 fn test_flat_to_sst_arrow_schema() {
1308 let metadata = build_test_region_metadata();
1309 let format = FlatWriteFormat::new(metadata, &FlatSchemaOptions::default());
1310 assert_eq!(
1311 &build_test_flat_sst_schema_with_field_ids(),
1312 format.arrow_schema()
1313 );
1314 }
1315
1316 fn input_columns_for_flat_batch(num_rows: usize) -> Vec<ArrayRef> {
1317 vec![
1318 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])), ]
1327 }
1328
1329 #[test]
1330 fn test_flat_convert_batch() {
1331 let metadata = build_test_region_metadata();
1332 let format = FlatWriteFormat::new(metadata, &FlatSchemaOptions::default());
1333
1334 let num_rows = 4;
1335 let columns: Vec<ArrayRef> = input_columns_for_flat_batch(num_rows);
1336 let batch =
1337 RecordBatch::try_new(build_test_flat_sst_schema_with_field_ids(), columns.clone())
1338 .unwrap();
1339 let expect_record =
1340 RecordBatch::try_new(build_test_flat_sst_schema_with_field_ids(), columns).unwrap();
1341
1342 let actual = format.convert_batch(&batch).unwrap();
1343 assert_eq!(expect_record, actual);
1344 }
1345
1346 #[test]
1347 fn test_flat_convert_with_override_sequence() {
1348 let metadata = build_test_region_metadata();
1349 let format = FlatWriteFormat::new(metadata, &FlatSchemaOptions::default())
1350 .with_override_sequence(Some(415411));
1351
1352 let num_rows = 4;
1353 let columns: Vec<ArrayRef> = input_columns_for_flat_batch(num_rows);
1354 let batch =
1355 RecordBatch::try_new(build_test_flat_sst_schema_with_field_ids(), columns).unwrap();
1356
1357 let expected_columns: Vec<ArrayRef> = vec![
1358 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])), ];
1367 let expected_record = RecordBatch::try_new(
1368 build_test_flat_sst_schema_with_field_ids(),
1369 expected_columns,
1370 )
1371 .unwrap();
1372
1373 let actual = format.convert_batch(&batch).unwrap();
1374 assert_eq!(expected_record, actual);
1375 }
1376
1377 #[test]
1378 fn test_flat_projection_indices() {
1379 let metadata = build_test_region_metadata();
1380 let read_format =
1385 FlatReadFormat::new(metadata.clone(), ReadColumns::new([3]), None, "test", false)
1386 .unwrap();
1387 assert_eq!(
1388 &[1, 4, 5, 6, 7],
1389 read_format.parquet_read_columns().root_indices()
1390 );
1391
1392 let read_format =
1394 FlatReadFormat::new(metadata.clone(), ReadColumns::new([4]), None, "test", false)
1395 .unwrap();
1396 assert_eq!(
1397 &[2, 4, 5, 6, 7],
1398 read_format.parquet_read_columns().root_indices()
1399 );
1400
1401 let read_format =
1403 FlatReadFormat::new(metadata.clone(), ReadColumns::new([5]), None, "test", false)
1404 .unwrap();
1405 assert_eq!(
1406 &[4, 5, 6, 7],
1407 read_format.parquet_read_columns().root_indices()
1408 );
1409
1410 let read_format =
1412 FlatReadFormat::new(metadata, ReadColumns::new([2, 1, 5]), None, "test", false)
1413 .unwrap();
1414 assert_eq!(
1415 &[0, 3, 4, 5, 6, 7],
1416 read_format.parquet_read_columns().root_indices()
1417 );
1418 }
1419
1420 #[test]
1421 fn test_flat_read_format_convert_batch() {
1422 let metadata = build_test_region_metadata();
1423 let mut format = FlatReadFormat::new(
1424 metadata,
1425 ReadColumns::new(std::iter::once(1)), Some(build_test_flat_sst_schema()),
1427 "test",
1428 false,
1429 )
1430 .unwrap();
1431
1432 let num_rows = 4;
1433 let original_sequence = 100u64;
1434 let override_sequence = 200u64;
1435
1436 let columns: Vec<ArrayRef> = input_columns_for_flat_batch(num_rows);
1438 let mut test_columns = columns.clone();
1439 test_columns[6] = Arc::new(UInt64Array::from(vec![original_sequence; num_rows]));
1441 let record_batch =
1442 RecordBatch::try_new(format.arrow_schema().clone(), test_columns).unwrap();
1443
1444 let result = format.convert_batch(record_batch.clone(), None).unwrap();
1446 let sequence_column = result.column(sequence_column_index(result.num_columns()));
1447 let sequence_array = sequence_column
1448 .as_any()
1449 .downcast_ref::<UInt64Array>()
1450 .unwrap();
1451
1452 let expected_original = UInt64Array::from(vec![original_sequence; num_rows]);
1453 assert_eq!(sequence_array, &expected_original);
1454
1455 format.set_override_sequence(Some(override_sequence));
1457 let override_sequence_array = format.new_override_sequence_array(num_rows).unwrap();
1458 let result = format
1459 .convert_batch(record_batch, Some(&override_sequence_array))
1460 .unwrap();
1461 let sequence_column = result.column(sequence_column_index(result.num_columns()));
1462 let sequence_array = sequence_column
1463 .as_any()
1464 .downcast_ref::<UInt64Array>()
1465 .unwrap();
1466
1467 let expected_override = UInt64Array::from(vec![override_sequence; num_rows]);
1468 assert_eq!(sequence_array, &expected_override);
1469 }
1470
1471 #[test]
1472 fn test_need_convert_to_flat() {
1473 let metadata = build_test_region_metadata();
1474
1475 let expected_columns = metadata.column_metadatas.len() + 3;
1478 let result =
1479 FlatReadFormat::is_legacy_format(&metadata, expected_columns, "test.parquet").unwrap();
1480 assert!(
1481 !result,
1482 "Should not need conversion when column counts match"
1483 );
1484
1485 let num_columns_without_pk = expected_columns - metadata.primary_key.len();
1488 let result =
1489 FlatReadFormat::is_legacy_format(&metadata, num_columns_without_pk, "test.parquet")
1490 .unwrap();
1491 assert!(
1492 result,
1493 "Should need conversion when primary key columns are missing"
1494 );
1495
1496 let too_many_columns = expected_columns + 1;
1498 let err = FlatReadFormat::is_legacy_format(&metadata, too_many_columns, "test.parquet")
1499 .unwrap_err();
1500 assert!(err.to_string().contains("Expected columns"), "{err:?}");
1501
1502 let wrong_diff_columns = expected_columns - 1; let err = FlatReadFormat::is_legacy_format(&metadata, wrong_diff_columns, "test.parquet")
1505 .unwrap_err();
1506 assert!(
1507 err.to_string().contains("Column number difference"),
1508 "{err:?}"
1509 );
1510 }
1511
1512 fn build_test_dense_pk_array(
1513 codec: &DensePrimaryKeyCodec,
1514 pk_values_per_row: &[&[Option<i64>]],
1515 ) -> Arc<PrimaryKeyArray> {
1516 let mut builder = PrimaryKeyArrayBuilder::with_capacity(pk_values_per_row.len(), 1024, 0);
1517
1518 for pk_values_row in pk_values_per_row {
1519 let values: Vec<ValueRef> = pk_values_row
1520 .iter()
1521 .map(|opt| match opt {
1522 Some(val) => ValueRef::Int64(*val),
1523 None => ValueRef::Null,
1524 })
1525 .collect();
1526
1527 let encoded = codec.encode(values.into_iter()).unwrap();
1528 builder.append_value(&encoded);
1529 }
1530
1531 Arc::new(builder.finish())
1532 }
1533
1534 fn build_test_sparse_region_metadata() -> RegionMetadataRef {
1535 let mut builder = RegionMetadataBuilder::new(RegionId::new(1, 1));
1536 builder
1537 .push_column_metadata(ColumnMetadata {
1538 column_schema: ColumnSchema::new(
1539 "__table_id",
1540 ConcreteDataType::uint32_datatype(),
1541 false,
1542 ),
1543 semantic_type: SemanticType::Tag,
1544 column_id: ReservedColumnId::table_id(),
1545 })
1546 .push_column_metadata(ColumnMetadata {
1547 column_schema: ColumnSchema::new(
1548 "__tsid",
1549 ConcreteDataType::uint64_datatype(),
1550 false,
1551 ),
1552 semantic_type: SemanticType::Tag,
1553 column_id: ReservedColumnId::tsid(),
1554 })
1555 .push_column_metadata(ColumnMetadata {
1556 column_schema: ColumnSchema::new("tag0", ConcreteDataType::string_datatype(), true),
1557 semantic_type: SemanticType::Tag,
1558 column_id: 1,
1559 })
1560 .push_column_metadata(ColumnMetadata {
1561 column_schema: ColumnSchema::new("tag1", ConcreteDataType::string_datatype(), true),
1562 semantic_type: SemanticType::Tag,
1563 column_id: 3,
1564 })
1565 .push_column_metadata(ColumnMetadata {
1566 column_schema: ColumnSchema::new(
1567 "field1",
1568 ConcreteDataType::int64_datatype(),
1569 true,
1570 ),
1571 semantic_type: SemanticType::Field,
1572 column_id: 4,
1573 })
1574 .push_column_metadata(ColumnMetadata {
1575 column_schema: ColumnSchema::new(
1576 "field0",
1577 ConcreteDataType::int64_datatype(),
1578 true,
1579 ),
1580 semantic_type: SemanticType::Field,
1581 column_id: 2,
1582 })
1583 .push_column_metadata(ColumnMetadata {
1584 column_schema: ColumnSchema::new(
1585 "ts",
1586 ConcreteDataType::timestamp_millisecond_datatype(),
1587 false,
1588 ),
1589 semantic_type: SemanticType::Timestamp,
1590 column_id: 5,
1591 })
1592 .primary_key(vec![
1593 ReservedColumnId::table_id(),
1594 ReservedColumnId::tsid(),
1595 1,
1596 3,
1597 ])
1598 .primary_key_encoding(PrimaryKeyEncoding::Sparse);
1599 Arc::new(builder.build().unwrap())
1600 }
1601
1602 fn build_test_sparse_pk_array(
1603 codec: &SparsePrimaryKeyCodec,
1604 pk_values_per_row: &[SparseTestRow],
1605 ) -> Arc<PrimaryKeyArray> {
1606 let mut builder = PrimaryKeyArrayBuilder::with_capacity(pk_values_per_row.len(), 1024, 0);
1607 for row in pk_values_per_row {
1608 let values = vec![
1609 (ReservedColumnId::table_id(), ValueRef::UInt32(row.table_id)),
1610 (ReservedColumnId::tsid(), ValueRef::UInt64(row.tsid)),
1611 (1, ValueRef::String(&row.tag0)),
1612 (3, ValueRef::String(&row.tag1)),
1613 ];
1614
1615 let mut buffer = Vec::new();
1616 codec.encode_value_refs(&values, &mut buffer).unwrap();
1617 builder.append_value(&buffer);
1618 }
1619
1620 Arc::new(builder.finish())
1621 }
1622
1623 #[derive(Clone)]
1624 struct SparseTestRow {
1625 table_id: u32,
1626 tsid: u64,
1627 tag0: String,
1628 tag1: String,
1629 }
1630
1631 #[test]
1632 fn test_flat_read_format_convert_format_with_dense_encoding() {
1633 let metadata = build_test_region_metadata();
1634
1635 let column_ids: Vec<_> = metadata
1636 .column_metadatas
1637 .iter()
1638 .map(|c| c.column_id)
1639 .collect();
1640 let format = FlatReadFormat::new(
1641 metadata.clone(),
1642 ReadColumns::new(column_ids),
1643 Some(build_test_arrow_schema()),
1644 "test",
1645 false,
1646 )
1647 .unwrap();
1648
1649 let num_rows = 4;
1650 let original_sequence = 100u64;
1651
1652 let pk_values_per_row = vec![
1654 &[Some(1i64), Some(1i64)][..]; num_rows ];
1656
1657 let codec = DensePrimaryKeyCodec::new(&metadata);
1659 let dense_pk_array = build_test_dense_pk_array(&codec, &pk_values_per_row);
1660 let columns: Vec<ArrayRef> = vec![
1661 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])), ];
1668
1669 let old_schema = build_test_arrow_schema();
1671 let record_batch = RecordBatch::try_new(old_schema, columns).unwrap();
1672
1673 let result = format.convert_batch(record_batch, None).unwrap();
1675
1676 let expected_columns: Vec<ArrayRef> = vec![
1678 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])), ];
1687 let expected_record_batch = RecordBatch::try_new(
1688 build_test_flat_sst_schema_with_field_ids(),
1689 expected_columns,
1690 )
1691 .unwrap();
1692
1693 assert_eq!(expected_record_batch, result);
1695 }
1696
1697 #[test]
1698 fn test_flat_read_format_convert_format_with_sparse_encoding() {
1699 let metadata = build_test_sparse_region_metadata();
1700
1701 let column_ids: Vec<_> = metadata
1702 .column_metadatas
1703 .iter()
1704 .map(|c| c.column_id)
1705 .collect();
1706 let format = FlatReadFormat::new(
1707 metadata.clone(),
1708 ReadColumns::new(column_ids.clone()),
1709 None,
1710 "test",
1711 false,
1712 )
1713 .unwrap();
1714
1715 let num_rows = 4;
1716 let original_sequence = 100u64;
1717
1718 let pk_test_rows = vec![
1720 SparseTestRow {
1721 table_id: 1,
1722 tsid: 123,
1723 tag0: "frontend".to_string(),
1724 tag1: "pod1".to_string(),
1725 };
1726 num_rows
1727 ];
1728
1729 let codec = SparsePrimaryKeyCodec::new(&metadata);
1730 let sparse_pk_array = build_test_sparse_pk_array(&codec, &pk_test_rows);
1731 let columns: Vec<ArrayRef> = vec![
1733 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])), ];
1740
1741 let old_schema = build_test_arrow_schema();
1743 let record_batch = RecordBatch::try_new(old_schema, columns).unwrap();
1744
1745 let result = format.convert_batch(record_batch.clone(), None).unwrap();
1747
1748 let tag0_array = Arc::new(DictionaryArray::new(
1750 UInt32Array::from(vec![0; num_rows]),
1751 Arc::new(StringArray::from(vec!["frontend"])),
1752 ));
1753 let tag1_array = Arc::new(DictionaryArray::new(
1754 UInt32Array::from(vec![0; num_rows]),
1755 Arc::new(StringArray::from(vec!["pod1"])),
1756 ));
1757 let expected_columns: Vec<ArrayRef> = vec![
1758 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])), ];
1769 let expected_schema = to_flat_sst_arrow_schema(&metadata, &FlatSchemaOptions::default());
1770 let expected_record_batch =
1771 RecordBatch::try_new(expected_schema, expected_columns).unwrap();
1772
1773 assert_eq!(expected_record_batch, result);
1775
1776 let format = FlatReadFormat::new(
1777 metadata.clone(),
1778 ReadColumns::new(column_ids),
1779 None,
1780 "test",
1781 true,
1782 )
1783 .unwrap();
1784 let result = format.convert_batch(record_batch.clone(), None).unwrap();
1786 assert_eq!(record_batch, result);
1787 }
1788
1789 #[test]
1790 fn test_convert_flat_batch() {
1791 let metadata = build_test_region_metadata();
1792 let write_format = PrimaryKeyWriteFormat::new(metadata);
1793
1794 let num_rows = 4;
1795 let flat_columns: Vec<ArrayRef> = input_columns_for_flat_batch(num_rows);
1797 let flat_batch = RecordBatch::try_new(build_test_flat_sst_schema(), flat_columns).unwrap();
1798
1799 let result = write_format.convert_flat_batch(&flat_batch, 2).unwrap();
1801
1802 let expected_columns: Vec<ArrayRef> = vec![
1804 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])), ];
1811 let expected = RecordBatch::try_new(build_test_arrow_schema(), expected_columns).unwrap();
1812
1813 assert_eq!(expected, result);
1814 }
1815
1816 #[test]
1817 fn test_convert_flat_batch_with_override_sequence() {
1818 let metadata = build_test_region_metadata();
1819 let write_format = PrimaryKeyWriteFormat::new(metadata).with_override_sequence(Some(999));
1820
1821 let num_rows = 4;
1822 let flat_columns: Vec<ArrayRef> = input_columns_for_flat_batch(num_rows);
1823 let flat_batch = RecordBatch::try_new(build_test_flat_sst_schema(), flat_columns).unwrap();
1824
1825 let result = write_format.convert_flat_batch(&flat_batch, 2).unwrap();
1826
1827 let expected_columns: Vec<ArrayRef> = vec![
1828 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])), ];
1835 let expected = RecordBatch::try_new(build_test_arrow_schema(), expected_columns).unwrap();
1836
1837 assert_eq!(expected, result);
1838 }
1839
1840 #[test]
1841 fn test_convert_flat_batch_no_tags() {
1842 let mut builder = RegionMetadataBuilder::new(RegionId::new(1, 1));
1844 builder
1845 .push_column_metadata(ColumnMetadata {
1846 column_schema: ColumnSchema::new(
1847 "field0",
1848 ConcreteDataType::int64_datatype(),
1849 true,
1850 ),
1851 semantic_type: SemanticType::Field,
1852 column_id: 1,
1853 })
1854 .push_column_metadata(ColumnMetadata {
1855 column_schema: ColumnSchema::new(
1856 "ts",
1857 ConcreteDataType::timestamp_millisecond_datatype(),
1858 false,
1859 ),
1860 semantic_type: SemanticType::Timestamp,
1861 column_id: 2,
1862 });
1863 let metadata = Arc::new(builder.build().unwrap());
1864 let write_format = PrimaryKeyWriteFormat::new(metadata);
1865
1866 let num_rows = 3;
1867 let sst_schema = write_format.arrow_schema().clone();
1869 let columns: Vec<ArrayRef> = vec![
1870 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])), ];
1876 let flat_batch = RecordBatch::try_new(sst_schema.clone(), columns.clone()).unwrap();
1877
1878 let result = write_format.convert_flat_batch(&flat_batch, 1).unwrap();
1880 let expected = RecordBatch::try_new(sst_schema, columns).unwrap();
1881
1882 assert_eq!(expected, result);
1883 }
1884}