1use std::collections::HashMap;
18use std::sync::Arc;
19
20use api::v1::SemanticType;
21use arrow_schema::DataType;
22use arrow_schema::extension::{EXTENSION_TYPE_NAME_KEY, ExtensionType};
23use common_base::readable_size::ReadableSize;
24use common_query::native_histogram::{
25 is_native_histogram_value_type, native_histogram_list_element_id, native_histogram_subfield_id,
26};
27use datatypes::arrow::datatypes::{
28 DataType as ArrowDataType, Field, FieldRef, Fields, Schema, SchemaRef,
29};
30use datatypes::arrow::record_batch::RecordBatch;
31use datatypes::extension::histogram::HistogramExtensionType;
32use datatypes::prelude::ConcreteDataType;
33use datatypes::timestamp::timestamp_array_to_primitive;
34use serde::{Deserialize, Serialize};
35use store_api::codec::PrimaryKeyEncoding;
36use store_api::metadata::RegionMetadata;
37use store_api::storage::consts::{
38 OP_TYPE_COLUMN_NAME, PRIMARY_KEY_COLUMN_NAME, SEQUENCE_COLUMN_NAME,
39};
40
41use crate::error::{InvalidNativeHistogramFieldIdSnafu, InvalidNativeHistogramSubfieldSnafu};
42use crate::sst::parquet::flat_format::time_index_column_index;
43
44pub mod file;
45pub mod file_purger;
46pub mod file_ref;
47pub mod index;
48pub mod location;
49pub mod parquet;
50pub(crate) mod primary_key;
51pub mod range_index;
52pub(crate) mod version;
53
54pub const DEFAULT_WRITE_BUFFER_SIZE: ReadableSize = ReadableSize::mb(8);
56
57pub const DEFAULT_WRITE_CONCURRENCY: usize = 8;
59
60#[derive(Debug, Default, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, strum::EnumString)]
62#[serde(rename_all = "snake_case")]
63#[strum(serialize_all = "snake_case")]
64pub enum FormatType {
65 #[default]
67 PrimaryKey,
68 Flat,
70}
71
72pub const PARQUET_FIELD_ID_KEY: &str = "PARQUET:field_id";
74
75pub fn with_field_id(mut field: Field, column_id: u32) -> Field {
80 field
81 .metadata_mut()
82 .insert(PARQUET_FIELD_ID_KEY.to_string(), column_id.to_string());
83 field
84}
85
86fn stamp_native_histogram_subfield_ids(field: &mut Field) -> crate::error::Result<()> {
102 if !is_native_histogram_value_type(&ConcreteDataType::from_arrow_type(field.data_type())) {
103 return Ok(());
104 }
105 let column_id = field
110 .metadata()
111 .get(PARQUET_FIELD_ID_KEY)
112 .and_then(|s| s.parse::<i32>().ok())
113 .ok_or_else(|| {
114 InvalidNativeHistogramFieldIdSnafu {
115 field_name: field.name().clone(),
116 }
117 .build()
118 })?;
119 field.metadata_mut().insert(
121 EXTENSION_TYPE_NAME_KEY.to_string(),
122 HistogramExtensionType::NAME.to_string(),
123 );
124 let ArrowDataType::Struct(children) = field.data_type() else {
125 return Ok(());
126 };
127 let new_children: crate::error::Result<Fields> = children
128 .iter()
129 .map(|child| {
130 let mut c = (**child).clone();
131 let id = native_histogram_subfield_id(column_id, c.name()).ok_or_else(|| {
136 InvalidNativeHistogramSubfieldSnafu {
137 column_id,
138 field_name: c.name().clone(),
139 }
140 .build()
141 })?;
142 c.metadata_mut()
144 .insert(PARQUET_FIELD_ID_KEY.to_string(), id.to_string());
145 if let ArrowDataType::List(elem) = c.data_type() {
147 let elem_id =
148 native_histogram_list_element_id(column_id, c.name()).ok_or_else(|| {
149 InvalidNativeHistogramSubfieldSnafu {
150 column_id,
151 field_name: c.name().clone(),
152 }
153 .build()
154 })?;
155 let mut new_elem = (**elem).clone();
156 new_elem
157 .metadata_mut()
158 .insert(PARQUET_FIELD_ID_KEY.to_string(), elem_id.to_string());
159 c.set_data_type(ArrowDataType::List(Arc::new(new_elem)));
160 }
161 Ok(Arc::new(c))
162 })
163 .collect();
164 field.set_data_type(ArrowDataType::Struct(new_children?));
165 Ok(())
166}
167
168pub fn maybe_wrap_schema(schema: &SchemaRef) -> crate::error::Result<SchemaRef> {
178 if !schema
181 .fields()
182 .iter()
183 .any(|f| matches!(f.data_type(), ArrowDataType::Struct(_)))
184 {
185 return Ok(schema.clone());
186 }
187 let new_fields: crate::error::Result<Vec<FieldRef>> = schema
188 .fields()
189 .iter()
190 .map(|f| {
191 let mut field = (**f).clone();
192 stamp_native_histogram_subfield_ids(&mut field)?;
193 Ok(Arc::new(field))
194 })
195 .collect();
196 Ok(Arc::new(Schema::new_with_metadata(
197 Fields::from(new_fields?),
198 schema.metadata().clone(),
199 )))
200}
201
202pub(crate) const INTERNAL_PARQUET_FIELD_ID_BASE: u32 = 1 << 30;
205
206pub(crate) const PRIMARY_KEY_PARQUET_FIELD_ID: u32 = INTERNAL_PARQUET_FIELD_ID_BASE;
208pub(crate) const SEQUENCE_PARQUET_FIELD_ID: u32 = INTERNAL_PARQUET_FIELD_ID_BASE + 1;
210pub(crate) const OP_TYPE_PARQUET_FIELD_ID: u32 = INTERNAL_PARQUET_FIELD_ID_BASE + 2;
212
213pub fn to_sst_arrow_schema(metadata: &RegionMetadata) -> SchemaRef {
215 let fields = Fields::from_iter(
216 metadata
217 .schema
218 .arrow_schema()
219 .fields()
220 .iter()
221 .zip(&metadata.column_metadatas)
222 .filter_map(|(field, column_meta)| {
223 if column_meta.semantic_type == SemanticType::Field {
224 Some(Arc::new(with_field_id(
225 (**field).clone(),
226 column_meta.column_id,
227 )))
228 } else {
229 None
231 }
232 })
233 .chain([Arc::new(with_field_id(
234 (*metadata.time_index_field()).clone(),
235 metadata.time_index_column().column_id,
236 ))])
237 .chain(internal_fields()),
238 );
239
240 Arc::new(Schema::new(fields))
241}
242
243pub struct FlatSchemaOptions {
245 pub raw_pk_columns: bool,
247 pub string_pk_use_dict: bool,
251 pub concretized_json_types: HashMap<String, DataType>,
254}
255
256impl Default for FlatSchemaOptions {
257 fn default() -> Self {
258 Self {
259 raw_pk_columns: true,
260 string_pk_use_dict: true,
261 concretized_json_types: HashMap::new(),
262 }
263 }
264}
265
266impl FlatSchemaOptions {
267 pub fn from_encoding(encoding: PrimaryKeyEncoding) -> Self {
269 if encoding == PrimaryKeyEncoding::Dense {
270 Self::default()
271 } else {
272 Self {
273 raw_pk_columns: false,
274 string_pk_use_dict: false,
275 concretized_json_types: HashMap::new(),
276 }
277 }
278 }
279}
280
281pub fn to_flat_sst_arrow_schema(
291 metadata: &RegionMetadata,
292 options: &FlatSchemaOptions,
293) -> SchemaRef {
294 let num_fields = flat_sst_arrow_schema_column_num(metadata, options);
295 let mut fields = Vec::with_capacity(num_fields);
296 let schema = metadata.schema.arrow_schema();
297 if options.raw_pk_columns {
298 for pk_id in &metadata.primary_key {
299 let pk_index = metadata.column_index_by_id(*pk_id).unwrap();
300 let column_id = metadata.column_metadatas[pk_index].column_id;
301 if options.string_pk_use_dict {
302 let old_field = &schema.fields[pk_index];
303 let new_field = tag_maybe_to_dictionary_field(
304 &metadata.column_metadatas[pk_index].column_schema.data_type,
305 old_field,
306 );
307 let new_field = concretize_json_type(new_field, options);
308 fields.push(Arc::new(with_field_id((*new_field).clone(), column_id)));
309 }
310 }
311 }
312 let remaining_fields = schema
313 .fields()
314 .iter()
315 .zip(&metadata.column_metadatas)
316 .filter_map(|(field, column_meta)| {
317 if column_meta.semantic_type == SemanticType::Field {
318 let field = concretize_json_type(field.clone(), options);
319 Some(Arc::new(with_field_id(
320 Arc::unwrap_or_clone(field),
321 column_meta.column_id,
322 )))
323 } else {
324 None
325 }
326 })
327 .chain([Arc::new(with_field_id(
328 (*metadata.time_index_field()).clone(),
329 metadata.time_index_column().column_id,
330 ))])
331 .chain(internal_fields());
332 for field in remaining_fields {
333 fields.push(field);
334 }
335
336 Arc::new(Schema::new(fields))
337}
338
339fn concretize_json_type(field: Arc<Field>, options: &FlatSchemaOptions) -> Arc<Field> {
340 if let Some(data_type) = options.concretized_json_types.get(field.name()) {
341 let mut field = Arc::unwrap_or_clone(field);
342 field.set_data_type(data_type.clone());
343 Arc::new(field)
344 } else {
345 field
346 }
347}
348
349pub fn flat_sst_arrow_schema_column_num(
351 metadata: &RegionMetadata,
352 options: &FlatSchemaOptions,
353) -> usize {
354 if options.raw_pk_columns {
355 metadata.column_metadatas.len() + 3
356 } else {
357 metadata.column_metadatas.len() + 3 - metadata.primary_key.len()
358 }
359}
360
361fn to_dictionary_field(field: &Field) -> Field {
363 let mut new_field = Field::new_dictionary(
364 field.name(),
365 datatypes::arrow::datatypes::DataType::UInt32,
366 field.data_type().clone(),
367 field.is_nullable(),
368 );
369
370 if let Some(field_id) = field.metadata().get(PARQUET_FIELD_ID_KEY) {
372 new_field
373 .metadata_mut()
374 .insert(PARQUET_FIELD_ID_KEY.to_string(), field_id.clone());
375 }
376
377 new_field
378}
379
380pub(crate) fn tag_maybe_to_dictionary_field(
382 data_type: &ConcreteDataType,
383 field: &Arc<Field>,
384) -> Arc<Field> {
385 if data_type.is_string() {
386 Arc::new(to_dictionary_field(field))
387 } else {
388 field.clone()
389 }
390}
391
392pub(crate) fn internal_fields() -> [FieldRef; 3] {
394 [
396 Arc::new(with_field_id(
397 Field::new_dictionary(
398 PRIMARY_KEY_COLUMN_NAME,
399 ArrowDataType::UInt32,
400 ArrowDataType::Binary,
401 false,
402 ),
403 PRIMARY_KEY_PARQUET_FIELD_ID,
404 )),
405 Arc::new(with_field_id(
406 Field::new(SEQUENCE_COLUMN_NAME, ArrowDataType::UInt64, false),
407 SEQUENCE_PARQUET_FIELD_ID,
408 )),
409 Arc::new(with_field_id(
410 Field::new(OP_TYPE_COLUMN_NAME, ArrowDataType::UInt8, false),
411 OP_TYPE_PARQUET_FIELD_ID,
412 )),
413 ]
414}
415
416pub(crate) fn override_pk_field_to_binary(schema: &SchemaRef) -> SchemaRef {
418 let new_fields = schema
419 .fields()
420 .iter()
421 .map(|field| {
422 if field.name() == PRIMARY_KEY_COLUMN_NAME {
423 let mut new_field = Field::new(
424 PRIMARY_KEY_COLUMN_NAME,
425 ArrowDataType::Binary,
426 field.is_nullable(),
427 );
428 if let Some(field_id) = field.metadata().get(PARQUET_FIELD_ID_KEY) {
431 new_field
432 .metadata_mut()
433 .insert(PARQUET_FIELD_ID_KEY.to_string(), field_id.clone());
434 }
435 Arc::new(new_field)
436 } else {
437 field.clone()
438 }
439 })
440 .collect::<Vec<_>>();
441 Arc::new(Schema::new(new_fields))
442}
443
444#[derive(Default)]
449pub(crate) struct SeriesEstimator {
450 last_timestamp: Option<i64>,
452 series_count: u64,
454}
455
456impl SeriesEstimator {
457 pub(crate) fn update_flat(&mut self, record_batch: &RecordBatch) {
461 let batch_rows = record_batch.num_rows();
462 if batch_rows == 0 {
463 return;
464 }
465
466 let time_index_pos = time_index_column_index(record_batch.num_columns());
467 let timestamps = record_batch.column(time_index_pos);
468 let Some((ts_values, _unit)) = timestamp_array_to_primitive(timestamps) else {
469 return;
470 };
471 let values = ts_values.values();
472
473 if let Some(last_ts) = self.last_timestamp {
475 if values[0] <= last_ts {
476 self.series_count += 1;
477 }
478 } else {
479 self.series_count = 1;
481 }
482
483 for i in 0..batch_rows - 1 {
485 if values[i] >= values[i + 1] {
488 self.series_count += 1;
489 }
490 }
491
492 self.last_timestamp = Some(values[batch_rows - 1]);
494 }
495
496 pub(crate) fn finish(&mut self) -> u64 {
498 self.last_timestamp = None;
499 let count = self.series_count;
500 self.series_count = 0;
501
502 count
503 }
504}
505
506#[cfg(test)]
507mod tests {
508 use std::sync::Arc;
509
510 use ::parquet::arrow::AsyncArrowWriter;
511 use ::parquet::arrow::arrow_reader::ParquetRecordBatchReaderBuilder;
512 use ::parquet::basic::LogicalType;
513 use ::parquet::variant::{VariantArray, VariantType, json_to_variant};
514 use common_query::prelude::greptime_native_histogram;
515 use datatypes::arrow::array::{
516 ArrayRef, BinaryArray, DictionaryArray, Int64Array, StringArray, StructArray,
517 TimestampMillisecondArray, UInt8Array, UInt32Array, UInt64Array,
518 };
519 use datatypes::arrow::datatypes::{DataType as ArrowDataType, Field, Schema, TimeUnit};
520 use datatypes::arrow::record_batch::RecordBatch;
521 use datatypes::extension::json::{Json2ExtensionType, Json2PhysicalLayout};
522 use datatypes::vectors::json::array::JsonArray;
523 use serde_json::json;
524
525 use super::*;
526
527 fn new_flat_record_batch(timestamps: &[i64]) -> RecordBatch {
528 let num_cols = 4; let time_index_pos = time_index_column_index(num_cols);
531 assert_eq!(time_index_pos, 0); let time_array = Arc::new(TimestampMillisecondArray::from(timestamps.to_vec()));
534 let pk_array = Arc::new(DictionaryArray::new(
535 UInt32Array::from(vec![0; timestamps.len()]),
536 Arc::new(BinaryArray::from(vec![b"test".as_slice()])),
537 ));
538 let seq_array = Arc::new(UInt64Array::from(vec![1; timestamps.len()]));
539 let op_array = Arc::new(UInt8Array::from(vec![1; timestamps.len()]));
540
541 let schema = Arc::new(Schema::new(vec![
542 Field::new(
543 "time",
544 ArrowDataType::Timestamp(TimeUnit::Millisecond, None),
545 false,
546 ),
547 Field::new_dictionary(
548 "__primary_key",
549 ArrowDataType::UInt32,
550 ArrowDataType::Binary,
551 false,
552 ),
553 Field::new("__sequence", ArrowDataType::UInt64, false),
554 Field::new("__op_type", ArrowDataType::UInt8, false),
555 ]));
556
557 RecordBatch::try_new(schema, vec![time_array, pk_array, seq_array, op_array]).unwrap()
558 }
559
560 #[test]
561 fn test_series_estimator_flat_empty_batch() {
562 let mut estimator = SeriesEstimator::default();
563 let record_batch = new_flat_record_batch(&[]);
564 estimator.update_flat(&record_batch);
565 assert_eq!(0, estimator.finish());
566 }
567
568 #[test]
569 fn test_series_estimator_flat_single_batch() {
570 let mut estimator = SeriesEstimator::default();
571 let record_batch = new_flat_record_batch(&[1, 2, 3]);
572 estimator.update_flat(&record_batch);
573 assert_eq!(1, estimator.finish());
574 }
575
576 #[test]
577 fn test_series_estimator_flat_series_boundary_within_batch() {
578 let mut estimator = SeriesEstimator::default();
579 let record_batch = new_flat_record_batch(&[1, 2, 3, 2, 4, 5]);
581 estimator.update_flat(&record_batch);
582 assert_eq!(2, estimator.finish());
584 }
585
586 #[test]
587 fn test_series_estimator_flat_multiple_boundaries_within_batch() {
588 let mut estimator = SeriesEstimator::default();
589 let record_batch = new_flat_record_batch(&[1, 2, 5, 4, 6, 3, 7]);
591 estimator.update_flat(&record_batch);
592 assert_eq!(3, estimator.finish());
593 }
594
595 #[test]
596 fn test_series_estimator_flat_equal_timestamps() {
597 let mut estimator = SeriesEstimator::default();
598 let record_batch = new_flat_record_batch(&[1, 2, 2, 3, 3, 3, 4]);
600 estimator.update_flat(&record_batch);
601 assert_eq!(4, estimator.finish());
603 }
604
605 #[test]
606 fn test_series_estimator_flat_multiple_batches_continuation() {
607 let mut estimator = SeriesEstimator::default();
608
609 let batch1 = new_flat_record_batch(&[1, 2, 3]);
611 estimator.update_flat(&batch1);
612
613 let batch2 = new_flat_record_batch(&[4, 5, 6]);
615 estimator.update_flat(&batch2);
616
617 assert_eq!(1, estimator.finish());
618 }
619
620 #[test]
621 fn test_series_estimator_flat_multiple_batches_new_series() {
622 let mut estimator = SeriesEstimator::default();
623
624 let batch1 = new_flat_record_batch(&[1, 2, 3]);
626 estimator.update_flat(&batch1);
627
628 let batch2 = new_flat_record_batch(&[2, 3, 4]);
630 estimator.update_flat(&batch2);
631
632 assert_eq!(2, estimator.finish());
633 }
634
635 #[test]
636 fn test_series_estimator_flat_boundary_at_batch_edge_equal() {
637 let mut estimator = SeriesEstimator::default();
638
639 let batch1 = new_flat_record_batch(&[1, 2, 5]);
641 estimator.update_flat(&batch1);
642
643 let batch2 = new_flat_record_batch(&[5, 6, 7]);
645 estimator.update_flat(&batch2);
646
647 assert_eq!(2, estimator.finish());
648 }
649
650 #[test]
651 fn test_series_estimator_flat_mixed_batches() {
652 let mut estimator = SeriesEstimator::default();
653
654 let batch1 = new_flat_record_batch(&[10, 20, 30]);
656 estimator.update_flat(&batch1);
657
658 let batch2 = new_flat_record_batch(&[5, 15, 10, 25]);
660 estimator.update_flat(&batch2);
661
662 let batch3 = new_flat_record_batch(&[30, 35]);
664 estimator.update_flat(&batch3);
665
666 assert_eq!(3, estimator.finish());
668 }
669
670 #[test]
671 fn test_series_estimator_flat_descending_timestamps() {
672 let mut estimator = SeriesEstimator::default();
673 let record_batch = new_flat_record_batch(&[10, 9, 8, 7, 6]);
675 estimator.update_flat(&record_batch);
676 assert_eq!(5, estimator.finish());
678 }
679
680 #[test]
681 fn test_series_estimator_flat_finish_resets_state() {
682 let mut estimator = SeriesEstimator::default();
683
684 let batch1 = new_flat_record_batch(&[1, 2, 3]);
685 estimator.update_flat(&batch1);
686
687 assert_eq!(1, estimator.finish());
688
689 let batch2 = new_flat_record_batch(&[4, 5, 6]);
691 estimator.update_flat(&batch2);
692
693 assert_eq!(1, estimator.finish());
694 }
695
696 fn histogram_field(name: &str, column_id: u32) -> Field {
699 use common_query::native_histogram::native_histogram_value_type;
700 use datatypes::data_type::DataType;
701 with_field_id(
702 Field::new(name, native_histogram_value_type().as_arrow_type(), true),
703 column_id,
704 )
705 }
706
707 fn assert_histogram_stamped(field: &Field, column_id: i32) {
711 use arrow_schema::extension::ExtensionType;
712 use common_query::native_histogram::{
713 native_histogram_list_element_id, native_histogram_subfield_id,
714 };
715 use datatypes::extension::histogram::HistogramExtensionType;
716
717 assert_eq!(
718 field
719 .metadata()
720 .get(arrow_schema::extension::EXTENSION_TYPE_NAME_KEY)
721 .map(|s| s.as_str()),
722 Some(HistogramExtensionType::NAME),
723 "histogram field must carry the greptime.histogram extension"
724 );
725 let ArrowDataType::Struct(children) = field.data_type() else {
726 panic!("expected a struct, got {:?}", field.data_type());
727 };
728 for child in children {
729 let expected = native_histogram_subfield_id(column_id, child.name())
730 .unwrap_or_else(|| panic!("no id for sub-field {}", child.name()));
731 let got: i32 = child
732 .metadata()
733 .get(PARQUET_FIELD_ID_KEY)
734 .unwrap_or_else(|| panic!("sub-field {} missing field id", child.name()))
735 .parse()
736 .unwrap();
737 assert_eq!(got, expected, "sub-field {} id", child.name());
738 if let ArrowDataType::List(elem) = child.data_type() {
739 let elem_expected =
740 native_histogram_list_element_id(column_id, child.name()).unwrap();
741 let elem_got: i32 = elem
742 .metadata()
743 .get(PARQUET_FIELD_ID_KEY)
744 .unwrap_or_else(|| panic!("list element of {} missing id", child.name()))
745 .parse()
746 .unwrap();
747 assert_eq!(
748 elem_got,
749 elem_expected,
750 "list element id of {}",
751 child.name()
752 );
753 }
754 }
755 }
756
757 #[test]
758 fn test_maybe_wrap_schema_native_histogram() {
759 let schema = Arc::new(Schema::new(vec![
760 Field::new(
761 "greptime_timestamp",
762 ArrowDataType::Timestamp(TimeUnit::Millisecond, None),
763 false,
764 ),
765 histogram_field(greptime_native_histogram(), 1),
766 ]));
767
768 let wrapped = maybe_wrap_schema(&schema).unwrap();
769 let hist = wrapped
770 .field_with_name(greptime_native_histogram())
771 .expect("histogram field present");
772 let ArrowDataType::Struct(children) = hist.data_type() else {
774 unreachable!()
775 };
776 assert_eq!(children.len(), 18);
777 assert_histogram_stamped(hist, 1);
778 }
779
780 #[test]
781 fn test_maybe_wrap_schema_multiple_histograms_disjoint_ids() {
782 use common_query::native_histogram::native_histogram_subfield_id;
786
787 let schema = Arc::new(Schema::new(vec![
788 histogram_field(greptime_native_histogram(), 1),
789 histogram_field(greptime_native_histogram(), 7),
790 ]));
791 let wrapped = maybe_wrap_schema(&schema).unwrap();
792 let h1 = &wrapped.fields()[0];
793 let h2 = &wrapped.fields()[1];
794 assert_histogram_stamped(h1, 1);
795 assert_histogram_stamped(h2, 7);
796 assert_ne!(
798 native_histogram_subfield_id(1, "sum"),
799 native_histogram_subfield_id(7, "sum")
800 );
801 }
802
803 #[test]
804 fn test_maybe_wrap_schema_recognizes_histogram_by_type() {
805 let schema = Arc::new(Schema::new(vec![histogram_field("custom_histogram", 5)]));
806
807 let wrapped = maybe_wrap_schema(&schema).unwrap();
808 let hist = wrapped.field_with_name("custom_histogram").unwrap();
809 assert_histogram_stamped(hist, 5);
810 }
811
812 #[test]
813 fn test_maybe_wrap_schema_plain_struct_not_stamped() {
814 use arrow_schema::extension::EXTENSION_TYPE_NAME_KEY;
815
816 let plain = ArrowDataType::Struct(
817 vec![
818 Arc::new(Field::new("a", ArrowDataType::Int32, true)),
819 Arc::new(Field::new("b", ArrowDataType::Utf8, true)),
820 ]
821 .into(),
822 );
823 let schema = Arc::new(Schema::new(vec![
824 Field::new(
825 "ts",
826 ArrowDataType::Timestamp(TimeUnit::Millisecond, None),
827 false,
828 ),
829 Field::new("data", plain, true),
830 ]));
831
832 let wrapped = maybe_wrap_schema(&schema).unwrap();
833 let data = wrapped.field_with_name("data").unwrap();
834 assert!(
835 data.metadata().get(EXTENSION_TYPE_NAME_KEY).is_none(),
836 "non-histogram struct must not get the extension"
837 );
838 if let ArrowDataType::Struct(children) = data.data_type() {
839 for child in children {
840 assert!(
841 child.metadata().get(PARQUET_FIELD_ID_KEY).is_none(),
842 "non-histogram sub-field {} must not get a field id",
843 child.name()
844 );
845 }
846 }
847 }
848
849 #[test]
850 fn test_maybe_wrap_schema_no_struct_unchanged() {
851 let schema: Arc<Schema> = Arc::new(Schema::new(vec![
852 Field::new(
853 "ts",
854 ArrowDataType::Timestamp(TimeUnit::Millisecond, None),
855 false,
856 ),
857 Field::new("v", ArrowDataType::Float64, true),
858 ]));
859 let wrapped = maybe_wrap_schema(&schema).unwrap();
860 assert!(
861 Arc::ptr_eq(&wrapped, &schema),
862 "a schema without any struct column must be returned unchanged"
863 );
864 }
865
866 fn parquet_footer_arrow_schema(schema: &SchemaRef) -> SchemaRef {
877 use ::parquet::arrow::ArrowWriter;
878 use ::parquet::arrow::arrow_reader::ParquetRecordBatchReaderBuilder;
879 use ::parquet::file::properties::WriterProperties;
880 use bytes::Bytes;
881
882 let wrapped = maybe_wrap_schema(schema).unwrap();
883 let mut bytes = Vec::new();
884 let props = WriterProperties::builder().build();
885 let mut writer = ArrowWriter::try_new(&mut bytes, wrapped.clone(), Some(props)).unwrap();
886 writer
887 .write(&RecordBatch::new_empty(wrapped.clone()))
888 .unwrap();
889 writer.close().unwrap();
890
891 ParquetRecordBatchReaderBuilder::try_new(Bytes::from(bytes))
892 .unwrap()
893 .schema()
894 .clone()
895 }
896
897 #[test]
898 fn test_maybe_wrap_schema_survives_parquet_roundtrip() {
899 let schema = Arc::new(Schema::new(vec![
903 Field::new(
904 "greptime_timestamp",
905 ArrowDataType::Timestamp(TimeUnit::Millisecond, None),
906 false,
907 ),
908 histogram_field(greptime_native_histogram(), 3),
909 ]));
910
911 let on_disk = parquet_footer_arrow_schema(&schema);
912 let hist = on_disk
913 .field_with_name(greptime_native_histogram())
914 .expect("histogram field present");
915 assert_histogram_stamped(hist, 3);
916 }
917
918 #[test]
919 fn test_parquet_roundtrip_recognizes_histogram_by_type() {
920 let schema = Arc::new(Schema::new(vec![
923 Field::new(
924 "ts",
925 ArrowDataType::Timestamp(TimeUnit::Millisecond, None),
926 false,
927 ),
928 histogram_field("custom_histogram", 9),
929 ]));
930
931 let on_disk = parquet_footer_arrow_schema(&schema);
932 let hist = on_disk.field_with_name("custom_histogram").unwrap();
933 assert_histogram_stamped(hist, 9);
934 }
935
936 #[test]
937 fn test_maybe_wrap_schema_overflows_return_error() {
938 let schema = Arc::new(Schema::new(vec![
943 Field::new(
944 "greptime_timestamp",
945 ArrowDataType::Timestamp(TimeUnit::Millisecond, None),
946 false,
947 ),
948 histogram_field(greptime_native_histogram(), 12_582_912),
949 ]));
950 let err = maybe_wrap_schema(&schema).unwrap_err();
951 assert!(
952 matches!(
953 err,
954 crate::error::Error::InvalidNativeHistogramSubfield { .. }
955 ),
956 "expected InvalidNativeHistogramSubfield, got {:?}",
957 err
958 );
959 }
960
961 #[test]
962 fn test_maybe_wrap_schema_missing_field_id_returns_error() {
963 use common_query::native_histogram::native_histogram_value_type;
969 use datatypes::data_type::DataType;
970
971 let field = Field::new(
972 greptime_native_histogram(),
973 native_histogram_value_type().as_arrow_type(),
974 true,
975 );
976 assert!(
977 field.metadata().get(PARQUET_FIELD_ID_KEY).is_none(),
978 "fixture must not carry a field id"
979 );
980 let schema = Arc::new(Schema::new(vec![field]));
981 let err = maybe_wrap_schema(&schema).unwrap_err();
982 assert!(
983 matches!(
984 err,
985 crate::error::Error::InvalidNativeHistogramFieldId { .. }
986 ),
987 "expected InvalidNativeHistogramFieldId, got {:?}",
988 err
989 );
990 }
991
992 #[test]
993 fn test_maybe_wrap_schema_field_id_above_i32_max_returns_error() {
994 let schema = Arc::new(Schema::new(vec![
999 Field::new(
1000 "greptime_timestamp",
1001 ArrowDataType::Timestamp(TimeUnit::Millisecond, None),
1002 false,
1003 ),
1004 histogram_field(greptime_native_histogram(), u32::MAX),
1005 ]));
1006 let err = maybe_wrap_schema(&schema).unwrap_err();
1007 assert!(
1008 matches!(
1009 err,
1010 crate::error::Error::InvalidNativeHistogramFieldId { .. }
1011 ),
1012 "expected InvalidNativeHistogramFieldId, got {:?}",
1013 err
1014 );
1015 }
1016
1017 fn json2_v2_test_type() -> ArrowDataType {
1018 ArrowDataType::Struct(
1019 vec![
1020 Arc::new(Field::new("active", ArrowDataType::Boolean, true)),
1021 Arc::new(Field::new("hot", ArrowDataType::Int64, true)),
1022 Arc::new(Field::new("name", ArrowDataType::Utf8, true)),
1023 ]
1024 .into(),
1025 )
1026 }
1027
1028 #[tokio::test]
1036 async fn test_nested_variant_survives_sst_writer_schema_roundtrip()
1037 -> Result<(), Box<dyn std::error::Error>> {
1038 let json: ArrayRef = Arc::new(StringArray::from(vec![
1039 Some(r#"{}"#),
1040 Some(r#"{"name":"Alice","active":true}"#),
1041 Some(r#"{"nested":{"count":42},"items":[1,"two",null]}"#),
1042 Some(r#"{"\u5b57\u6bb5":"\u503c"}"#),
1043 None,
1044 ]));
1045 let remainder = json_to_variant(&json)?;
1046 let remainder_field = remainder.field("!__remainder__!");
1047 let remainder_array = ArrayRef::from(remainder);
1048 let hot_field = Field::new("hot", ArrowDataType::Int64, true);
1049 let data_array = Arc::new(StructArray::new(
1050 vec![remainder_field.clone(), hot_field.clone()].into(),
1051 vec![
1052 remainder_array,
1053 Arc::new(Int64Array::from(vec![
1054 Some(1),
1055 Some(2),
1056 Some(3),
1057 Some(4),
1058 None,
1059 ])),
1060 ],
1061 None,
1062 ));
1063 let data_field = Field::new(
1064 "data",
1065 ArrowDataType::Struct(vec![remainder_field, hot_field].into()),
1066 true,
1067 )
1068 .with_extension_type(Json2ExtensionType::default());
1069 let schema = Arc::new(Schema::new(vec![data_field]));
1070 let source = RecordBatch::try_new(schema.clone(), vec![data_array])?;
1071
1072 let wrapped = maybe_wrap_schema(&schema)?;
1073 let mut buffer = Vec::new();
1074 let mut writer = AsyncArrowWriter::try_new(&mut buffer, wrapped, None)?;
1075 writer.write(&source).await?;
1076 writer.close().await?;
1077
1078 let builder = ParquetRecordBatchReaderBuilder::try_new(bytes::Bytes::from(buffer))?;
1079 let parquet_remainder =
1080 &builder.parquet_schema().root_schema().get_fields()[0].get_fields()[0];
1081 assert_eq!(
1082 parquet_remainder.get_basic_info().logical_type_ref(),
1083 Some(&LogicalType::variant(None))
1084 );
1085
1086 let ArrowDataType::Struct(children) = builder.schema().field_with_name("data")?.data_type()
1087 else {
1088 unreachable!();
1089 };
1090 assert!(children[0].has_valid_extension_type::<VariantType>());
1091
1092 let mut reader = builder.build()?;
1093 let result = reader.next().unwrap()?;
1094 assert_eq!(source, result);
1095 let result_field = result.schema().field(0).clone();
1096 let result = result
1097 .column(0)
1098 .as_any()
1099 .downcast_ref::<StructArray>()
1100 .unwrap();
1101 VariantArray::try_new(result.column(0))?;
1102 let result: ArrayRef = Arc::new(result.clone());
1103 let result =
1104 JsonArray::from(&result).project_to_v2(&result_field, &json2_v2_test_type())?;
1105 assert_eq!(
1106 json!({"active": true, "hot": 2, "name": "Alice"}),
1107 JsonArray::from(&result).try_get_value(1)?
1108 );
1109 Ok(())
1110 }
1111
1112 #[test]
1114 fn test_read_json2_v2_fixture() -> Result<(), Box<dyn std::error::Error>> {
1115 let bytes = bytes::Bytes::from_static(include_bytes!("../test-data/json2-v2.parquet"));
1116 let builder = ParquetRecordBatchReaderBuilder::try_new(bytes)?;
1117 let field = builder.schema().field(0).clone();
1118 assert!(Json2PhysicalLayout::try_from_root(&field)?.is_version_2());
1119
1120 let batch = builder.build()?.next().unwrap()?;
1121 let data = batch
1122 .column(0)
1123 .as_any()
1124 .downcast_ref::<StructArray>()
1125 .unwrap();
1126 VariantArray::try_new(data.column(0))?;
1127 let data: ArrayRef = Arc::new(data.clone());
1128 let data = JsonArray::from(&data).project_to_v2(&field, &json2_v2_test_type())?;
1129 assert_eq!(
1130 json!({"active": true, "hot": 2, "name": "Alice"}),
1131 JsonArray::from(&data).try_get_value(1)?
1132 );
1133 Ok(())
1134 }
1135}