1use std::collections::{HashMap, HashSet};
18use std::sync::Arc;
19use std::time::{Duration, Instant};
20
21use api::helper::{ColumnDataTypeWrapper, to_grpc_value};
22use api::v1::bulk_wal_entry::Body;
23use api::v1::{ArrowIpc, BulkWalEntry, Mutation, OpType};
24use bytes::Bytes;
25use common_grpc::flight::{FlightDecoder, FlightEncoder, FlightMessage};
26use common_recordbatch::DfRecordBatch as RecordBatch;
27use common_time::Timestamp;
28use datafusion_common::Column;
29use datafusion_common::pruning::PruningStatistics;
30use datafusion_expr::utils::expr_to_columns;
31use datatypes::arrow;
32use datatypes::arrow::array::{
33 Array, ArrayRef, BinaryArray, BooleanArray, DictionaryArray, StringDictionaryBuilder,
34 TimestampMicrosecondArray, TimestampMillisecondArray, TimestampNanosecondArray,
35 TimestampSecondArray, UInt8Array, UInt32Array, UInt64Array,
36};
37use datatypes::arrow::compute::{SortColumn, SortOptions, concat_batches};
38use datatypes::arrow::datatypes::{
39 DataType as ArrowDataType, Field, Schema, SchemaRef, TimeUnit, UInt32Type,
40};
41use datatypes::data_type::DataType;
42use datatypes::extension::json::is_json2_extension_type;
43use datatypes::prelude::{MutableVector, Vector};
44use datatypes::value::ValueRef;
45use datatypes::vectors::Helper;
46use mito_codec::key_values::{KeyValue, KeyValues};
47use mito_codec::row_converter::{PrimaryKeyCodec, SortField, build_primary_key_codec_with_fields};
48use parquet::arrow::ArrowWriter;
49use parquet::basic::{Compression, ZstdLevel};
50use parquet::file::metadata::ParquetMetaData;
51use parquet::file::properties::WriterProperties;
52use smallvec::SmallVec;
53use snafu::{OptionExt, ResultExt};
54use store_api::codec::PrimaryKeyEncoding;
55use store_api::metadata::{RegionMetadata, RegionMetadataRef};
56use store_api::storage::consts::PRIMARY_KEY_COLUMN_NAME;
57use store_api::storage::{ColumnId, FileId, SequenceNumber, SequenceRange};
58
59use crate::error::{
60 self, ColumnNotFoundSnafu, ComputeArrowSnafu, CreateDefaultSnafu, DataTypeMismatchSnafu,
61 EncodeMemtableSnafu, EncodeSnafu, InvalidMetadataSnafu, InvalidRequestSnafu,
62 NewRecordBatchSnafu, Result,
63};
64use crate::memtable::bulk::context::{BulkIterContext, BulkIterContextRef};
65use crate::memtable::bulk::json_align::Json2Aligner;
66use crate::memtable::bulk::part_reader::EncodedBulkPartIter;
67use crate::memtable::time_series::{ValueBuilder, Values};
68use crate::memtable::{BoxedRecordBatchIterator, MemScanMetrics, MemtableStats};
69use crate::sst::SeriesEstimator;
70use crate::sst::index::IndexOutput;
71use crate::sst::parquet::flat_format::primary_key_column_index;
72use crate::sst::parquet::format::{PrimaryKeyArray, PrimaryKeyArrayBuilder};
73use crate::sst::parquet::{PARQUET_METADATA_KEY, SstInfo};
74
75const INIT_DICT_VALUE_CAPACITY: usize = 8;
76
77#[derive(Clone)]
79pub struct BulkPart {
80 pub batch: RecordBatch,
81 pub max_timestamp: i64,
82 pub min_timestamp: i64,
83 pub sequence: u64,
84 pub timestamp_index: usize,
85 pub raw_data: Option<ArrowIpc>,
86}
87
88impl TryFrom<BulkWalEntry> for BulkPart {
89 type Error = error::Error;
90
91 fn try_from(value: BulkWalEntry) -> std::result::Result<Self, Self::Error> {
92 match value.body.expect("Entry payload should be present") {
93 Body::ArrowIpc(ipc) => {
94 let mut decoder = FlightDecoder::try_from_schema_bytes(&ipc.schema)
95 .context(error::ConvertBulkWalEntrySnafu)?;
96 let batch = decoder
97 .try_decode_record_batch(&ipc.data_header, &ipc.payload)
98 .context(error::ConvertBulkWalEntrySnafu)?;
99 Ok(Self {
100 batch,
101 max_timestamp: value.max_ts,
102 min_timestamp: value.min_ts,
103 sequence: value.sequence,
104 timestamp_index: value.timestamp_index as usize,
105 raw_data: Some(ipc),
106 })
107 }
108 }
109 }
110}
111
112impl TryFrom<&BulkPart> for BulkWalEntry {
113 type Error = error::Error;
114
115 fn try_from(value: &BulkPart) -> Result<Self> {
116 if let Some(ipc) = &value.raw_data {
117 Ok(BulkWalEntry {
118 sequence: value.sequence,
119 max_ts: value.max_timestamp,
120 min_ts: value.min_timestamp,
121 timestamp_index: value.timestamp_index as u32,
122 body: Some(Body::ArrowIpc(ipc.clone())),
123 })
124 } else {
125 let mut encoder = FlightEncoder::default();
126 let schema_bytes = encoder
127 .encode_schema(value.batch.schema().as_ref())
128 .data_header;
129 let [rb_data] = encoder
130 .encode(FlightMessage::RecordBatch(value.batch.clone()))
131 .try_into()
132 .map_err(|_| {
133 error::UnsupportedOperationSnafu {
134 err_msg: "create BulkWalEntry from RecordBatch with dictionary arrays",
135 }
136 .build()
137 })?;
138 Ok(BulkWalEntry {
139 sequence: value.sequence,
140 max_ts: value.max_timestamp,
141 min_ts: value.min_timestamp,
142 timestamp_index: value.timestamp_index as u32,
143 body: Some(Body::ArrowIpc(ArrowIpc {
144 schema: schema_bytes,
145 data_header: rb_data.data_header,
146 payload: rb_data.data_body,
147 })),
148 })
149 }
150 }
151}
152
153impl BulkPart {
154 pub(crate) fn schema(&self) -> SchemaRef {
155 self.batch.schema()
156 }
157
158 pub(crate) fn estimated_size(&self) -> usize {
159 record_batch_estimated_size(&self.batch)
160 }
161
162 pub fn estimated_series_count(&self) -> usize {
165 let pk_column_idx = primary_key_column_index(self.batch.num_columns());
166 let pk_column = self.batch.column(pk_column_idx);
167 if let Some(dict_array) = pk_column.as_any().downcast_ref::<PrimaryKeyArray>() {
168 dict_array.values().len()
169 } else {
170 0
171 }
172 }
173
174 pub fn to_memtable_stats(&self, region_metadata: &RegionMetadataRef) -> MemtableStats {
176 let ts_type = region_metadata.time_index_type();
177 let min_ts = ts_type.create_timestamp(self.min_timestamp);
178 let max_ts = ts_type.create_timestamp(self.max_timestamp);
179
180 MemtableStats {
181 estimated_bytes: self.estimated_size(),
182 time_range: Some((min_ts, max_ts)),
183 num_rows: self.num_rows(),
184 num_ranges: 1,
185 max_sequence: self.sequence,
186 series_count: self.estimated_series_count(),
187 }
188 }
189
190 pub fn fill_missing_columns(&mut self, region_metadata: &RegionMetadata) -> Result<()> {
200 let batch_schema = self.batch.schema();
202 let batch_columns: HashSet<_> = batch_schema
203 .fields()
204 .iter()
205 .map(|f| f.name().as_str())
206 .collect();
207
208 let mut columns_to_fill = Vec::new();
210 for column_meta in ®ion_metadata.column_metadatas {
211 if !batch_columns.contains(column_meta.column_schema.name.as_str()) {
214 columns_to_fill.push(column_meta);
215 }
216 }
217
218 if columns_to_fill.is_empty() {
219 return Ok(());
220 }
221
222 let num_rows = self.batch.num_rows();
223
224 let mut new_columns = Vec::new();
225 let mut new_fields = Vec::new();
226
227 new_fields.extend(batch_schema.fields().iter().cloned());
229 new_columns.extend_from_slice(self.batch.columns());
230
231 let region_id = region_metadata.region_id;
232 for column_meta in columns_to_fill {
234 let default_vector = column_meta
235 .column_schema
236 .create_default_vector(num_rows)
237 .context(CreateDefaultSnafu {
238 region_id,
239 column: &column_meta.column_schema.name,
240 })?
241 .with_context(|| InvalidRequestSnafu {
242 region_id,
243 reason: format!(
244 "column {} does not have default value",
245 column_meta.column_schema.name
246 ),
247 })?;
248 let arrow_array = default_vector.to_arrow_array();
249 column_meta.column_schema.data_type.as_arrow_type();
250
251 new_fields.push(Arc::new(Field::new(
252 column_meta.column_schema.name.clone(),
253 column_meta.column_schema.data_type.as_arrow_type(),
254 column_meta.column_schema.is_nullable(),
255 )));
256 new_columns.push(arrow_array);
257 }
258
259 let new_schema = Arc::new(Schema::new(new_fields));
261 let new_batch =
262 RecordBatch::try_new(new_schema, new_columns).context(NewRecordBatchSnafu)?;
263
264 self.batch = new_batch;
266
267 self.raw_data = None;
272
273 Ok(())
274 }
275
276 pub(crate) fn to_mutation(&self, region_metadata: &RegionMetadataRef) -> Result<Mutation> {
278 let vectors = region_metadata
279 .schema
280 .column_schemas()
281 .iter()
282 .map(|col| match self.batch.column_by_name(&col.name) {
283 None => Ok(None),
284 Some(col) => Helper::try_into_vector(col).map(Some),
285 })
286 .collect::<datatypes::error::Result<Vec<_>>>()
287 .context(error::ComputeVectorSnafu)?;
288
289 let rows = (0..self.num_rows())
290 .map(|row_idx| {
291 let values = (0..self.batch.num_columns())
292 .map(|col_idx| {
293 if let Some(v) = &vectors[col_idx] {
294 to_grpc_value(v.get(row_idx))
295 } else {
296 api::v1::Value { value_data: None }
297 }
298 })
299 .collect::<Vec<_>>();
300 api::v1::Row { values }
301 })
302 .collect::<Vec<_>>();
303
304 let schema = region_metadata
305 .column_metadatas
306 .iter()
307 .map(|c| {
308 let data_type_wrapper =
309 ColumnDataTypeWrapper::try_from(c.column_schema.data_type.clone())?;
310 Ok(api::v1::ColumnSchema {
311 column_name: c.column_schema.name.clone(),
312 datatype: data_type_wrapper.datatype() as i32,
313 semantic_type: c.semantic_type as i32,
314 ..Default::default()
315 })
316 })
317 .collect::<api::error::Result<Vec<_>>>()
318 .context(error::ConvertColumnDataTypeSnafu {
319 reason: "failed to convert region metadata to column schema",
320 })?;
321
322 let rows = api::v1::Rows { schema, rows };
323
324 Ok(Mutation {
325 op_type: OpType::Put as i32,
326 sequence: self.sequence,
327 rows: Some(rows),
328 write_hint: None,
329 })
330 }
331
332 pub fn timestamps(&self) -> &ArrayRef {
333 self.batch.column(self.timestamp_index)
334 }
335
336 pub fn num_rows(&self) -> usize {
337 self.batch.num_rows()
338 }
339}
340
341pub struct UnorderedPart {
344 parts: Vec<BulkPart>,
346 total_rows: usize,
348 total_bytes: usize,
350 min_timestamp: i64,
352 max_timestamp: i64,
354 max_sequence: u64,
356 threshold: usize,
358 compact_threshold: usize,
360}
361
362impl Default for UnorderedPart {
363 fn default() -> Self {
364 Self::new()
365 }
366}
367
368impl UnorderedPart {
369 pub fn new() -> Self {
371 Self {
372 parts: Vec::new(),
373 total_rows: 0,
374 total_bytes: 0,
375 min_timestamp: i64::MAX,
376 max_timestamp: i64::MIN,
377 max_sequence: 0,
378 threshold: 1024,
379 compact_threshold: 4096,
380 }
381 }
382
383 pub fn set_threshold(&mut self, threshold: usize) {
385 self.threshold = threshold;
386 }
387
388 pub fn set_compact_threshold(&mut self, compact_threshold: usize) {
390 self.compact_threshold = compact_threshold;
391 }
392
393 pub fn threshold(&self) -> usize {
395 self.threshold
396 }
397
398 pub fn compact_threshold(&self) -> usize {
400 self.compact_threshold
401 }
402
403 pub fn should_accept(&self, num_rows: usize) -> bool {
405 num_rows < self.threshold
406 }
407
408 pub fn should_compact(&self) -> bool {
410 self.total_rows >= self.compact_threshold
411 }
412
413 pub(super) fn estimated_bytes(&self) -> usize {
415 self.total_bytes
416 }
417
418 pub fn push(&mut self, part: BulkPart) {
420 self.total_rows += part.num_rows();
421 self.total_bytes = self.total_bytes.saturating_add(part.estimated_size());
422 self.min_timestamp = self.min_timestamp.min(part.min_timestamp);
423 self.max_timestamp = self.max_timestamp.max(part.max_timestamp);
424 self.max_sequence = self.max_sequence.max(part.sequence);
425 self.parts.push(part);
426 }
427
428 pub fn num_rows(&self) -> usize {
430 self.total_rows
431 }
432
433 pub fn is_empty(&self) -> bool {
435 self.parts.is_empty()
436 }
437
438 pub fn num_parts(&self) -> usize {
440 self.parts.len()
441 }
442
443 pub fn concat_and_sort(&self) -> Result<Option<RecordBatch>> {
446 if self.parts.is_empty() {
447 return Ok(None);
448 }
449
450 if self.parts.len() == 1 {
451 return Ok(Some(self.parts[0].batch.clone()));
453 }
454
455 let schema = self.parts[0].batch.schema();
457 let concatenated = if schema.fields().iter().any(is_json2_extension_type) {
458 let aligner = Json2Aligner::try_new(self.parts.iter().map(|part| part.batch.schema()))?;
459 let aligned_batches =
460 aligner.align_batches(self.parts.iter().map(|part| part.batch.clone()))?;
461 concat_batches(aligner.schema(), &aligned_batches).context(ComputeArrowSnafu)?
462 } else {
463 concat_batches(&schema, self.parts.iter().map(|x| &x.batch))
464 .context(ComputeArrowSnafu)?
465 };
466
467 let sorted_batch = sort_primary_key_record_batch(&concatenated)?;
469
470 Ok(Some(sorted_batch))
471 }
472
473 pub fn to_bulk_part(&self) -> Result<Option<BulkPart>> {
476 let Some(sorted_batch) = self.concat_and_sort()? else {
477 return Ok(None);
478 };
479
480 let timestamp_index = self.parts[0].timestamp_index;
481
482 Ok(Some(BulkPart {
483 batch: sorted_batch,
484 max_timestamp: self.max_timestamp,
485 min_timestamp: self.min_timestamp,
486 sequence: self.max_sequence,
487 timestamp_index,
488 raw_data: None,
489 }))
490 }
491
492 pub fn clear(&mut self) {
494 self.parts.clear();
495 self.total_rows = 0;
496 self.total_bytes = 0;
497 self.min_timestamp = i64::MAX;
498 self.max_timestamp = i64::MIN;
499 self.max_sequence = 0;
500 }
501}
502
503pub fn record_batch_estimated_size(batch: &RecordBatch) -> usize {
505 batch
506 .columns()
507 .iter()
508 .map(|c| c.to_data().get_slice_memory_size().unwrap_or(0))
510 .sum()
511}
512
513enum PrimaryKeyColumnBuilder {
515 StringDict(StringDictionaryBuilder<UInt32Type>),
517 Vector(Box<dyn MutableVector>),
519}
520
521impl PrimaryKeyColumnBuilder {
522 fn push_value_ref(&mut self, value: ValueRef) -> Result<()> {
524 match self {
525 PrimaryKeyColumnBuilder::StringDict(builder) => {
526 if let Some(s) = value.try_into_string().context(DataTypeMismatchSnafu)? {
527 builder.append_value(s);
529 } else {
530 builder.append_null();
531 }
532 }
533 PrimaryKeyColumnBuilder::Vector(builder) => {
534 builder.push_value_ref(&value);
535 }
536 }
537 Ok(())
538 }
539
540 fn into_arrow_array(self) -> ArrayRef {
542 match self {
543 PrimaryKeyColumnBuilder::StringDict(mut builder) => Arc::new(builder.finish()),
544 PrimaryKeyColumnBuilder::Vector(mut builder) => builder.to_vector().to_arrow_array(),
545 }
546 }
547}
548
549pub struct BulkPartConverter {
551 schema: SchemaRef,
553 primary_key_codec: Arc<dyn PrimaryKeyCodec>,
555 key_buf: Vec<u8>,
557 key_array_builder: PrimaryKeyArrayBuilder,
559 value_builder: ValueBuilder,
561 primary_key_column_builders: Vec<PrimaryKeyColumnBuilder>,
564
565 max_ts: i64,
567 min_ts: i64,
569 max_sequence: SequenceNumber,
571}
572
573impl BulkPartConverter {
574 pub fn new(
579 region_metadata: &RegionMetadataRef,
580 schema: SchemaRef,
581 capacity: usize,
582 primary_key_codec: Arc<dyn PrimaryKeyCodec>,
583 store_primary_key_columns: bool,
584 ) -> Self {
585 debug_assert_eq!(
586 region_metadata.primary_key_encoding,
587 primary_key_codec.encoding()
588 );
589
590 let primary_key_column_builders = if store_primary_key_columns
591 && region_metadata.primary_key_encoding != PrimaryKeyEncoding::Sparse
592 {
593 new_primary_key_column_builders(region_metadata, capacity)
594 } else {
595 Vec::new()
596 };
597
598 Self {
599 schema,
600 primary_key_codec,
601 key_buf: Vec::new(),
602 key_array_builder: PrimaryKeyArrayBuilder::new(),
603 value_builder: ValueBuilder::new(region_metadata, capacity),
604 primary_key_column_builders,
605 min_ts: i64::MAX,
606 max_ts: i64::MIN,
607 max_sequence: SequenceNumber::MIN,
608 }
609 }
610
611 pub fn append_key_values(&mut self, key_values: &KeyValues) -> Result<()> {
613 for kv in key_values.iter() {
614 self.append_key_value(&kv)?;
615 }
616
617 Ok(())
618 }
619
620 fn append_key_value(&mut self, kv: &KeyValue) -> Result<()> {
624 if self.primary_key_codec.encoding() == PrimaryKeyEncoding::Sparse {
626 let mut primary_keys = kv.primary_keys();
629 if let Some(encoded) = primary_keys
630 .next()
631 .context(ColumnNotFoundSnafu {
632 column: PRIMARY_KEY_COLUMN_NAME,
633 })?
634 .try_into_binary()
635 .context(DataTypeMismatchSnafu)?
636 {
637 self.key_array_builder
638 .append(encoded)
639 .context(ComputeArrowSnafu)?;
640 } else {
641 self.key_array_builder
642 .append("")
643 .context(ComputeArrowSnafu)?;
644 }
645 } else {
646 self.key_buf.clear();
648 self.primary_key_codec
649 .encode_key_value(kv, &mut self.key_buf)
650 .context(EncodeSnafu)?;
651 self.key_array_builder
652 .append(&self.key_buf)
653 .context(ComputeArrowSnafu)?;
654 };
655
656 if !self.primary_key_column_builders.is_empty() {
658 for (builder, pk_value) in self
659 .primary_key_column_builders
660 .iter_mut()
661 .zip(kv.primary_keys())
662 {
663 builder.push_value_ref(pk_value)?;
664 }
665 }
666
667 self.value_builder.push(
669 kv.timestamp(),
670 kv.sequence(),
671 kv.op_type() as u8,
672 kv.fields(),
673 );
674
675 let ts = kv
678 .timestamp()
679 .try_into_timestamp()
680 .unwrap()
681 .unwrap()
682 .value();
683 self.min_ts = self.min_ts.min(ts);
684 self.max_ts = self.max_ts.max(ts);
685 self.max_sequence = self.max_sequence.max(kv.sequence());
686
687 Ok(())
688 }
689
690 pub fn convert(mut self) -> Result<BulkPart> {
694 let values = Values::from(self.value_builder);
695 let mut columns =
696 Vec::with_capacity(4 + values.fields.len() + self.primary_key_column_builders.len());
697
698 for builder in self.primary_key_column_builders {
700 columns.push(builder.into_arrow_array());
701 }
702 columns.extend(values.fields.iter().map(|field| field.to_arrow_array()));
704 let timestamp_index = columns.len();
706 columns.push(values.timestamp.to_arrow_array());
707 let pk_array = self.key_array_builder.finish();
709 columns.push(Arc::new(pk_array));
710 columns.push(values.sequence.to_arrow_array());
712 columns.push(values.op_type.to_arrow_array());
713
714 let schema = align_schema_with_json_array(self.schema, &columns);
717 let batch = RecordBatch::try_new(schema, columns).context(NewRecordBatchSnafu)?;
718 let batch = sort_primary_key_record_batch(&batch)?;
720
721 Ok(BulkPart {
722 batch,
723 max_timestamp: self.max_ts,
724 min_timestamp: self.min_ts,
725 sequence: self.max_sequence,
726 timestamp_index,
727 raw_data: None,
728 })
729 }
730}
731
732fn align_schema_with_json_array(schema: SchemaRef, columns: &[ArrayRef]) -> SchemaRef {
733 if schema.fields().iter().all(|f| !is_json2_extension_type(f)) {
734 return schema;
735 }
736
737 let mut fields = Vec::with_capacity(schema.fields().len());
738 for (field, array) in schema.fields().iter().zip(columns) {
739 if !is_json2_extension_type(field) {
740 fields.push(field.clone());
741 continue;
742 }
743
744 let mut field = field.as_ref().clone();
745 field.set_data_type(array.data_type().clone());
746 fields.push(Arc::new(field));
747 }
748
749 Arc::new(Schema::new_with_metadata(fields, schema.metadata().clone()))
750}
751
752fn new_primary_key_column_builders(
753 metadata: &RegionMetadata,
754 capacity: usize,
755) -> Vec<PrimaryKeyColumnBuilder> {
756 metadata
757 .primary_key_columns()
758 .map(|col| {
759 if col.column_schema.data_type.is_string() {
760 PrimaryKeyColumnBuilder::StringDict(StringDictionaryBuilder::with_capacity(
761 capacity,
762 INIT_DICT_VALUE_CAPACITY,
763 capacity,
764 ))
765 } else {
766 PrimaryKeyColumnBuilder::Vector(
767 col.column_schema.data_type.create_mutable_vector(capacity),
768 )
769 }
770 })
771 .collect()
772}
773
774pub fn sort_primary_key_record_batch(batch: &RecordBatch) -> Result<RecordBatch> {
776 let indices = if let Some(indices) = try_sort_primary_key_indices(batch)? {
777 indices
778 } else {
779 lexsort_primary_key_indices(batch)?
780 };
781
782 datatypes::arrow::compute::take_record_batch(batch, &indices).context(ComputeArrowSnafu)
783}
784
785fn try_sort_primary_key_indices(batch: &RecordBatch) -> Result<Option<UInt32Array>> {
789 let total_columns = batch.num_columns();
790 let timestamp = batch.column(total_columns - 4);
791 let primary_key = batch.column(total_columns - 3);
792 let sequence = batch.column(total_columns - 2);
793
794 let Some(primary_key) = primary_key
795 .as_any()
796 .downcast_ref::<DictionaryArray<UInt32Type>>()
797 else {
798 return Ok(None);
799 };
800 let Some(sequence) = sequence.as_any().downcast_ref::<UInt64Array>() else {
801 return Ok(None);
802 };
803
804 if batch.num_rows() > u32::MAX as usize
805 || primary_key.null_count() != 0
806 || primary_key.values().data_type() != &ArrowDataType::Binary
807 || primary_key.values().null_count() != 0
808 || primary_key.values().len() > batch.num_rows()
809 || timestamp.null_count() != 0
810 || sequence.null_count() != 0
811 {
812 return Ok(None);
813 }
814
815 let indices = match timestamp.data_type() {
816 ArrowDataType::Timestamp(TimeUnit::Second, _) => timestamp
817 .as_any()
818 .downcast_ref::<TimestampSecondArray>()
819 .map(|timestamp| sort_primary_key_indices(primary_key, timestamp.values(), sequence)),
820 ArrowDataType::Timestamp(TimeUnit::Millisecond, _) => timestamp
821 .as_any()
822 .downcast_ref::<TimestampMillisecondArray>()
823 .map(|timestamp| sort_primary_key_indices(primary_key, timestamp.values(), sequence)),
824 ArrowDataType::Timestamp(TimeUnit::Microsecond, _) => timestamp
825 .as_any()
826 .downcast_ref::<TimestampMicrosecondArray>()
827 .map(|timestamp| sort_primary_key_indices(primary_key, timestamp.values(), sequence)),
828 ArrowDataType::Timestamp(TimeUnit::Nanosecond, _) => timestamp
829 .as_any()
830 .downcast_ref::<TimestampNanosecondArray>()
831 .map(|timestamp| sort_primary_key_indices(primary_key, timestamp.values(), sequence)),
832 _ => None,
833 };
834
835 indices.transpose()
836}
837
838fn sort_primary_key_indices(
841 primary_key: &PrimaryKeyArray,
842 timestamps: &[i64],
843 sequences: &UInt64Array,
844) -> Result<UInt32Array> {
845 debug_assert_eq!(primary_key.len(), timestamps.len());
846 debug_assert_eq!(primary_key.len(), sequences.len());
847 debug_assert_eq!(primary_key.null_count(), 0);
848 debug_assert_eq!(primary_key.values().null_count(), 0);
849 debug_assert_eq!(sequences.null_count(), 0);
850
851 let value_ranks = datatypes::arrow::compute::rank(
855 primary_key.values().as_ref(),
856 Some(SortOptions {
857 descending: false,
858 nulls_first: true,
859 }),
860 )
861 .context(ComputeArrowSnafu)?;
862
863 let mut counts = vec![0usize; value_ranks.len() + 1];
866 for &key in primary_key.keys().values() {
867 counts[value_ranks[key as usize] as usize] += 1;
868 }
869
870 let mut bucket_offsets = vec![0usize; counts.len()];
871 let mut offset = 0;
872 for (bucket_offset, &count) in bucket_offsets.iter_mut().zip(&counts) {
873 *bucket_offset = offset;
874 offset += count;
875 }
876
877 let mut next_offsets = bucket_offsets.clone();
879 let mut indices = vec![0u32; primary_key.len()];
880 for (row, &key) in primary_key.keys().values().iter().enumerate() {
881 let rank = value_ranks[key as usize] as usize;
882 indices[next_offsets[rank]] = row as u32;
883 next_offsets[rank] += 1;
884 }
885
886 let sequences = sequences.values();
887 for (&start, count) in bucket_offsets.iter().zip(counts) {
888 let end = start + count;
889 let bucket = &mut indices[start..end];
890 if !bucket.is_sorted_by(|left, right| {
891 compare_time_sequence(timestamps, sequences, *left, *right).is_le()
892 }) {
893 bucket.sort_unstable_by(|left, right| {
894 compare_time_sequence(timestamps, sequences, *left, *right)
895 });
896 }
897 }
898
899 Ok(UInt32Array::from(indices))
900}
901
902#[inline]
903fn compare_time_sequence(
904 timestamps: &[i64],
905 sequences: &[u64],
906 left: u32,
907 right: u32,
908) -> std::cmp::Ordering {
909 let left = left as usize;
910 let right = right as usize;
911 timestamps[left]
912 .cmp(×tamps[right])
913 .then_with(|| sequences[right].cmp(&sequences[left]))
914}
915
916fn lexsort_primary_key_indices(batch: &RecordBatch) -> Result<UInt32Array> {
917 let total_columns = batch.num_columns();
918 let sort_columns = vec![
919 SortColumn {
921 values: batch.column(total_columns - 3).clone(),
922 options: Some(SortOptions {
923 descending: false,
924 nulls_first: true,
925 }),
926 },
927 SortColumn {
929 values: batch.column(total_columns - 4).clone(),
930 options: Some(SortOptions {
931 descending: false,
932 nulls_first: true,
933 }),
934 },
935 SortColumn {
937 values: batch.column(total_columns - 2).clone(),
938 options: Some(SortOptions {
939 descending: true,
940 nulls_first: true,
941 }),
942 },
943 ];
944
945 datatypes::arrow::compute::lexsort_to_indices(&sort_columns, None).context(ComputeArrowSnafu)
946}
947
948pub fn convert_bulk_part(
974 part: BulkPart,
975 region_metadata: &RegionMetadataRef,
976 primary_key_codec: Arc<dyn PrimaryKeyCodec>,
977 schema: SchemaRef,
978 store_primary_key_columns: bool,
979) -> Result<Option<BulkPart>> {
980 if part.num_rows() == 0 {
981 return Ok(None);
982 }
983
984 let num_rows = part.num_rows();
985 let is_sparse = region_metadata.primary_key_encoding == PrimaryKeyEncoding::Sparse;
986
987 let input_schema = part.batch.schema();
989 let column_indices: HashMap<&str, usize> = input_schema
990 .fields()
991 .iter()
992 .enumerate()
993 .map(|(idx, field)| (field.name().as_str(), idx))
994 .collect();
995
996 let mut output_columns = Vec::new();
998
999 let pk_array = if is_sparse {
1001 None
1004 } else {
1005 let pk_vectors: Result<Vec<_>> = region_metadata
1007 .primary_key_columns()
1008 .map(|col_meta| {
1009 let col_idx = column_indices
1010 .get(col_meta.column_schema.name.as_str())
1011 .context(ColumnNotFoundSnafu {
1012 column: &col_meta.column_schema.name,
1013 })?;
1014 let col = part.batch.column(*col_idx);
1015 Helper::try_into_vector(col).context(error::ComputeVectorSnafu)
1016 })
1017 .collect();
1018 let pk_vectors = pk_vectors?;
1019
1020 let mut key_array_builder = PrimaryKeyArrayBuilder::new();
1021 let mut encode_buf = Vec::new();
1022
1023 for row_idx in 0..num_rows {
1024 encode_buf.clear();
1025
1026 let pk_values_with_ids: Vec<_> = region_metadata
1028 .primary_key
1029 .iter()
1030 .zip(pk_vectors.iter())
1031 .map(|(col_id, vector)| (*col_id, vector.get_ref(row_idx)))
1032 .collect();
1033
1034 primary_key_codec
1036 .encode_value_refs(&pk_values_with_ids, &mut encode_buf)
1037 .context(EncodeSnafu)?;
1038
1039 key_array_builder
1040 .append(&encode_buf)
1041 .context(ComputeArrowSnafu)?;
1042 }
1043
1044 Some(key_array_builder.finish())
1045 };
1046
1047 if store_primary_key_columns && !is_sparse {
1049 for col_meta in region_metadata.primary_key_columns() {
1050 let col_idx = column_indices
1051 .get(col_meta.column_schema.name.as_str())
1052 .context(ColumnNotFoundSnafu {
1053 column: &col_meta.column_schema.name,
1054 })?;
1055 let col = part.batch.column(*col_idx);
1056
1057 let col = if col_meta.column_schema.data_type.is_string() {
1059 let target_type = ArrowDataType::Dictionary(
1060 Box::new(ArrowDataType::UInt32),
1061 Box::new(ArrowDataType::Utf8),
1062 );
1063 arrow::compute::cast(col, &target_type).context(ComputeArrowSnafu)?
1064 } else {
1065 col.clone()
1066 };
1067 output_columns.push(col);
1068 }
1069 }
1070
1071 for col_meta in region_metadata.field_columns() {
1073 let col_idx = column_indices
1074 .get(col_meta.column_schema.name.as_str())
1075 .context(ColumnNotFoundSnafu {
1076 column: &col_meta.column_schema.name,
1077 })?;
1078 output_columns.push(part.batch.column(*col_idx).clone());
1079 }
1080
1081 let new_timestamp_index = output_columns.len();
1083 let ts_col_idx = column_indices
1084 .get(
1085 region_metadata
1086 .time_index_column()
1087 .column_schema
1088 .name
1089 .as_str(),
1090 )
1091 .context(ColumnNotFoundSnafu {
1092 column: ®ion_metadata.time_index_column().column_schema.name,
1093 })?;
1094 output_columns.push(part.batch.column(*ts_col_idx).clone());
1095
1096 let pk_dictionary = if let Some(pk_dict_array) = pk_array {
1098 Arc::new(pk_dict_array) as ArrayRef
1099 } else {
1100 let pk_col_idx =
1101 column_indices
1102 .get(PRIMARY_KEY_COLUMN_NAME)
1103 .context(ColumnNotFoundSnafu {
1104 column: PRIMARY_KEY_COLUMN_NAME,
1105 })?;
1106 let col = part.batch.column(*pk_col_idx);
1107
1108 let target_type = ArrowDataType::Dictionary(
1110 Box::new(ArrowDataType::UInt32),
1111 Box::new(ArrowDataType::Binary),
1112 );
1113 arrow::compute::cast(col, &target_type).context(ComputeArrowSnafu)?
1114 };
1115 output_columns.push(pk_dictionary);
1116
1117 let sequence_array = UInt64Array::from(vec![part.sequence; num_rows]);
1118 output_columns.push(Arc::new(sequence_array) as ArrayRef);
1119
1120 let op_type_array = UInt8Array::from(vec![OpType::Put as u8; num_rows]);
1121 output_columns.push(Arc::new(op_type_array) as ArrayRef);
1122
1123 let batch = RecordBatch::try_new(schema, output_columns).context(NewRecordBatchSnafu)?;
1124
1125 let sorted_batch = sort_primary_key_record_batch(&batch)?;
1127
1128 Ok(Some(BulkPart {
1129 batch: sorted_batch,
1130 max_timestamp: part.max_timestamp,
1131 min_timestamp: part.min_timestamp,
1132 sequence: part.sequence,
1133 timestamp_index: new_timestamp_index,
1134 raw_data: None,
1135 }))
1136}
1137
1138#[derive(Debug, Clone)]
1139pub struct EncodedBulkPart {
1140 data: Bytes,
1141 metadata: BulkPartMeta,
1142 schema: SchemaRef,
1144}
1145
1146impl EncodedBulkPart {
1147 pub fn new(data: Bytes, metadata: BulkPartMeta, schema: SchemaRef) -> Self {
1148 Self {
1149 data,
1150 metadata,
1151 schema,
1152 }
1153 }
1154
1155 pub fn metadata(&self) -> &BulkPartMeta {
1156 &self.metadata
1157 }
1158
1159 pub(crate) fn schema(&self) -> SchemaRef {
1160 self.schema.clone()
1161 }
1162
1163 pub(crate) fn size_bytes(&self) -> usize {
1165 self.data.len()
1166 }
1167
1168 pub fn data(&self) -> &Bytes {
1170 &self.data
1171 }
1172
1173 pub fn to_memtable_stats(&self) -> MemtableStats {
1175 let meta = &self.metadata;
1176 let ts_type = meta.region_metadata.time_index_type();
1177 let min_ts = ts_type.create_timestamp(meta.min_timestamp);
1178 let max_ts = ts_type.create_timestamp(meta.max_timestamp);
1179
1180 MemtableStats {
1181 estimated_bytes: self.size_bytes(),
1182 time_range: Some((min_ts, max_ts)),
1183 num_rows: meta.num_rows,
1184 num_ranges: 1,
1185 max_sequence: meta.max_sequence,
1186 series_count: meta.num_series as usize,
1187 }
1188 }
1189
1190 pub(crate) fn to_sst_info(&self, file_id: FileId) -> SstInfo {
1198 let unit = self.metadata.region_metadata.time_index_type().unit();
1199 let max_row_group_uncompressed_size: u64 = self
1200 .metadata
1201 .parquet_metadata
1202 .row_groups()
1203 .iter()
1204 .map(|rg| {
1205 rg.columns()
1206 .iter()
1207 .map(|c| c.uncompressed_size() as u64)
1208 .sum::<u64>()
1209 })
1210 .max()
1211 .unwrap_or(0);
1212 SstInfo {
1213 file_id,
1214 time_range: (
1215 Timestamp::new(self.metadata.min_timestamp, unit),
1216 Timestamp::new(self.metadata.max_timestamp, unit),
1217 ),
1218 file_size: self.data.len() as u64,
1219 max_row_group_uncompressed_size,
1220 num_rows: self.metadata.num_rows,
1221 num_row_groups: self.metadata.parquet_metadata.num_row_groups() as u64,
1222 file_metadata: Some(self.metadata.parquet_metadata.clone()),
1223 index_metadata: IndexOutput::default(),
1224 num_series: self.metadata.num_series,
1225 }
1226 }
1227
1228 pub(crate) fn read(
1229 &self,
1230 context: BulkIterContextRef,
1231 sequence: Option<SequenceRange>,
1232 mem_scan_metrics: Option<MemScanMetrics>,
1233 ) -> Result<Option<BoxedRecordBatchIterator>> {
1234 let skip_fields_for_pruning = context.pre_filter_mode().skip_fields();
1236
1237 let row_groups_to_read =
1239 context.row_groups_to_read(&self.metadata.parquet_metadata, skip_fields_for_pruning);
1240
1241 if row_groups_to_read.is_empty() {
1242 return Ok(None);
1244 }
1245
1246 let iter = EncodedBulkPartIter::try_new(
1247 self,
1248 context,
1249 row_groups_to_read,
1250 sequence,
1251 mem_scan_metrics,
1252 )?;
1253 Ok(Some(Box::new(iter) as BoxedRecordBatchIterator))
1254 }
1255}
1256
1257#[derive(Debug, Clone)]
1259pub struct BulkPartMeta {
1260 pub num_rows: usize,
1262 pub max_timestamp: i64,
1264 pub min_timestamp: i64,
1266 pub parquet_metadata: Arc<ParquetMetaData>,
1268 pub region_metadata: RegionMetadataRef,
1270 pub num_series: u64,
1272 pub max_sequence: u64,
1274}
1275
1276#[derive(Default, Debug)]
1278pub struct BulkPartEncodeMetrics {
1279 pub iter_cost: Duration,
1281 pub write_cost: Duration,
1283 pub raw_size: usize,
1285 pub encoded_size: usize,
1287 pub num_rows: usize,
1289}
1290
1291pub struct BulkPartEncoder {
1292 metadata: RegionMetadataRef,
1293 writer_props: Option<WriterProperties>,
1294}
1295
1296impl BulkPartEncoder {
1297 pub fn new(metadata: RegionMetadataRef, row_group_size: usize) -> Result<BulkPartEncoder> {
1298 let json = metadata.to_json().context(InvalidMetadataSnafu)?;
1300 let key_value_meta =
1301 parquet::file::metadata::KeyValue::new(PARQUET_METADATA_KEY.to_string(), json);
1302
1303 let writer_props = Some(
1305 WriterProperties::builder()
1306 .set_key_value_metadata(Some(vec![key_value_meta]))
1307 .set_write_batch_size(row_group_size)
1308 .set_max_row_group_row_count(Some(row_group_size))
1309 .set_compression(Compression::ZSTD(ZstdLevel::default()))
1310 .set_column_index_truncate_length(None)
1311 .set_statistics_truncate_length(None)
1312 .build(),
1313 );
1314
1315 Ok(Self {
1316 metadata,
1317 writer_props,
1318 })
1319 }
1320}
1321
1322impl BulkPartEncoder {
1323 pub fn encode_record_batch_iter(
1325 &self,
1326 iter: BoxedRecordBatchIterator,
1327 arrow_schema: SchemaRef,
1328 min_timestamp: i64,
1329 max_timestamp: i64,
1330 max_sequence: u64,
1331 metrics: &mut BulkPartEncodeMetrics,
1332 ) -> Result<Option<EncodedBulkPart>> {
1333 let mut buf = Vec::with_capacity(4096);
1334 let mut writer =
1335 ArrowWriter::try_new(&mut buf, arrow_schema.clone(), self.writer_props.clone())
1336 .context(EncodeMemtableSnafu)?;
1337 let mut total_rows = 0;
1338 let mut series_estimator = SeriesEstimator::default();
1339
1340 let mut iter_start = Instant::now();
1342 for batch_result in iter {
1343 metrics.iter_cost += iter_start.elapsed();
1344 let batch = batch_result?;
1345 if batch.num_rows() == 0 {
1346 continue;
1347 }
1348
1349 series_estimator.update_flat(&batch);
1350 metrics.raw_size += record_batch_estimated_size(&batch);
1351 let write_start = Instant::now();
1352 writer.write(&batch).context(EncodeMemtableSnafu)?;
1353 metrics.write_cost += write_start.elapsed();
1354 total_rows += batch.num_rows();
1355 iter_start = Instant::now();
1356 }
1357 metrics.iter_cost += iter_start.elapsed();
1358
1359 if total_rows == 0 {
1360 return Ok(None);
1361 }
1362
1363 let close_start = Instant::now();
1364 let file_metadata = writer.close().context(EncodeMemtableSnafu)?;
1365 metrics.write_cost += close_start.elapsed();
1366 metrics.encoded_size += buf.len();
1367 metrics.num_rows += total_rows;
1368
1369 let buf = Bytes::from(buf);
1370 let parquet_metadata = Arc::new(file_metadata);
1371 let num_series = series_estimator.finish();
1372
1373 Ok(Some(EncodedBulkPart {
1374 data: buf,
1375 metadata: BulkPartMeta {
1376 num_rows: total_rows,
1377 max_timestamp,
1378 min_timestamp,
1379 parquet_metadata,
1380 region_metadata: self.metadata.clone(),
1381 num_series,
1382 max_sequence,
1383 },
1384 schema: arrow_schema,
1385 }))
1386 }
1387
1388 pub fn encode_part(&self, part: &BulkPart) -> Result<Option<EncodedBulkPart>> {
1390 if part.batch.num_rows() == 0 {
1391 return Ok(None);
1392 }
1393
1394 let mut buf = Vec::with_capacity(4096);
1395 let arrow_schema = part.batch.schema();
1396
1397 let file_metadata = {
1398 let mut writer =
1399 ArrowWriter::try_new(&mut buf, arrow_schema.clone(), self.writer_props.clone())
1400 .context(EncodeMemtableSnafu)?;
1401 writer.write(&part.batch).context(EncodeMemtableSnafu)?;
1402 writer.finish().context(EncodeMemtableSnafu)?
1403 };
1404
1405 let buf = Bytes::from(buf);
1406 let parquet_metadata = Arc::new(file_metadata);
1407
1408 Ok(Some(EncodedBulkPart {
1409 data: buf,
1410 metadata: BulkPartMeta {
1411 num_rows: part.batch.num_rows(),
1412 max_timestamp: part.max_timestamp,
1413 min_timestamp: part.min_timestamp,
1414 parquet_metadata,
1415 region_metadata: self.metadata.clone(),
1416 num_series: part.estimated_series_count() as u64,
1417 max_sequence: part.sequence,
1418 },
1419 schema: arrow_schema,
1420 }))
1421 }
1422}
1423
1424#[derive(Debug, Clone)]
1430struct BatchStats {
1431 num_batches: usize,
1433 first_tag_id: ColumnId,
1435 min_values: ArrayRef,
1437 max_values: ArrayRef,
1439}
1440
1441impl BatchStats {
1442 fn compute(batches: &[RecordBatch], metadata: &RegionMetadata) -> Option<Self> {
1447 let first_tag_id = *metadata.primary_key.first()?;
1452 let first_tag_column = metadata.column_by_id(first_tag_id)?;
1453 let data_type = &first_tag_column.column_schema.data_type;
1454
1455 let converter = build_primary_key_codec_with_fields(
1456 metadata.primary_key_encoding,
1457 [(first_tag_id, SortField::new(data_type.clone()))].into_iter(),
1458 );
1459 let pk_index = primary_key_column_index(batches.first()?.num_columns());
1460
1461 let mut min_builder = data_type.create_mutable_vector(batches.len());
1462 let mut max_builder = data_type.create_mutable_vector(batches.len());
1463
1464 for batch in batches {
1465 match Self::extract_first_tag_bounds(batch, pk_index, &*converter) {
1466 Some((min_val, max_val)) => {
1467 min_builder.push_value_ref(&min_val.as_value_ref());
1468 max_builder.push_value_ref(&max_val.as_value_ref());
1469 }
1470 None => {
1471 min_builder.push_null();
1472 max_builder.push_null();
1473 }
1474 }
1475 }
1476
1477 Some(Self {
1478 num_batches: batches.len(),
1479 first_tag_id,
1480 min_values: min_builder.to_vector().to_arrow_array(),
1481 max_values: max_builder.to_vector().to_arrow_array(),
1482 })
1483 }
1484
1485 fn extract_first_tag_bounds(
1487 batch: &RecordBatch,
1488 pk_index: usize,
1489 converter: &dyn PrimaryKeyCodec,
1490 ) -> Option<(datatypes::value::Value, datatypes::value::Value)> {
1491 if batch.num_rows() == 0 {
1492 return None;
1493 }
1494
1495 let pk_dict = batch
1496 .column(pk_index)
1497 .as_any()
1498 .downcast_ref::<PrimaryKeyArray>()?;
1499 let pk_values = pk_dict.values().as_any().downcast_ref::<BinaryArray>()?;
1500
1501 let keys = pk_dict.keys();
1502 let min_key = keys.value(0);
1503 let max_key = keys.value(batch.num_rows() - 1);
1504 let min_bytes = pk_values.value(min_key as usize);
1505 let max_bytes = pk_values.value(max_key as usize);
1506
1507 Some((
1508 converter.decode_leftmost(min_bytes).ok()??,
1509 converter.decode_leftmost(max_bytes).ok()??,
1510 ))
1511 }
1512}
1513
1514struct BatchPruningStats<'a> {
1519 stats: &'a BatchStats,
1520 metadata: &'a RegionMetadataRef,
1521}
1522
1523impl PruningStatistics for BatchPruningStats<'_> {
1524 fn min_values(&self, column: &Column) -> Option<ArrayRef> {
1525 let col = self.metadata.column_by_name(&column.name)?;
1526 if col.column_id == self.stats.first_tag_id {
1527 Some(self.stats.min_values.clone())
1528 } else {
1529 None
1530 }
1531 }
1532
1533 fn max_values(&self, column: &Column) -> Option<ArrayRef> {
1534 let col = self.metadata.column_by_name(&column.name)?;
1535 if col.column_id == self.stats.first_tag_id {
1536 Some(self.stats.max_values.clone())
1537 } else {
1538 None
1539 }
1540 }
1541
1542 fn num_containers(&self) -> usize {
1543 self.stats.num_batches
1544 }
1545
1546 fn null_counts(&self, _column: &Column) -> Option<ArrayRef> {
1547 None
1548 }
1549
1550 fn row_counts(&self, _column: &Column) -> Option<ArrayRef> {
1551 None
1552 }
1553
1554 fn contained(
1555 &self,
1556 _column: &Column,
1557 _values: &std::collections::HashSet<datafusion_common::ScalarValue>,
1558 ) -> Option<BooleanArray> {
1559 None
1560 }
1561}
1562
1563fn predicate_references_column(predicate: &table::predicate::Predicate, column_name: &str) -> bool {
1565 let mut columns = HashSet::new();
1566 for expr in predicate.exprs() {
1567 let _ = expr_to_columns(expr, &mut columns);
1568 }
1569 columns.iter().any(|col| col.name == column_name)
1570}
1571
1572pub(crate) fn should_prune_bulk_part(
1576 batch: &RecordBatch,
1577 context: &BulkIterContext,
1578 metadata: &RegionMetadata,
1579) -> bool {
1580 let predicate = match &context.predicate {
1581 Some(p) => p,
1582 None => return false,
1583 };
1584 let first_tag_id = match metadata.primary_key.first() {
1587 Some(id) => *id,
1588 None => return false,
1589 };
1590 let first_tag_name = &metadata
1592 .column_by_id(first_tag_id)
1593 .unwrap()
1594 .column_schema
1595 .name;
1596 if !predicate_references_column(predicate, first_tag_name) {
1597 return false;
1598 }
1599 let stats = match BatchStats::compute(std::slice::from_ref(batch), metadata) {
1600 Some(s) => s,
1601 None => return false,
1602 };
1603 let region_meta = context.read_format().metadata();
1604 let pruning_stats = BatchPruningStats {
1605 stats: &stats,
1606 metadata: region_meta,
1607 };
1608 let mask = predicate.prune_with_stats(&pruning_stats, region_meta.schema.arrow_schema());
1609 !mask.first().copied().unwrap_or(true)
1610}
1611
1612#[derive(Debug, Clone)]
1618pub struct MultiBulkPart {
1619 batches: SmallVec<[RecordBatch; 4]>,
1621 total_rows: usize,
1623 max_timestamp: i64,
1625 min_timestamp: i64,
1627 max_sequence: SequenceNumber,
1629 series_count: usize,
1631 batch_stats: Option<BatchStats>,
1634}
1635
1636impl MultiBulkPart {
1637 pub fn from_bulk_part(part: BulkPart, metadata: &RegionMetadata) -> Self {
1639 let num_rows = part.num_rows();
1640 let series_count = part.estimated_series_count();
1641 let batch_stats = BatchStats::compute(std::slice::from_ref(&part.batch), metadata);
1642 let mut batches = SmallVec::new();
1643 batches.push(part.batch);
1644
1645 Self {
1646 batches,
1647 total_rows: num_rows,
1648 max_timestamp: part.max_timestamp,
1649 min_timestamp: part.min_timestamp,
1650 max_sequence: part.sequence,
1651 series_count,
1652 batch_stats,
1653 }
1654 }
1655
1656 pub fn new(
1669 batches: Vec<RecordBatch>,
1670 min_timestamp: i64,
1671 max_timestamp: i64,
1672 max_sequence: SequenceNumber,
1673 series_count: usize,
1674 metadata: &RegionMetadata,
1675 ) -> Self {
1676 assert!(!batches.is_empty(), "batches must not be empty");
1677
1678 let total_rows = batches.iter().map(|b| b.num_rows()).sum();
1679 let batch_stats = BatchStats::compute(&batches, metadata);
1680
1681 Self {
1682 batches: SmallVec::from_vec(batches),
1683 total_rows,
1684 max_timestamp,
1685 min_timestamp,
1686 max_sequence,
1687 series_count,
1688 batch_stats,
1689 }
1690 }
1691
1692 pub fn num_rows(&self) -> usize {
1694 self.total_rows
1695 }
1696
1697 pub(crate) fn schemas(&self) -> impl Iterator<Item = SchemaRef> + '_ {
1698 self.batches.iter().map(|batch| batch.schema())
1699 }
1700
1701 pub fn min_timestamp(&self) -> i64 {
1703 self.min_timestamp
1704 }
1705
1706 pub fn max_timestamp(&self) -> i64 {
1708 self.max_timestamp
1709 }
1710
1711 pub fn max_sequence(&self) -> SequenceNumber {
1713 self.max_sequence
1714 }
1715
1716 pub fn series_count(&self) -> usize {
1718 self.series_count
1719 }
1720
1721 pub fn num_batches(&self) -> usize {
1723 self.batches.len()
1724 }
1725
1726 pub(crate) fn estimated_size(&self) -> usize {
1728 self.batches.iter().map(record_batch_estimated_size).sum()
1729 }
1730
1731 pub(crate) fn read(
1737 &self,
1738 context: BulkIterContextRef,
1739 sequence: Option<SequenceRange>,
1740 mem_scan_metrics: Option<MemScanMetrics>,
1741 ) -> Result<Option<BoxedRecordBatchIterator>> {
1742 if self.batches.is_empty() {
1743 return Ok(None);
1744 }
1745
1746 let batches_to_read = self.prune_batches(&context);
1747
1748 if batches_to_read.is_empty() {
1749 return Ok(None);
1750 }
1751
1752 let iter = crate::memtable::bulk::part_reader::BulkPartBatchIter::new(
1753 batches_to_read,
1754 context,
1755 sequence,
1756 self.series_count,
1757 mem_scan_metrics,
1758 );
1759 Ok(Some(Box::new(iter) as BoxedRecordBatchIterator))
1760 }
1761
1762 fn prune_batches(&self, context: &BulkIterContextRef) -> Vec<RecordBatch> {
1765 if let Some(stats) = &self.batch_stats
1766 && let Some(predicate) = &context.predicate
1767 {
1768 let region_meta = context.read_format().metadata();
1769 let pruning_stats = BatchPruningStats {
1770 stats,
1771 metadata: region_meta,
1772 };
1773 let mask =
1774 predicate.prune_with_stats(&pruning_stats, region_meta.schema.arrow_schema());
1775 self.batches
1776 .iter()
1777 .zip(mask.iter())
1778 .filter_map(
1779 |(batch, &selected)| {
1780 if selected { Some(batch.clone()) } else { None }
1781 },
1782 )
1783 .collect()
1784 } else {
1785 self.batches.iter().cloned().collect()
1786 }
1787 }
1788
1789 pub fn to_memtable_stats(&self, region_metadata: &RegionMetadataRef) -> MemtableStats {
1791 let ts_type = region_metadata.time_index_type();
1792 let min_ts = ts_type.create_timestamp(self.min_timestamp);
1793 let max_ts = ts_type.create_timestamp(self.max_timestamp);
1794
1795 MemtableStats {
1796 estimated_bytes: self.estimated_size(),
1797 time_range: Some((min_ts, max_ts)),
1798 num_rows: self.num_rows(),
1799 num_ranges: 1,
1800 max_sequence: self.max_sequence,
1801 series_count: self.series_count,
1802 }
1803 }
1804}
1805
1806#[cfg(test)]
1807mod tests {
1808 use api::v1::{Row, SemanticType, WriteHint};
1809 use datafusion_common::ScalarValue;
1810 use datatypes::arrow::array::{
1811 BinaryArray, DictionaryArray, Float64Array, TimestampMillisecondArray,
1812 };
1813 use datatypes::arrow::datatypes::UInt32Type;
1814 use datatypes::prelude::{ConcreteDataType, Value};
1815 use datatypes::schema::ColumnSchema;
1816 use mito_codec::row_converter::build_primary_key_codec;
1817 use rand::rngs::StdRng;
1818 use rand::{Rng, SeedableRng};
1819 use store_api::metadata::{ColumnMetadata, RegionMetadataBuilder};
1820 use store_api::storage::RegionId;
1821 use store_api::storage::consts::ReservedColumnId;
1822 use table::predicate::Predicate;
1823
1824 use super::*;
1825 use crate::memtable::bulk::context::BulkIterContext;
1826 use crate::sst::{FlatSchemaOptions, to_flat_sst_arrow_schema};
1827 use crate::test_util::memtable_util::{build_key_values_with_ts_seq_values, metadata_for_test};
1828
1829 struct MutationInput<'a> {
1830 k0: &'a str,
1831 k1: u32,
1832 timestamps: &'a [i64],
1833 v1: &'a [Option<f64>],
1834 sequence: u64,
1835 }
1836
1837 #[test]
1838 fn test_sort_primary_key_record_batch_with_duplicate_dictionary_values() {
1839 let primary_key = DictionaryArray::try_new(
1842 UInt32Array::from(vec![0, 1, 0, 2, 2]),
1843 Arc::new(BinaryArray::from_vec(vec![
1844 b"series_b".as_slice(),
1845 b"series_a".as_slice(),
1846 b"series_a".as_slice(),
1847 ])),
1848 )
1849 .unwrap();
1850 let timestamps = TimestampMillisecondArray::from(vec![20, 30, 10, 30, 10]);
1851 let sequences = UInt64Array::from(vec![5, 6, 4, 8, 9]);
1852 let row_ids = UInt32Array::from_iter_values(0..5);
1853 let op_types = UInt8Array::from_value(OpType::Put as u8, 5);
1854 let schema = Arc::new(Schema::new(vec![
1855 Field::new("row_id", ArrowDataType::UInt32, false),
1856 Field::new(
1857 "ts",
1858 ArrowDataType::Timestamp(TimeUnit::Millisecond, None),
1859 false,
1860 ),
1861 Field::new(
1862 PRIMARY_KEY_COLUMN_NAME,
1863 ArrowDataType::Dictionary(
1864 Box::new(ArrowDataType::UInt32),
1865 Box::new(ArrowDataType::Binary),
1866 ),
1867 false,
1868 ),
1869 Field::new("__sequence", ArrowDataType::UInt64, false),
1870 Field::new("__op_type", ArrowDataType::UInt8, false),
1871 ]));
1872 let batch = RecordBatch::try_new(
1873 schema,
1874 vec![
1875 Arc::new(row_ids),
1876 Arc::new(timestamps),
1877 Arc::new(primary_key),
1878 Arc::new(sequences),
1879 Arc::new(op_types),
1880 ],
1881 )
1882 .unwrap();
1883
1884 let expected_indices = lexsort_primary_key_indices(&batch).unwrap();
1885 let expected =
1886 datatypes::arrow::compute::take_record_batch(&batch, &expected_indices).unwrap();
1887 let actual = sort_primary_key_record_batch(&batch).unwrap();
1888
1889 let expected_row_ids = expected
1890 .column(0)
1891 .as_any()
1892 .downcast_ref::<UInt32Array>()
1893 .unwrap();
1894 let actual_row_ids = actual
1895 .column(0)
1896 .as_any()
1897 .downcast_ref::<UInt32Array>()
1898 .unwrap();
1899 assert_eq!(expected_row_ids, actual_row_ids);
1900 assert_eq!(actual_row_ids.values(), &[4, 3, 1, 2, 0]);
1901 }
1902
1903 #[test]
1904 fn test_sort_primary_key_record_batch_against_lexsort() {
1905 let mut rng = StdRng::seed_from_u64(0x5eed);
1906 for _ in 0..100 {
1907 let num_rows = rng.random_range(0..128);
1908 let keys = UInt32Array::from_iter_values((0..num_rows).map(|_| rng.random_range(0..8)));
1909 let primary_key = DictionaryArray::try_new(
1910 keys,
1911 Arc::new(BinaryArray::from_vec(vec![
1912 b"d".as_slice(),
1913 b"a".as_slice(),
1914 b"c".as_slice(),
1915 b"b".as_slice(),
1916 b"a".as_slice(),
1917 b"d".as_slice(),
1918 b"b".as_slice(),
1919 b"c".as_slice(),
1920 ])),
1921 )
1922 .unwrap();
1923 let timestamps = TimestampMillisecondArray::from_iter_values(
1924 (0..num_rows).map(|_| rng.random_range(-1000..1000)),
1925 );
1926 let sequences = UInt64Array::from_iter_values(0..num_rows as u64);
1928 let row_ids = UInt32Array::from_iter_values(0..num_rows as u32);
1929 let op_types = UInt8Array::from_value(OpType::Put as u8, num_rows);
1930 let schema = Arc::new(Schema::new(vec![
1931 Field::new("row_id", ArrowDataType::UInt32, false),
1932 Field::new(
1933 "ts",
1934 ArrowDataType::Timestamp(TimeUnit::Millisecond, None),
1935 false,
1936 ),
1937 Field::new(
1938 PRIMARY_KEY_COLUMN_NAME,
1939 ArrowDataType::Dictionary(
1940 Box::new(ArrowDataType::UInt32),
1941 Box::new(ArrowDataType::Binary),
1942 ),
1943 false,
1944 ),
1945 Field::new("__sequence", ArrowDataType::UInt64, false),
1946 Field::new("__op_type", ArrowDataType::UInt8, false),
1947 ]));
1948 let batch = RecordBatch::try_new(
1949 schema,
1950 vec![
1951 Arc::new(row_ids),
1952 Arc::new(timestamps),
1953 Arc::new(primary_key),
1954 Arc::new(sequences),
1955 Arc::new(op_types),
1956 ],
1957 )
1958 .unwrap();
1959
1960 let expected_indices = lexsort_primary_key_indices(&batch).unwrap();
1961 let expected =
1962 datatypes::arrow::compute::take_record_batch(&batch, &expected_indices).unwrap();
1963 let actual = sort_primary_key_record_batch(&batch).unwrap();
1964 assert_eq!(expected.column(0), actual.column(0));
1965 }
1966 }
1967
1968 #[test]
1969 fn test_unordered_part_tracks_estimated_bytes() {
1970 let mut part = UnorderedPart::new();
1971 let bulk_part = BulkPart {
1972 batch: RecordBatch::new_empty(Arc::new(arrow::datatypes::Schema::empty())),
1973 max_timestamp: 0,
1974 min_timestamp: 0,
1975 sequence: 0,
1976 timestamp_index: 0,
1977 raw_data: None,
1978 };
1979 let estimated_size = bulk_part.estimated_size();
1980
1981 part.push(bulk_part);
1982 assert_eq!(estimated_size, part.estimated_bytes());
1983 part.clear();
1984 assert_eq!(0, part.estimated_bytes());
1985 assert!(part.is_empty());
1986 }
1987
1988 #[test]
1989 fn test_unordered_part_should_accept() {
1990 let mut part = UnorderedPart::new();
1991 part.set_threshold(10);
1992 assert!(part.should_accept(9));
1993 assert!(!part.should_accept(10));
1994 }
1995
1996 fn encode(input: &[MutationInput]) -> EncodedBulkPart {
1997 let metadata = metadata_for_test();
1998 let kvs = input
1999 .iter()
2000 .map(|m| {
2001 build_key_values_with_ts_seq_values(
2002 &metadata,
2003 m.k0.to_string(),
2004 m.k1,
2005 m.timestamps.iter().copied(),
2006 m.v1.iter().copied(),
2007 m.sequence,
2008 )
2009 })
2010 .collect::<Vec<_>>();
2011 let schema = to_flat_sst_arrow_schema(&metadata, &FlatSchemaOptions::default());
2012 let primary_key_codec = build_primary_key_codec(&metadata);
2013 let mut converter = BulkPartConverter::new(&metadata, schema, 64, primary_key_codec, true);
2014 for kv in kvs {
2015 converter.append_key_values(&kv).unwrap();
2016 }
2017 let part = converter.convert().unwrap();
2018 let encoder = BulkPartEncoder::new(metadata, 1024).unwrap();
2019 encoder.encode_part(&part).unwrap().unwrap()
2020 }
2021
2022 #[test]
2023 fn test_write_and_read_part_projection() {
2024 let part = encode(&[
2025 MutationInput {
2026 k0: "a",
2027 k1: 0,
2028 timestamps: &[1],
2029 v1: &[Some(0.1)],
2030 sequence: 0,
2031 },
2032 MutationInput {
2033 k0: "b",
2034 k1: 0,
2035 timestamps: &[1],
2036 v1: &[Some(0.0)],
2037 sequence: 0,
2038 },
2039 MutationInput {
2040 k0: "a",
2041 k1: 0,
2042 timestamps: &[2],
2043 v1: &[Some(0.2)],
2044 sequence: 1,
2045 },
2046 ]);
2047
2048 let projection = &[4u32];
2049 let reader = part
2050 .read(
2051 Arc::new(
2052 BulkIterContext::new(
2053 part.metadata.region_metadata.clone(),
2054 Some(projection.as_slice()),
2055 None,
2056 false,
2057 crate::sst::parquet::DEFAULT_READ_BATCH_SIZE,
2058 )
2059 .unwrap(),
2060 ),
2061 None,
2062 None,
2063 )
2064 .unwrap()
2065 .expect("expect at least one row group");
2066
2067 let mut total_rows_read = 0;
2068 let mut field: Vec<f64> = vec![];
2069 for res in reader {
2070 let batch = res.unwrap();
2071 assert_eq!(5, batch.num_columns());
2072 field.extend_from_slice(
2073 batch
2074 .column(0)
2075 .as_any()
2076 .downcast_ref::<Float64Array>()
2077 .unwrap()
2078 .values(),
2079 );
2080 total_rows_read += batch.num_rows();
2081 }
2082 assert_eq!(3, total_rows_read);
2083 assert_eq!(vec![0.1, 0.2, 0.0], field);
2084 }
2085
2086 fn prepare(key_values: Vec<(&str, u32, (i64, i64), u64)>) -> EncodedBulkPart {
2087 let metadata = metadata_for_test();
2088 let kvs = key_values
2089 .into_iter()
2090 .map(|(k0, k1, (start, end), sequence)| {
2091 let ts = start..end;
2092 let v1 = (start..end).map(|_| None);
2093 build_key_values_with_ts_seq_values(&metadata, k0.to_string(), k1, ts, v1, sequence)
2094 })
2095 .collect::<Vec<_>>();
2096 let schema = to_flat_sst_arrow_schema(&metadata, &FlatSchemaOptions::default());
2097 let primary_key_codec = build_primary_key_codec(&metadata);
2098 let mut converter = BulkPartConverter::new(&metadata, schema, 64, primary_key_codec, true);
2099 for kv in kvs {
2100 converter.append_key_values(&kv).unwrap();
2101 }
2102 let part = converter.convert().unwrap();
2103 let encoder = BulkPartEncoder::new(metadata, 1024).unwrap();
2104 encoder.encode_part(&part).unwrap().unwrap()
2105 }
2106
2107 fn check_prune_row_group(
2108 part: &EncodedBulkPart,
2109 predicate: Option<Predicate>,
2110 expected_rows: usize,
2111 ) {
2112 let context = Arc::new(
2113 BulkIterContext::new(
2114 part.metadata.region_metadata.clone(),
2115 None,
2116 predicate,
2117 false,
2118 crate::sst::parquet::DEFAULT_READ_BATCH_SIZE,
2119 )
2120 .unwrap(),
2121 );
2122 let reader = part
2123 .read(context, None, None)
2124 .unwrap()
2125 .expect("expect at least one row group");
2126 let mut total_rows_read = 0;
2127 for res in reader {
2128 let batch = res.unwrap();
2129 total_rows_read += batch.num_rows();
2130 }
2131 assert_eq!(expected_rows, total_rows_read);
2133 }
2134
2135 #[test]
2136 fn test_prune_row_groups() {
2137 let part = prepare(vec![
2138 ("a", 0, (0, 40), 1),
2139 ("a", 1, (0, 60), 1),
2140 ("b", 0, (0, 100), 2),
2141 ("b", 1, (100, 180), 3),
2142 ("b", 1, (180, 210), 4),
2143 ]);
2144
2145 let context = Arc::new(
2146 BulkIterContext::new(
2147 part.metadata.region_metadata.clone(),
2148 None,
2149 Some(Predicate::new(vec![datafusion_expr::col("ts").eq(
2150 datafusion_expr::lit(ScalarValue::TimestampMillisecond(Some(300), None)),
2151 )])),
2152 false,
2153 crate::sst::parquet::DEFAULT_READ_BATCH_SIZE,
2154 )
2155 .unwrap(),
2156 );
2157 assert!(part.read(context, None, None).unwrap().is_none());
2158
2159 check_prune_row_group(&part, None, 310);
2160
2161 check_prune_row_group(
2162 &part,
2163 Some(Predicate::new(vec![
2164 datafusion_expr::col("k0").eq(datafusion_expr::lit("a")),
2165 datafusion_expr::col("k1").eq(datafusion_expr::lit(0u32)),
2166 ])),
2167 40,
2168 );
2169
2170 check_prune_row_group(
2171 &part,
2172 Some(Predicate::new(vec![
2173 datafusion_expr::col("k0").eq(datafusion_expr::lit("a")),
2174 datafusion_expr::col("k1").eq(datafusion_expr::lit(1u32)),
2175 ])),
2176 60,
2177 );
2178
2179 check_prune_row_group(
2180 &part,
2181 Some(Predicate::new(vec![
2182 datafusion_expr::col("k0").eq(datafusion_expr::lit("a")),
2183 ])),
2184 100,
2185 );
2186
2187 check_prune_row_group(
2188 &part,
2189 Some(Predicate::new(vec![
2190 datafusion_expr::col("k0").eq(datafusion_expr::lit("b")),
2191 datafusion_expr::col("k1").eq(datafusion_expr::lit(0u32)),
2192 ])),
2193 100,
2194 );
2195
2196 check_prune_row_group(
2198 &part,
2199 Some(Predicate::new(vec![
2200 datafusion_expr::col("v0").eq(datafusion_expr::lit(150i64)),
2201 ])),
2202 1,
2203 );
2204 }
2205
2206 #[test]
2207 fn test_bulk_part_converter_append_and_convert() {
2208 let metadata = metadata_for_test();
2209 let capacity = 100;
2210 let primary_key_codec = build_primary_key_codec(&metadata);
2211 let schema = to_flat_sst_arrow_schema(
2212 &metadata,
2213 &FlatSchemaOptions::from_encoding(metadata.primary_key_encoding),
2214 );
2215
2216 let mut converter =
2217 BulkPartConverter::new(&metadata, schema, capacity, primary_key_codec, true);
2218
2219 let key_values1 = build_key_values_with_ts_seq_values(
2220 &metadata,
2221 "key1".to_string(),
2222 1u32,
2223 vec![1000, 2000].into_iter(),
2224 vec![Some(1.0), Some(2.0)].into_iter(),
2225 1,
2226 );
2227
2228 let key_values2 = build_key_values_with_ts_seq_values(
2229 &metadata,
2230 "key2".to_string(),
2231 2u32,
2232 vec![1500].into_iter(),
2233 vec![Some(3.0)].into_iter(),
2234 2,
2235 );
2236
2237 converter.append_key_values(&key_values1).unwrap();
2238 converter.append_key_values(&key_values2).unwrap();
2239
2240 let bulk_part = converter.convert().unwrap();
2241
2242 assert_eq!(bulk_part.num_rows(), 3);
2243 assert_eq!(bulk_part.min_timestamp, 1000);
2244 assert_eq!(bulk_part.max_timestamp, 2000);
2245 assert_eq!(bulk_part.sequence, 2);
2246 assert_eq!(bulk_part.timestamp_index, bulk_part.batch.num_columns() - 4);
2247
2248 let schema = bulk_part.batch.schema();
2251 let field_names: Vec<&str> = schema.fields().iter().map(|f| f.name().as_str()).collect();
2252 assert_eq!(
2253 field_names,
2254 vec![
2255 "k0",
2256 "k1",
2257 "v0",
2258 "v1",
2259 "ts",
2260 "__primary_key",
2261 "__sequence",
2262 "__op_type"
2263 ]
2264 );
2265 }
2266
2267 #[test]
2268 fn test_bulk_part_converter_sorting() {
2269 let metadata = metadata_for_test();
2270 let capacity = 100;
2271 let primary_key_codec = build_primary_key_codec(&metadata);
2272 let schema = to_flat_sst_arrow_schema(
2273 &metadata,
2274 &FlatSchemaOptions::from_encoding(metadata.primary_key_encoding),
2275 );
2276
2277 let mut converter =
2278 BulkPartConverter::new(&metadata, schema, capacity, primary_key_codec, true);
2279
2280 let key_values1 = build_key_values_with_ts_seq_values(
2281 &metadata,
2282 "z_key".to_string(),
2283 3u32,
2284 vec![3000].into_iter(),
2285 vec![Some(3.0)].into_iter(),
2286 3,
2287 );
2288
2289 let key_values2 = build_key_values_with_ts_seq_values(
2290 &metadata,
2291 "a_key".to_string(),
2292 1u32,
2293 vec![1000].into_iter(),
2294 vec![Some(1.0)].into_iter(),
2295 1,
2296 );
2297
2298 let key_values3 = build_key_values_with_ts_seq_values(
2299 &metadata,
2300 "m_key".to_string(),
2301 2u32,
2302 vec![2000].into_iter(),
2303 vec![Some(2.0)].into_iter(),
2304 2,
2305 );
2306
2307 converter.append_key_values(&key_values1).unwrap();
2308 converter.append_key_values(&key_values2).unwrap();
2309 converter.append_key_values(&key_values3).unwrap();
2310
2311 let bulk_part = converter.convert().unwrap();
2312
2313 assert_eq!(bulk_part.num_rows(), 3);
2314
2315 let ts_column = bulk_part.batch.column(bulk_part.timestamp_index);
2316 let seq_column = bulk_part.batch.column(bulk_part.batch.num_columns() - 2);
2317
2318 let ts_array = ts_column
2319 .as_any()
2320 .downcast_ref::<TimestampMillisecondArray>()
2321 .unwrap();
2322 let seq_array = seq_column.as_any().downcast_ref::<UInt64Array>().unwrap();
2323
2324 assert_eq!(ts_array.values(), &[1000, 2000, 3000]);
2325 assert_eq!(seq_array.values(), &[1, 2, 3]);
2326
2327 let schema = bulk_part.batch.schema();
2329 let field_names: Vec<&str> = schema.fields().iter().map(|f| f.name().as_str()).collect();
2330 assert_eq!(
2331 field_names,
2332 vec![
2333 "k0",
2334 "k1",
2335 "v0",
2336 "v1",
2337 "ts",
2338 "__primary_key",
2339 "__sequence",
2340 "__op_type"
2341 ]
2342 );
2343 }
2344
2345 #[test]
2346 fn test_bulk_part_converter_empty() {
2347 let metadata = metadata_for_test();
2348 let capacity = 10;
2349 let primary_key_codec = build_primary_key_codec(&metadata);
2350 let schema = to_flat_sst_arrow_schema(
2351 &metadata,
2352 &FlatSchemaOptions::from_encoding(metadata.primary_key_encoding),
2353 );
2354
2355 let converter =
2356 BulkPartConverter::new(&metadata, schema, capacity, primary_key_codec, true);
2357
2358 let bulk_part = converter.convert().unwrap();
2359
2360 assert_eq!(bulk_part.num_rows(), 0);
2361 assert_eq!(bulk_part.min_timestamp, i64::MAX);
2362 assert_eq!(bulk_part.max_timestamp, i64::MIN);
2363 assert_eq!(bulk_part.sequence, SequenceNumber::MIN);
2364
2365 let schema = bulk_part.batch.schema();
2367 let field_names: Vec<&str> = schema.fields().iter().map(|f| f.name().as_str()).collect();
2368 assert_eq!(
2369 field_names,
2370 vec![
2371 "k0",
2372 "k1",
2373 "v0",
2374 "v1",
2375 "ts",
2376 "__primary_key",
2377 "__sequence",
2378 "__op_type"
2379 ]
2380 );
2381 }
2382
2383 #[test]
2384 fn test_bulk_part_converter_without_primary_key_columns() {
2385 let metadata = metadata_for_test();
2386 let primary_key_codec = build_primary_key_codec(&metadata);
2387 let schema = to_flat_sst_arrow_schema(
2388 &metadata,
2389 &FlatSchemaOptions {
2390 raw_pk_columns: false,
2391 string_pk_use_dict: true,
2392 ..Default::default()
2393 },
2394 );
2395
2396 let capacity = 100;
2397 let mut converter =
2398 BulkPartConverter::new(&metadata, schema, capacity, primary_key_codec, false);
2399
2400 let key_values1 = build_key_values_with_ts_seq_values(
2401 &metadata,
2402 "key1".to_string(),
2403 1u32,
2404 vec![1000, 2000].into_iter(),
2405 vec![Some(1.0), Some(2.0)].into_iter(),
2406 1,
2407 );
2408
2409 let key_values2 = build_key_values_with_ts_seq_values(
2410 &metadata,
2411 "key2".to_string(),
2412 2u32,
2413 vec![1500].into_iter(),
2414 vec![Some(3.0)].into_iter(),
2415 2,
2416 );
2417
2418 converter.append_key_values(&key_values1).unwrap();
2419 converter.append_key_values(&key_values2).unwrap();
2420
2421 let bulk_part = converter.convert().unwrap();
2422
2423 assert_eq!(bulk_part.num_rows(), 3);
2424 assert_eq!(bulk_part.min_timestamp, 1000);
2425 assert_eq!(bulk_part.max_timestamp, 2000);
2426 assert_eq!(bulk_part.sequence, 2);
2427 assert_eq!(bulk_part.timestamp_index, bulk_part.batch.num_columns() - 4);
2428
2429 let schema = bulk_part.batch.schema();
2431 let field_names: Vec<&str> = schema.fields().iter().map(|f| f.name().as_str()).collect();
2432 assert_eq!(
2433 field_names,
2434 vec!["v0", "v1", "ts", "__primary_key", "__sequence", "__op_type"]
2435 );
2436 }
2437
2438 #[allow(clippy::too_many_arguments)]
2439 fn build_key_values_with_sparse_encoding(
2440 metadata: &RegionMetadataRef,
2441 primary_key_codec: &Arc<dyn PrimaryKeyCodec>,
2442 table_id: u32,
2443 tsid: u64,
2444 k0: String,
2445 k1: String,
2446 timestamps: impl Iterator<Item = i64>,
2447 values: impl Iterator<Item = Option<f64>>,
2448 sequence: SequenceNumber,
2449 ) -> KeyValues {
2450 let pk_values = vec![
2452 (ReservedColumnId::table_id(), Value::UInt32(table_id)),
2453 (ReservedColumnId::tsid(), Value::UInt64(tsid)),
2454 (0, Value::String(k0.clone().into())),
2455 (1, Value::String(k1.clone().into())),
2456 ];
2457 let mut encoded_key = Vec::new();
2458 primary_key_codec
2459 .encode_values(&pk_values, &mut encoded_key)
2460 .unwrap();
2461 assert!(!encoded_key.is_empty());
2462
2463 let column_schema = vec![
2465 api::v1::ColumnSchema {
2466 column_name: PRIMARY_KEY_COLUMN_NAME.to_string(),
2467 datatype: api::helper::ColumnDataTypeWrapper::try_from(
2468 ConcreteDataType::binary_datatype(),
2469 )
2470 .unwrap()
2471 .datatype() as i32,
2472 semantic_type: api::v1::SemanticType::Tag as i32,
2473 ..Default::default()
2474 },
2475 api::v1::ColumnSchema {
2476 column_name: "ts".to_string(),
2477 datatype: api::helper::ColumnDataTypeWrapper::try_from(
2478 ConcreteDataType::timestamp_millisecond_datatype(),
2479 )
2480 .unwrap()
2481 .datatype() as i32,
2482 semantic_type: api::v1::SemanticType::Timestamp as i32,
2483 ..Default::default()
2484 },
2485 api::v1::ColumnSchema {
2486 column_name: "v0".to_string(),
2487 datatype: api::helper::ColumnDataTypeWrapper::try_from(
2488 ConcreteDataType::int64_datatype(),
2489 )
2490 .unwrap()
2491 .datatype() as i32,
2492 semantic_type: api::v1::SemanticType::Field as i32,
2493 ..Default::default()
2494 },
2495 api::v1::ColumnSchema {
2496 column_name: "v1".to_string(),
2497 datatype: api::helper::ColumnDataTypeWrapper::try_from(
2498 ConcreteDataType::float64_datatype(),
2499 )
2500 .unwrap()
2501 .datatype() as i32,
2502 semantic_type: api::v1::SemanticType::Field as i32,
2503 ..Default::default()
2504 },
2505 ];
2506
2507 let rows = timestamps
2508 .zip(values)
2509 .map(|(ts, v)| Row {
2510 values: vec![
2511 api::v1::Value {
2512 value_data: Some(api::v1::value::ValueData::BinaryValue(
2513 encoded_key.clone(),
2514 )),
2515 },
2516 api::v1::Value {
2517 value_data: Some(api::v1::value::ValueData::TimestampMillisecondValue(ts)),
2518 },
2519 api::v1::Value {
2520 value_data: Some(api::v1::value::ValueData::I64Value(ts)),
2521 },
2522 api::v1::Value {
2523 value_data: v.map(api::v1::value::ValueData::F64Value),
2524 },
2525 ],
2526 })
2527 .collect();
2528
2529 let mutation = api::v1::Mutation {
2530 op_type: 1,
2531 sequence,
2532 rows: Some(api::v1::Rows {
2533 schema: column_schema,
2534 rows,
2535 }),
2536 write_hint: Some(WriteHint {
2537 primary_key_encoding: api::v1::PrimaryKeyEncoding::Sparse.into(),
2538 }),
2539 };
2540 KeyValues::new(metadata.as_ref(), mutation).unwrap()
2541 }
2542
2543 #[test]
2544 fn test_bulk_part_converter_sparse_primary_key_encoding() {
2545 use api::v1::SemanticType;
2546 use datatypes::schema::ColumnSchema;
2547 use store_api::metadata::{ColumnMetadata, RegionMetadataBuilder};
2548 use store_api::storage::RegionId;
2549
2550 let mut builder = RegionMetadataBuilder::new(RegionId::new(123, 456));
2551 builder
2552 .push_column_metadata(ColumnMetadata {
2553 column_schema: ColumnSchema::new("k0", ConcreteDataType::string_datatype(), false),
2554 semantic_type: SemanticType::Tag,
2555 column_id: 0,
2556 })
2557 .push_column_metadata(ColumnMetadata {
2558 column_schema: ColumnSchema::new("k1", ConcreteDataType::string_datatype(), false),
2559 semantic_type: SemanticType::Tag,
2560 column_id: 1,
2561 })
2562 .push_column_metadata(ColumnMetadata {
2563 column_schema: ColumnSchema::new(
2564 "ts",
2565 ConcreteDataType::timestamp_millisecond_datatype(),
2566 false,
2567 ),
2568 semantic_type: SemanticType::Timestamp,
2569 column_id: 2,
2570 })
2571 .push_column_metadata(ColumnMetadata {
2572 column_schema: ColumnSchema::new("v0", ConcreteDataType::int64_datatype(), true),
2573 semantic_type: SemanticType::Field,
2574 column_id: 3,
2575 })
2576 .push_column_metadata(ColumnMetadata {
2577 column_schema: ColumnSchema::new("v1", ConcreteDataType::float64_datatype(), true),
2578 semantic_type: SemanticType::Field,
2579 column_id: 4,
2580 })
2581 .primary_key(vec![0, 1])
2582 .primary_key_encoding(PrimaryKeyEncoding::Sparse);
2583 let metadata = Arc::new(builder.build().unwrap());
2584
2585 let primary_key_codec = build_primary_key_codec(&metadata);
2586 let schema = to_flat_sst_arrow_schema(
2587 &metadata,
2588 &FlatSchemaOptions::from_encoding(metadata.primary_key_encoding),
2589 );
2590
2591 assert_eq!(metadata.primary_key_encoding, PrimaryKeyEncoding::Sparse);
2592 assert_eq!(primary_key_codec.encoding(), PrimaryKeyEncoding::Sparse);
2593
2594 let capacity = 100;
2595 let mut converter =
2596 BulkPartConverter::new(&metadata, schema, capacity, primary_key_codec.clone(), true);
2597
2598 let key_values1 = build_key_values_with_sparse_encoding(
2599 &metadata,
2600 &primary_key_codec,
2601 2048u32, 100u64, "key11".to_string(),
2604 "key21".to_string(),
2605 vec![1000, 2000].into_iter(),
2606 vec![Some(1.0), Some(2.0)].into_iter(),
2607 1,
2608 );
2609
2610 let key_values2 = build_key_values_with_sparse_encoding(
2611 &metadata,
2612 &primary_key_codec,
2613 4096u32, 200u64, "key12".to_string(),
2616 "key22".to_string(),
2617 vec![1500].into_iter(),
2618 vec![Some(3.0)].into_iter(),
2619 2,
2620 );
2621
2622 converter.append_key_values(&key_values1).unwrap();
2623 converter.append_key_values(&key_values2).unwrap();
2624
2625 let bulk_part = converter.convert().unwrap();
2626
2627 assert_eq!(bulk_part.num_rows(), 3);
2628 assert_eq!(bulk_part.min_timestamp, 1000);
2629 assert_eq!(bulk_part.max_timestamp, 2000);
2630 assert_eq!(bulk_part.sequence, 2);
2631 assert_eq!(bulk_part.timestamp_index, bulk_part.batch.num_columns() - 4);
2632
2633 let schema = bulk_part.batch.schema();
2637 let field_names: Vec<&str> = schema.fields().iter().map(|f| f.name().as_str()).collect();
2638 assert_eq!(
2639 field_names,
2640 vec!["v0", "v1", "ts", "__primary_key", "__sequence", "__op_type"]
2641 );
2642
2643 let primary_key_column = bulk_part.batch.column_by_name("__primary_key").unwrap();
2645 let dict_array = primary_key_column
2646 .as_any()
2647 .downcast_ref::<DictionaryArray<UInt32Type>>()
2648 .unwrap();
2649
2650 assert!(!dict_array.is_empty());
2652 assert_eq!(dict_array.len(), 3); let values = dict_array
2656 .values()
2657 .as_any()
2658 .downcast_ref::<BinaryArray>()
2659 .unwrap();
2660 for i in 0..values.len() {
2661 assert!(
2662 !values.value(i).is_empty(),
2663 "Encoded primary key should not be empty"
2664 );
2665 }
2666 }
2667
2668 #[test]
2669 fn test_convert_bulk_part_empty() {
2670 let metadata = metadata_for_test();
2671 let schema = to_flat_sst_arrow_schema(
2672 &metadata,
2673 &FlatSchemaOptions::from_encoding(metadata.primary_key_encoding),
2674 );
2675 let primary_key_codec = build_primary_key_codec(&metadata);
2676
2677 let empty_batch = RecordBatch::new_empty(schema.clone());
2679 let empty_part = BulkPart {
2680 batch: empty_batch,
2681 max_timestamp: 0,
2682 min_timestamp: 0,
2683 sequence: 0,
2684 timestamp_index: 0,
2685 raw_data: None,
2686 };
2687
2688 let result =
2689 convert_bulk_part(empty_part, &metadata, primary_key_codec, schema, true).unwrap();
2690 assert!(result.is_none());
2691 }
2692
2693 #[test]
2694 fn test_convert_bulk_part_dense_with_pk_columns() {
2695 let metadata = metadata_for_test();
2696 let primary_key_codec = build_primary_key_codec(&metadata);
2697
2698 let k0_array = Arc::new(arrow::array::StringArray::from(vec![
2699 "key1", "key2", "key1",
2700 ]));
2701 let k1_array = Arc::new(arrow::array::UInt32Array::from(vec![1, 2, 1]));
2702 let v0_array = Arc::new(arrow::array::Int64Array::from(vec![100, 200, 300]));
2703 let v1_array = Arc::new(arrow::array::Float64Array::from(vec![1.0, 2.0, 3.0]));
2704 let ts_array = Arc::new(TimestampMillisecondArray::from(vec![1000, 2000, 1500]));
2705
2706 let input_schema = Arc::new(Schema::new(vec![
2707 Field::new("k0", ArrowDataType::Utf8, false),
2708 Field::new("k1", ArrowDataType::UInt32, false),
2709 Field::new("v0", ArrowDataType::Int64, true),
2710 Field::new("v1", ArrowDataType::Float64, true),
2711 Field::new(
2712 "ts",
2713 ArrowDataType::Timestamp(arrow::datatypes::TimeUnit::Millisecond, None),
2714 false,
2715 ),
2716 ]));
2717
2718 let input_batch = RecordBatch::try_new(
2719 input_schema,
2720 vec![k0_array, k1_array, v0_array, v1_array, ts_array],
2721 )
2722 .unwrap();
2723
2724 let part = BulkPart {
2725 batch: input_batch,
2726 max_timestamp: 2000,
2727 min_timestamp: 1000,
2728 sequence: 5,
2729 timestamp_index: 4,
2730 raw_data: None,
2731 };
2732
2733 let output_schema = to_flat_sst_arrow_schema(
2734 &metadata,
2735 &FlatSchemaOptions::from_encoding(metadata.primary_key_encoding),
2736 );
2737
2738 let result = convert_bulk_part(
2739 part,
2740 &metadata,
2741 primary_key_codec,
2742 output_schema,
2743 true, )
2745 .unwrap();
2746
2747 let converted = result.unwrap();
2748
2749 assert_eq!(converted.num_rows(), 3);
2750 assert_eq!(converted.max_timestamp, 2000);
2751 assert_eq!(converted.min_timestamp, 1000);
2752 assert_eq!(converted.sequence, 5);
2753
2754 let schema = converted.batch.schema();
2755 let field_names: Vec<&str> = schema.fields().iter().map(|f| f.name().as_str()).collect();
2756 assert_eq!(
2757 field_names,
2758 vec![
2759 "k0",
2760 "k1",
2761 "v0",
2762 "v1",
2763 "ts",
2764 "__primary_key",
2765 "__sequence",
2766 "__op_type"
2767 ]
2768 );
2769
2770 let k0_col = converted.batch.column_by_name("k0").unwrap();
2771 assert!(matches!(
2772 k0_col.data_type(),
2773 ArrowDataType::Dictionary(_, _)
2774 ));
2775
2776 let pk_col = converted.batch.column_by_name("__primary_key").unwrap();
2777 let dict_array = pk_col
2778 .as_any()
2779 .downcast_ref::<DictionaryArray<UInt32Type>>()
2780 .unwrap();
2781 let keys = dict_array.keys();
2782
2783 assert_eq!(keys.len(), 3);
2784 }
2785
2786 #[test]
2787 fn test_convert_bulk_part_dense_without_pk_columns() {
2788 let metadata = metadata_for_test();
2789 let primary_key_codec = build_primary_key_codec(&metadata);
2790
2791 let k0_array = Arc::new(arrow::array::StringArray::from(vec!["key1", "key2"]));
2793 let k1_array = Arc::new(arrow::array::UInt32Array::from(vec![1, 2]));
2794 let v0_array = Arc::new(arrow::array::Int64Array::from(vec![100, 200]));
2795 let v1_array = Arc::new(arrow::array::Float64Array::from(vec![1.0, 2.0]));
2796 let ts_array = Arc::new(TimestampMillisecondArray::from(vec![1000, 2000]));
2797
2798 let input_schema = Arc::new(Schema::new(vec![
2799 Field::new("k0", ArrowDataType::Utf8, false),
2800 Field::new("k1", ArrowDataType::UInt32, false),
2801 Field::new("v0", ArrowDataType::Int64, true),
2802 Field::new("v1", ArrowDataType::Float64, true),
2803 Field::new(
2804 "ts",
2805 ArrowDataType::Timestamp(arrow::datatypes::TimeUnit::Millisecond, None),
2806 false,
2807 ),
2808 ]));
2809
2810 let input_batch = RecordBatch::try_new(
2811 input_schema,
2812 vec![k0_array, k1_array, v0_array, v1_array, ts_array],
2813 )
2814 .unwrap();
2815
2816 let part = BulkPart {
2817 batch: input_batch,
2818 max_timestamp: 2000,
2819 min_timestamp: 1000,
2820 sequence: 3,
2821 timestamp_index: 4,
2822 raw_data: None,
2823 };
2824
2825 let output_schema = to_flat_sst_arrow_schema(
2826 &metadata,
2827 &FlatSchemaOptions {
2828 raw_pk_columns: false,
2829 string_pk_use_dict: true,
2830 ..Default::default()
2831 },
2832 );
2833
2834 let result = convert_bulk_part(
2835 part,
2836 &metadata,
2837 primary_key_codec,
2838 output_schema,
2839 false, )
2841 .unwrap();
2842
2843 let converted = result.unwrap();
2844
2845 assert_eq!(converted.num_rows(), 2);
2846 assert_eq!(converted.max_timestamp, 2000);
2847 assert_eq!(converted.min_timestamp, 1000);
2848 assert_eq!(converted.sequence, 3);
2849
2850 let schema = converted.batch.schema();
2852 let field_names: Vec<&str> = schema.fields().iter().map(|f| f.name().as_str()).collect();
2853 assert_eq!(
2854 field_names,
2855 vec!["v0", "v1", "ts", "__primary_key", "__sequence", "__op_type"]
2856 );
2857
2858 let pk_col = converted.batch.column_by_name("__primary_key").unwrap();
2860 assert!(matches!(
2861 pk_col.data_type(),
2862 ArrowDataType::Dictionary(_, _)
2863 ));
2864 }
2865
2866 #[test]
2867 fn test_convert_bulk_part_sparse_encoding() {
2868 let mut builder = RegionMetadataBuilder::new(RegionId::new(123, 456));
2869 builder
2870 .push_column_metadata(ColumnMetadata {
2871 column_schema: ColumnSchema::new("k0", ConcreteDataType::string_datatype(), false),
2872 semantic_type: SemanticType::Tag,
2873 column_id: 0,
2874 })
2875 .push_column_metadata(ColumnMetadata {
2876 column_schema: ColumnSchema::new("k1", ConcreteDataType::string_datatype(), false),
2877 semantic_type: SemanticType::Tag,
2878 column_id: 1,
2879 })
2880 .push_column_metadata(ColumnMetadata {
2881 column_schema: ColumnSchema::new(
2882 "ts",
2883 ConcreteDataType::timestamp_millisecond_datatype(),
2884 false,
2885 ),
2886 semantic_type: SemanticType::Timestamp,
2887 column_id: 2,
2888 })
2889 .push_column_metadata(ColumnMetadata {
2890 column_schema: ColumnSchema::new("v0", ConcreteDataType::int64_datatype(), true),
2891 semantic_type: SemanticType::Field,
2892 column_id: 3,
2893 })
2894 .push_column_metadata(ColumnMetadata {
2895 column_schema: ColumnSchema::new("v1", ConcreteDataType::float64_datatype(), true),
2896 semantic_type: SemanticType::Field,
2897 column_id: 4,
2898 })
2899 .primary_key(vec![0, 1])
2900 .primary_key_encoding(PrimaryKeyEncoding::Sparse);
2901 let metadata = Arc::new(builder.build().unwrap());
2902
2903 let primary_key_codec = build_primary_key_codec(&metadata);
2904
2905 let pk_array = Arc::new(arrow::array::BinaryArray::from(vec![
2907 b"encoded_key_1".as_slice(),
2908 b"encoded_key_2".as_slice(),
2909 ]));
2910 let v0_array = Arc::new(arrow::array::Int64Array::from(vec![100, 200]));
2911 let v1_array = Arc::new(arrow::array::Float64Array::from(vec![1.0, 2.0]));
2912 let ts_array = Arc::new(TimestampMillisecondArray::from(vec![1000, 2000]));
2913
2914 let input_schema = Arc::new(Schema::new(vec![
2915 Field::new("__primary_key", ArrowDataType::Binary, false),
2916 Field::new("v0", ArrowDataType::Int64, true),
2917 Field::new("v1", ArrowDataType::Float64, true),
2918 Field::new(
2919 "ts",
2920 ArrowDataType::Timestamp(arrow::datatypes::TimeUnit::Millisecond, None),
2921 false,
2922 ),
2923 ]));
2924
2925 let input_batch =
2926 RecordBatch::try_new(input_schema, vec![pk_array, v0_array, v1_array, ts_array])
2927 .unwrap();
2928
2929 let part = BulkPart {
2930 batch: input_batch,
2931 max_timestamp: 2000,
2932 min_timestamp: 1000,
2933 sequence: 7,
2934 timestamp_index: 3,
2935 raw_data: None,
2936 };
2937
2938 let output_schema = to_flat_sst_arrow_schema(
2939 &metadata,
2940 &FlatSchemaOptions::from_encoding(metadata.primary_key_encoding),
2941 );
2942
2943 let result = convert_bulk_part(
2944 part,
2945 &metadata,
2946 primary_key_codec,
2947 output_schema,
2948 true, )
2950 .unwrap();
2951
2952 let converted = result.unwrap();
2953
2954 assert_eq!(converted.num_rows(), 2);
2955 assert_eq!(converted.max_timestamp, 2000);
2956 assert_eq!(converted.min_timestamp, 1000);
2957 assert_eq!(converted.sequence, 7);
2958
2959 let schema = converted.batch.schema();
2961 let field_names: Vec<&str> = schema.fields().iter().map(|f| f.name().as_str()).collect();
2962 assert_eq!(
2963 field_names,
2964 vec!["v0", "v1", "ts", "__primary_key", "__sequence", "__op_type"]
2965 );
2966
2967 let pk_col = converted.batch.column_by_name("__primary_key").unwrap();
2969 assert!(matches!(
2970 pk_col.data_type(),
2971 ArrowDataType::Dictionary(_, _)
2972 ));
2973 }
2974
2975 #[test]
2976 fn test_convert_bulk_part_sorting_with_multiple_series() {
2977 let metadata = metadata_for_test();
2978 let primary_key_codec = build_primary_key_codec(&metadata);
2979
2980 let k0_array = Arc::new(arrow::array::StringArray::from(vec![
2982 "series_b", "series_a", "series_b", "series_a",
2983 ]));
2984 let k1_array = Arc::new(arrow::array::UInt32Array::from(vec![2, 1, 2, 1]));
2985 let v0_array = Arc::new(arrow::array::Int64Array::from(vec![200, 100, 400, 300]));
2986 let v1_array = Arc::new(arrow::array::Float64Array::from(vec![2.0, 1.0, 4.0, 3.0]));
2987 let ts_array = Arc::new(TimestampMillisecondArray::from(vec![
2988 2000, 1000, 4000, 3000,
2989 ]));
2990
2991 let input_schema = Arc::new(Schema::new(vec![
2992 Field::new("k0", ArrowDataType::Utf8, false),
2993 Field::new("k1", ArrowDataType::UInt32, false),
2994 Field::new("v0", ArrowDataType::Int64, true),
2995 Field::new("v1", ArrowDataType::Float64, true),
2996 Field::new(
2997 "ts",
2998 ArrowDataType::Timestamp(arrow::datatypes::TimeUnit::Millisecond, None),
2999 false,
3000 ),
3001 ]));
3002
3003 let input_batch = RecordBatch::try_new(
3004 input_schema,
3005 vec![k0_array, k1_array, v0_array, v1_array, ts_array],
3006 )
3007 .unwrap();
3008
3009 let part = BulkPart {
3010 batch: input_batch,
3011 max_timestamp: 4000,
3012 min_timestamp: 1000,
3013 sequence: 10,
3014 timestamp_index: 4,
3015 raw_data: None,
3016 };
3017
3018 let output_schema = to_flat_sst_arrow_schema(
3019 &metadata,
3020 &FlatSchemaOptions::from_encoding(metadata.primary_key_encoding),
3021 );
3022
3023 let result =
3024 convert_bulk_part(part, &metadata, primary_key_codec, output_schema, true).unwrap();
3025
3026 let converted = result.unwrap();
3027
3028 assert_eq!(converted.num_rows(), 4);
3029
3030 let ts_col = converted.batch.column(converted.timestamp_index);
3032 let ts_array = ts_col
3033 .as_any()
3034 .downcast_ref::<TimestampMillisecondArray>()
3035 .unwrap();
3036
3037 let timestamps: Vec<i64> = ts_array.values().to_vec();
3041 assert_eq!(timestamps, vec![1000, 3000, 2000, 4000]);
3042 }
3043
3044 fn build_converted_bulk_part(inputs: &[MutationInput]) -> BulkPart {
3046 let metadata = metadata_for_test();
3047 let kvs = inputs
3048 .iter()
3049 .map(|m| {
3050 build_key_values_with_ts_seq_values(
3051 &metadata,
3052 m.k0.to_string(),
3053 m.k1,
3054 m.timestamps.iter().copied(),
3055 m.v1.iter().copied(),
3056 m.sequence,
3057 )
3058 })
3059 .collect::<Vec<_>>();
3060 let schema = to_flat_sst_arrow_schema(&metadata, &FlatSchemaOptions::default());
3061 let primary_key_codec = build_primary_key_codec(&metadata);
3062 let mut converter = BulkPartConverter::new(&metadata, schema, 64, primary_key_codec, true);
3063 for kv in kvs {
3064 converter.append_key_values(&kv).unwrap();
3065 }
3066 converter.convert().unwrap()
3067 }
3068
3069 fn build_multi_bulk_part(groups: &[&[MutationInput]]) -> (MultiBulkPart, RegionMetadataRef) {
3071 let metadata = metadata_for_test();
3072 let mut all_batches = Vec::new();
3073 let mut min_ts = i64::MAX;
3074 let mut max_ts = i64::MIN;
3075 let mut max_seq = 0u64;
3076
3077 for inputs in groups {
3078 let part = build_converted_bulk_part(inputs);
3079 min_ts = min_ts.min(part.min_timestamp);
3080 max_ts = max_ts.max(part.max_timestamp);
3081 max_seq = max_seq.max(part.sequence);
3082 all_batches.push(part.batch);
3083 }
3084
3085 let multi = MultiBulkPart::new(
3086 all_batches,
3087 min_ts,
3088 max_ts,
3089 max_seq,
3090 groups.len(),
3091 &metadata,
3092 );
3093 (multi, metadata)
3094 }
3095
3096 #[test]
3097 fn test_multi_bulk_part_prune_batches() {
3098 let (multi, metadata) = build_multi_bulk_part(&[
3100 &[MutationInput {
3101 k0: "a",
3102 k1: 0,
3103 timestamps: &[1, 2],
3104 v1: &[Some(1.0), Some(2.0)],
3105 sequence: 0,
3106 }],
3107 &[MutationInput {
3108 k0: "m",
3109 k1: 0,
3110 timestamps: &[3, 4],
3111 v1: &[Some(3.0), Some(4.0)],
3112 sequence: 1,
3113 }],
3114 &[MutationInput {
3115 k0: "z",
3116 k1: 0,
3117 timestamps: &[5, 6],
3118 v1: &[Some(5.0), Some(6.0)],
3119 sequence: 2,
3120 }],
3121 ]);
3122 assert_eq!(multi.num_rows(), 6);
3123 assert_eq!(multi.num_batches(), 3);
3124
3125 let context = Arc::new(
3127 BulkIterContext::new(
3128 metadata.clone(),
3129 None,
3130 Some(Predicate::new(vec![
3131 datafusion_expr::col("k0").eq(datafusion_expr::lit("m")),
3132 ])),
3133 false,
3134 crate::sst::parquet::DEFAULT_READ_BATCH_SIZE,
3135 )
3136 .unwrap(),
3137 );
3138 let reader = multi
3139 .read(context, None, None)
3140 .unwrap()
3141 .expect("should have results");
3142 let total_rows: usize = reader.map(|r| r.unwrap().num_rows()).sum();
3143 assert_eq!(total_rows, 2);
3144
3145 let context = Arc::new(
3147 BulkIterContext::new(
3148 metadata.clone(),
3149 None,
3150 Some(Predicate::new(vec![
3151 datafusion_expr::col("k0").eq(datafusion_expr::lit("nonexistent")),
3152 ])),
3153 false,
3154 crate::sst::parquet::DEFAULT_READ_BATCH_SIZE,
3155 )
3156 .unwrap(),
3157 );
3158 assert!(multi.read(context, None, None).unwrap().is_none());
3159
3160 let context = Arc::new(
3162 BulkIterContext::new(
3163 metadata.clone(),
3164 None,
3165 None,
3166 false,
3167 crate::sst::parquet::DEFAULT_READ_BATCH_SIZE,
3168 )
3169 .unwrap(),
3170 );
3171 let reader = multi
3172 .read(context, None, None)
3173 .unwrap()
3174 .expect("should have results");
3175 let total_rows: usize = reader.map(|r| r.unwrap().num_rows()).sum();
3176 assert_eq!(total_rows, 6);
3177 }
3178
3179 #[test]
3180 fn test_fill_missing_columns_resets_raw_data() {
3181 let metadata = metadata_for_test();
3182
3183 let k0_array = Arc::new(arrow::array::StringArray::from(vec!["key1", "key2"]));
3186 let k1_array = Arc::new(arrow::array::UInt32Array::from(vec![1u32, 2]));
3187 let v0_array = Arc::new(arrow::array::Int64Array::from(vec![100i64, 200]));
3188 let ts_array = Arc::new(TimestampMillisecondArray::from(vec![1000i64, 2000]));
3189 let input_schema = Arc::new(Schema::new(vec![
3190 Field::new("k0", ArrowDataType::Utf8, false),
3191 Field::new("k1", ArrowDataType::UInt32, false),
3192 Field::new("v0", ArrowDataType::Int64, true),
3193 Field::new(
3194 "ts",
3195 ArrowDataType::Timestamp(arrow::datatypes::TimeUnit::Millisecond, None),
3196 false,
3197 ),
3198 ]));
3199 let batch =
3200 RecordBatch::try_new(input_schema, vec![k0_array, k1_array, v0_array, ts_array])
3201 .unwrap();
3202
3203 let mut encoder = FlightEncoder::default();
3205 let schema_bytes = encoder.encode_schema(batch.schema().as_ref()).data_header;
3206 let [flight_data] = encoder
3207 .encode(FlightMessage::RecordBatch(batch.clone()))
3208 .try_into()
3209 .unwrap();
3210 let mut part = BulkPart {
3211 batch,
3212 max_timestamp: 2000,
3213 min_timestamp: 1000,
3214 sequence: 5,
3215 timestamp_index: 3,
3216 raw_data: Some(ArrowIpc {
3217 schema: schema_bytes,
3218 data_header: flight_data.data_header,
3219 payload: flight_data.data_body,
3220 }),
3221 };
3222
3223 part.fill_missing_columns(&metadata).unwrap();
3224 assert!(part.batch.column_by_name("v1").is_some());
3225 assert!(part.raw_data.is_none());
3229
3230 let entry = BulkWalEntry::try_from(&part).unwrap();
3232 let replayed = BulkPart::try_from(entry).unwrap();
3233 assert_eq!(2, replayed.num_rows());
3234 assert!(replayed.batch.column_by_name("v1").is_some());
3235 }
3236}