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