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