1pub mod batch_adapter;
18pub mod compat;
19pub mod dedup;
20pub mod flat_dedup;
21pub mod flat_merge;
22pub mod flat_projection;
23pub mod last_row;
24pub mod projection;
25pub(crate) mod prune;
26pub(crate) mod pruner;
27pub mod range;
28#[cfg(feature = "test")]
29pub mod range_cache;
30#[cfg(not(feature = "test"))]
31pub(crate) mod range_cache;
32pub mod read_columns;
33pub mod scan_region;
34pub mod scan_util;
35pub(crate) mod seq_scan;
36pub(crate) mod series_candidate;
37pub(crate) mod series_reader;
38pub mod series_scan;
39pub mod stream;
40pub(crate) mod unordered_scan;
41
42use std::collections::HashMap;
43use std::sync::Arc;
44use std::time::Duration;
45
46use api::v1::OpType;
47use arrow_schema::SchemaRef;
48use common_time::Timestamp;
49use datafusion_common::arrow::array::UInt8Array;
50use datatypes::arrow;
51use datatypes::arrow::array::{Array, ArrayRef};
52use datatypes::arrow::compute::SortOptions;
53use datatypes::arrow::record_batch::RecordBatch;
54use datatypes::arrow::row::{RowConverter, SortField};
55use datatypes::prelude::{ConcreteDataType, DataType, ScalarVector};
56use datatypes::scalars::ScalarVectorBuilder;
57use datatypes::types::TimestampType;
58use datatypes::value::{Value, ValueRef};
59use datatypes::vectors::{
60 BooleanVector, Helper, TimestampMicrosecondVector, TimestampMillisecondVector,
61 TimestampMillisecondVectorBuilder, TimestampNanosecondVector, TimestampSecondVector,
62 UInt8Vector, UInt8VectorBuilder, UInt32Vector, UInt64Vector, UInt64VectorBuilder, Vector,
63 VectorRef,
64};
65use futures::TryStreamExt;
66use futures::stream::BoxStream;
67use mito_codec::row_converter::{CompositeValues, DensePrimaryKeyCodec, PrimaryKeyCodec};
68use snafu::{OptionExt, ResultExt, ensure};
69use store_api::storage::{ColumnId, SequenceNumber, SequenceRange};
70
71use crate::error::{
72 ComputeArrowSnafu, ComputeVectorSnafu, ConvertVectorSnafu, DecodeSnafu, InvalidBatchSnafu,
73 Result,
74};
75use crate::memtable::BoxedRecordBatchIterator;
76
77pub(crate) fn timestamp_array_to_i64_slice(arr: &ArrayRef) -> &[i64] {
78 use datatypes::arrow::array::{
79 TimestampMicrosecondArray, TimestampMillisecondArray, TimestampNanosecondArray,
80 TimestampSecondArray,
81 };
82 use datatypes::arrow::datatypes::{DataType, TimeUnit};
83
84 match arr.data_type() {
85 DataType::Timestamp(t, _) => match t {
86 TimeUnit::Second => arr
87 .as_any()
88 .downcast_ref::<TimestampSecondArray>()
89 .unwrap()
90 .values(),
91 TimeUnit::Millisecond => arr
92 .as_any()
93 .downcast_ref::<TimestampMillisecondArray>()
94 .unwrap()
95 .values(),
96 TimeUnit::Microsecond => arr
97 .as_any()
98 .downcast_ref::<TimestampMicrosecondArray>()
99 .unwrap()
100 .values(),
101 TimeUnit::Nanosecond => arr
102 .as_any()
103 .downcast_ref::<TimestampNanosecondArray>()
104 .unwrap()
105 .values(),
106 },
107 _ => unreachable!(),
108 }
109}
110
111#[derive(Debug, PartialEq, Clone)]
116pub struct Batch {
117 primary_key: Vec<u8>,
119 pk_values: Option<CompositeValues>,
121 dense_pk_cache: Option<Box<DensePkCache>>,
123 timestamps: VectorRef,
125 sequences: Arc<UInt64Vector>,
129 op_types: Arc<UInt8Vector>,
133 fields: Vec<BatchColumn>,
135 fields_idx: Option<HashMap<ColumnId, usize>>,
137}
138
139#[derive(Debug, PartialEq, Clone)]
141struct DensePkCache {
142 offsets: Vec<usize>,
143 values: Vec<Option<Value>>,
144}
145
146impl Batch {
147 pub fn new(
149 primary_key: Vec<u8>,
150 timestamps: VectorRef,
151 sequences: Arc<UInt64Vector>,
152 op_types: Arc<UInt8Vector>,
153 fields: Vec<BatchColumn>,
154 ) -> Result<Batch> {
155 BatchBuilder::with_required_columns(primary_key, timestamps, sequences, op_types)
156 .with_fields(fields)
157 .build()
158 }
159
160 pub fn with_fields(self, fields: Vec<BatchColumn>) -> Result<Batch> {
162 Batch::new(
163 self.primary_key,
164 self.timestamps,
165 self.sequences,
166 self.op_types,
167 fields,
168 )
169 }
170
171 pub fn primary_key(&self) -> &[u8] {
173 &self.primary_key
174 }
175
176 pub fn pk_values(&self) -> Option<&CompositeValues> {
178 self.pk_values.as_ref()
179 }
180
181 pub fn set_pk_values(&mut self, pk_values: CompositeValues) {
183 self.pk_values = Some(pk_values);
184 self.dense_pk_cache = None;
185 }
186
187 #[cfg(any(test, feature = "test"))]
189 pub fn remove_pk_values(&mut self) {
190 self.pk_values = None;
191 self.dense_pk_cache = None;
192 }
193
194 pub fn fields(&self) -> &[BatchColumn] {
196 &self.fields
197 }
198
199 pub fn timestamps(&self) -> &VectorRef {
201 &self.timestamps
202 }
203
204 pub fn sequences(&self) -> &Arc<UInt64Vector> {
206 &self.sequences
207 }
208
209 pub fn op_types(&self) -> &Arc<UInt8Vector> {
211 &self.op_types
212 }
213
214 pub fn num_rows(&self) -> usize {
216 self.sequences.len()
219 }
220
221 #[allow(dead_code)]
223 pub(crate) fn empty() -> Self {
224 Self {
225 primary_key: vec![],
226 pk_values: None,
227 dense_pk_cache: None,
228 timestamps: Arc::new(TimestampMillisecondVectorBuilder::with_capacity(0).finish()),
229 sequences: Arc::new(UInt64VectorBuilder::with_capacity(0).finish()),
230 op_types: Arc::new(UInt8VectorBuilder::with_capacity(0).finish()),
231 fields: vec![],
232 fields_idx: None,
233 }
234 }
235
236 pub fn is_empty(&self) -> bool {
238 self.num_rows() == 0
239 }
240
241 pub fn first_timestamp(&self) -> Option<Timestamp> {
243 if self.timestamps.is_empty() {
244 return None;
245 }
246
247 Some(self.get_timestamp(0))
248 }
249
250 pub fn last_timestamp(&self) -> Option<Timestamp> {
252 if self.timestamps.is_empty() {
253 return None;
254 }
255
256 Some(self.get_timestamp(self.timestamps.len() - 1))
257 }
258
259 pub fn first_sequence(&self) -> Option<SequenceNumber> {
261 if self.sequences.is_empty() {
262 return None;
263 }
264
265 Some(self.get_sequence(0))
266 }
267
268 pub fn last_sequence(&self) -> Option<SequenceNumber> {
270 if self.sequences.is_empty() {
271 return None;
272 }
273
274 Some(self.get_sequence(self.sequences.len() - 1))
275 }
276
277 pub fn set_primary_key(&mut self, primary_key: Vec<u8>) {
282 self.primary_key = primary_key;
283 self.dense_pk_cache = None;
284 }
285
286 pub fn slice(&self, offset: usize, length: usize) -> Batch {
291 let fields = self
292 .fields
293 .iter()
294 .map(|column| BatchColumn {
295 column_id: column.column_id,
296 data: column.data.slice(offset, length),
297 })
298 .collect();
299 Batch {
301 primary_key: self.primary_key.clone(),
304 pk_values: self.pk_values.clone(),
305 dense_pk_cache: self.dense_pk_cache.clone(),
306 timestamps: self.timestamps.slice(offset, length),
307 sequences: Arc::new(self.sequences.get_slice(offset, length)),
308 op_types: Arc::new(self.op_types.get_slice(offset, length)),
309 fields,
310 fields_idx: self.fields_idx.clone(),
311 }
312 }
313
314 pub fn concat(mut batches: Vec<Batch>) -> Result<Batch> {
318 ensure!(
319 !batches.is_empty(),
320 InvalidBatchSnafu {
321 reason: "empty batches",
322 }
323 );
324 if batches.len() == 1 {
325 return Ok(batches.pop().unwrap());
327 }
328
329 let primary_key = std::mem::take(&mut batches[0].primary_key);
330 let first = &batches[0];
331 ensure!(
333 batches
334 .iter()
335 .skip(1)
336 .all(|b| b.primary_key() == primary_key),
337 InvalidBatchSnafu {
338 reason: "batches have different primary key",
339 }
340 );
341 for b in batches.iter().skip(1) {
342 ensure!(
343 b.fields.len() == first.fields.len(),
344 InvalidBatchSnafu {
345 reason: "batches have different field num",
346 }
347 );
348 for (l, r) in b.fields.iter().zip(&first.fields) {
349 ensure!(
350 l.column_id == r.column_id,
351 InvalidBatchSnafu {
352 reason: "batches have different fields",
353 }
354 );
355 }
356 }
357
358 let mut builder = BatchBuilder::new(primary_key);
360 let array = concat_arrays(batches.iter().map(|b| b.timestamps().to_arrow_array()))?;
362 builder.timestamps_array(array)?;
363 let array = concat_arrays(batches.iter().map(|b| b.sequences().to_arrow_array()))?;
364 builder.sequences_array(array)?;
365 let array = concat_arrays(batches.iter().map(|b| b.op_types().to_arrow_array()))?;
366 builder.op_types_array(array)?;
367 for (i, batch_column) in first.fields.iter().enumerate() {
368 let array = concat_arrays(batches.iter().map(|b| b.fields()[i].data.to_arrow_array()))?;
369 builder.push_field_array(batch_column.column_id, array)?;
370 }
371
372 builder.build()
373 }
374
375 pub fn filter_deleted(&mut self) -> Result<()> {
377 let array = self.op_types.as_arrow();
379 let rhs = UInt8Array::new_scalar(OpType::Delete as u8);
381 let predicate =
382 arrow::compute::kernels::cmp::neq(array, &rhs).context(ComputeArrowSnafu)?;
383 self.filter(&BooleanVector::from(predicate))
384 }
385
386 pub fn filter(&mut self, predicate: &BooleanVector) -> Result<()> {
389 self.timestamps = self
390 .timestamps
391 .filter(predicate)
392 .context(ComputeVectorSnafu)?;
393 self.sequences = Arc::new(
394 UInt64Vector::try_from_arrow_array(
395 arrow::compute::filter(self.sequences.as_arrow(), predicate.as_boolean_array())
396 .context(ComputeArrowSnafu)?,
397 )
398 .unwrap(),
399 );
400 self.op_types = Arc::new(
401 UInt8Vector::try_from_arrow_array(
402 arrow::compute::filter(self.op_types.as_arrow(), predicate.as_boolean_array())
403 .context(ComputeArrowSnafu)?,
404 )
405 .unwrap(),
406 );
407 for batch_column in &mut self.fields {
408 batch_column.data = batch_column
409 .data
410 .filter(predicate)
411 .context(ComputeVectorSnafu)?;
412 }
413
414 Ok(())
415 }
416
417 pub fn filter_by_sequence(&mut self, sequence: Option<SequenceRange>) -> Result<()> {
419 let seq_range = match sequence {
420 None => return Ok(()),
421 Some(seq_range) => {
422 let (Some(first), Some(last)) = (self.first_sequence(), self.last_sequence())
423 else {
424 return Ok(());
425 };
426 let is_subset = match seq_range {
427 SequenceRange::Gt { min } => min < first,
428 SequenceRange::LtEq { max } => max >= last,
429 SequenceRange::GtLtEq { min, max } => min < first && max >= last,
430 };
431 if is_subset {
432 return Ok(());
433 }
434 seq_range
435 }
436 };
437
438 let seqs = self.sequences.as_arrow();
439 let predicate = seq_range.filter(seqs).context(ComputeArrowSnafu)?;
440
441 let predicate = BooleanVector::from(predicate);
442 self.filter(&predicate)?;
443
444 Ok(())
445 }
446
447 pub fn sort(&mut self, dedup: bool) -> Result<()> {
454 let converter = RowConverter::new(vec![
457 SortField::new(self.timestamps.data_type().as_arrow_type()),
458 SortField::new_with_options(
459 self.sequences.data_type().as_arrow_type(),
460 SortOptions {
461 descending: true,
462 ..Default::default()
463 },
464 ),
465 ])
466 .context(ComputeArrowSnafu)?;
467 let columns = [
469 self.timestamps.to_arrow_array(),
470 self.sequences.to_arrow_array(),
471 ];
472 let rows = converter.convert_columns(&columns).unwrap();
473 let mut to_sort: Vec<_> = rows.iter().enumerate().collect();
474
475 let was_sorted = to_sort.is_sorted_by_key(|x| x.1);
476 if !was_sorted {
477 to_sort.sort_unstable_by_key(|x| x.1);
478 }
479
480 let num_rows = to_sort.len();
481 if dedup {
482 to_sort.dedup_by(|left, right| {
484 debug_assert_eq!(18, left.1.as_ref().len());
485 debug_assert_eq!(18, right.1.as_ref().len());
486 let (left_key, right_key) = (left.1.as_ref(), right.1.as_ref());
487 left_key[..TIMESTAMP_KEY_LEN] == right_key[..TIMESTAMP_KEY_LEN]
489 });
490 }
491 let no_dedup = to_sort.len() == num_rows;
492
493 if was_sorted && no_dedup {
494 return Ok(());
495 }
496 let indices = UInt32Vector::from_iter_values(to_sort.iter().map(|v| v.0 as u32));
497 self.take_in_place(&indices)
498 }
499
500 pub(crate) fn merge_last_non_null(&mut self) -> Result<()> {
508 let num_rows = self.num_rows();
509 if num_rows < 2 {
510 return Ok(());
511 }
512
513 let Some(timestamps) = self.timestamps_native() else {
514 return Ok(());
515 };
516
517 let mut has_dup = false;
519 let mut group_count = 1;
520 for i in 1..num_rows {
521 has_dup |= timestamps[i] == timestamps[i - 1];
522 group_count += (timestamps[i] != timestamps[i - 1]) as usize;
523 }
524 if !has_dup {
525 return Ok(());
526 }
527
528 let num_fields = self.fields.len();
529 let op_types = self.op_types.as_arrow().values();
530
531 let mut base_indices: Vec<u32> = Vec::with_capacity(group_count);
532 let mut field_indices: Vec<Vec<u32>> = (0..num_fields)
533 .map(|_| Vec::with_capacity(group_count))
534 .collect();
535
536 let mut start = 0;
537 while start < num_rows {
538 let ts = timestamps[start];
539 let mut end = start + 1;
540 while end < num_rows && timestamps[end] == ts {
541 end += 1;
542 }
543
544 let group_pos = base_indices.len();
545 base_indices.push(start as u32);
546
547 if num_fields > 0 {
548 for idx in &mut field_indices {
550 idx.push(start as u32);
551 }
552
553 let base_deleted = op_types[start] == OpType::Delete as u8;
554 if !base_deleted {
555 let mut missing_fields = Vec::new();
558 for (field_idx, col) in self.fields.iter().enumerate() {
559 if col.data.is_null(start) {
560 missing_fields.push(field_idx);
561 }
562 }
563
564 if !missing_fields.is_empty() {
565 for row_idx in (start + 1)..end {
566 if op_types[row_idx] == OpType::Delete as u8 {
567 break;
568 }
569
570 missing_fields.retain(|&field_idx| {
571 if self.fields[field_idx].data.is_null(row_idx) {
572 true
573 } else {
574 field_indices[field_idx][group_pos] = row_idx as u32;
575 false
576 }
577 });
578
579 if missing_fields.is_empty() {
580 break;
581 }
582 }
583 }
584 }
585 }
586
587 start = end;
588 }
589
590 let base_indices = UInt32Vector::from_vec(base_indices);
591 self.timestamps = self
592 .timestamps
593 .take(&base_indices)
594 .context(ComputeVectorSnafu)?;
595 let array = arrow::compute::take(self.sequences.as_arrow(), base_indices.as_arrow(), None)
596 .context(ComputeArrowSnafu)?;
597 self.sequences = Arc::new(UInt64Vector::try_from_arrow_array(array).unwrap());
599 let array = arrow::compute::take(self.op_types.as_arrow(), base_indices.as_arrow(), None)
600 .context(ComputeArrowSnafu)?;
601 self.op_types = Arc::new(UInt8Vector::try_from_arrow_array(array).unwrap());
603
604 for (field_idx, batch_column) in self.fields.iter_mut().enumerate() {
605 let idx = UInt32Vector::from_vec(std::mem::take(&mut field_indices[field_idx]));
606 batch_column.data = batch_column.data.take(&idx).context(ComputeVectorSnafu)?;
607 }
608
609 Ok(())
610 }
611
612 pub fn memory_size(&self) -> usize {
614 let mut size = std::mem::size_of::<Self>();
615 size += self.primary_key.len();
616 size += self.timestamps.memory_size();
617 size += self.sequences.memory_size();
618 size += self.op_types.memory_size();
619 for batch_column in &self.fields {
620 size += batch_column.data.memory_size();
621 }
622 size
623 }
624
625 pub(crate) fn timestamps_native(&self) -> Option<&[i64]> {
627 if self.timestamps.is_empty() {
628 return None;
629 }
630
631 let values = match self.timestamps.data_type() {
632 ConcreteDataType::Timestamp(TimestampType::Second(_)) => self
633 .timestamps
634 .as_any()
635 .downcast_ref::<TimestampSecondVector>()
636 .unwrap()
637 .as_arrow()
638 .values(),
639 ConcreteDataType::Timestamp(TimestampType::Millisecond(_)) => self
640 .timestamps
641 .as_any()
642 .downcast_ref::<TimestampMillisecondVector>()
643 .unwrap()
644 .as_arrow()
645 .values(),
646 ConcreteDataType::Timestamp(TimestampType::Microsecond(_)) => self
647 .timestamps
648 .as_any()
649 .downcast_ref::<TimestampMicrosecondVector>()
650 .unwrap()
651 .as_arrow()
652 .values(),
653 ConcreteDataType::Timestamp(TimestampType::Nanosecond(_)) => self
654 .timestamps
655 .as_any()
656 .downcast_ref::<TimestampNanosecondVector>()
657 .unwrap()
658 .as_arrow()
659 .values(),
660 other => panic!("timestamps in a Batch has other type {:?}", other),
661 };
662
663 Some(values)
664 }
665
666 fn take_in_place(&mut self, indices: &UInt32Vector) -> Result<()> {
668 self.timestamps = self.timestamps.take(indices).context(ComputeVectorSnafu)?;
669 let array = arrow::compute::take(self.sequences.as_arrow(), indices.as_arrow(), None)
670 .context(ComputeArrowSnafu)?;
671 self.sequences = Arc::new(UInt64Vector::try_from_arrow_array(array).unwrap());
673 let array = arrow::compute::take(self.op_types.as_arrow(), indices.as_arrow(), None)
674 .context(ComputeArrowSnafu)?;
675 self.op_types = Arc::new(UInt8Vector::try_from_arrow_array(array).unwrap());
676 for batch_column in &mut self.fields {
677 batch_column.data = batch_column
678 .data
679 .take(indices)
680 .context(ComputeVectorSnafu)?;
681 }
682
683 Ok(())
684 }
685
686 fn get_timestamp(&self, index: usize) -> Timestamp {
691 match self.timestamps.get_ref(index) {
692 ValueRef::Timestamp(timestamp) => timestamp,
693
694 value => panic!("{:?} is not a timestamp", value),
696 }
697 }
698
699 pub(crate) fn get_sequence(&self, index: usize) -> SequenceNumber {
704 self.sequences.get_data(index).unwrap()
706 }
707
708 pub(crate) fn ensure_dense_pk_decoded(&mut self, codec: &DensePrimaryKeyCodec) -> Result<()> {
711 if self.pk_values.is_none() {
712 let mut values = Vec::with_capacity(codec.num_fields());
715 for value in codec.decode_dense_iter(&self.primary_key) {
716 values.push(value.context(DecodeSnafu)?);
717 }
718 self.set_pk_values(CompositeValues::Dense(values));
719 }
720 Ok(())
721 }
722
723 pub fn pk_col_value(
729 &mut self,
730 codec: &dyn PrimaryKeyCodec,
731 col_idx_in_pk: usize,
732 column_id: ColumnId,
733 ) -> Result<Option<&Value>> {
734 if self.pk_values.is_none() {
735 if let Some(codec) = codec.as_dense() {
736 if col_idx_in_pk >= codec.num_fields() {
737 return Ok(None);
738 }
739 let cache = self.dense_pk_cache.get_or_insert_with(|| {
740 Box::new(DensePkCache {
741 offsets: Vec::new(),
742 values: vec![None; codec.num_fields()],
743 })
744 });
745 if cache.values[col_idx_in_pk].is_none() {
746 let value = codec
747 .decode_value_at(&self.primary_key, col_idx_in_pk, &mut cache.offsets)
748 .context(DecodeSnafu)?;
749 cache.values[col_idx_in_pk] = Some(value);
750 }
751 return Ok(cache.values[col_idx_in_pk].as_ref());
752 }
753 self.pk_values = Some(codec.decode(&self.primary_key).context(DecodeSnafu)?);
754 }
755
756 let pk_values = self.pk_values.as_ref().unwrap();
757 Ok(match pk_values {
758 CompositeValues::Dense(values) => values.get(col_idx_in_pk).map(|(_, v)| v),
759 CompositeValues::Sparse(values) => values.get(&column_id),
760 })
761 }
762
763 pub fn field_col_value(&mut self, column_id: ColumnId) -> Option<&BatchColumn> {
767 if self.fields_idx.is_none() {
768 self.fields_idx = Some(
769 self.fields
770 .iter()
771 .enumerate()
772 .map(|(i, c)| (c.column_id, i))
773 .collect(),
774 );
775 }
776
777 self.fields_idx
778 .as_ref()
779 .unwrap()
780 .get(&column_id)
781 .map(|&idx| &self.fields[idx])
782 }
783}
784
785const TIMESTAMP_KEY_LEN: usize = 9;
787
788fn concat_arrays(iter: impl Iterator<Item = ArrayRef>) -> Result<ArrayRef> {
790 let arrays: Vec<_> = iter.collect();
791 let dyn_arrays: Vec<_> = arrays.iter().map(|array| array.as_ref()).collect();
792 arrow::compute::concat(&dyn_arrays).context(ComputeArrowSnafu)
793}
794
795#[derive(Debug, PartialEq, Eq, Clone)]
797pub struct BatchColumn {
798 pub column_id: ColumnId,
800 pub data: VectorRef,
802}
803
804pub struct BatchBuilder {
806 primary_key: Vec<u8>,
807 timestamps: Option<VectorRef>,
808 sequences: Option<Arc<UInt64Vector>>,
809 op_types: Option<Arc<UInt8Vector>>,
810 fields: Vec<BatchColumn>,
811}
812
813impl BatchBuilder {
814 pub fn new(primary_key: Vec<u8>) -> BatchBuilder {
816 BatchBuilder {
817 primary_key,
818 timestamps: None,
819 sequences: None,
820 op_types: None,
821 fields: Vec::new(),
822 }
823 }
824
825 pub fn with_required_columns(
827 primary_key: Vec<u8>,
828 timestamps: VectorRef,
829 sequences: Arc<UInt64Vector>,
830 op_types: Arc<UInt8Vector>,
831 ) -> BatchBuilder {
832 BatchBuilder {
833 primary_key,
834 timestamps: Some(timestamps),
835 sequences: Some(sequences),
836 op_types: Some(op_types),
837 fields: Vec::new(),
838 }
839 }
840
841 pub fn with_fields(mut self, fields: Vec<BatchColumn>) -> Self {
843 self.fields = fields;
844 self
845 }
846
847 pub fn push_field(&mut self, column: BatchColumn) -> &mut Self {
849 self.fields.push(column);
850 self
851 }
852
853 pub fn push_field_array(&mut self, column_id: ColumnId, array: ArrayRef) -> Result<&mut Self> {
855 let vector = Helper::try_into_vector(array).context(ConvertVectorSnafu)?;
856 self.fields.push(BatchColumn {
857 column_id,
858 data: vector,
859 });
860
861 Ok(self)
862 }
863
864 pub fn timestamps_array(&mut self, array: ArrayRef) -> Result<&mut Self> {
866 let vector = Helper::try_into_vector(array).context(ConvertVectorSnafu)?;
867 ensure!(
868 vector.data_type().is_timestamp(),
869 InvalidBatchSnafu {
870 reason: format!("{:?} is not a timestamp type", vector.data_type()),
871 }
872 );
873
874 self.timestamps = Some(vector);
875 Ok(self)
876 }
877
878 pub fn sequences_array(&mut self, array: ArrayRef) -> Result<&mut Self> {
880 ensure!(
881 *array.data_type() == arrow::datatypes::DataType::UInt64,
882 InvalidBatchSnafu {
883 reason: "sequence array is not UInt64 type",
884 }
885 );
886 let vector = Arc::new(UInt64Vector::try_from_arrow_array(array).unwrap());
888 self.sequences = Some(vector);
889
890 Ok(self)
891 }
892
893 pub fn op_types_array(&mut self, array: ArrayRef) -> Result<&mut Self> {
895 ensure!(
896 *array.data_type() == arrow::datatypes::DataType::UInt8,
897 InvalidBatchSnafu {
898 reason: "sequence array is not UInt8 type",
899 }
900 );
901 let vector = Arc::new(UInt8Vector::try_from_arrow_array(array).unwrap());
903 self.op_types = Some(vector);
904
905 Ok(self)
906 }
907
908 pub fn build(self) -> Result<Batch> {
910 let timestamps = self.timestamps.context(InvalidBatchSnafu {
911 reason: "missing timestamps",
912 })?;
913 let sequences = self.sequences.context(InvalidBatchSnafu {
914 reason: "missing sequences",
915 })?;
916 let op_types = self.op_types.context(InvalidBatchSnafu {
917 reason: "missing op_types",
918 })?;
919 assert_eq!(0, timestamps.null_count());
922 assert_eq!(0, sequences.null_count());
923 assert_eq!(0, op_types.null_count());
924
925 let ts_len = timestamps.len();
926 ensure!(
927 sequences.len() == ts_len,
928 InvalidBatchSnafu {
929 reason: format!(
930 "sequence have different len {} != {}",
931 sequences.len(),
932 ts_len
933 ),
934 }
935 );
936 ensure!(
937 op_types.len() == ts_len,
938 InvalidBatchSnafu {
939 reason: format!(
940 "op type have different len {} != {}",
941 op_types.len(),
942 ts_len
943 ),
944 }
945 );
946 for column in &self.fields {
947 ensure!(
948 column.data.len() == ts_len,
949 InvalidBatchSnafu {
950 reason: format!(
951 "column {} has different len {} != {}",
952 column.column_id,
953 column.data.len(),
954 ts_len
955 ),
956 }
957 );
958 }
959
960 Ok(Batch {
961 primary_key: self.primary_key,
962 pk_values: None,
963 dense_pk_cache: None,
964 timestamps,
965 sequences,
966 op_types,
967 fields: self.fields,
968 fields_idx: None,
969 })
970 }
971}
972
973impl From<Batch> for BatchBuilder {
974 fn from(batch: Batch) -> Self {
975 Self {
976 primary_key: batch.primary_key,
977 timestamps: Some(batch.timestamps),
978 sequences: Some(batch.sequences),
979 op_types: Some(batch.op_types),
980 fields: batch.fields,
981 }
982 }
983}
984
985pub struct FlatSource {
987 schema: SchemaRef,
988 inner: FlatSourceInner,
989}
990
991impl FlatSource {
992 pub fn new_iter(schema: SchemaRef, iter: BoxedRecordBatchIterator) -> Self {
994 Self {
995 schema,
996 inner: FlatSourceInner::Iter(iter),
997 }
998 }
999
1000 pub fn new_stream(schema: SchemaRef, stream: BoxedRecordBatchStream) -> Self {
1002 Self {
1003 schema,
1004 inner: FlatSourceInner::Stream(stream),
1005 }
1006 }
1007
1008 pub(crate) fn schema(&self) -> &SchemaRef {
1009 &self.schema
1010 }
1011
1012 pub async fn next_batch(&mut self) -> Result<Option<RecordBatch>> {
1013 self.inner.next_batch().await
1014 }
1015
1016 #[cfg(test)]
1017 pub(crate) fn take_iter(self) -> BoxedRecordBatchIterator {
1018 match self.inner {
1019 FlatSourceInner::Iter(iter) => iter,
1020 FlatSourceInner::Stream(_) => unreachable!(),
1021 }
1022 }
1023}
1024
1025enum FlatSourceInner {
1026 Iter(BoxedRecordBatchIterator),
1028 Stream(BoxedRecordBatchStream),
1030}
1031
1032impl FlatSourceInner {
1033 pub async fn next_batch(&mut self) -> Result<Option<RecordBatch>> {
1035 match self {
1036 Self::Iter(iter) => iter.next().transpose(),
1037 Self::Stream(stream) => stream.try_next().await,
1038 }
1039 }
1040}
1041
1042pub type BoxedRecordBatchStream = BoxStream<'static, Result<RecordBatch>>;
1044
1045#[derive(Debug, Default)]
1047pub(crate) struct ScannerMetrics {
1048 scan_cost: Duration,
1050 yield_cost: Duration,
1052 num_batches: usize,
1054 num_rows: usize,
1056}
1057
1058#[cfg(test)]
1059mod tests {
1060 use datatypes::arrow::array::{TimestampMillisecondArray, UInt8Array, UInt64Array};
1061 use mito_codec::row_converter::{self, build_primary_key_codec_with_fields};
1062 use store_api::codec::PrimaryKeyEncoding;
1063 use store_api::storage::consts::ReservedColumnId;
1064
1065 use super::*;
1066 use crate::error::Error;
1067 use crate::test_util::new_batch_builder;
1068
1069 #[test]
1070 fn dense_pk_columns_are_lazy_and_reset_with_the_key() {
1071 use mito_codec::row_converter::{DensePrimaryKeyCodec, PrimaryKeyCodecExt, SortField};
1072
1073 let codec = DensePrimaryKeyCodec::with_fields(vec![
1074 (7, SortField::new(ConcreteDataType::string_datatype())),
1075 (3, SortField::new(ConcreteDataType::int64_datatype())),
1076 (9, SortField::new(ConcreteDataType::string_datatype())),
1077 ]);
1078 let values = [
1079 Value::from("䏿–‡\0abcdefgh"),
1080 Value::Int64(-42),
1081 Value::Null,
1082 ];
1083 let key = codec
1084 .encode(values.iter().map(Value::as_value_ref))
1085 .unwrap();
1086 let mut batch = new_batch_without_fields(&[1, 2], &[1, 1], &[OpType::Put, OpType::Put]);
1087 batch.set_primary_key(key);
1088 for pos in [2, 0, 1, 2] {
1089 assert_eq!(
1090 batch.pk_col_value(&codec, pos, [7, 3, 9][pos]).unwrap(),
1091 Some(&values[pos])
1092 );
1093 }
1094 assert!(batch.pk_col_value(&codec, 3, 99).unwrap().is_none());
1095 let mut sliced = batch.slice(1, 1);
1096 assert_eq!(sliced.pk_col_value(&codec, 0, 7).unwrap(), Some(&values[0]));
1097
1098 batch.set_primary_key(vec![0, 1]);
1100 assert_eq!(
1101 batch.pk_col_value(&codec, 0, 7).unwrap(),
1102 Some(&Value::Null)
1103 );
1104 assert!(batch.pk_col_value(&codec, 1, 3).is_err());
1105 batch.set_primary_key(
1106 codec
1107 .encode(
1108 [
1109 ValueRef::String(""),
1110 ValueRef::Int64(8),
1111 ValueRef::String("new"),
1112 ]
1113 .into_iter(),
1114 )
1115 .unwrap(),
1116 );
1117 assert_eq!(
1118 batch.pk_col_value(&codec, 2, 9).unwrap(),
1119 Some(&Value::from("new"))
1120 );
1121 assert_eq!(
1122 batch.pk_col_value(&codec, 0, 7).unwrap(),
1123 Some(&Value::from(""))
1124 );
1125
1126 batch.set_pk_values(CompositeValues::Dense(vec![(7, Value::from("default"))]));
1128 batch.ensure_dense_pk_decoded(&codec).unwrap();
1129 assert_eq!(
1130 batch.pk_col_value(&codec, 0, 7).unwrap(),
1131 Some(&Value::from("default"))
1132 );
1133 assert!(batch.pk_col_value(&codec, 1, 3).unwrap().is_none());
1134 batch.remove_pk_values();
1135 batch.ensure_dense_pk_decoded(&codec).unwrap();
1136 assert_eq!(
1137 batch.pk_col_value(&codec, 1, 3).unwrap(),
1138 Some(&Value::Int64(8))
1139 );
1140 }
1141
1142 fn new_batch(
1143 timestamps: &[i64],
1144 sequences: &[u64],
1145 op_types: &[OpType],
1146 field: &[u64],
1147 ) -> Batch {
1148 new_batch_builder(b"test", timestamps, sequences, op_types, 1, field)
1149 .build()
1150 .unwrap()
1151 }
1152
1153 fn new_batch_with_u64_fields(
1154 timestamps: &[i64],
1155 sequences: &[u64],
1156 op_types: &[OpType],
1157 fields: &[(ColumnId, &[Option<u64>])],
1158 ) -> Batch {
1159 assert_eq!(timestamps.len(), sequences.len());
1160 assert_eq!(timestamps.len(), op_types.len());
1161 for (_, values) in fields {
1162 assert_eq!(timestamps.len(), values.len());
1163 }
1164
1165 let mut builder = BatchBuilder::new(b"test".to_vec());
1166 builder
1167 .timestamps_array(Arc::new(TimestampMillisecondArray::from_iter_values(
1168 timestamps.iter().copied(),
1169 )))
1170 .unwrap()
1171 .sequences_array(Arc::new(UInt64Array::from_iter_values(
1172 sequences.iter().copied(),
1173 )))
1174 .unwrap()
1175 .op_types_array(Arc::new(UInt8Array::from_iter_values(
1176 op_types.iter().map(|v| *v as u8),
1177 )))
1178 .unwrap();
1179
1180 for (col_id, values) in fields {
1181 builder
1182 .push_field_array(*col_id, Arc::new(UInt64Array::from(values.to_vec())))
1183 .unwrap();
1184 }
1185
1186 builder.build().unwrap()
1187 }
1188
1189 fn new_batch_without_fields(
1190 timestamps: &[i64],
1191 sequences: &[u64],
1192 op_types: &[OpType],
1193 ) -> Batch {
1194 assert_eq!(timestamps.len(), sequences.len());
1195 assert_eq!(timestamps.len(), op_types.len());
1196
1197 let mut builder = BatchBuilder::new(b"test".to_vec());
1198 builder
1199 .timestamps_array(Arc::new(TimestampMillisecondArray::from_iter_values(
1200 timestamps.iter().copied(),
1201 )))
1202 .unwrap()
1203 .sequences_array(Arc::new(UInt64Array::from_iter_values(
1204 sequences.iter().copied(),
1205 )))
1206 .unwrap()
1207 .op_types_array(Arc::new(UInt8Array::from_iter_values(
1208 op_types.iter().map(|v| *v as u8),
1209 )))
1210 .unwrap();
1211
1212 builder.build().unwrap()
1213 }
1214
1215 #[test]
1216 fn test_empty_batch() {
1217 let batch = Batch::empty();
1218 assert!(batch.is_empty());
1219 assert_eq!(None, batch.first_timestamp());
1220 assert_eq!(None, batch.last_timestamp());
1221 assert_eq!(None, batch.first_sequence());
1222 assert_eq!(None, batch.last_sequence());
1223 assert!(batch.timestamps_native().is_none());
1224 }
1225
1226 #[test]
1227 fn test_first_last_one() {
1228 let batch = new_batch(&[1], &[2], &[OpType::Put], &[4]);
1229 assert_eq!(
1230 Timestamp::new_millisecond(1),
1231 batch.first_timestamp().unwrap()
1232 );
1233 assert_eq!(
1234 Timestamp::new_millisecond(1),
1235 batch.last_timestamp().unwrap()
1236 );
1237 assert_eq!(2, batch.first_sequence().unwrap());
1238 assert_eq!(2, batch.last_sequence().unwrap());
1239 }
1240
1241 #[test]
1242 fn test_first_last_multiple() {
1243 let batch = new_batch(
1244 &[1, 2, 3],
1245 &[11, 12, 13],
1246 &[OpType::Put, OpType::Put, OpType::Put],
1247 &[21, 22, 23],
1248 );
1249 assert_eq!(
1250 Timestamp::new_millisecond(1),
1251 batch.first_timestamp().unwrap()
1252 );
1253 assert_eq!(
1254 Timestamp::new_millisecond(3),
1255 batch.last_timestamp().unwrap()
1256 );
1257 assert_eq!(11, batch.first_sequence().unwrap());
1258 assert_eq!(13, batch.last_sequence().unwrap());
1259 }
1260
1261 #[test]
1262 fn test_slice() {
1263 let batch = new_batch(
1264 &[1, 2, 3, 4],
1265 &[11, 12, 13, 14],
1266 &[OpType::Put, OpType::Delete, OpType::Put, OpType::Put],
1267 &[21, 22, 23, 24],
1268 );
1269 let batch = batch.slice(1, 2);
1270 let expect = new_batch(
1271 &[2, 3],
1272 &[12, 13],
1273 &[OpType::Delete, OpType::Put],
1274 &[22, 23],
1275 );
1276 assert_eq!(expect, batch);
1277 }
1278
1279 #[test]
1280 fn test_timestamps_native() {
1281 let batch = new_batch(
1282 &[1, 2, 3, 4],
1283 &[11, 12, 13, 14],
1284 &[OpType::Put, OpType::Delete, OpType::Put, OpType::Put],
1285 &[21, 22, 23, 24],
1286 );
1287 assert_eq!(&[1, 2, 3, 4], batch.timestamps_native().unwrap());
1288 }
1289
1290 #[test]
1291 fn test_concat_empty() {
1292 let err = Batch::concat(vec![]).unwrap_err();
1293 assert!(
1294 matches!(err, Error::InvalidBatch { .. }),
1295 "unexpected err: {err}"
1296 );
1297 }
1298
1299 #[test]
1300 fn test_concat_one() {
1301 let batch = new_batch(&[], &[], &[], &[]);
1302 let actual = Batch::concat(vec![batch.clone()]).unwrap();
1303 assert_eq!(batch, actual);
1304
1305 let batch = new_batch(&[1, 2], &[11, 12], &[OpType::Put, OpType::Put], &[21, 22]);
1306 let actual = Batch::concat(vec![batch.clone()]).unwrap();
1307 assert_eq!(batch, actual);
1308 }
1309
1310 #[test]
1311 fn test_concat_multiple() {
1312 let batches = vec![
1313 new_batch(&[1, 2], &[11, 12], &[OpType::Put, OpType::Put], &[21, 22]),
1314 new_batch(
1315 &[3, 4, 5],
1316 &[13, 14, 15],
1317 &[OpType::Put, OpType::Delete, OpType::Put],
1318 &[23, 24, 25],
1319 ),
1320 new_batch(&[], &[], &[], &[]),
1321 new_batch(&[6], &[16], &[OpType::Put], &[26]),
1322 ];
1323 let batch = Batch::concat(batches).unwrap();
1324 let expect = new_batch(
1325 &[1, 2, 3, 4, 5, 6],
1326 &[11, 12, 13, 14, 15, 16],
1327 &[
1328 OpType::Put,
1329 OpType::Put,
1330 OpType::Put,
1331 OpType::Delete,
1332 OpType::Put,
1333 OpType::Put,
1334 ],
1335 &[21, 22, 23, 24, 25, 26],
1336 );
1337 assert_eq!(expect, batch);
1338 }
1339
1340 #[test]
1341 fn test_concat_different() {
1342 let batch1 = new_batch(&[1], &[1], &[OpType::Put], &[1]);
1343 let mut batch2 = new_batch(&[2], &[2], &[OpType::Put], &[2]);
1344 batch2.primary_key = b"hello".to_vec();
1345 let err = Batch::concat(vec![batch1, batch2]).unwrap_err();
1346 assert!(
1347 matches!(err, Error::InvalidBatch { .. }),
1348 "unexpected err: {err}"
1349 );
1350 }
1351
1352 #[test]
1353 fn test_concat_different_fields() {
1354 let batch1 = new_batch(&[1], &[1], &[OpType::Put], &[1]);
1355 let fields = vec![
1356 batch1.fields()[0].clone(),
1357 BatchColumn {
1358 column_id: 2,
1359 data: Arc::new(UInt64Vector::from_slice([2])),
1360 },
1361 ];
1362 let batch2 = batch1.clone().with_fields(fields).unwrap();
1364 let err = Batch::concat(vec![batch1.clone(), batch2]).unwrap_err();
1365 assert!(
1366 matches!(err, Error::InvalidBatch { .. }),
1367 "unexpected err: {err}"
1368 );
1369
1370 let fields = vec![BatchColumn {
1372 column_id: 2,
1373 data: Arc::new(UInt64Vector::from_slice([2])),
1374 }];
1375 let batch2 = batch1.clone().with_fields(fields).unwrap();
1376 let err = Batch::concat(vec![batch1, batch2]).unwrap_err();
1377 assert!(
1378 matches!(err, Error::InvalidBatch { .. }),
1379 "unexpected err: {err}"
1380 );
1381 }
1382
1383 #[test]
1384 fn test_filter_deleted_empty() {
1385 let mut batch = new_batch(&[], &[], &[], &[]);
1386 batch.filter_deleted().unwrap();
1387 assert!(batch.is_empty());
1388 }
1389
1390 #[test]
1391 fn test_filter_deleted() {
1392 let mut batch = new_batch(
1393 &[1, 2, 3, 4],
1394 &[11, 12, 13, 14],
1395 &[OpType::Delete, OpType::Put, OpType::Delete, OpType::Put],
1396 &[21, 22, 23, 24],
1397 );
1398 batch.filter_deleted().unwrap();
1399 let expect = new_batch(&[2, 4], &[12, 14], &[OpType::Put, OpType::Put], &[22, 24]);
1400 assert_eq!(expect, batch);
1401
1402 let mut batch = new_batch(
1403 &[1, 2, 3, 4],
1404 &[11, 12, 13, 14],
1405 &[OpType::Put, OpType::Put, OpType::Put, OpType::Put],
1406 &[21, 22, 23, 24],
1407 );
1408 let expect = batch.clone();
1409 batch.filter_deleted().unwrap();
1410 assert_eq!(expect, batch);
1411 }
1412
1413 #[test]
1414 fn test_filter_by_sequence() {
1415 let mut batch = new_batch(
1417 &[1, 2, 3, 4],
1418 &[11, 12, 13, 14],
1419 &[OpType::Put, OpType::Put, OpType::Put, OpType::Put],
1420 &[21, 22, 23, 24],
1421 );
1422 batch
1423 .filter_by_sequence(Some(SequenceRange::LtEq { max: 13 }))
1424 .unwrap();
1425 let expect = new_batch(
1426 &[1, 2, 3],
1427 &[11, 12, 13],
1428 &[OpType::Put, OpType::Put, OpType::Put],
1429 &[21, 22, 23],
1430 );
1431 assert_eq!(expect, batch);
1432
1433 let mut batch = new_batch(
1435 &[1, 2, 3, 4],
1436 &[11, 12, 13, 14],
1437 &[OpType::Put, OpType::Delete, OpType::Put, OpType::Put],
1438 &[21, 22, 23, 24],
1439 );
1440
1441 batch
1442 .filter_by_sequence(Some(SequenceRange::LtEq { max: 10 }))
1443 .unwrap();
1444 assert!(batch.is_empty());
1445
1446 let mut batch = new_batch(
1448 &[1, 2, 3, 4],
1449 &[11, 12, 13, 14],
1450 &[OpType::Put, OpType::Delete, OpType::Put, OpType::Put],
1451 &[21, 22, 23, 24],
1452 );
1453 let expect = batch.clone();
1454 batch.filter_by_sequence(None).unwrap();
1455 assert_eq!(expect, batch);
1456
1457 let mut batch = new_batch(&[], &[], &[], &[]);
1459 batch
1460 .filter_by_sequence(Some(SequenceRange::LtEq { max: 10 }))
1461 .unwrap();
1462 assert!(batch.is_empty());
1463
1464 let mut batch = new_batch(&[], &[], &[], &[]);
1466 batch.filter_by_sequence(None).unwrap();
1467 assert!(batch.is_empty());
1468
1469 let mut batch = new_batch(
1471 &[1, 2, 3, 4],
1472 &[11, 12, 13, 14],
1473 &[OpType::Put, OpType::Put, OpType::Put, OpType::Put],
1474 &[21, 22, 23, 24],
1475 );
1476 batch
1477 .filter_by_sequence(Some(SequenceRange::Gt { min: 12 }))
1478 .unwrap();
1479 let expect = new_batch(&[3, 4], &[13, 14], &[OpType::Put, OpType::Put], &[23, 24]);
1480 assert_eq!(expect, batch);
1481
1482 let mut batch = new_batch(
1484 &[1, 2, 3, 4],
1485 &[11, 12, 13, 14],
1486 &[OpType::Put, OpType::Delete, OpType::Put, OpType::Put],
1487 &[21, 22, 23, 24],
1488 );
1489 batch
1490 .filter_by_sequence(Some(SequenceRange::Gt { min: 20 }))
1491 .unwrap();
1492 assert!(batch.is_empty());
1493
1494 let mut batch = new_batch(
1496 &[1, 2, 3, 4, 5],
1497 &[11, 12, 13, 14, 15],
1498 &[
1499 OpType::Put,
1500 OpType::Put,
1501 OpType::Put,
1502 OpType::Put,
1503 OpType::Put,
1504 ],
1505 &[21, 22, 23, 24, 25],
1506 );
1507 batch
1508 .filter_by_sequence(Some(SequenceRange::GtLtEq { min: 12, max: 14 }))
1509 .unwrap();
1510 let expect = new_batch(&[3, 4], &[13, 14], &[OpType::Put, OpType::Put], &[23, 24]);
1511 assert_eq!(expect, batch);
1512
1513 let mut batch = new_batch(
1515 &[1, 2, 3, 4, 5],
1516 &[11, 12, 13, 14, 15],
1517 &[
1518 OpType::Put,
1519 OpType::Delete,
1520 OpType::Put,
1521 OpType::Delete,
1522 OpType::Put,
1523 ],
1524 &[21, 22, 23, 24, 25],
1525 );
1526 batch
1527 .filter_by_sequence(Some(SequenceRange::GtLtEq { min: 11, max: 13 }))
1528 .unwrap();
1529 let expect = new_batch(
1530 &[2, 3],
1531 &[12, 13],
1532 &[OpType::Delete, OpType::Put],
1533 &[22, 23],
1534 );
1535 assert_eq!(expect, batch);
1536
1537 let mut batch = new_batch(
1539 &[1, 2, 3, 4],
1540 &[11, 12, 13, 14],
1541 &[OpType::Put, OpType::Put, OpType::Put, OpType::Put],
1542 &[21, 22, 23, 24],
1543 );
1544 batch
1545 .filter_by_sequence(Some(SequenceRange::GtLtEq { min: 20, max: 25 }))
1546 .unwrap();
1547 assert!(batch.is_empty());
1548 }
1549
1550 #[test]
1551 fn test_merge_last_non_null_no_dup() {
1552 let mut batch = new_batch_with_u64_fields(
1553 &[1, 2],
1554 &[2, 1],
1555 &[OpType::Put, OpType::Put],
1556 &[(1, &[Some(10), None]), (2, &[Some(100), Some(200)])],
1557 );
1558 let expect = batch.clone();
1559 batch.merge_last_non_null().unwrap();
1560 assert_eq!(expect, batch);
1561 }
1562
1563 #[test]
1564 fn test_merge_last_non_null_fill_null_fields() {
1565 let mut batch = new_batch_with_u64_fields(
1567 &[1, 1, 1],
1568 &[3, 2, 1],
1569 &[OpType::Put, OpType::Put, OpType::Put],
1570 &[
1571 (1, &[None, Some(10), Some(11)]),
1572 (2, &[Some(100), Some(200), Some(300)]),
1573 ],
1574 );
1575 batch.merge_last_non_null().unwrap();
1576
1577 let expect = new_batch_with_u64_fields(
1580 &[1],
1581 &[3],
1582 &[OpType::Put],
1583 &[(1, &[Some(10)]), (2, &[Some(100)])],
1584 );
1585 assert_eq!(expect, batch);
1586 }
1587
1588 #[test]
1589 fn test_merge_last_non_null_stop_at_delete_row() {
1590 let mut batch = new_batch_with_u64_fields(
1593 &[1, 1, 1],
1594 &[3, 2, 1],
1595 &[OpType::Put, OpType::Delete, OpType::Put],
1596 &[
1597 (1, &[None, Some(10), Some(11)]),
1598 (2, &[Some(100), Some(200), Some(300)]),
1599 ],
1600 );
1601 batch.merge_last_non_null().unwrap();
1602
1603 let expect = new_batch_with_u64_fields(
1604 &[1],
1605 &[3],
1606 &[OpType::Put],
1607 &[(1, &[None]), (2, &[Some(100)])],
1608 );
1609 assert_eq!(expect, batch);
1610 }
1611
1612 #[test]
1613 fn test_merge_last_non_null_base_delete_no_merge() {
1614 let mut batch = new_batch_with_u64_fields(
1615 &[1, 1],
1616 &[3, 2],
1617 &[OpType::Delete, OpType::Put],
1618 &[(1, &[None, Some(10)]), (2, &[None, Some(200)])],
1619 );
1620 batch.merge_last_non_null().unwrap();
1621
1622 let expect =
1624 new_batch_with_u64_fields(&[1], &[3], &[OpType::Delete], &[(1, &[None]), (2, &[None])]);
1625 assert_eq!(expect, batch);
1626 }
1627
1628 #[test]
1629 fn test_merge_last_non_null_multiple_timestamp_groups() {
1630 let mut batch = new_batch_with_u64_fields(
1631 &[1, 1, 2, 3, 3],
1632 &[5, 4, 3, 2, 1],
1633 &[
1634 OpType::Put,
1635 OpType::Put,
1636 OpType::Put,
1637 OpType::Put,
1638 OpType::Put,
1639 ],
1640 &[
1641 (1, &[None, Some(10), Some(20), None, Some(30)]),
1642 (2, &[Some(100), Some(110), Some(120), None, Some(130)]),
1643 ],
1644 );
1645 batch.merge_last_non_null().unwrap();
1646
1647 let expect = new_batch_with_u64_fields(
1648 &[1, 2, 3],
1649 &[5, 3, 2],
1650 &[OpType::Put, OpType::Put, OpType::Put],
1651 &[
1652 (1, &[Some(10), Some(20), Some(30)]),
1653 (2, &[Some(100), Some(120), Some(130)]),
1654 ],
1655 );
1656 assert_eq!(expect, batch);
1657 }
1658
1659 #[test]
1660 fn test_merge_last_non_null_no_fields() {
1661 let mut batch = new_batch_without_fields(
1662 &[1, 1, 2],
1663 &[3, 2, 1],
1664 &[OpType::Put, OpType::Put, OpType::Put],
1665 );
1666 batch.merge_last_non_null().unwrap();
1667
1668 let expect = new_batch_without_fields(&[1, 2], &[3, 1], &[OpType::Put, OpType::Put]);
1669 assert_eq!(expect, batch);
1670 }
1671
1672 #[test]
1673 fn test_filter() {
1674 let mut batch = new_batch(
1676 &[1, 2, 3, 4],
1677 &[11, 12, 13, 14],
1678 &[OpType::Put, OpType::Put, OpType::Put, OpType::Put],
1679 &[21, 22, 23, 24],
1680 );
1681 let predicate = BooleanVector::from_vec(vec![false, false, true, true]);
1682 batch.filter(&predicate).unwrap();
1683 let expect = new_batch(&[3, 4], &[13, 14], &[OpType::Put, OpType::Put], &[23, 24]);
1684 assert_eq!(expect, batch);
1685
1686 let mut batch = new_batch(
1688 &[1, 2, 3, 4],
1689 &[11, 12, 13, 14],
1690 &[OpType::Put, OpType::Delete, OpType::Put, OpType::Put],
1691 &[21, 22, 23, 24],
1692 );
1693 let predicate = BooleanVector::from_vec(vec![false, false, true, true]);
1694 batch.filter(&predicate).unwrap();
1695 let expect = new_batch(&[3, 4], &[13, 14], &[OpType::Put, OpType::Put], &[23, 24]);
1696 assert_eq!(expect, batch);
1697
1698 let predicate = BooleanVector::from_vec(vec![false, false]);
1700 batch.filter(&predicate).unwrap();
1701 assert!(batch.is_empty());
1702 }
1703
1704 #[test]
1705 fn test_sort_and_dedup() {
1706 let original = new_batch(
1707 &[2, 3, 1, 4, 5, 2],
1708 &[1, 2, 3, 4, 5, 6],
1709 &[
1710 OpType::Put,
1711 OpType::Put,
1712 OpType::Put,
1713 OpType::Put,
1714 OpType::Put,
1715 OpType::Put,
1716 ],
1717 &[21, 22, 23, 24, 25, 26],
1718 );
1719
1720 let mut batch = original.clone();
1721 batch.sort(true).unwrap();
1722 assert_eq!(
1724 new_batch(
1725 &[1, 2, 3, 4, 5],
1726 &[3, 6, 2, 4, 5],
1727 &[
1728 OpType::Put,
1729 OpType::Put,
1730 OpType::Put,
1731 OpType::Put,
1732 OpType::Put,
1733 ],
1734 &[23, 26, 22, 24, 25],
1735 ),
1736 batch
1737 );
1738
1739 let mut batch = original.clone();
1740 batch.sort(false).unwrap();
1741
1742 assert_eq!(
1744 new_batch(
1745 &[1, 2, 2, 3, 4, 5],
1746 &[3, 6, 1, 2, 4, 5],
1747 &[
1748 OpType::Put,
1749 OpType::Put,
1750 OpType::Put,
1751 OpType::Put,
1752 OpType::Put,
1753 OpType::Put,
1754 ],
1755 &[23, 26, 21, 22, 24, 25],
1756 ),
1757 batch
1758 );
1759
1760 let original = new_batch(
1761 &[2, 2, 1],
1762 &[1, 6, 1],
1763 &[OpType::Delete, OpType::Put, OpType::Put],
1764 &[21, 22, 23],
1765 );
1766
1767 let mut batch = original.clone();
1768 batch.sort(true).unwrap();
1769 let expect = new_batch(&[1, 2], &[1, 6], &[OpType::Put, OpType::Put], &[23, 22]);
1770 assert_eq!(expect, batch);
1771
1772 let mut batch = original.clone();
1773 batch.sort(false).unwrap();
1774 let expect = new_batch(
1775 &[1, 2, 2],
1776 &[1, 6, 1],
1777 &[OpType::Put, OpType::Put, OpType::Delete],
1778 &[23, 22, 21],
1779 );
1780 assert_eq!(expect, batch);
1781 }
1782
1783 #[test]
1784 fn test_get_value() {
1785 let encodings = [PrimaryKeyEncoding::Dense, PrimaryKeyEncoding::Sparse];
1786
1787 for encoding in encodings {
1788 let codec = build_primary_key_codec_with_fields(
1789 encoding,
1790 [
1791 (
1792 ReservedColumnId::table_id(),
1793 row_converter::SortField::new(ConcreteDataType::uint32_datatype()),
1794 ),
1795 (
1796 ReservedColumnId::tsid(),
1797 row_converter::SortField::new(ConcreteDataType::uint64_datatype()),
1798 ),
1799 (
1800 100,
1801 row_converter::SortField::new(ConcreteDataType::string_datatype()),
1802 ),
1803 (
1804 200,
1805 row_converter::SortField::new(ConcreteDataType::string_datatype()),
1806 ),
1807 ]
1808 .into_iter(),
1809 );
1810
1811 let values = [
1812 Value::UInt32(1000),
1813 Value::UInt64(2000),
1814 Value::String("abcdefgh".into()),
1815 Value::String("zyxwvu".into()),
1816 ];
1817 let mut buf = vec![];
1818 codec
1819 .encode_values(
1820 &[
1821 (ReservedColumnId::table_id(), values[0].clone()),
1822 (ReservedColumnId::tsid(), values[1].clone()),
1823 (100, values[2].clone()),
1824 (200, values[3].clone()),
1825 ],
1826 &mut buf,
1827 )
1828 .unwrap();
1829
1830 let field_col_id = 2;
1831 let mut batch = new_batch_builder(
1832 &buf,
1833 &[1, 2, 3],
1834 &[1, 1, 1],
1835 &[OpType::Put, OpType::Put, OpType::Put],
1836 field_col_id,
1837 &[42, 43, 44],
1838 )
1839 .build()
1840 .unwrap();
1841
1842 let v = batch
1843 .pk_col_value(&*codec, 0, ReservedColumnId::table_id())
1844 .unwrap()
1845 .unwrap();
1846 assert_eq!(values[0], *v);
1847
1848 let v = batch
1849 .pk_col_value(&*codec, 1, ReservedColumnId::tsid())
1850 .unwrap()
1851 .unwrap();
1852 assert_eq!(values[1], *v);
1853
1854 let v = batch.pk_col_value(&*codec, 2, 100).unwrap().unwrap();
1855 assert_eq!(values[2], *v);
1856
1857 let v = batch.pk_col_value(&*codec, 3, 200).unwrap().unwrap();
1858 assert_eq!(values[3], *v);
1859
1860 let v = batch.field_col_value(field_col_id).unwrap();
1861 assert_eq!(v.data.get(0), Value::UInt64(42));
1862 assert_eq!(v.data.get(1), Value::UInt64(43));
1863 assert_eq!(v.data.get(2), Value::UInt64(44));
1864 }
1865 }
1866}