Skip to main content

mito2/read/
flat_dedup.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//! Dedup implementation for flat format.
16
17use std::ops::Range;
18use std::sync::Arc;
19use std::time::Instant;
20
21use api::v1::OpType;
22use async_stream::try_stream;
23use common_telemetry::debug;
24use datatypes::arrow::array::{
25    Array, ArrayRef, BinaryArray, BooleanArray, BooleanBufferBuilder, UInt8Array, UInt64Array,
26    make_comparator,
27};
28use datatypes::arrow::buffer::BooleanBuffer;
29use datatypes::arrow::compute::kernels::cmp::distinct;
30use datatypes::arrow::compute::kernels::partition::{Partitions, partition};
31use datatypes::arrow::compute::kernels::take::take;
32use datatypes::arrow::compute::{
33    SortOptions, TakeOptions, concat_batches, filter_record_batch, take_record_batch,
34};
35use datatypes::arrow::error::ArrowError;
36use datatypes::arrow::record_batch::RecordBatch;
37use futures::{Stream, TryStreamExt};
38use snafu::ResultExt;
39
40use crate::error::{ComputeArrowSnafu, NewRecordBatchSnafu, Result};
41use crate::metrics::MERGE_FILTER_ROWS_TOTAL;
42use crate::read::dedup::{DedupMetrics, DedupMetricsReport};
43use crate::read::timestamp_array_to_i64_slice;
44use crate::sst::parquet::flat_format::{
45    op_type_column_index, primary_key_column_index, time_index_column_index,
46};
47use crate::sst::parquet::format::{FIXED_POS_COLUMN_NUM, PrimaryKeyArray};
48
49/// An iterator to dedup sorted batches from an iterator based on the dedup strategy.
50pub struct FlatDedupIterator<I, S> {
51    iter: I,
52    strategy: S,
53    metrics: DedupMetrics,
54}
55
56impl<I, S> FlatDedupIterator<I, S> {
57    /// Creates a new dedup iterator.
58    pub fn new(iter: I, strategy: S) -> Self {
59        Self {
60            iter,
61            strategy,
62            metrics: DedupMetrics::default(),
63        }
64    }
65}
66
67impl<I: Iterator<Item = Result<RecordBatch>>, S: RecordBatchDedupStrategy> FlatDedupIterator<I, S> {
68    /// Returns the next deduplicated batch.
69    fn fetch_next_batch(&mut self) -> Result<Option<RecordBatch>> {
70        while let Some(batch) = self.iter.next().transpose()? {
71            if let Some(batch) = self.strategy.push_batch(batch, &mut self.metrics)? {
72                return Ok(Some(batch));
73            }
74        }
75
76        self.strategy.finish(&mut self.metrics)
77    }
78}
79
80impl<I: Iterator<Item = Result<RecordBatch>>, S: RecordBatchDedupStrategy> Iterator
81    for FlatDedupIterator<I, S>
82{
83    type Item = Result<RecordBatch>;
84
85    fn next(&mut self) -> Option<Self::Item> {
86        self.fetch_next_batch().transpose()
87    }
88}
89
90/// An async reader to dedup sorted record batches from a stream based on the dedup strategy.
91pub struct FlatDedupReader<I, S> {
92    stream: I,
93    strategy: S,
94    metrics: DedupMetrics,
95    /// Optional metrics reporter.
96    metrics_reporter: Option<Arc<dyn DedupMetricsReport>>,
97}
98
99impl<I, S> FlatDedupReader<I, S> {
100    /// Creates a new dedup reader.
101    pub fn new(
102        stream: I,
103        strategy: S,
104        metrics_reporter: Option<Arc<dyn DedupMetricsReport>>,
105    ) -> Self {
106        Self {
107            stream,
108            strategy,
109            metrics: DedupMetrics::default(),
110            metrics_reporter,
111        }
112    }
113}
114
115impl<I: Stream<Item = Result<RecordBatch>> + Unpin, S: RecordBatchDedupStrategy>
116    FlatDedupReader<I, S>
117{
118    /// Returns the next deduplicated batch.
119    async fn fetch_next_batch(&mut self) -> Result<Option<RecordBatch>> {
120        while let Some(batch) = self.stream.try_next().await? {
121            if let Some(batch) = self.strategy.push_batch(batch, &mut self.metrics)? {
122                self.metrics.maybe_report(&self.metrics_reporter);
123                return Ok(Some(batch));
124            }
125        }
126
127        let result = self.strategy.finish(&mut self.metrics)?;
128        self.metrics.maybe_report(&self.metrics_reporter);
129        Ok(result)
130    }
131
132    /// Converts the reader into a stream.
133    pub fn into_stream(mut self) -> impl Stream<Item = Result<RecordBatch>> {
134        try_stream! {
135            while let Some(batch) = self.fetch_next_batch().await? {
136                yield batch;
137            }
138        }
139    }
140}
141
142impl<I, S> Drop for FlatDedupReader<I, S> {
143    fn drop(&mut self) {
144        debug!("Flat dedup reader finished, metrics: {:?}", self.metrics);
145
146        MERGE_FILTER_ROWS_TOTAL
147            .with_label_values(&["dedup"])
148            .inc_by(self.metrics.num_unselected_rows as u64);
149        MERGE_FILTER_ROWS_TOTAL
150            .with_label_values(&["delete"])
151            .inc_by(self.metrics.num_deleted_rows as u64);
152
153        // Report any remaining metrics.
154        if let Some(reporter) = &self.metrics_reporter {
155            reporter.report(&mut self.metrics);
156        }
157    }
158}
159
160/// Strategy to remove duplicate rows from sorted record batches.
161pub trait RecordBatchDedupStrategy: Send {
162    /// Pushes a batch to the dedup strategy.
163    /// Returns a batch if the strategy ensures there is no duplications based on
164    /// the input batch.
165    fn push_batch(
166        &mut self,
167        batch: RecordBatch,
168        metrics: &mut DedupMetrics,
169    ) -> Result<Option<RecordBatch>>;
170
171    /// Finishes the deduplication process and returns any remaining batch.
172    ///
173    /// Users must ensure that `push_batch` is called for all batches before
174    /// calling this method.
175    fn finish(&mut self, metrics: &mut DedupMetrics) -> Result<Option<RecordBatch>>;
176}
177
178/// Dedup strategy that keeps the row with latest sequence of each key.
179pub struct FlatLastRow {
180    /// Meta of the last row in the previous batch that has the same key
181    /// as the batch to push.
182    prev_batch: Option<BatchLastRow>,
183    /// Filter deleted rows.
184    filter_deleted: bool,
185}
186
187impl FlatLastRow {
188    /// Creates a new strategy with the given `filter_deleted` flag.
189    pub fn new(filter_deleted: bool) -> Self {
190        Self {
191            prev_batch: None,
192            filter_deleted,
193        }
194    }
195
196    /// Remove duplications from the batch without considering previous rows.
197    fn dedup_one_batch(batch: RecordBatch) -> Result<RecordBatch> {
198        let num_rows = batch.num_rows();
199        if num_rows < 2 {
200            return Ok(batch);
201        }
202
203        let num_columns = batch.num_columns();
204        let timestamps = batch.column(time_index_column_index(num_columns));
205        // Checks duplications based on the timestamp.
206        let mask = find_boundaries(timestamps).context(ComputeArrowSnafu)?;
207        if mask.count_set_bits() == num_rows - 1 {
208            // Fast path: No duplication.
209            return Ok(batch);
210        }
211
212        // The batch has duplicated timestamps, but it doesn't mean it must
213        // has duplicated rows.
214        // Partitions the batch by the primary key and time index.
215        let columns: Vec<_> = [
216            primary_key_column_index(num_columns),
217            time_index_column_index(num_columns),
218        ]
219        .iter()
220        .map(|index| batch.column(*index).clone())
221        .collect();
222        let partitions = partition(&columns).context(ComputeArrowSnafu)?;
223
224        Self::dedup_by_partitions(batch, &partitions)
225    }
226
227    /// Remove duplications for each partition.
228    fn dedup_by_partitions(batch: RecordBatch, partitions: &Partitions) -> Result<RecordBatch> {
229        let ranges = partitions.ranges();
230        // Each range at least has 1 row.
231        let num_duplications: usize = ranges.iter().map(|r| r.end - r.start - 1).sum();
232        if num_duplications == 0 {
233            // Fast path, no duplications.
234            return Ok(batch);
235        }
236
237        // Always takes the first row in each range.
238        let take_indices: UInt64Array = ranges.iter().map(|r| Some(r.start as u64)).collect();
239        take_record_batch(&batch, &take_indices).context(ComputeArrowSnafu)
240    }
241}
242
243impl RecordBatchDedupStrategy for FlatLastRow {
244    fn push_batch(
245        &mut self,
246        batch: RecordBatch,
247        metrics: &mut DedupMetrics,
248    ) -> Result<Option<RecordBatch>> {
249        let start = Instant::now();
250
251        if batch.num_rows() == 0 {
252            return Ok(None);
253        }
254
255        // Dedup current batch to ensure no duplication before we checking the previous row.
256        let row_before_dedup = batch.num_rows();
257        let mut batch = Self::dedup_one_batch(batch)?;
258
259        if let Some(prev_batch) = &self.prev_batch {
260            // If we have previous batch.
261            if prev_batch.is_last_row_duplicated(&batch) {
262                // Duplicated with the last batch, skip the first row.
263                batch = batch.slice(1, batch.num_rows() - 1);
264            }
265        }
266        metrics.num_unselected_rows += row_before_dedup - batch.num_rows();
267
268        let Some(batch_last_row) = BatchLastRow::try_new(batch.clone()) else {
269            // The batch after dedup is empty.
270            // We don't need to update `prev_batch` because they have the same
271            // key and timestamp.
272            metrics.dedup_cost += start.elapsed();
273            return Ok(None);
274        };
275
276        // Store current batch to `prev_batch` so we could compare the next batch
277        // with this batch. We store batch before filtering it as rows with `OpType::Delete`
278        // would be removed from the batch after filter, then we may store an incorrect `last row`
279        // of previous batch.
280        // Safety: We checked the batch is not empty before.
281        self.prev_batch = Some(batch_last_row);
282
283        // Filters deleted rows at last.
284        let result = maybe_filter_deleted(batch, self.filter_deleted, metrics);
285
286        metrics.dedup_cost += start.elapsed();
287
288        result
289    }
290
291    fn finish(&mut self, _metrics: &mut DedupMetrics) -> Result<Option<RecordBatch>> {
292        Ok(None)
293    }
294}
295
296/// Dedup strategy that keeps the last non-null field for the same key.
297pub struct FlatLastNonNull {
298    /// The start index of field columns:
299    field_column_start: usize,
300    /// Filter deleted rows.
301    filter_deleted: bool,
302    /// Buffered batch to check whether the next batch have duplicated rows with this batch.
303    /// Fields in the last row of this batch may be updated by the next batch.
304    /// The buffered batch should contain no duplication.
305    buffer: Option<BatchLastRow>,
306    /// Whether the last row range contains a delete operation.
307    /// If so, we don't need to update null fields.
308    contains_delete: bool,
309}
310
311impl RecordBatchDedupStrategy for FlatLastNonNull {
312    fn push_batch(
313        &mut self,
314        batch: RecordBatch,
315        metrics: &mut DedupMetrics,
316    ) -> Result<Option<RecordBatch>> {
317        let start = Instant::now();
318
319        if batch.num_rows() == 0 {
320            return Ok(None);
321        }
322
323        let row_before_dedup = batch.num_rows();
324
325        let Some(buffer) = self.buffer.take() else {
326            // If the buffer is None, dedup the batch, put the batch into the buffer and return.
327            // There is no previous batch with the same key, we can pass contains_delete as false.
328            let (record_batch, contains_delete) =
329                Self::dedup_one_batch(batch, self.field_column_start, false)?;
330            metrics.num_unselected_rows += row_before_dedup - record_batch.num_rows();
331            self.buffer = BatchLastRow::try_new(record_batch);
332            self.contains_delete = contains_delete;
333
334            metrics.dedup_cost += start.elapsed();
335            return Ok(None);
336        };
337
338        if !buffer.is_last_row_duplicated(&batch) {
339            // The first row of batch has different key from the buffer.
340            // We can replace the buffer with the new batch.
341            // Dedup the batch.
342            // There is no previous batch with the same key, we can pass contains_delete as false.
343            let (record_batch, contains_delete) =
344                Self::dedup_one_batch(batch, self.field_column_start, false)?;
345            metrics.num_unselected_rows += row_before_dedup - record_batch.num_rows();
346            debug_assert!(record_batch.num_rows() > 0);
347            self.buffer = BatchLastRow::try_new(record_batch);
348            self.contains_delete = contains_delete;
349
350            let result = maybe_filter_deleted(buffer.last_batch, self.filter_deleted, metrics);
351            metrics.dedup_cost += start.elapsed();
352            return result;
353        }
354
355        // The next batch has duplicated rows.
356        // We can return rows except the last row in the buffer.
357        let output = if buffer.last_batch.num_rows() > 1 {
358            let dedup_batch = buffer.last_batch.slice(0, buffer.last_batch.num_rows() - 1);
359            debug_assert_eq!(buffer.last_batch.num_rows() - 1, dedup_batch.num_rows());
360
361            maybe_filter_deleted(dedup_batch, self.filter_deleted, metrics)?
362        } else {
363            None
364        };
365        let last_row = buffer.last_batch.slice(buffer.last_batch.num_rows() - 1, 1);
366
367        // We concat the last row with the next batch.
368        let schema = batch.schema();
369        let merged = concat_batches(&schema, &[last_row, batch]).context(ComputeArrowSnafu)?;
370        let merged_row_count = merged.num_rows();
371        // Dedup the merged batch and update the buffer.
372        let (record_batch, contains_delete) =
373            Self::dedup_one_batch(merged, self.field_column_start, self.contains_delete)?;
374        metrics.num_unselected_rows += merged_row_count - record_batch.num_rows();
375        debug_assert!(record_batch.num_rows() > 0);
376        self.buffer = BatchLastRow::try_new(record_batch);
377        self.contains_delete = contains_delete;
378
379        metrics.dedup_cost += start.elapsed();
380
381        Ok(output)
382    }
383
384    fn finish(&mut self, metrics: &mut DedupMetrics) -> Result<Option<RecordBatch>> {
385        let Some(buffer) = self.buffer.take() else {
386            return Ok(None);
387        };
388
389        let start = Instant::now();
390
391        let result = maybe_filter_deleted(buffer.last_batch, self.filter_deleted, metrics);
392
393        metrics.dedup_cost += start.elapsed();
394
395        result
396    }
397}
398
399impl FlatLastNonNull {
400    /// Creates a new strategy with the given `filter_deleted` flag.
401    pub fn new(field_column_start: usize, filter_deleted: bool) -> Self {
402        Self {
403            field_column_start,
404            filter_deleted,
405            buffer: None,
406            contains_delete: false,
407        }
408    }
409
410    /// Remove duplications from the batch without considering the previous and next rows.
411    /// Returns a tuple containing the deduplicated batch and a boolean indicating whether the last range contains deleted rows.
412    fn dedup_one_batch(
413        batch: RecordBatch,
414        field_column_start: usize,
415        prev_batch_contains_delete: bool,
416    ) -> Result<(RecordBatch, bool)> {
417        // Get op type array for checking delete operations
418        let op_type_column = batch
419            .column(op_type_column_index(batch.num_columns()))
420            .clone();
421        let op_types = op_type_column
422            .as_any()
423            .downcast_ref::<UInt8Array>()
424            .unwrap();
425        let num_rows = batch.num_rows();
426        if num_rows < 2 {
427            let contains_delete = if num_rows > 0 {
428                op_types.value(0) == OpType::Delete as u8
429            } else {
430                false
431            };
432            return Ok((batch, contains_delete));
433        }
434
435        let num_columns = batch.num_columns();
436        let timestamps = batch.column(time_index_column_index(num_columns));
437        // Checks duplications based on the timestamp.
438        let mask = find_boundaries(timestamps).context(ComputeArrowSnafu)?;
439        if mask.count_set_bits() == num_rows - 1 {
440            let contains_delete = op_types.value(num_rows - 1) == OpType::Delete as u8;
441            // Fast path: No duplication.
442            return Ok((batch, contains_delete));
443        }
444
445        // The batch has duplicated timestamps, but it doesn't mean it must
446        // has duplicated rows.
447        // Partitions the batch by the primary key and time index.
448        let columns: Vec<_> = [
449            primary_key_column_index(num_columns),
450            time_index_column_index(num_columns),
451        ]
452        .iter()
453        .map(|index| batch.column(*index).clone())
454        .collect();
455        let partitions = partition(&columns).context(ComputeArrowSnafu)?;
456
457        Self::dedup_by_partitions(
458            batch,
459            &partitions,
460            field_column_start,
461            op_types,
462            prev_batch_contains_delete,
463        )
464    }
465
466    /// Remove depulications for each partition.
467    /// Returns a tuple containing the deduplicated batch and a boolean indicating whether the last range contains deleted rows.
468    fn dedup_by_partitions(
469        batch: RecordBatch,
470        partitions: &Partitions,
471        field_column_start: usize,
472        op_types: &UInt8Array,
473        first_range_contains_delete: bool,
474    ) -> Result<(RecordBatch, bool)> {
475        let ranges = partitions.ranges();
476        let contains_delete = Self::last_range_has_delete(&ranges, op_types);
477
478        // Each range at least has 1 row.
479        let num_duplications: usize = ranges.iter().map(|r| r.end - r.start - 1).sum();
480        if num_duplications == 0 {
481            // Fast path, no duplication.
482            return Ok((batch, contains_delete));
483        }
484
485        let field_column_end = batch.num_columns() - FIXED_POS_COLUMN_NUM;
486        let take_options = Some(TakeOptions {
487            check_bounds: false,
488        });
489        // Always takes the first value for non-field columns in each range.
490        let non_field_indices: UInt64Array = ranges.iter().map(|r| Some(r.start as u64)).collect();
491        let new_columns = batch
492            .columns()
493            .iter()
494            .enumerate()
495            .map(|(col_idx, column)| {
496                if col_idx >= field_column_start && col_idx < field_column_end {
497                    let field_indices = Self::compute_field_indices(
498                        &ranges,
499                        column,
500                        op_types,
501                        first_range_contains_delete,
502                    );
503                    take(column, &field_indices, take_options.clone()).context(ComputeArrowSnafu)
504                } else {
505                    take(column, &non_field_indices, take_options.clone())
506                        .context(ComputeArrowSnafu)
507                }
508            })
509            .collect::<Result<Vec<ArrayRef>>>()?;
510
511        let record_batch =
512            RecordBatch::try_new(batch.schema(), new_columns).context(NewRecordBatchSnafu)?;
513        Ok((record_batch, contains_delete))
514    }
515
516    /// Returns an array of indices of the latest non null value for
517    /// each input range.
518    /// If all values in a range are null, the returned index is unspecific.
519    /// Stops when encountering a delete operation and ignores all subsequent rows.
520    fn compute_field_indices(
521        ranges: &[Range<usize>],
522        field_array: &ArrayRef,
523        op_types: &UInt8Array,
524        first_range_contains_delete: bool,
525    ) -> UInt64Array {
526        ranges
527            .iter()
528            .enumerate()
529            .map(|(range_idx, r)| {
530                let mut value_index = r.start as u64;
531                if range_idx == 0 && first_range_contains_delete {
532                    return Some(value_index);
533                }
534
535                // Iterate through the range to find the first valid non-null value
536                // but stop if we encounter a delete operation.
537                for i in r.clone() {
538                    if op_types.value(i) == OpType::Delete as u8 {
539                        break;
540                    }
541                    if field_array.is_valid(i) {
542                        value_index = i as u64;
543                        break;
544                    }
545                }
546
547                Some(value_index)
548            })
549            .collect()
550    }
551
552    /// Checks whether the last range contains a delete operation.
553    fn last_range_has_delete(ranges: &[Range<usize>], op_types: &UInt8Array) -> bool {
554        if let Some(last_range) = ranges.last() {
555            last_range
556                .clone()
557                .any(|i| op_types.value(i) == OpType::Delete as u8)
558        } else {
559            false
560        }
561    }
562}
563
564/// State of the batch with the last row for dedup.
565struct BatchLastRow {
566    /// The record batch that contains the last row.
567    /// It must has at least one row.
568    last_batch: RecordBatch,
569    /// Primary keys of the last batch.
570    primary_key: PrimaryKeyArray,
571    /// Last timestamp value.
572    timestamp: i64,
573}
574
575impl BatchLastRow {
576    /// Returns a new [BatchLastRow] if the record batch is not empty.
577    fn try_new(record_batch: RecordBatch) -> Option<Self> {
578        if record_batch.num_rows() > 0 {
579            let num_columns = record_batch.num_columns();
580            let primary_key = record_batch
581                .column(primary_key_column_index(num_columns))
582                .as_any()
583                .downcast_ref::<PrimaryKeyArray>()
584                .unwrap()
585                .clone();
586            let timestamp_array = record_batch.column(time_index_column_index(num_columns));
587            let timestamp = timestamp_value(timestamp_array, timestamp_array.len() - 1);
588
589            Some(Self {
590                last_batch: record_batch,
591                primary_key,
592                timestamp,
593            })
594        } else {
595            None
596        }
597    }
598
599    /// Returns true if the first row of the input `batch` is duplicated with the last row.
600    fn is_last_row_duplicated(&self, batch: &RecordBatch) -> bool {
601        if batch.num_rows() == 0 {
602            return false;
603        }
604
605        // The first timestamp in the batch.
606        let batch_timestamp = timestamp_value(
607            batch.column(time_index_column_index(batch.num_columns())),
608            0,
609        );
610        if batch_timestamp != self.timestamp {
611            return false;
612        }
613
614        let last_key = primary_key_at(&self.primary_key, self.last_batch.num_rows() - 1);
615        let primary_key = batch
616            .column(primary_key_column_index(batch.num_columns()))
617            .as_any()
618            .downcast_ref::<PrimaryKeyArray>()
619            .unwrap();
620        // Primary key of the first row in the batch.
621        let batch_key = primary_key_at(primary_key, 0);
622
623        last_key == batch_key
624    }
625}
626
627// TODO(yingwen): We only compares timestamp arrays, we can modify this function
628// to simplify the comparator.
629// Port from https://github.com/apache/arrow-rs/blob/55.0.0/arrow-ord/src/partition.rs#L155-L168
630/// Returns a mask with bits set whenever the value or nullability changes
631fn find_boundaries(v: &dyn Array) -> Result<BooleanBuffer, ArrowError> {
632    let slice_len = v.len() - 1;
633    let v1 = v.slice(0, slice_len);
634    let v2 = v.slice(1, slice_len);
635
636    if !v.data_type().is_nested() {
637        return Ok(distinct(&v1, &v2)?.values().clone());
638    }
639    // Given that we're only comparing values, null ordering in the input or
640    // sort options do not matter.
641    let cmp = make_comparator(&v1, &v2, SortOptions::default())?;
642    Ok((0..slice_len).map(|i| !cmp(i, i).is_eq()).collect())
643}
644
645/// Filters deleted rows from the record batch if `filter_deleted` is true.
646fn maybe_filter_deleted(
647    record_batch: RecordBatch,
648    filter_deleted: bool,
649    metrics: &mut DedupMetrics,
650) -> Result<Option<RecordBatch>> {
651    if !filter_deleted {
652        return Ok(Some(record_batch));
653    }
654    let batch = filter_deleted_from_batch(record_batch, metrics)?;
655    // Skips empty batches.
656    if batch.num_rows() == 0 {
657        return Ok(None);
658    }
659    Ok(Some(batch))
660}
661
662/// Removes deleted rows from the batch and updates metrics.
663fn filter_deleted_from_batch(
664    batch: RecordBatch,
665    metrics: &mut DedupMetrics,
666) -> Result<RecordBatch> {
667    let num_rows = batch.num_rows();
668    let op_type_column = batch.column(op_type_column_index(batch.num_columns()));
669    // Safety: The column should be op type.
670    let op_types = op_type_column
671        .as_any()
672        .downcast_ref::<UInt8Array>()
673        .unwrap();
674    let has_delete = op_types
675        .values()
676        .iter()
677        .any(|op_type| *op_type != OpType::Put as u8);
678    if !has_delete {
679        return Ok(batch);
680    }
681
682    let mut builder = BooleanBufferBuilder::new(op_types.len());
683    for op_type in op_types.values() {
684        if *op_type == OpType::Delete as u8 {
685            builder.append(false);
686        } else {
687            builder.append(true);
688        }
689    }
690    let predicate = BooleanArray::new(builder.into(), None);
691    let new_batch = filter_record_batch(&batch, &predicate).context(ComputeArrowSnafu)?;
692    let num_deleted = num_rows - new_batch.num_rows();
693    metrics.num_deleted_rows += num_deleted;
694    metrics.num_unselected_rows += num_deleted;
695
696    Ok(new_batch)
697}
698
699/// Gets the primary key at `index`.
700fn primary_key_at(primary_key: &PrimaryKeyArray, index: usize) -> &[u8] {
701    let key = primary_key.keys().value(index);
702    let binary_values = primary_key
703        .values()
704        .as_any()
705        .downcast_ref::<BinaryArray>()
706        .unwrap();
707    binary_values.value(key as usize)
708}
709
710/// Gets the timestamp value from the timestamp array.
711///
712/// # Panics
713/// Panics if the array is not a timestamp array or
714/// the index is out of bound.
715pub(crate) fn timestamp_value(array: &ArrayRef, idx: usize) -> i64 {
716    timestamp_array_to_i64_slice(array)[idx]
717}
718
719#[cfg(test)]
720mod tests {
721    use std::sync::Arc;
722
723    use api::v1::OpType;
724    use datatypes::arrow::array::{
725        ArrayRef, BinaryDictionaryBuilder, Int64Array, StringDictionaryBuilder,
726        TimestampMillisecondArray, UInt8Array, UInt64Array,
727    };
728    use datatypes::arrow::datatypes::{DataType, Field, Schema, SchemaRef, TimeUnit, UInt32Type};
729    use datatypes::arrow::record_batch::RecordBatch;
730
731    use super::*;
732
733    /// Creates a test RecordBatch in flat format with given parameters.
734    fn new_record_batch(
735        primary_keys: &[&[u8]],
736        timestamps: &[i64],
737        sequences: &[u64],
738        op_types: &[OpType],
739        fields: &[u64],
740    ) -> RecordBatch {
741        let num_rows = timestamps.len();
742        debug_assert_eq!(sequences.len(), num_rows);
743        debug_assert_eq!(op_types.len(), num_rows);
744        debug_assert_eq!(fields.len(), num_rows);
745        debug_assert_eq!(primary_keys.len(), num_rows);
746
747        let columns: Vec<ArrayRef> = vec![
748            // k0 column (primary key as string dictionary)
749            build_test_pk_string_dict_array(primary_keys),
750            // field0 column
751            Arc::new(Int64Array::from_iter(
752                fields.iter().map(|v| Some(*v as i64)),
753            )),
754            // ts column (time index)
755            Arc::new(TimestampMillisecondArray::from_iter_values(
756                timestamps.iter().copied(),
757            )),
758            // __primary_key column
759            build_test_pk_array(primary_keys),
760            // __sequence column
761            Arc::new(UInt64Array::from_iter_values(sequences.iter().copied())),
762            // __op_type column
763            Arc::new(UInt8Array::from_iter_values(
764                op_types.iter().map(|v| *v as u8),
765            )),
766        ];
767
768        RecordBatch::try_new(build_test_flat_schema(), columns).unwrap()
769    }
770
771    /// Creates a test RecordBatch in flat format with multiple fields for testing FlatLastNonNull.
772    fn new_record_batch_multi_fields(
773        primary_keys: &[&[u8]],
774        timestamps: &[i64],
775        sequences: &[u64],
776        op_types: &[OpType],
777        fields: &[(Option<u64>, Option<u64>)],
778    ) -> RecordBatch {
779        let num_rows = timestamps.len();
780        debug_assert_eq!(sequences.len(), num_rows);
781        debug_assert_eq!(op_types.len(), num_rows);
782        debug_assert_eq!(fields.len(), num_rows);
783        debug_assert_eq!(primary_keys.len(), num_rows);
784
785        let columns: Vec<ArrayRef> = vec![
786            // k0 column (primary key as string dictionary)
787            build_test_pk_string_dict_array(primary_keys),
788            // field0 column
789            Arc::new(Int64Array::from_iter(
790                fields.iter().map(|field| field.0.map(|v| v as i64)),
791            )),
792            // field1 column
793            Arc::new(Int64Array::from_iter(
794                fields.iter().map(|field| field.1.map(|v| v as i64)),
795            )),
796            // ts column (time index)
797            Arc::new(TimestampMillisecondArray::from_iter_values(
798                timestamps.iter().copied(),
799            )),
800            // __primary_key column
801            build_test_pk_array(primary_keys),
802            // __sequence column
803            Arc::new(UInt64Array::from_iter_values(sequences.iter().copied())),
804            // __op_type column
805            Arc::new(UInt8Array::from_iter_values(
806                op_types.iter().map(|v| *v as u8),
807            )),
808        ];
809
810        RecordBatch::try_new(build_test_multi_field_schema(), columns).unwrap()
811    }
812
813    /// Creates a test string dictionary primary key array for given primary keys.
814    fn build_test_pk_string_dict_array(primary_keys: &[&[u8]]) -> ArrayRef {
815        let mut builder = StringDictionaryBuilder::<UInt32Type>::new();
816        for &pk in primary_keys {
817            let pk_str = std::str::from_utf8(pk).unwrap();
818            builder.append(pk_str).unwrap();
819        }
820        Arc::new(builder.finish())
821    }
822
823    /// Creates a test primary key array for given primary keys.
824    fn build_test_pk_array(primary_keys: &[&[u8]]) -> ArrayRef {
825        let mut builder = BinaryDictionaryBuilder::<UInt32Type>::new();
826        for &pk in primary_keys {
827            builder.append(pk).unwrap();
828        }
829        Arc::new(builder.finish())
830    }
831
832    /// Builds the arrow schema for test flat format.
833    fn build_test_flat_schema() -> SchemaRef {
834        let fields = vec![
835            Field::new(
836                "k0",
837                DataType::Dictionary(Box::new(DataType::UInt32), Box::new(DataType::Utf8)),
838                false,
839            ),
840            Field::new("field0", DataType::Int64, true),
841            Field::new(
842                "ts",
843                DataType::Timestamp(TimeUnit::Millisecond, None),
844                false,
845            ),
846            Field::new(
847                "__primary_key",
848                DataType::Dictionary(Box::new(DataType::UInt32), Box::new(DataType::Binary)),
849                false,
850            ),
851            Field::new("__sequence", DataType::UInt64, false),
852            Field::new("__op_type", DataType::UInt8, false),
853        ];
854        Arc::new(Schema::new(fields))
855    }
856
857    /// Builds the arrow schema for test flat format with multiple fields.
858    fn build_test_multi_field_schema() -> SchemaRef {
859        let fields = vec![
860            Field::new(
861                "k0",
862                DataType::Dictionary(Box::new(DataType::UInt32), Box::new(DataType::Utf8)),
863                false,
864            ),
865            Field::new("field0", DataType::Int64, true),
866            Field::new("field1", DataType::Int64, true),
867            Field::new(
868                "ts",
869                DataType::Timestamp(TimeUnit::Millisecond, None),
870                false,
871            ),
872            Field::new(
873                "__primary_key",
874                DataType::Dictionary(Box::new(DataType::UInt32), Box::new(DataType::Binary)),
875                false,
876            ),
877            Field::new("__sequence", DataType::UInt64, false),
878            Field::new("__op_type", DataType::UInt8, false),
879        ];
880        Arc::new(Schema::new(fields))
881    }
882
883    /// Asserts that two RecordBatch vectors are equal.
884    fn check_record_batches_equal(expected: &[RecordBatch], actual: &[RecordBatch]) {
885        for (i, (exp, act)) in expected.iter().zip(actual.iter()).enumerate() {
886            assert_eq!(exp, act, "RecordBatch {} differs", i);
887        }
888        assert_eq!(
889            expected.len(),
890            actual.len(),
891            "Number of batches don't match"
892        );
893    }
894
895    /// Helper function to collect iterator results.
896    fn collect_iterator_results<I>(iter: I) -> Vec<RecordBatch>
897    where
898        I: Iterator<Item = Result<RecordBatch>>,
899    {
900        iter.map(|result| result.unwrap()).collect()
901    }
902
903    #[test]
904    fn test_flat_last_row_no_duplications() {
905        let input = vec![
906            new_record_batch(
907                &[b"k1", b"k1"],
908                &[1, 2],
909                &[11, 12],
910                &[OpType::Put, OpType::Put],
911                &[21, 22],
912            ),
913            new_record_batch(&[b"k1"], &[3], &[13], &[OpType::Put], &[23]),
914            new_record_batch(
915                &[b"k2", b"k2"],
916                &[1, 2],
917                &[111, 112],
918                &[OpType::Put, OpType::Put],
919                &[31, 32],
920            ),
921        ];
922
923        // Test with filter_deleted = true
924        let iter = input.clone().into_iter().map(Ok);
925        let mut dedup_iter = FlatDedupIterator::new(iter, FlatLastRow::new(true));
926        let result = collect_iterator_results(&mut dedup_iter);
927        check_record_batches_equal(&input, &result);
928        assert_eq!(0, dedup_iter.metrics.num_unselected_rows);
929        assert_eq!(0, dedup_iter.metrics.num_deleted_rows);
930
931        // Test with filter_deleted = false
932        let iter = input.clone().into_iter().map(Ok);
933        let mut dedup_iter = FlatDedupIterator::new(iter, FlatLastRow::new(false));
934        let result = collect_iterator_results(&mut dedup_iter);
935        check_record_batches_equal(&input, &result);
936        assert_eq!(0, dedup_iter.metrics.num_unselected_rows);
937        assert_eq!(0, dedup_iter.metrics.num_deleted_rows);
938    }
939
940    #[test]
941    fn test_flat_last_row_duplications() {
942        let input = vec![
943            new_record_batch(
944                &[b"k1", b"k1"],
945                &[1, 2],
946                &[13, 11],
947                &[OpType::Put, OpType::Put],
948                &[11, 12],
949            ),
950            // empty batch.
951            new_record_batch(&[], &[], &[], &[], &[]),
952            // Duplicate with the previous batch.
953            new_record_batch(
954                &[b"k1", b"k1", b"k1"],
955                &[2, 3, 4],
956                &[10, 13, 13],
957                &[OpType::Put, OpType::Put, OpType::Delete],
958                &[2, 13, 14],
959            ),
960            new_record_batch(
961                &[b"k2", b"k2"],
962                &[1, 2],
963                &[20, 20],
964                &[OpType::Put, OpType::Delete],
965                &[101, 0],
966            ),
967            new_record_batch(&[b"k2"], &[2], &[19], &[OpType::Put], &[102]),
968            new_record_batch(&[b"k3"], &[2], &[20], &[OpType::Put], &[202]),
969            // This batch won't increase the deleted rows count as it
970            // is filtered out by the previous batch.
971            new_record_batch(&[b"k3"], &[2], &[19], &[OpType::Delete], &[0]),
972        ];
973
974        // Test with filter_deleted = true
975        let expected_filter_deleted = vec![
976            new_record_batch(
977                &[b"k1", b"k1"],
978                &[1, 2],
979                &[13, 11],
980                &[OpType::Put, OpType::Put],
981                &[11, 12],
982            ),
983            new_record_batch(&[b"k1"], &[3], &[13], &[OpType::Put], &[13]),
984            new_record_batch(&[b"k2"], &[1], &[20], &[OpType::Put], &[101]),
985            new_record_batch(&[b"k3"], &[2], &[20], &[OpType::Put], &[202]),
986        ];
987
988        let iter = input.clone().into_iter().map(Ok);
989        let mut dedup_iter = FlatDedupIterator::new(iter, FlatLastRow::new(true));
990        let result = collect_iterator_results(&mut dedup_iter);
991        check_record_batches_equal(&expected_filter_deleted, &result);
992        assert_eq!(5, dedup_iter.metrics.num_unselected_rows);
993        assert_eq!(2, dedup_iter.metrics.num_deleted_rows);
994
995        // Test with filter_deleted = false
996        let expected_no_filter = vec![
997            new_record_batch(
998                &[b"k1", b"k1"],
999                &[1, 2],
1000                &[13, 11],
1001                &[OpType::Put, OpType::Put],
1002                &[11, 12],
1003            ),
1004            new_record_batch(
1005                &[b"k1", b"k1"],
1006                &[3, 4],
1007                &[13, 13],
1008                &[OpType::Put, OpType::Delete],
1009                &[13, 14],
1010            ),
1011            new_record_batch(
1012                &[b"k2", b"k2"],
1013                &[1, 2],
1014                &[20, 20],
1015                &[OpType::Put, OpType::Delete],
1016                &[101, 0],
1017            ),
1018            new_record_batch(&[b"k3"], &[2], &[20], &[OpType::Put], &[202]),
1019        ];
1020
1021        let iter = input.clone().into_iter().map(Ok);
1022        let mut dedup_iter = FlatDedupIterator::new(iter, FlatLastRow::new(false));
1023        let result = collect_iterator_results(&mut dedup_iter);
1024        check_record_batches_equal(&expected_no_filter, &result);
1025        assert_eq!(3, dedup_iter.metrics.num_unselected_rows);
1026        assert_eq!(0, dedup_iter.metrics.num_deleted_rows);
1027    }
1028
1029    #[test]
1030    fn exact_sequence_filter_component_composition_covers_tombstones_and_last_non_null() {
1031        use store_api::storage::SequenceRange;
1032
1033        use crate::read::scan_util::filter_flat_batch_by_sequence;
1034
1035        let filter = |batch| {
1036            filter_flat_batch_by_sequence(
1037                batch,
1038                Some(SequenceRange::GtLtEq { min: 1, max: 3 }),
1039                true,
1040            )
1041            .unwrap()
1042            .unwrap()
1043        };
1044
1045        // The shared merge presents duplicate rows newest-first: the in-range
1046        // tombstone (seq 3) reaches shared dedup and suppresses the eligible
1047        // put (seq 2) across batches.
1048        let output = FlatDedupIterator::new(
1049            vec![
1050                Ok(filter(new_record_batch_multi_fields(
1051                    &[b"series"],
1052                    &[1000],
1053                    &[3],
1054                    &[OpType::Delete],
1055                    &[(None, None)],
1056                ))),
1057                Ok(filter(new_record_batch_multi_fields(
1058                    &[b"series"],
1059                    &[1000],
1060                    &[2],
1061                    &[OpType::Put],
1062                    &[(Some(10), None)],
1063                ))),
1064            ]
1065            .into_iter(),
1066            FlatLastNonNull::new(1, true),
1067        )
1068        .collect::<Result<Vec<_>>>()
1069        .unwrap();
1070        assert!(output.is_empty());
1071
1072        // A newest outside-range tombstone is removed before shared dedup and
1073        // cannot suppress the eligible put.
1074        let output = FlatDedupIterator::new(
1075            vec![Ok(filter(new_record_batch_multi_fields(
1076                &[b"series", b"series"],
1077                &[1000, 1000],
1078                &[4, 2],
1079                &[OpType::Delete, OpType::Put],
1080                &[(None, None), (Some(10), None)],
1081            )))]
1082            .into_iter(),
1083            FlatLastNonNull::new(1, true),
1084        )
1085        .collect::<Result<Vec<_>>>()
1086        .unwrap();
1087        assert_eq!(1, output[0].num_rows());
1088
1089        // A newest outside-range non-null value is removed before
1090        // LastNonNull, so it cannot fill the eligible row's null field.
1091        let output = FlatDedupIterator::new(
1092            vec![Ok(filter(new_record_batch_multi_fields(
1093                &[b"series", b"series"],
1094                &[1000, 1000],
1095                &[4, 2],
1096                &[OpType::Put, OpType::Put],
1097                &[(None, Some(40)), (Some(10), None)],
1098            )))]
1099            .into_iter(),
1100            FlatLastNonNull::new(1, true),
1101        )
1102        .collect::<Result<Vec<_>>>()
1103        .unwrap();
1104        assert_eq!(1, output[0].num_rows());
1105        assert!(output[0].column(2).is_null(0));
1106    }
1107
1108    #[test]
1109    fn test_flat_last_non_null_no_duplications() {
1110        let input = vec![
1111            new_record_batch(
1112                &[b"k1", b"k1"],
1113                &[1, 2],
1114                &[11, 12],
1115                &[OpType::Put, OpType::Put],
1116                &[21, 22],
1117            ),
1118            new_record_batch(&[b"k1"], &[3], &[13], &[OpType::Put], &[23]),
1119            new_record_batch(
1120                &[b"k2", b"k2"],
1121                &[1, 2],
1122                &[111, 112],
1123                &[OpType::Put, OpType::Put],
1124                &[31, 32],
1125            ),
1126        ];
1127
1128        // Test with filter_deleted = true
1129        let iter = input.clone().into_iter().map(Ok);
1130        let mut dedup_iter = FlatDedupIterator::new(iter, FlatLastNonNull::new(1, true));
1131        let result = collect_iterator_results(&mut dedup_iter);
1132        check_record_batches_equal(&input, &result);
1133        assert_eq!(0, dedup_iter.metrics.num_unselected_rows);
1134        assert_eq!(0, dedup_iter.metrics.num_deleted_rows);
1135
1136        // Test with filter_deleted = false
1137        let iter = input.clone().into_iter().map(Ok);
1138        let mut dedup_iter = FlatDedupIterator::new(iter, FlatLastNonNull::new(1, false));
1139        let result = collect_iterator_results(&mut dedup_iter);
1140        check_record_batches_equal(&input, &result);
1141        assert_eq!(0, dedup_iter.metrics.num_unselected_rows);
1142        assert_eq!(0, dedup_iter.metrics.num_deleted_rows);
1143    }
1144
1145    #[test]
1146    fn test_flat_last_non_null_field_merging() {
1147        let input = vec![
1148            new_record_batch_multi_fields(
1149                &[b"k1", b"k1"],
1150                &[1, 2],
1151                &[13, 11],
1152                &[OpType::Put, OpType::Put],
1153                &[(Some(11), Some(11)), (None, None)],
1154            ),
1155            // empty batch
1156            new_record_batch_multi_fields(&[], &[], &[], &[], &[]),
1157            // Duplicate with the previous batch - should merge fields
1158            new_record_batch_multi_fields(
1159                &[b"k1"],
1160                &[2],
1161                &[10],
1162                &[OpType::Put],
1163                &[(Some(12), None)],
1164            ),
1165            new_record_batch_multi_fields(
1166                &[b"k1", b"k1", b"k1"],
1167                &[2, 3, 4],
1168                &[10, 13, 13],
1169                &[OpType::Put, OpType::Put, OpType::Delete],
1170                &[(Some(2), Some(22)), (Some(13), None), (None, Some(14))],
1171            ),
1172            new_record_batch_multi_fields(
1173                &[b"k2", b"k2"],
1174                &[1, 2],
1175                &[20, 20],
1176                &[OpType::Put, OpType::Delete],
1177                &[(Some(101), Some(101)), (None, None)],
1178            ),
1179            new_record_batch_multi_fields(
1180                &[b"k2"],
1181                &[2],
1182                &[19],
1183                &[OpType::Put],
1184                &[(Some(102), Some(102))],
1185            ),
1186            new_record_batch_multi_fields(
1187                &[b"k3"],
1188                &[2],
1189                &[20],
1190                &[OpType::Put],
1191                &[(Some(202), Some(202))],
1192            ),
1193            // This batch won't increase the deleted rows count as it
1194            // is filtered out by the previous batch. (All fields are null).
1195            new_record_batch_multi_fields(
1196                &[b"k3"],
1197                &[2],
1198                &[19],
1199                &[OpType::Delete],
1200                &[(None, None)],
1201            ),
1202        ];
1203
1204        // Test with filter_deleted = true
1205        let expected_filter_deleted = vec![
1206            new_record_batch_multi_fields(
1207                &[b"k1"],
1208                &[1],
1209                &[13],
1210                &[OpType::Put],
1211                &[(Some(11), Some(11))],
1212            ),
1213            new_record_batch_multi_fields(
1214                &[b"k1", b"k1"],
1215                &[2, 3],
1216                &[11, 13],
1217                &[OpType::Put, OpType::Put],
1218                &[(Some(12), Some(22)), (Some(13), None)],
1219            ),
1220            new_record_batch_multi_fields(
1221                &[b"k2"],
1222                &[1],
1223                &[20],
1224                &[OpType::Put],
1225                &[(Some(101), Some(101))],
1226            ),
1227            new_record_batch_multi_fields(
1228                &[b"k3"],
1229                &[2],
1230                &[20],
1231                &[OpType::Put],
1232                &[(Some(202), Some(202))],
1233            ),
1234        ];
1235
1236        let iter = input.clone().into_iter().map(Ok);
1237        let mut dedup_iter = FlatDedupIterator::new(iter, FlatLastNonNull::new(1, true));
1238        let result = collect_iterator_results(&mut dedup_iter);
1239        check_record_batches_equal(&expected_filter_deleted, &result);
1240        assert_eq!(6, dedup_iter.metrics.num_unselected_rows);
1241        assert_eq!(2, dedup_iter.metrics.num_deleted_rows);
1242
1243        // Test with filter_deleted = false
1244        let expected_no_filter = vec![
1245            new_record_batch_multi_fields(
1246                &[b"k1"],
1247                &[1],
1248                &[13],
1249                &[OpType::Put],
1250                &[(Some(11), Some(11))],
1251            ),
1252            new_record_batch_multi_fields(
1253                &[b"k1", b"k1", b"k1"],
1254                &[2, 3, 4],
1255                &[11, 13, 13],
1256                &[OpType::Put, OpType::Put, OpType::Delete],
1257                &[(Some(12), Some(22)), (Some(13), None), (None, Some(14))],
1258            ),
1259            new_record_batch_multi_fields(
1260                &[b"k2"],
1261                &[1],
1262                &[20],
1263                &[OpType::Put],
1264                &[(Some(101), Some(101))],
1265            ),
1266            new_record_batch_multi_fields(
1267                &[b"k2"],
1268                &[2],
1269                &[20],
1270                &[OpType::Delete],
1271                &[(None, None)],
1272            ),
1273            new_record_batch_multi_fields(
1274                &[b"k3"],
1275                &[2],
1276                &[20],
1277                &[OpType::Put],
1278                &[(Some(202), Some(202))],
1279            ),
1280        ];
1281
1282        let iter = input.clone().into_iter().map(Ok);
1283        let mut dedup_iter = FlatDedupIterator::new(iter, FlatLastNonNull::new(1, false));
1284        let result = collect_iterator_results(&mut dedup_iter);
1285        check_record_batches_equal(&expected_no_filter, &result);
1286        assert_eq!(4, dedup_iter.metrics.num_unselected_rows);
1287        assert_eq!(0, dedup_iter.metrics.num_deleted_rows);
1288    }
1289
1290    #[test]
1291    fn test_flat_last_non_null_skip_merge_no_null() {
1292        let input = vec![
1293            new_record_batch_multi_fields(
1294                &[b"k1", b"k1"],
1295                &[1, 2],
1296                &[13, 11],
1297                &[OpType::Put, OpType::Put],
1298                &[(Some(11), Some(11)), (Some(12), Some(12))],
1299            ),
1300            new_record_batch_multi_fields(
1301                &[b"k1"],
1302                &[2],
1303                &[10],
1304                &[OpType::Put],
1305                &[(None, Some(22))],
1306            ),
1307            new_record_batch_multi_fields(
1308                &[b"k1", b"k1"],
1309                &[2, 3],
1310                &[9, 13],
1311                &[OpType::Put, OpType::Put],
1312                &[(Some(32), None), (Some(13), Some(13))],
1313            ),
1314        ];
1315
1316        let expected = vec![
1317            new_record_batch_multi_fields(
1318                &[b"k1"],
1319                &[1],
1320                &[13],
1321                &[OpType::Put],
1322                &[(Some(11), Some(11))],
1323            ),
1324            new_record_batch_multi_fields(
1325                &[b"k1", b"k1"],
1326                &[2, 3],
1327                &[11, 13],
1328                &[OpType::Put, OpType::Put],
1329                &[(Some(12), Some(12)), (Some(13), Some(13))],
1330            ),
1331        ];
1332
1333        let iter = input.into_iter().map(Ok);
1334        let mut dedup_iter = FlatDedupIterator::new(iter, FlatLastNonNull::new(1, true));
1335        let result = collect_iterator_results(&mut dedup_iter);
1336        check_record_batches_equal(&expected, &result);
1337        assert_eq!(2, dedup_iter.metrics.num_unselected_rows);
1338        assert_eq!(0, dedup_iter.metrics.num_deleted_rows);
1339    }
1340
1341    #[test]
1342    fn test_flat_last_non_null_merge_null() {
1343        let input = vec![
1344            new_record_batch_multi_fields(
1345                &[b"k1", b"k1"],
1346                &[1, 2],
1347                &[13, 11],
1348                &[OpType::Put, OpType::Put],
1349                &[(Some(11), Some(11)), (None, None)],
1350            ),
1351            new_record_batch_multi_fields(
1352                &[b"k1"],
1353                &[2],
1354                &[10],
1355                &[OpType::Put],
1356                &[(None, Some(22))],
1357            ),
1358            new_record_batch_multi_fields(
1359                &[b"k1"],
1360                &[3],
1361                &[13],
1362                &[OpType::Put],
1363                &[(Some(33), None)],
1364            ),
1365        ];
1366
1367        let expected = vec![
1368            new_record_batch_multi_fields(
1369                &[b"k1"],
1370                &[1],
1371                &[13],
1372                &[OpType::Put],
1373                &[(Some(11), Some(11))],
1374            ),
1375            new_record_batch_multi_fields(
1376                &[b"k1"],
1377                &[2],
1378                &[11],
1379                &[OpType::Put],
1380                &[(None, Some(22))],
1381            ),
1382            new_record_batch_multi_fields(
1383                &[b"k1"],
1384                &[3],
1385                &[13],
1386                &[OpType::Put],
1387                &[(Some(33), None)],
1388            ),
1389        ];
1390
1391        let iter = input.into_iter().map(Ok);
1392        let mut dedup_iter = FlatDedupIterator::new(iter, FlatLastNonNull::new(1, true));
1393        let result = collect_iterator_results(&mut dedup_iter);
1394        check_record_batches_equal(&expected, &result);
1395        assert_eq!(1, dedup_iter.metrics.num_unselected_rows);
1396        assert_eq!(0, dedup_iter.metrics.num_deleted_rows);
1397    }
1398
1399    /// Helper function to check dedup strategy behavior directly.
1400    fn check_flat_dedup_strategy(
1401        input: &[RecordBatch],
1402        strategy: &mut dyn RecordBatchDedupStrategy,
1403        expect: &[RecordBatch],
1404    ) {
1405        let mut actual = Vec::new();
1406        let mut metrics = DedupMetrics::default();
1407        for batch in input {
1408            if let Some(out) = strategy.push_batch(batch.clone(), &mut metrics).unwrap() {
1409                actual.push(out);
1410            }
1411        }
1412        if let Some(out) = strategy.finish(&mut metrics).unwrap() {
1413            actual.push(out);
1414        }
1415
1416        check_record_batches_equal(expect, &actual);
1417    }
1418
1419    #[test]
1420    fn test_flat_last_non_null_strategy_delete_last() {
1421        let input = vec![
1422            new_record_batch_multi_fields(
1423                &[b"k1"],
1424                &[1],
1425                &[6],
1426                &[OpType::Put],
1427                &[(Some(11), None)],
1428            ),
1429            new_record_batch_multi_fields(
1430                &[b"k1", b"k1"],
1431                &[1, 2],
1432                &[1, 7],
1433                &[OpType::Put, OpType::Put],
1434                &[(Some(1), None), (Some(22), Some(222))],
1435            ),
1436            new_record_batch_multi_fields(
1437                &[b"k1"],
1438                &[2],
1439                &[4],
1440                &[OpType::Put],
1441                &[(Some(12), None)],
1442            ),
1443            new_record_batch_multi_fields(
1444                &[b"k2", b"k2"],
1445                &[2, 3],
1446                &[2, 5],
1447                &[OpType::Put, OpType::Delete],
1448                &[(None, None), (Some(13), None)],
1449            ),
1450            new_record_batch_multi_fields(&[b"k2"], &[3], &[3], &[OpType::Put], &[(None, Some(3))]),
1451        ];
1452
1453        let mut strategy = FlatLastNonNull::new(1, true);
1454        check_flat_dedup_strategy(
1455            &input,
1456            &mut strategy,
1457            &[
1458                new_record_batch_multi_fields(
1459                    &[b"k1"],
1460                    &[1],
1461                    &[6],
1462                    &[OpType::Put],
1463                    &[(Some(11), None)],
1464                ),
1465                new_record_batch_multi_fields(
1466                    &[b"k1"],
1467                    &[2],
1468                    &[7],
1469                    &[OpType::Put],
1470                    &[(Some(22), Some(222))],
1471                ),
1472                new_record_batch_multi_fields(
1473                    &[b"k2"],
1474                    &[2],
1475                    &[2],
1476                    &[OpType::Put],
1477                    &[(None, None)],
1478                ),
1479            ],
1480        );
1481    }
1482
1483    #[test]
1484    fn test_flat_last_non_null_strategy_delete_one() {
1485        let input = vec![
1486            new_record_batch_multi_fields(&[b"k1"], &[1], &[1], &[OpType::Delete], &[(None, None)]),
1487            new_record_batch_multi_fields(
1488                &[b"k2"],
1489                &[1],
1490                &[6],
1491                &[OpType::Put],
1492                &[(Some(11), None)],
1493            ),
1494        ];
1495
1496        let mut strategy = FlatLastNonNull::new(1, true);
1497        check_flat_dedup_strategy(
1498            &input,
1499            &mut strategy,
1500            &[new_record_batch_multi_fields(
1501                &[b"k2"],
1502                &[1],
1503                &[6],
1504                &[OpType::Put],
1505                &[(Some(11), None)],
1506            )],
1507        );
1508    }
1509
1510    #[test]
1511    fn test_flat_last_non_null_strategy_delete_all() {
1512        let input = vec![
1513            new_record_batch_multi_fields(&[b"k1"], &[1], &[1], &[OpType::Delete], &[(None, None)]),
1514            new_record_batch_multi_fields(
1515                &[b"k2"],
1516                &[1],
1517                &[6],
1518                &[OpType::Delete],
1519                &[(Some(11), None)],
1520            ),
1521        ];
1522
1523        let mut strategy = FlatLastNonNull::new(1, true);
1524        check_flat_dedup_strategy(&input, &mut strategy, &[]);
1525    }
1526
1527    #[test]
1528    fn test_flat_last_non_null_strategy_same_batch() {
1529        let input = vec![
1530            new_record_batch_multi_fields(
1531                &[b"k1"],
1532                &[1],
1533                &[6],
1534                &[OpType::Put],
1535                &[(Some(11), None)],
1536            ),
1537            new_record_batch_multi_fields(
1538                &[b"k1", b"k1"],
1539                &[1, 2],
1540                &[1, 7],
1541                &[OpType::Put, OpType::Put],
1542                &[(Some(1), None), (Some(22), Some(222))],
1543            ),
1544            new_record_batch_multi_fields(
1545                &[b"k1"],
1546                &[2],
1547                &[4],
1548                &[OpType::Put],
1549                &[(Some(12), None)],
1550            ),
1551            new_record_batch_multi_fields(
1552                &[b"k1", b"k1"],
1553                &[2, 3],
1554                &[2, 5],
1555                &[OpType::Put, OpType::Put],
1556                &[(None, None), (Some(13), None)],
1557            ),
1558            new_record_batch_multi_fields(&[b"k1"], &[3], &[3], &[OpType::Put], &[(None, Some(3))]),
1559        ];
1560
1561        let mut strategy = FlatLastNonNull::new(1, true);
1562        check_flat_dedup_strategy(
1563            &input,
1564            &mut strategy,
1565            &[
1566                new_record_batch_multi_fields(
1567                    &[b"k1"],
1568                    &[1],
1569                    &[6],
1570                    &[OpType::Put],
1571                    &[(Some(11), None)],
1572                ),
1573                new_record_batch_multi_fields(
1574                    &[b"k1"],
1575                    &[2],
1576                    &[7],
1577                    &[OpType::Put],
1578                    &[(Some(22), Some(222))],
1579                ),
1580                new_record_batch_multi_fields(
1581                    &[b"k1"],
1582                    &[3],
1583                    &[5],
1584                    &[OpType::Put],
1585                    &[(Some(13), Some(3))],
1586                ),
1587            ],
1588        );
1589    }
1590
1591    #[test]
1592    fn test_flat_last_non_null_strategy_delete_middle() {
1593        let input = vec![
1594            new_record_batch_multi_fields(
1595                &[b"k1"],
1596                &[1],
1597                &[7],
1598                &[OpType::Put],
1599                &[(Some(11), None)],
1600            ),
1601            new_record_batch_multi_fields(&[b"k1"], &[1], &[4], &[OpType::Delete], &[(None, None)]),
1602            new_record_batch_multi_fields(
1603                &[b"k1"],
1604                &[1],
1605                &[1],
1606                &[OpType::Put],
1607                &[(Some(12), Some(1))],
1608            ),
1609            new_record_batch_multi_fields(
1610                &[b"k1"],
1611                &[2],
1612                &[8],
1613                &[OpType::Put],
1614                &[(Some(21), None)],
1615            ),
1616            new_record_batch_multi_fields(&[b"k1"], &[2], &[5], &[OpType::Delete], &[(None, None)]),
1617            new_record_batch_multi_fields(
1618                &[b"k1"],
1619                &[2],
1620                &[2],
1621                &[OpType::Put],
1622                &[(Some(22), Some(2))],
1623            ),
1624            new_record_batch_multi_fields(
1625                &[b"k1"],
1626                &[3],
1627                &[9],
1628                &[OpType::Put],
1629                &[(Some(31), None)],
1630            ),
1631            new_record_batch_multi_fields(&[b"k1"], &[3], &[6], &[OpType::Delete], &[(None, None)]),
1632            new_record_batch_multi_fields(
1633                &[b"k1"],
1634                &[3],
1635                &[3],
1636                &[OpType::Put],
1637                &[(Some(32), Some(3))],
1638            ),
1639        ];
1640
1641        let mut strategy = FlatLastNonNull::new(1, true);
1642        check_flat_dedup_strategy(
1643            &input,
1644            &mut strategy,
1645            &[
1646                new_record_batch_multi_fields(
1647                    &[b"k1"],
1648                    &[1],
1649                    &[7],
1650                    &[OpType::Put],
1651                    &[(Some(11), None)],
1652                ),
1653                new_record_batch_multi_fields(
1654                    &[b"k1"],
1655                    &[2],
1656                    &[8],
1657                    &[OpType::Put],
1658                    &[(Some(21), None)],
1659                ),
1660                new_record_batch_multi_fields(
1661                    &[b"k1"],
1662                    &[3],
1663                    &[9],
1664                    &[OpType::Put],
1665                    &[(Some(31), None)],
1666                ),
1667            ],
1668        );
1669    }
1670}