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::prelude::DataType;
29use datatypes::value::Value;
30use datatypes::vectors::VectorRef;
31use mito_codec::row_converter::{
32 CompositeValues, PrimaryKeyCodec, SortField, build_primary_key_codec,
33 build_primary_key_codec_with_fields,
34};
35use snafu::{OptionExt, ResultExt, ensure};
36use store_api::codec::PrimaryKeyEncoding;
37use store_api::metadata::{RegionMetadata, RegionMetadataRef};
38use store_api::storage::ColumnId;
39
40use crate::error::{
41 CompatReaderSnafu, ComputeArrowSnafu, CreateDefaultSnafu, DecodeSnafu, EncodeSnafu,
42 NewRecordBatchSnafu, Result, UnexpectedSnafu, UnsupportedOperationSnafu,
43};
44use crate::read::flat_projection::{FlatProjectionMapper, flat_projected_columns};
45use crate::sst::parquet::flat_format::{FlatReadFormat, primary_key_column_index};
46use crate::sst::parquet::format::{INTERNAL_COLUMN_NUM, PrimaryKeyArray};
47use crate::sst::{internal_fields, tag_maybe_to_dictionary_field, with_field_id};
48
49pub(crate) fn has_same_columns_and_pk_encoding(
52 projection_mapper: &FlatProjectionMapper,
53 read_format: &FlatReadFormat,
54 compaction: bool,
55) -> bool {
56 let left = projection_mapper.metadata();
57 let right = read_format.metadata();
58 if left.primary_key_encoding != right.primary_key_encoding {
59 return false;
60 }
61
62 if left.column_metadatas.len() != right.column_metadatas.len() {
63 return false;
64 }
65
66 for (left_col, right_col) in left.column_metadatas.iter().zip(&right.column_metadatas) {
67 if left_col.column_id != right_col.column_id {
68 return false;
69 }
70 debug_assert_eq!(left_col.semantic_type, right_col.semantic_type);
71 }
72
73 &projection_mapper.input_arrow_schema(compaction) == read_format.arrow_schema()
74}
75
76pub(crate) struct FlatCompatBatch {
78 index_or_defaults: Vec<IndexOrDefault>,
80 arrow_schema: SchemaRef,
82 compat_pk: FlatCompatPrimaryKey,
84}
85
86impl FlatCompatBatch {
87 pub(crate) fn try_new(
93 mapper: &FlatProjectionMapper,
94 read_format: &FlatReadFormat,
95 compaction: bool,
96 ) -> Result<Option<Self>> {
97 let actual = read_format.metadata();
98 let actual_schema = flat_projected_columns(
99 actual,
100 read_format.format_projection(),
101 read_format.json_target_types(),
102 );
103
104 let expect_schema = mapper.batch_schema();
105 if expect_schema == actual_schema
106 && actual.primary_key == mapper.metadata().primary_key
107 && actual.primary_key_encoding == mapper.metadata().primary_key_encoding
108 {
109 return Ok(None);
112 }
113
114 if actual.primary_key_encoding == PrimaryKeyEncoding::Sparse && compaction {
115 return FlatCompatBatch::try_new_compact_sparse(mapper, actual);
117 }
118
119 let (index_or_defaults, fields) =
120 Self::compute_index_and_fields(&actual_schema, expect_schema, mapper.metadata())?;
121
122 let compat_pk = FlatCompatPrimaryKey::new(mapper.metadata(), actual)?;
123
124 Ok(Some(Self {
125 index_or_defaults,
126 arrow_schema: Arc::new(Schema::new(fields)),
127 compat_pk,
128 }))
129 }
130
131 fn compute_index_and_fields(
132 actual_schema: &[(ColumnId, ConcreteDataType)],
133 expect_schema: &[(ColumnId, ConcreteDataType)],
134 expect_metadata: &RegionMetadata,
135 ) -> Result<(Vec<IndexOrDefault>, Vec<FieldRef>)> {
136 let actual_schema_index: HashMap<_, _> = actual_schema
138 .iter()
139 .enumerate()
140 .map(|(idx, (column_id, data_type))| (*column_id, (idx, data_type)))
141 .collect();
142
143 let mut index_or_defaults = Vec::with_capacity(expect_schema.len());
144 let mut fields = Vec::with_capacity(expect_schema.len());
145 for (column_id, expect_data_type) in expect_schema {
146 let column_index = expect_metadata.column_index_by_id(*column_id).unwrap();
148 let expect_column = &expect_metadata.column_metadatas[column_index];
149 let column_field = &expect_metadata.schema.arrow_schema().fields()[column_index];
150 if expect_column.semantic_type == SemanticType::Tag {
152 let field = tag_maybe_to_dictionary_field(
153 &expect_column.column_schema.data_type,
154 column_field,
155 );
156 fields.push(Arc::new(with_field_id(
157 (*field).clone(),
158 expect_column.column_id,
159 )));
160 } else {
161 let field = with_field_id(
162 Arc::unwrap_or_clone(column_field.clone()),
163 expect_column.column_id,
164 )
165 .with_data_type(expect_data_type.as_arrow_type());
166 fields.push(Arc::new(field))
167 };
168
169 if let Some((index, actual_data_type)) = actual_schema_index.get(column_id) {
170 let mut cast_type = None;
171
172 if expect_data_type != *actual_data_type {
174 ensure!(
175 !expect_data_type.is_json2() && !actual_data_type.is_json2(),
176 CompatReaderSnafu {
177 region_id: expect_metadata.region_id,
178 reason: format!(
179 "JSON2 column '{}' must be aligned before FlatCompatBatch, actual: {}, expected: {}",
180 expect_column.column_schema.name,
181 actual_data_type,
182 expect_data_type,
183 ),
184 }
185 );
186 cast_type = Some(expect_data_type.clone());
187 }
188 index_or_defaults.push(IndexOrDefault::Index {
190 pos: *index,
191 cast_type,
192 });
193 } else {
194 let default_vector = expect_column
196 .column_schema
197 .create_default_vector(1)
198 .context(CreateDefaultSnafu {
199 region_id: expect_metadata.region_id,
200 column: &expect_column.column_schema.name,
201 })?
202 .with_context(|| CompatReaderSnafu {
203 region_id: expect_metadata.region_id,
204 reason: format!(
205 "column {} does not have a default value to read",
206 expect_column.column_schema.name
207 ),
208 })?;
209 index_or_defaults.push(IndexOrDefault::DefaultValue {
210 default_vector,
211 semantic_type: expect_column.semantic_type,
212 });
213 };
214 }
215 fields.extend_from_slice(&internal_fields());
216
217 Ok((index_or_defaults, fields))
218 }
219
220 fn try_new_compact_sparse(
221 mapper: &FlatProjectionMapper,
222 actual: &RegionMetadataRef,
223 ) -> Result<Option<Self>> {
224 ensure!(
227 mapper.metadata().primary_key_encoding == PrimaryKeyEncoding::Sparse,
228 UnsupportedOperationSnafu {
229 err_msg: "Flat format doesn't support converting sparse encoding back to dense encoding"
230 }
231 );
232
233 let actual_schema: Vec<_> = actual
236 .field_columns()
237 .chain([actual.time_index_column()])
238 .map(|col| (col.column_id, col.column_schema.data_type.clone()))
239 .collect();
240 let expect_schema: Vec<_> = mapper
241 .metadata()
242 .field_columns()
243 .chain([mapper.metadata().time_index_column()])
244 .map(|col| (col.column_id, col.column_schema.data_type.clone()))
245 .collect();
246
247 let (index_or_defaults, fields) =
248 Self::compute_index_and_fields(&actual_schema, &expect_schema, mapper.metadata())?;
249
250 let compat_pk = FlatCompatPrimaryKey::default();
251
252 Ok(Some(Self {
253 index_or_defaults,
254 arrow_schema: Arc::new(Schema::new(fields)),
255 compat_pk,
256 }))
257 }
258
259 pub(crate) fn compat(&self, batch: RecordBatch) -> Result<RecordBatch> {
261 let len = batch.num_rows();
262 let columns = self
263 .index_or_defaults
264 .iter()
265 .map(|index_or_default| match index_or_default {
266 IndexOrDefault::Index { pos, cast_type } => {
267 let old_column = batch.column(*pos);
268
269 if let Some(ty) = cast_type {
270 let casted =
271 datatypes::arrow::compute::cast(old_column, &ty.as_arrow_type())
272 .context(ComputeArrowSnafu)?;
273 Ok(casted)
274 } else {
275 Ok(old_column.clone())
276 }
277 }
278 IndexOrDefault::DefaultValue {
279 default_vector,
280 semantic_type,
281 } => repeat_vector(default_vector, len, *semantic_type == SemanticType::Tag),
282 })
283 .chain(
284 batch.columns()[batch.num_columns() - INTERNAL_COLUMN_NUM..]
286 .iter()
287 .map(|col| Ok(col.clone())),
288 )
289 .collect::<Result<Vec<_>>>()?;
290
291 let mut columns = columns;
292 let primary_key_index = primary_key_column_index(columns.len());
293 columns[primary_key_index] = self.compat_primary_key(&columns[primary_key_index])?;
294
295 RecordBatch::try_new(self.arrow_schema.clone(), columns).context(NewRecordBatchSnafu)
296 }
297
298 pub(crate) fn compat_primary_key(&self, primary_key: &ArrayRef) -> Result<ArrayRef> {
300 self.compat_pk.compat(primary_key)
301 }
302}
303
304fn repeat_vector(vector: &VectorRef, to_len: usize, is_tag: bool) -> Result<ArrayRef> {
306 assert_eq!(1, vector.len());
307 let data_type = vector.data_type();
308 if is_tag && data_type.is_string() {
309 let values = vector.to_arrow_array();
310 if values.is_null(0) {
311 let keys = UInt32Array::new_null(to_len);
313 Ok(Arc::new(DictionaryArray::new(keys, values.slice(0, 0))))
314 } else {
315 let keys = UInt32Array::from_value(0, to_len);
316 Ok(Arc::new(DictionaryArray::new(keys, values)))
317 }
318 } else {
319 let keys = UInt32Array::from_value(0, to_len);
320 take(
321 &vector.to_arrow_array(),
322 &keys,
323 Some(TakeOptions {
324 check_bounds: false,
325 }),
326 )
327 .context(ComputeArrowSnafu)
328 }
329}
330
331fn is_primary_key_same(expect: &RegionMetadata, actual: &RegionMetadata) -> Result<bool> {
333 ensure!(
334 actual.primary_key.len() <= expect.primary_key.len(),
335 CompatReaderSnafu {
336 region_id: expect.region_id,
337 reason: format!(
338 "primary key has more columns {} than expect {}",
339 actual.primary_key.len(),
340 expect.primary_key.len()
341 ),
342 }
343 );
344 ensure!(
345 actual.primary_key == expect.primary_key[..actual.primary_key.len()],
346 CompatReaderSnafu {
347 region_id: expect.region_id,
348 reason: format!(
349 "primary key has different prefix, expect: {:?}, actual: {:?}",
350 expect.primary_key, actual.primary_key
351 ),
352 }
353 );
354
355 Ok(actual.primary_key.len() == expect.primary_key.len())
356}
357
358#[derive(Debug)]
360enum IndexOrDefault {
361 Index {
363 pos: usize,
364 cast_type: Option<ConcreteDataType>,
365 },
366 DefaultValue {
368 default_vector: VectorRef,
370 semantic_type: SemanticType,
372 },
373}
374
375struct FlatRewritePrimaryKey {
377 codec: Arc<dyn PrimaryKeyCodec>,
379 metadata: RegionMetadataRef,
381 old_codec: Arc<dyn PrimaryKeyCodec>,
384}
385
386impl FlatRewritePrimaryKey {
387 fn new(
388 expect: &RegionMetadataRef,
389 actual: &RegionMetadataRef,
390 ) -> Option<FlatRewritePrimaryKey> {
391 if expect.primary_key_encoding == actual.primary_key_encoding {
392 return None;
393 }
394 let codec = build_primary_key_codec(expect);
395 let old_codec = build_primary_key_codec(actual);
396
397 Some(FlatRewritePrimaryKey {
398 codec,
399 metadata: expect.clone(),
400 old_codec,
401 })
402 }
403
404 fn rewrite_key(
407 &self,
408 append_values: &[(ColumnId, Value)],
409 primary_key: &ArrayRef,
410 ) -> Result<ArrayRef> {
411 if let Some(old_pk_dict_array) = primary_key.as_any().downcast_ref::<PrimaryKeyArray>() {
412 let old_pk_values_array = old_pk_dict_array
413 .values()
414 .as_any()
415 .downcast_ref::<BinaryArray>()
416 .context(UnexpectedSnafu {
417 reason: "Primary-key dictionary values are not binary",
418 })?;
419 let new_pk_values_array =
420 Arc::new(self.rewrite_values(append_values, old_pk_values_array)?);
421 return Ok(Arc::new(PrimaryKeyArray::new(
422 old_pk_dict_array.keys().clone(),
423 new_pk_values_array,
424 )));
425 }
426
427 let old_pk_values_array =
428 primary_key
429 .as_any()
430 .downcast_ref::<BinaryArray>()
431 .context(UnexpectedSnafu {
432 reason: format!(
433 "Primary-key column is neither binary nor dictionary, got {:?}",
434 primary_key.data_type()
435 ),
436 })?;
437 Ok(Arc::new(
438 self.rewrite_values(append_values, old_pk_values_array)?,
439 ))
440 }
441
442 fn rewrite_values(
443 &self,
444 append_values: &[(ColumnId, Value)],
445 old_pk_values_array: &BinaryArray,
446 ) -> Result<BinaryArray> {
447 let mut builder = BinaryBuilder::with_capacity(
448 old_pk_values_array.len(),
449 old_pk_values_array.value_data().len(),
450 );
451
452 let mut buffer = Vec::with_capacity(
454 old_pk_values_array.value_data().len() / old_pk_values_array.len().max(1),
455 );
456 let mut column_id_values = Vec::new();
457 for value in old_pk_values_array.iter() {
459 let Some(old_pk) = value else {
460 builder.append_null();
461 continue;
462 };
463 let mut pk_values = self.old_codec.decode(old_pk).context(DecodeSnafu)?;
465 pk_values.extend(append_values);
466
467 buffer.clear();
468 column_id_values.clear();
469 match pk_values {
471 CompositeValues::Dense(dense_values) => {
472 self.codec
473 .encode_values(dense_values.as_slice(), &mut buffer)
474 .context(EncodeSnafu)?;
475 }
476 CompositeValues::Sparse(sparse_values) => {
477 for id in &self.metadata.primary_key {
478 let value = sparse_values.get_or_null(*id);
479 column_id_values.push((*id, value.clone()));
480 }
481 self.codec
482 .encode_values(&column_id_values, &mut buffer)
483 .context(EncodeSnafu)?;
484 }
485 }
486 builder.append_value(&buffer);
487 }
488 Ok(builder.finish())
489 }
490}
491
492#[derive(Default)]
494struct FlatCompatPrimaryKey {
495 rewriter: Option<FlatRewritePrimaryKey>,
497 converter: Option<Arc<dyn PrimaryKeyCodec>>,
499 values: Vec<(ColumnId, Value)>,
501}
502
503impl FlatCompatPrimaryKey {
504 fn new(expect: &RegionMetadataRef, actual: &RegionMetadataRef) -> Result<Self> {
505 let rewriter = FlatRewritePrimaryKey::new(expect, actual);
506
507 if is_primary_key_same(expect, actual)? {
508 return Ok(Self {
509 rewriter,
510 converter: None,
511 values: Vec::new(),
512 });
513 }
514
515 let to_add = &expect.primary_key[actual.primary_key.len()..];
517 let mut values = Vec::with_capacity(to_add.len());
518 let mut fields = Vec::with_capacity(to_add.len());
519 for column_id in to_add {
520 let column = expect.column_by_id(*column_id).unwrap();
522 fields.push((
523 *column_id,
524 SortField::new(column.column_schema.data_type.clone()),
525 ));
526 let default_value = column
527 .column_schema
528 .create_default()
529 .context(CreateDefaultSnafu {
530 region_id: expect.region_id,
531 column: &column.column_schema.name,
532 })?
533 .with_context(|| CompatReaderSnafu {
534 region_id: expect.region_id,
535 reason: format!(
536 "key column {} does not have a default value to read",
537 column.column_schema.name
538 ),
539 })?;
540 values.push((*column_id, default_value));
541 }
542 debug_assert!(!fields.is_empty());
544
545 let converter = Some(build_primary_key_codec_with_fields(
547 expect.primary_key_encoding,
548 fields.into_iter(),
549 ));
550
551 Ok(Self {
552 rewriter,
553 converter,
554 values,
555 })
556 }
557
558 fn compat(&self, primary_key: &ArrayRef) -> Result<ArrayRef> {
560 if let Some(rewriter) = &self.rewriter {
561 return rewriter.rewrite_key(&self.values, primary_key);
563 }
564
565 self.append_key(primary_key)
566 }
567
568 fn append_key(&self, primary_key: &ArrayRef) -> Result<ArrayRef> {
570 let Some(converter) = &self.converter else {
571 return Ok(primary_key.clone());
572 };
573
574 if let Some(old_pk_dict_array) = primary_key.as_any().downcast_ref::<PrimaryKeyArray>() {
575 let old_pk_values_array = old_pk_dict_array
576 .values()
577 .as_any()
578 .downcast_ref::<BinaryArray>()
579 .context(UnexpectedSnafu {
580 reason: "Primary-key dictionary values are not binary",
581 })?;
582 let new_pk_values_array =
583 Arc::new(self.append_values(old_pk_values_array, converter.as_ref())?);
584 return Ok(Arc::new(PrimaryKeyArray::new(
585 old_pk_dict_array.keys().clone(),
586 new_pk_values_array,
587 )));
588 }
589
590 let old_pk_values_array =
591 primary_key
592 .as_any()
593 .downcast_ref::<BinaryArray>()
594 .context(UnexpectedSnafu {
595 reason: format!(
596 "Primary-key column is neither binary nor dictionary, got {:?}",
597 primary_key.data_type()
598 ),
599 })?;
600 Ok(Arc::new(
601 self.append_values(old_pk_values_array, converter.as_ref())?,
602 ))
603 }
604
605 fn append_values(
606 &self,
607 old_pk_values_array: &BinaryArray,
608 converter: &dyn PrimaryKeyCodec,
609 ) -> Result<BinaryArray> {
610 let mut builder = BinaryBuilder::with_capacity(
611 old_pk_values_array.len(),
612 old_pk_values_array.value_data().len()
613 + converter.estimated_size().unwrap_or_default() * old_pk_values_array.len(),
614 );
615
616 let mut buffer = Vec::with_capacity(
618 old_pk_values_array.value_data().len() / old_pk_values_array.len().max(1)
619 + converter.estimated_size().unwrap_or_default(),
620 );
621
622 for value in old_pk_values_array.iter() {
624 let Some(old_pk) = value else {
625 builder.append_null();
626 continue;
627 };
628
629 buffer.clear();
630 buffer.extend_from_slice(old_pk);
631 converter
632 .encode_values(&self.values, &mut buffer)
633 .context(EncodeSnafu)?;
634
635 builder.append_value(&buffer);
636 }
637
638 Ok(builder.finish())
639 }
640}
641
642#[cfg(test)]
643mod tests {
644 use std::collections::BTreeMap;
645 use std::sync::Arc;
646
647 use api::v1::{OpType, SemanticType};
648 use datatypes::arrow::array::{
649 ArrayRef, BinaryArray, BinaryDictionaryBuilder, Int64Array, StringDictionaryBuilder,
650 TimestampMicrosecondArray, TimestampMillisecondArray, UInt8Array, UInt64Array,
651 };
652 use datatypes::arrow::datatypes::UInt32Type;
653 use datatypes::arrow::record_batch::RecordBatch;
654 use datatypes::prelude::ConcreteDataType;
655 use datatypes::schema::ColumnSchema;
656 use datatypes::types::json_type::JsonNativeType;
657 use datatypes::value::ValueRef;
658 use mito_codec::row_converter::{
659 DensePrimaryKeyCodec, PrimaryKeyCodecExt, SparsePrimaryKeyCodec,
660 };
661 use store_api::codec::PrimaryKeyEncoding;
662 use store_api::metadata::{ColumnMetadata, RegionMetadataBuilder};
663 use store_api::storage::RegionId;
664
665 use super::*;
666 use crate::read::flat_projection::FlatProjectionMapper;
667 use crate::read::read_columns::ReadColumns;
668 use crate::sst::parquet::flat_format::FlatReadFormat;
669 use crate::sst::{FlatSchemaOptions, to_flat_sst_arrow_schema};
670
671 fn new_metadata(
673 semantic_types: &[(ColumnId, SemanticType, ConcreteDataType)],
674 primary_key: &[ColumnId],
675 ) -> RegionMetadata {
676 let mut builder = RegionMetadataBuilder::new(RegionId::new(1, 1));
677 for (id, semantic_type, data_type) in semantic_types {
678 let column_schema = match semantic_type {
679 SemanticType::Tag => {
680 ColumnSchema::new(format!("tag_{id}"), data_type.clone(), true)
681 }
682 SemanticType::Field => {
683 ColumnSchema::new(format!("field_{id}"), data_type.clone(), true)
684 }
685 SemanticType::Timestamp => ColumnSchema::new("ts", data_type.clone(), false),
686 };
687
688 builder.push_column_metadata(ColumnMetadata {
689 column_schema,
690 semantic_type: *semantic_type,
691 column_id: *id,
692 });
693 }
694 builder.primary_key(primary_key.to_vec());
695 builder.build().unwrap()
696 }
697
698 fn encode_key(keys: &[Option<&str>]) -> Vec<u8> {
700 let fields = (0..keys.len())
701 .map(|_| (0, SortField::new(ConcreteDataType::string_datatype())))
702 .collect();
703 let converter = DensePrimaryKeyCodec::with_fields(fields);
704 let row = keys.iter().map(|str_opt| match str_opt {
705 Some(v) => ValueRef::String(v),
706 None => ValueRef::Null,
707 });
708
709 converter.encode(row).unwrap()
710 }
711
712 fn encode_sparse_key(keys: &[(ColumnId, Option<&str>)]) -> Vec<u8> {
714 let fields = (0..keys.len())
715 .map(|_| (1, SortField::new(ConcreteDataType::string_datatype())))
716 .collect();
717 let converter = SparsePrimaryKeyCodec::with_fields(fields);
718 let row = keys
719 .iter()
720 .map(|(id, str_opt)| match str_opt {
721 Some(v) => (*id, ValueRef::String(v)),
722 None => (*id, ValueRef::Null),
723 })
724 .collect::<Vec<_>>();
725 let mut buffer = vec![];
726 converter.encode_value_refs(&row, &mut buffer).unwrap();
727 buffer
728 }
729
730 fn build_flat_test_pk_array(primary_keys: &[&[u8]]) -> ArrayRef {
732 let mut builder = BinaryDictionaryBuilder::<UInt32Type>::new();
733 for &pk in primary_keys {
734 builder.append(pk).unwrap();
735 }
736 Arc::new(builder.finish())
737 }
738
739 #[test]
740 fn test_flat_compat_batch_with_missing_columns() {
741 let actual_metadata = Arc::new(new_metadata(
742 &[
743 (
744 0,
745 SemanticType::Timestamp,
746 ConcreteDataType::timestamp_millisecond_datatype(),
747 ),
748 (1, SemanticType::Tag, ConcreteDataType::string_datatype()),
749 (2, SemanticType::Field, ConcreteDataType::int64_datatype()),
750 ],
751 &[1],
752 ));
753
754 let expected_metadata = Arc::new(new_metadata(
755 &[
756 (
757 0,
758 SemanticType::Timestamp,
759 ConcreteDataType::timestamp_millisecond_datatype(),
760 ),
761 (1, SemanticType::Tag, ConcreteDataType::string_datatype()),
762 (2, SemanticType::Field, ConcreteDataType::int64_datatype()),
763 (3, SemanticType::Field, ConcreteDataType::int64_datatype()),
765 ],
766 &[1],
767 ));
768
769 let mapper = FlatProjectionMapper::all(&expected_metadata).unwrap();
770 let read_format = FlatReadFormat::new(
771 actual_metadata.clone(),
772 ReadColumns::new([0, 1, 2, 3]),
773 None,
774 "test",
775 false,
776 )
777 .unwrap();
778
779 let compat_batch = FlatCompatBatch::try_new(&mapper, &read_format, false)
780 .unwrap()
781 .unwrap();
782
783 let mut tag_builder = StringDictionaryBuilder::<UInt32Type>::new();
784 tag_builder.append_value("tag1");
785 tag_builder.append_value("tag1");
786 let tag_dict_array = Arc::new(tag_builder.finish());
787
788 let k1 = encode_key(&[Some("tag1")]);
789 let input_columns: Vec<ArrayRef> = vec![
790 tag_dict_array.clone(),
791 Arc::new(Int64Array::from(vec![100, 200])),
792 Arc::new(TimestampMillisecondArray::from_iter_values([1000, 2000])),
793 build_flat_test_pk_array(&[&k1, &k1]),
794 Arc::new(UInt64Array::from_iter_values([1, 2])),
795 Arc::new(UInt8Array::from_iter_values([
796 OpType::Put as u8,
797 OpType::Put as u8,
798 ])),
799 ];
800 let input_schema =
801 to_flat_sst_arrow_schema(&actual_metadata, &FlatSchemaOptions::default());
802 let input_batch = RecordBatch::try_new(input_schema, input_columns).unwrap();
803
804 let result = compat_batch.compat(input_batch).unwrap();
805
806 let expected_schema =
807 to_flat_sst_arrow_schema(&expected_metadata, &FlatSchemaOptions::default());
808
809 let expected_columns: Vec<ArrayRef> = vec![
810 tag_dict_array.clone(),
811 Arc::new(Int64Array::from(vec![100, 200])),
812 Arc::new(Int64Array::from(vec![None::<i64>, None::<i64>])),
813 Arc::new(TimestampMillisecondArray::from_iter_values([1000, 2000])),
814 build_flat_test_pk_array(&[&k1, &k1]),
815 Arc::new(UInt64Array::from_iter_values([1, 2])),
816 Arc::new(UInt8Array::from_iter_values([
817 OpType::Put as u8,
818 OpType::Put as u8,
819 ])),
820 ];
821 let expected_batch = RecordBatch::try_new(expected_schema, expected_columns).unwrap();
822
823 assert_eq!(expected_batch, result);
824 }
825
826 #[test]
827 fn test_flat_compat_batch_uses_projected_json2_type() -> Result<()> {
828 let json2 = ConcreteDataType::json2(JsonNativeType::object());
829 let actual_metadata = Arc::new(new_metadata(
830 &[
831 (
832 0,
833 SemanticType::Timestamp,
834 ConcreteDataType::timestamp_millisecond_datatype(),
835 ),
836 (1, SemanticType::Field, json2.clone()),
837 ],
838 &[],
839 ));
840 let expected_metadata = Arc::new(new_metadata(
841 &[
842 (
843 0,
844 SemanticType::Timestamp,
845 ConcreteDataType::timestamp_millisecond_datatype(),
846 ),
847 (1, SemanticType::Field, json2),
848 (2, SemanticType::Field, ConcreteDataType::int64_datatype()),
849 ],
850 &[],
851 ));
852 let read_columns = ReadColumns::new([0, 1, 2])
853 .with_json_target_types(BTreeMap::from([(1, JsonNativeType::Variant)]));
854 let mapper = FlatProjectionMapper::new_with_read_columns(
855 &expected_metadata,
856 vec![0, 1, 2],
857 read_columns.clone(),
858 )?;
859 let read_format = FlatReadFormat::new(actual_metadata, read_columns, None, "test", false)?;
860
861 let compat = FlatCompatBatch::try_new(&mapper, &read_format, false)?.unwrap();
862 let json_index = mapper
863 .batch_schema()
864 .iter()
865 .position(|(id, _)| *id == 1)
866 .unwrap();
867 assert!(matches!(
868 &compat.index_or_defaults[json_index],
869 IndexOrDefault::Index {
870 cast_type: None,
871 ..
872 }
873 ));
874 Ok(())
875 }
876
877 #[test]
878 fn test_flat_compat_batch_with_read_projection_superset() {
879 let actual_metadata = Arc::new(new_metadata(
880 &[
881 (
882 0,
883 SemanticType::Timestamp,
884 ConcreteDataType::timestamp_millisecond_datatype(),
885 ),
886 (1, SemanticType::Tag, ConcreteDataType::string_datatype()),
887 (2, SemanticType::Field, ConcreteDataType::int64_datatype()),
888 ],
889 &[1],
890 ));
891
892 let expected_metadata = Arc::new(new_metadata(
893 &[
894 (
895 0,
896 SemanticType::Timestamp,
897 ConcreteDataType::timestamp_millisecond_datatype(),
898 ),
899 (1, SemanticType::Tag, ConcreteDataType::string_datatype()),
900 (2, SemanticType::Field, ConcreteDataType::int64_datatype()),
901 (3, SemanticType::Field, ConcreteDataType::int64_datatype()),
903 ],
904 &[1],
905 ));
906
907 let mapper = FlatProjectionMapper::new_with_read_columns(
908 &expected_metadata,
909 vec![1, 2],
910 ReadColumns::new([1, 2, 3]),
911 )
912 .unwrap();
913 let read_format = FlatReadFormat::new(
914 actual_metadata.clone(),
915 ReadColumns::new([1, 2, 3]),
916 None,
917 "test",
918 false,
919 )
920 .unwrap();
921
922 let compat_batch = FlatCompatBatch::try_new(&mapper, &read_format, false)
923 .unwrap()
924 .unwrap();
925
926 let mut tag_builder = StringDictionaryBuilder::<UInt32Type>::new();
927 tag_builder.append_value("tag1");
928 tag_builder.append_value("tag1");
929 let tag_dict_array = Arc::new(tag_builder.finish());
930
931 let k1 = encode_key(&[Some("tag1")]);
932 let input_columns: Vec<ArrayRef> = vec![
933 tag_dict_array.clone(),
934 Arc::new(Int64Array::from(vec![100, 200])),
935 Arc::new(TimestampMillisecondArray::from_iter_values([1000, 2000])),
936 build_flat_test_pk_array(&[&k1, &k1]),
937 Arc::new(UInt64Array::from_iter_values([1, 2])),
938 Arc::new(UInt8Array::from_iter_values([
939 OpType::Put as u8,
940 OpType::Put as u8,
941 ])),
942 ];
943 let input_schema =
944 to_flat_sst_arrow_schema(&actual_metadata, &FlatSchemaOptions::default());
945 let input_batch = RecordBatch::try_new(input_schema, input_columns).unwrap();
946
947 let result = compat_batch.compat(input_batch).unwrap();
948
949 let expected_schema =
950 to_flat_sst_arrow_schema(&expected_metadata, &FlatSchemaOptions::default());
951 let expected_columns: Vec<ArrayRef> = vec![
952 tag_dict_array.clone(),
953 Arc::new(Int64Array::from(vec![100, 200])),
954 Arc::new(Int64Array::from(vec![None::<i64>, None::<i64>])),
955 Arc::new(TimestampMillisecondArray::from_iter_values([1000, 2000])),
956 build_flat_test_pk_array(&[&k1, &k1]),
957 Arc::new(UInt64Array::from_iter_values([1, 2])),
958 Arc::new(UInt8Array::from_iter_values([
959 OpType::Put as u8,
960 OpType::Put as u8,
961 ])),
962 ];
963 let expected_batch = RecordBatch::try_new(expected_schema, expected_columns).unwrap();
964
965 assert_eq!(expected_batch, result);
966 }
967
968 #[test]
969 fn test_flat_compat_batch_with_different_pk_encoding() {
970 let mut actual_metadata = new_metadata(
971 &[
972 (
973 0,
974 SemanticType::Timestamp,
975 ConcreteDataType::timestamp_millisecond_datatype(),
976 ),
977 (1, SemanticType::Tag, ConcreteDataType::string_datatype()),
978 (2, SemanticType::Field, ConcreteDataType::int64_datatype()),
979 ],
980 &[1],
981 );
982 actual_metadata.primary_key_encoding = PrimaryKeyEncoding::Dense;
983 let actual_metadata = Arc::new(actual_metadata);
984
985 let mut expected_metadata = new_metadata(
986 &[
987 (
988 0,
989 SemanticType::Timestamp,
990 ConcreteDataType::timestamp_millisecond_datatype(),
991 ),
992 (1, SemanticType::Tag, ConcreteDataType::string_datatype()),
993 (2, SemanticType::Field, ConcreteDataType::int64_datatype()),
994 (3, SemanticType::Tag, ConcreteDataType::string_datatype()),
995 ],
996 &[1, 3],
997 );
998 expected_metadata.primary_key_encoding = PrimaryKeyEncoding::Sparse;
999 let expected_metadata = Arc::new(expected_metadata);
1000
1001 let mapper = FlatProjectionMapper::all(&expected_metadata).unwrap();
1002 let read_format = FlatReadFormat::new(
1003 actual_metadata.clone(),
1004 ReadColumns::new([0, 1, 2, 3]),
1005 None,
1006 "test",
1007 false,
1008 )
1009 .unwrap();
1010
1011 let compat_batch = FlatCompatBatch::try_new(&mapper, &read_format, false)
1012 .unwrap()
1013 .unwrap();
1014
1015 let mut tag1_builder = StringDictionaryBuilder::<UInt32Type>::new();
1017 tag1_builder.append_value("tag1");
1018 tag1_builder.append_value("tag1");
1019 let tag1_dict_array = Arc::new(tag1_builder.finish());
1020
1021 let k1 = encode_key(&[Some("tag1")]);
1022 let input_columns: Vec<ArrayRef> = vec![
1023 tag1_dict_array.clone(),
1024 Arc::new(Int64Array::from(vec![100, 200])),
1025 Arc::new(TimestampMillisecondArray::from_iter_values([1000, 2000])),
1026 build_flat_test_pk_array(&[&k1, &k1]),
1027 Arc::new(UInt64Array::from_iter_values([1, 2])),
1028 Arc::new(UInt8Array::from_iter_values([
1029 OpType::Put as u8,
1030 OpType::Put as u8,
1031 ])),
1032 ];
1033 let input_schema =
1034 to_flat_sst_arrow_schema(&actual_metadata, &FlatSchemaOptions::default());
1035 let input_batch = RecordBatch::try_new(input_schema, input_columns).unwrap();
1036
1037 let result = compat_batch.compat(input_batch).unwrap();
1038
1039 let sparse_k1 = encode_sparse_key(&[(1, Some("tag1")), (3, None)]);
1040 let mut null_tag_builder = StringDictionaryBuilder::<UInt32Type>::new();
1041 null_tag_builder.append_nulls(2);
1042 let null_tag_dict_array = Arc::new(null_tag_builder.finish());
1043 let expected_columns: Vec<ArrayRef> = vec![
1044 tag1_dict_array.clone(),
1045 null_tag_dict_array,
1046 Arc::new(Int64Array::from(vec![100, 200])),
1047 Arc::new(TimestampMillisecondArray::from_iter_values([1000, 2000])),
1048 build_flat_test_pk_array(&[&sparse_k1, &sparse_k1]),
1049 Arc::new(UInt64Array::from_iter_values([1, 2])),
1050 Arc::new(UInt8Array::from_iter_values([
1051 OpType::Put as u8,
1052 OpType::Put as u8,
1053 ])),
1054 ];
1055 let output_schema =
1056 to_flat_sst_arrow_schema(&expected_metadata, &FlatSchemaOptions::default());
1057 let expected_batch = RecordBatch::try_new(output_schema, expected_columns).unwrap();
1058
1059 assert_eq!(expected_batch, result);
1060 }
1061
1062 #[test]
1063 fn test_compat_primary_key_with_different_encoding_only() {
1064 let mut actual_metadata = new_metadata(
1065 &[
1066 (
1067 0,
1068 SemanticType::Timestamp,
1069 ConcreteDataType::timestamp_millisecond_datatype(),
1070 ),
1071 (1, SemanticType::Tag, ConcreteDataType::string_datatype()),
1072 (2, SemanticType::Field, ConcreteDataType::int64_datatype()),
1073 ],
1074 &[1],
1075 );
1076 actual_metadata.primary_key_encoding = PrimaryKeyEncoding::Dense;
1077 let actual_metadata = Arc::new(actual_metadata);
1078
1079 let mut expected_metadata = (*actual_metadata).clone();
1080 expected_metadata.primary_key_encoding = PrimaryKeyEncoding::Sparse;
1081 let expected_metadata = Arc::new(expected_metadata);
1082
1083 let mapper = FlatProjectionMapper::all(&expected_metadata).unwrap();
1084 let read_format = FlatReadFormat::new(
1085 actual_metadata,
1086 ReadColumns::new([0, 1, 2]),
1087 None,
1088 "test",
1089 false,
1090 )
1091 .unwrap();
1092 let compat = FlatCompatBatch::try_new(&mapper, &read_format, false)
1093 .unwrap()
1094 .unwrap();
1095
1096 let dense_key = encode_key(&[Some("tag1")]);
1097 let sparse_key = encode_sparse_key(&[(1, Some("tag1"))]);
1098
1099 let dictionary_key = build_flat_test_pk_array(&[&dense_key, &dense_key]);
1100 let result = compat.compat_primary_key(&dictionary_key).unwrap();
1101 let result = result.as_any().downcast_ref::<PrimaryKeyArray>().unwrap();
1102 let values = result
1103 .values()
1104 .as_any()
1105 .downcast_ref::<BinaryArray>()
1106 .unwrap();
1107 assert_eq!(values.value(result.keys().value(0) as usize), sparse_key);
1108 assert_eq!(values.value(result.keys().value(1) as usize), sparse_key);
1109
1110 let binary_key: ArrayRef = Arc::new(BinaryArray::from(vec![
1111 Some(dense_key.as_slice()),
1112 Some(dense_key.as_slice()),
1113 ]));
1114 let result = compat.compat_primary_key(&binary_key).unwrap();
1115 let result = result.as_any().downcast_ref::<BinaryArray>().unwrap();
1116 assert_eq!(result.value(0), sparse_key);
1117 assert_eq!(result.value(1), sparse_key);
1118 }
1119
1120 #[test]
1123 fn test_flat_compat_batch_with_time_index_unit_change() {
1124 let actual_metadata = Arc::new(new_metadata(
1125 &[
1126 (
1127 0,
1128 SemanticType::Timestamp,
1129 ConcreteDataType::timestamp_millisecond_datatype(),
1130 ),
1131 (1, SemanticType::Tag, ConcreteDataType::string_datatype()),
1132 (2, SemanticType::Field, ConcreteDataType::int64_datatype()),
1133 ],
1134 &[1],
1135 ));
1136
1137 let mut expected_metadata = (*actual_metadata).clone();
1138 for column in expected_metadata.column_metadatas.iter_mut() {
1139 if column.semantic_type == SemanticType::Timestamp {
1140 column.column_schema.data_type = ConcreteDataType::timestamp_microsecond_datatype();
1141 }
1142 }
1143 let expected_metadata = Arc::new(
1145 RegionMetadataBuilder::from_existing(expected_metadata)
1146 .build()
1147 .unwrap(),
1148 );
1149
1150 let same_unit_metadata = Arc::new((*actual_metadata).clone());
1152 let mapper = FlatProjectionMapper::all(&same_unit_metadata).unwrap();
1153 let read_format = FlatReadFormat::new(
1154 actual_metadata.clone(),
1155 ReadColumns::new([0, 1, 2]),
1156 None,
1157 "test",
1158 false,
1159 )
1160 .unwrap();
1161 assert!(
1162 FlatCompatBatch::try_new(&mapper, &read_format, false)
1163 .unwrap()
1164 .is_none()
1165 );
1166
1167 let mapper = FlatProjectionMapper::all(&expected_metadata).unwrap();
1169 let read_format = FlatReadFormat::new(
1170 actual_metadata.clone(),
1171 ReadColumns::new([0, 1, 2]),
1172 None,
1173 "test",
1174 false,
1175 )
1176 .unwrap();
1177 let compat_batch = FlatCompatBatch::try_new(&mapper, &read_format, false)
1178 .unwrap()
1179 .unwrap();
1180
1181 let mut tag_builder = StringDictionaryBuilder::<UInt32Type>::new();
1182 tag_builder.append_value("tag1");
1183 tag_builder.append_value("tag1");
1184 let tag_dict_array = Arc::new(tag_builder.finish());
1185
1186 let k1 = encode_key(&[Some("tag1")]);
1187 let input_columns: Vec<ArrayRef> = vec![
1188 tag_dict_array.clone(),
1189 Arc::new(Int64Array::from(vec![100, 200])),
1190 Arc::new(TimestampMillisecondArray::from_iter_values([1000, 2000])),
1191 build_flat_test_pk_array(&[&k1, &k1]),
1192 Arc::new(UInt64Array::from_iter_values([1, 2])),
1193 Arc::new(UInt8Array::from_iter_values([
1194 OpType::Put as u8,
1195 OpType::Put as u8,
1196 ])),
1197 ];
1198 let input_schema =
1199 to_flat_sst_arrow_schema(&actual_metadata, &FlatSchemaOptions::default());
1200 let input_batch = RecordBatch::try_new(input_schema, input_columns).unwrap();
1201
1202 let result = compat_batch.compat(input_batch).unwrap();
1203
1204 let expected_columns: Vec<ArrayRef> = vec![
1205 tag_dict_array.clone(),
1206 Arc::new(Int64Array::from(vec![100, 200])),
1207 Arc::new(TimestampMicrosecondArray::from_iter_values([
1209 1_000_000, 2_000_000,
1210 ])),
1211 build_flat_test_pk_array(&[&k1, &k1]),
1212 Arc::new(UInt64Array::from_iter_values([1, 2])),
1213 Arc::new(UInt8Array::from_iter_values([
1214 OpType::Put as u8,
1215 OpType::Put as u8,
1216 ])),
1217 ];
1218 let expected_schema =
1219 to_flat_sst_arrow_schema(&expected_metadata, &FlatSchemaOptions::default());
1220 let expected_batch = RecordBatch::try_new(expected_schema, expected_columns).unwrap();
1221
1222 assert_eq!(expected_batch, result);
1223 }
1224
1225 #[test]
1228 fn test_flat_compat_batch_time_index_unit_change_with_added_column() {
1229 let actual_metadata = Arc::new(new_metadata(
1230 &[
1231 (
1232 0,
1233 SemanticType::Timestamp,
1234 ConcreteDataType::timestamp_millisecond_datatype(),
1235 ),
1236 (1, SemanticType::Tag, ConcreteDataType::string_datatype()),
1237 (2, SemanticType::Field, ConcreteDataType::int64_datatype()),
1238 ],
1239 &[1],
1240 ));
1241
1242 let mut expected_metadata = new_metadata(
1243 &[
1244 (
1245 0,
1246 SemanticType::Timestamp,
1247 ConcreteDataType::timestamp_microsecond_datatype(),
1248 ),
1249 (1, SemanticType::Tag, ConcreteDataType::string_datatype()),
1250 (2, SemanticType::Field, ConcreteDataType::int64_datatype()),
1251 (3, SemanticType::Field, ConcreteDataType::int64_datatype()),
1252 ],
1253 &[1],
1254 );
1255 expected_metadata.primary_key_encoding = PrimaryKeyEncoding::Dense;
1256 let expected_metadata = Arc::new(expected_metadata);
1257
1258 let mapper = FlatProjectionMapper::all(&expected_metadata).unwrap();
1259 let read_format = FlatReadFormat::new(
1260 actual_metadata.clone(),
1261 ReadColumns::new([0, 1, 2, 3]),
1262 None,
1263 "test",
1264 false,
1265 )
1266 .unwrap();
1267
1268 let compat_batch = FlatCompatBatch::try_new(&mapper, &read_format, false)
1269 .unwrap()
1270 .unwrap();
1271
1272 let mut tag_builder = StringDictionaryBuilder::<UInt32Type>::new();
1273 tag_builder.append_value("tag1");
1274 tag_builder.append_value("tag1");
1275 let tag_dict_array = Arc::new(tag_builder.finish());
1276
1277 let k1 = encode_key(&[Some("tag1")]);
1278 let input_columns: Vec<ArrayRef> = vec![
1279 tag_dict_array.clone(),
1280 Arc::new(Int64Array::from(vec![100, 200])),
1281 Arc::new(TimestampMillisecondArray::from_iter_values([1000, 2000])),
1282 build_flat_test_pk_array(&[&k1, &k1]),
1283 Arc::new(UInt64Array::from_iter_values([1, 2])),
1284 Arc::new(UInt8Array::from_iter_values([
1285 OpType::Put as u8,
1286 OpType::Put as u8,
1287 ])),
1288 ];
1289 let input_schema =
1290 to_flat_sst_arrow_schema(&actual_metadata, &FlatSchemaOptions::default());
1291 let input_batch = RecordBatch::try_new(input_schema, input_columns).unwrap();
1292
1293 let result = compat_batch.compat(input_batch).unwrap();
1294
1295 let expected_columns: Vec<ArrayRef> = vec![
1296 tag_dict_array.clone(),
1297 Arc::new(Int64Array::from(vec![100, 200])),
1298 Arc::new(Int64Array::from(vec![None::<i64>, None::<i64>])),
1299 Arc::new(TimestampMicrosecondArray::from_iter_values([
1300 1_000_000, 2_000_000,
1301 ])),
1302 build_flat_test_pk_array(&[&k1, &k1]),
1303 Arc::new(UInt64Array::from_iter_values([1, 2])),
1304 Arc::new(UInt8Array::from_iter_values([
1305 OpType::Put as u8,
1306 OpType::Put as u8,
1307 ])),
1308 ];
1309 let expected_schema =
1310 to_flat_sst_arrow_schema(&expected_metadata, &FlatSchemaOptions::default());
1311 let expected_batch = RecordBatch::try_new(expected_schema, expected_columns).unwrap();
1312
1313 assert_eq!(expected_batch, result);
1314 }
1315
1316 #[test]
1319 fn test_flat_compat_batch_compact_sparse_time_index_unit_change() {
1320 let mut actual_metadata = new_metadata(
1321 &[
1322 (
1323 0,
1324 SemanticType::Timestamp,
1325 ConcreteDataType::timestamp_millisecond_datatype(),
1326 ),
1327 (2, SemanticType::Field, ConcreteDataType::int64_datatype()),
1328 ],
1329 &[],
1330 );
1331 actual_metadata.primary_key_encoding = PrimaryKeyEncoding::Sparse;
1332 let actual_metadata = Arc::new(actual_metadata);
1333
1334 let mut expected_metadata = (*actual_metadata).clone();
1335 for column in expected_metadata.column_metadatas.iter_mut() {
1336 if column.semantic_type == SemanticType::Timestamp {
1337 column.column_schema.data_type = ConcreteDataType::timestamp_microsecond_datatype();
1338 }
1339 }
1340 let expected_metadata = Arc::new(
1341 RegionMetadataBuilder::from_existing(expected_metadata)
1342 .build()
1343 .unwrap(),
1344 );
1345
1346 let mapper = FlatProjectionMapper::all(&expected_metadata).unwrap();
1347 let read_format = FlatReadFormat::new(
1348 actual_metadata.clone(),
1349 ReadColumns::new([0, 2]),
1350 None,
1351 "test",
1352 true,
1353 )
1354 .unwrap();
1355
1356 let compat_batch = FlatCompatBatch::try_new(&mapper, &read_format, true)
1357 .unwrap()
1358 .unwrap();
1359
1360 let sparse_k1 = encode_sparse_key(&[]);
1361 let input_columns: Vec<ArrayRef> = vec![
1362 Arc::new(Int64Array::from(vec![100, 200])),
1363 Arc::new(TimestampMillisecondArray::from_iter_values([1000, 2000])),
1364 build_flat_test_pk_array(&[&sparse_k1, &sparse_k1]),
1365 Arc::new(UInt64Array::from_iter_values([1, 2])),
1366 Arc::new(UInt8Array::from_iter_values([
1367 OpType::Put as u8,
1368 OpType::Put as u8,
1369 ])),
1370 ];
1371 let input_schema =
1372 to_flat_sst_arrow_schema(&actual_metadata, &FlatSchemaOptions::default());
1373 let input_batch = RecordBatch::try_new(input_schema, input_columns).unwrap();
1374
1375 let result = compat_batch.compat(input_batch).unwrap();
1376
1377 let expected_columns: Vec<ArrayRef> = vec![
1378 Arc::new(Int64Array::from(vec![100, 200])),
1379 Arc::new(TimestampMicrosecondArray::from_iter_values([
1380 1_000_000, 2_000_000,
1381 ])),
1382 build_flat_test_pk_array(&[&sparse_k1, &sparse_k1]),
1383 Arc::new(UInt64Array::from_iter_values([1, 2])),
1384 Arc::new(UInt8Array::from_iter_values([
1385 OpType::Put as u8,
1386 OpType::Put as u8,
1387 ])),
1388 ];
1389 let output_schema =
1390 to_flat_sst_arrow_schema(&expected_metadata, &FlatSchemaOptions::default());
1391 let expected_batch = RecordBatch::try_new(output_schema, expected_columns).unwrap();
1392
1393 assert_eq!(expected_batch, result);
1394 }
1395
1396 #[test]
1397 fn test_flat_compat_batch_compact_sparse() {
1398 let mut actual_metadata = new_metadata(
1399 &[
1400 (
1401 0,
1402 SemanticType::Timestamp,
1403 ConcreteDataType::timestamp_millisecond_datatype(),
1404 ),
1405 (2, SemanticType::Field, ConcreteDataType::int64_datatype()),
1406 ],
1407 &[],
1408 );
1409 actual_metadata.primary_key_encoding = PrimaryKeyEncoding::Sparse;
1410 let actual_metadata = Arc::new(actual_metadata);
1411
1412 let mut expected_metadata = new_metadata(
1413 &[
1414 (
1415 0,
1416 SemanticType::Timestamp,
1417 ConcreteDataType::timestamp_millisecond_datatype(),
1418 ),
1419 (2, SemanticType::Field, ConcreteDataType::int64_datatype()),
1420 (3, SemanticType::Field, ConcreteDataType::int64_datatype()),
1421 ],
1422 &[],
1423 );
1424 expected_metadata.primary_key_encoding = PrimaryKeyEncoding::Sparse;
1425 let expected_metadata = Arc::new(expected_metadata);
1426
1427 let mapper = FlatProjectionMapper::all(&expected_metadata).unwrap();
1428 let read_format = FlatReadFormat::new(
1429 actual_metadata.clone(),
1430 ReadColumns::new([0, 2, 3]),
1431 None,
1432 "test",
1433 true,
1434 )
1435 .unwrap();
1436
1437 let compat_batch = FlatCompatBatch::try_new(&mapper, &read_format, true)
1438 .unwrap()
1439 .unwrap();
1440
1441 let sparse_k1 = encode_sparse_key(&[]);
1442 let input_columns: Vec<ArrayRef> = vec![
1443 Arc::new(Int64Array::from(vec![100, 200])),
1444 Arc::new(TimestampMillisecondArray::from_iter_values([1000, 2000])),
1445 build_flat_test_pk_array(&[&sparse_k1, &sparse_k1]),
1446 Arc::new(UInt64Array::from_iter_values([1, 2])),
1447 Arc::new(UInt8Array::from_iter_values([
1448 OpType::Put as u8,
1449 OpType::Put as u8,
1450 ])),
1451 ];
1452 let input_schema =
1453 to_flat_sst_arrow_schema(&actual_metadata, &FlatSchemaOptions::default());
1454 let input_batch = RecordBatch::try_new(input_schema, input_columns).unwrap();
1455
1456 let result = compat_batch.compat(input_batch).unwrap();
1457
1458 let expected_columns: Vec<ArrayRef> = vec![
1459 Arc::new(Int64Array::from(vec![100, 200])),
1460 Arc::new(Int64Array::from(vec![None::<i64>, None::<i64>])),
1461 Arc::new(TimestampMillisecondArray::from_iter_values([1000, 2000])),
1462 build_flat_test_pk_array(&[&sparse_k1, &sparse_k1]),
1463 Arc::new(UInt64Array::from_iter_values([1, 2])),
1464 Arc::new(UInt8Array::from_iter_values([
1465 OpType::Put as u8,
1466 OpType::Put as u8,
1467 ])),
1468 ];
1469 let output_schema =
1470 to_flat_sst_arrow_schema(&expected_metadata, &FlatSchemaOptions::default());
1471 let expected_batch = RecordBatch::try_new(output_schema, expected_columns).unwrap();
1472
1473 assert_eq!(expected_batch, result);
1474 }
1475}