Skip to main content

mito2/
read.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15//! Common structs and utilities for reading data.
16
17pub 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/// Storage internal representation of a batch of rows for a primary key (time series).
112///
113/// Rows are sorted by primary key, timestamp, sequence desc, op_type desc. Fields
114/// always keep the same relative order as fields in [RegionMetadata](store_api::metadata::RegionMetadata).
115#[derive(Debug, PartialEq, Clone)]
116pub struct Batch {
117    /// Primary key encoded in a comparable form.
118    primary_key: Vec<u8>,
119    /// Possibly decoded `primary_key` values. Some places would decode it in advance.
120    pk_values: Option<CompositeValues>,
121    /// Lazily decoded Dense columns, shared by consumers such as index creators.
122    dense_pk_cache: Option<Box<DensePkCache>>,
123    /// Timestamps of rows, should be sorted and not null.
124    timestamps: VectorRef,
125    /// Sequences of rows
126    ///
127    /// UInt64 type, not null.
128    sequences: Arc<UInt64Vector>,
129    /// Op types of rows
130    ///
131    /// UInt8 type, not null.
132    op_types: Arc<UInt8Vector>,
133    /// Fields organized in columnar format.
134    fields: Vec<BatchColumn>,
135    /// Cache for field index lookup.
136    fields_idx: Option<HashMap<ColumnId, usize>>,
137}
138
139/// Only used when no fully decoded primary key was supplied by the reader.
140#[derive(Debug, PartialEq, Clone)]
141struct DensePkCache {
142    offsets: Vec<usize>,
143    values: Vec<Option<Value>>,
144}
145
146impl Batch {
147    /// Creates a new batch.
148    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    /// Tries to set fields for the batch.
161    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    /// Returns primary key of the batch.
172    pub fn primary_key(&self) -> &[u8] {
173        &self.primary_key
174    }
175
176    /// Returns possibly decoded primary-key values.
177    pub fn pk_values(&self) -> Option<&CompositeValues> {
178        self.pk_values.as_ref()
179    }
180
181    /// Sets possibly decoded primary-key values.
182    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    /// Removes possibly decoded primary-key values. For testing only.
188    #[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    /// Returns fields in the batch.
195    pub fn fields(&self) -> &[BatchColumn] {
196        &self.fields
197    }
198
199    /// Returns timestamps of the batch.
200    pub fn timestamps(&self) -> &VectorRef {
201        &self.timestamps
202    }
203
204    /// Returns sequences of the batch.
205    pub fn sequences(&self) -> &Arc<UInt64Vector> {
206        &self.sequences
207    }
208
209    /// Returns op types of the batch.
210    pub fn op_types(&self) -> &Arc<UInt8Vector> {
211        &self.op_types
212    }
213
214    /// Returns the number of rows in the batch.
215    pub fn num_rows(&self) -> usize {
216        // All vectors have the same length. We use the length of sequences vector
217        // since it has static type.
218        self.sequences.len()
219    }
220
221    /// Create an empty [`Batch`].
222    #[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    /// Returns true if the number of rows in the batch is 0.
237    pub fn is_empty(&self) -> bool {
238        self.num_rows() == 0
239    }
240
241    /// Returns the first timestamp in the batch or `None` if the batch is empty.
242    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    /// Returns the last timestamp in the batch or `None` if the batch is empty.
251    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    /// Returns the first sequence in the batch or `None` if the batch is empty.
260    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    /// Returns the last sequence in the batch or `None` if the batch is empty.
269    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    /// Replaces the primary key of the batch.
278    ///
279    /// Notice that this [Batch] also contains a maybe-exist `pk_values`.
280    /// Be sure to update that field as well.
281    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    /// Slice the batch, returning a new batch.
287    ///
288    /// # Panics
289    /// Panics if `offset + length > self.num_rows()`.
290    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        // We skip using the builder to avoid validating the batch again.
300        Batch {
301            // Now we need to clone the primary key. We could try `Bytes` if
302            // this becomes a bottleneck.
303            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    /// Takes `batches` and concat them into one batch.
315    ///
316    /// All `batches` must have the same primary key.
317    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            // Now we own the `batches` so we could pop it directly.
326            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        // We took the primary key from the first batch so we don't use `first.primary_key()`.
332        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        // We take the primary key from the first batch.
359        let mut builder = BatchBuilder::new(primary_key);
360        // Concat timestamps, sequences, op_types, fields.
361        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    /// Removes rows whose op type is delete.
376    pub fn filter_deleted(&mut self) -> Result<()> {
377        // Safety: op type column is not null.
378        let array = self.op_types.as_arrow();
379        // Find rows with non-delete op type.
380        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    // Applies the `predicate` to the batch.
387    // Safety: We know the array type so we unwrap on casting.
388    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    /// Filters rows by the given `sequence`. Only preserves rows with sequence less than or equal to `sequence`.
418    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    /// Sorts rows in the batch. If `dedup` is true, it also removes
448    /// duplicated rows according to primary keys.
449    ///
450    /// It orders rows by timestamp, sequence desc and only keep the latest
451    /// row for the same timestamp. It doesn't consider op type as sequence
452    /// should already provide uniqueness for a row.
453    pub fn sort(&mut self, dedup: bool) -> Result<()> {
454        // If building a converter each time is costly, we may allow passing a
455        // converter.
456        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        // Columns to sort.
468        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            // Dedup by timestamps.
483            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                // We only compare the timestamp part and ignore sequence.
488                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    /// Merges duplicated timestamps in the batch by keeping the latest non-null field values.
501    ///
502    /// Rows must already be sorted by timestamp (ascending) and sequence (descending).
503    ///
504    /// This method deduplicates rows with the same timestamp (keeping the first row in each
505    /// timestamp range as the base row) and fills null fields from subsequent rows until all
506    /// fields are filled or a delete operation is encountered.
507    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        // Fast path: check if there are any duplicate timestamps.
518        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                // Default: take the base row for all fields.
549                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                    // Track fields that are null in the base row and try to fill them from older
556                    // rows in the same timestamp range.
557                    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        // Safety: We know the array and vector type.
598        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        // Safety: We know the array and vector type.
602        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    /// Returns the estimated memory size of the batch.
613    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    /// Returns timestamps in a native slice or `None` if the batch is empty.
626    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    /// Takes the batch in place.
667    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        // Safety: we know the array and vector type.
672        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    /// Gets a timestamp at given `index`.
687    ///
688    /// # Panics
689    /// Panics if `index` is out-of-bound or the timestamp vector returns null.
690    fn get_timestamp(&self, index: usize) -> Timestamp {
691        match self.timestamps.get_ref(index) {
692            ValueRef::Timestamp(timestamp) => timestamp,
693
694            // We have check the data type is timestamp compatible in the [BatchBuilder] so it's safe to panic.
695            value => panic!("{:?} is not a timestamp", value),
696        }
697    }
698
699    /// Gets a sequence at given `index`.
700    ///
701    /// # Panics
702    /// Panics if `index` is out-of-bound or the sequence vector returns null.
703    pub(crate) fn get_sequence(&self, index: usize) -> SequenceNumber {
704        // Safety: sequences is not null so it actually returns Some.
705        self.sequences.get_data(index).unwrap()
706    }
707
708    /// Prepares a shared full-key cache when all Dense PK columns are needed.
709    /// Explicitly supplied decoded/defaulted values take precedence.
710    pub(crate) fn ensure_dense_pk_decoded(&mut self, codec: &DensePrimaryKeyCodec) -> Result<()> {
711        if self.pk_values.is_none() {
712            // A fallible iterator has no nonzero lower size hint. Reserve the
713            // known field count instead of growing its collected Vec per key.
714            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    /// Returns the value of the column in the primary key.
724    ///
725    /// Reuses predecoded values when available. Otherwise Dense keys decode only
726    /// the requested column, sharing offsets and values across callers; Sparse
727    /// keys cache the full decode.
728    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    /// Returns values of the field in the batch.
764    ///
765    /// Lazily caches the field index.
766    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
785/// Len of timestamp in arrow row format.
786const TIMESTAMP_KEY_LEN: usize = 9;
787
788/// Helper function to concat arrays from `iter`.
789fn 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/// A column in a [Batch].
796#[derive(Debug, PartialEq, Eq, Clone)]
797pub struct BatchColumn {
798    /// Id of the column.
799    pub column_id: ColumnId,
800    /// Data of the column.
801    pub data: VectorRef,
802}
803
804/// Builder to build [Batch].
805pub 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    /// Creates a new [BatchBuilder] with primary key.
815    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    /// Creates a new [BatchBuilder] with all required columns.
826    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    /// Set all field columns.
842    pub fn with_fields(mut self, fields: Vec<BatchColumn>) -> Self {
843        self.fields = fields;
844        self
845    }
846
847    /// Push a field column.
848    pub fn push_field(&mut self, column: BatchColumn) -> &mut Self {
849        self.fields.push(column);
850        self
851    }
852
853    /// Push an array as a field.
854    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    /// Try to set an array as timestamps.
865    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    /// Try to set an array as sequences.
879    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        // Safety: The cast must success as we have ensured it is uint64 type.
887        let vector = Arc::new(UInt64Vector::try_from_arrow_array(array).unwrap());
888        self.sequences = Some(vector);
889
890        Ok(self)
891    }
892
893    /// Try to set an array as op types.
894    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        // Safety: The cast must success as we have ensured it is uint64 type.
902        let vector = Arc::new(UInt8Vector::try_from_arrow_array(array).unwrap());
903        self.op_types = Some(vector);
904
905        Ok(self)
906    }
907
908    /// Builds the [Batch].
909    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        // Our storage format ensure these columns are not nullable so
920        // we use assert here.
921        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
985/// Async [RecordBatch] reader and iterator wrapper for flat format.
986pub struct FlatSource {
987    schema: SchemaRef,
988    inner: FlatSourceInner,
989}
990
991impl FlatSource {
992    /// Create a [FlatSource] from a [BoxedRecordBatchIterator] and its schema.
993    pub fn new_iter(schema: SchemaRef, iter: BoxedRecordBatchIterator) -> Self {
994        Self {
995            schema,
996            inner: FlatSourceInner::Iter(iter),
997        }
998    }
999
1000    /// Create a [FlatSource] from a [BoxedRecordBatchStream] and its schema.
1001    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    /// Source from a [BoxedRecordBatchIterator].
1027    Iter(BoxedRecordBatchIterator),
1028    /// Source from a [BoxedRecordBatchStream].
1029    Stream(BoxedRecordBatchStream),
1030}
1031
1032impl FlatSourceInner {
1033    /// Returns next [RecordBatch] from this data source.
1034    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
1042/// Pointer to a stream that yields [RecordBatch].
1043pub type BoxedRecordBatchStream = BoxStream<'static, Result<RecordBatch>>;
1044
1045/// Local metrics for scanners.
1046#[derive(Debug, Default)]
1047pub(crate) struct ScannerMetrics {
1048    /// Duration to scan data.
1049    scan_cost: Duration,
1050    /// Duration while waiting for `yield`.
1051    yield_cost: Duration,
1052    /// Number of batches returned.
1053    num_batches: usize,
1054    /// Number of rows returned.
1055    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        // A corrupt unrequested suffix must not force whole-key decoding.
1099        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        // Schema compatibility may supply already decoded/defaulted values.
1127        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        // Batch 2 has more fields.
1363        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        // Batch 2 has different field.
1371        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        // Filters put only.
1416        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        // Filters to empty.
1434        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        // None filter.
1447        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        // Filter a empty batch
1458        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        // Filter a empty batch with None
1465        let mut batch = new_batch(&[], &[], &[], &[]);
1466        batch.filter_by_sequence(None).unwrap();
1467        assert!(batch.is_empty());
1468
1469        // Test From variant - exclusive lower bound
1470        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        // Test From variant with no matches
1483        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        // Test Range variant - exclusive lower bound, inclusive upper bound
1495        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        // Test Range variant with mixed operations
1514        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        // Test Range variant with no matches
1538        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        // Rows are already sorted by timestamp asc and sequence desc.
1566        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        // Field 1 is filled from the first older row (seq=2). Field 2 keeps the base value.
1578        // Filled fields must not be overwritten by even older duplicates.
1579        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        // A delete row in older duplicates should stop filling to avoid resurrecting values before
1591        // deletion.
1592        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        // Base row is delete, keep it as is and don't merge fields from older rows.
1623        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        // Filters put only.
1675        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        // Filters deletion.
1687        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        // Filters to empty.
1699        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        // It should only keep one timestamp 2.
1723        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        // It should only keep one timestamp 2.
1743        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}