1use 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
49pub struct FlatDedupIterator<I, S> {
51 iter: I,
52 strategy: S,
53 metrics: DedupMetrics,
54}
55
56impl<I, S> FlatDedupIterator<I, S> {
57 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 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
90pub struct FlatDedupReader<I, S> {
92 stream: I,
93 strategy: S,
94 metrics: DedupMetrics,
95 metrics_reporter: Option<Arc<dyn DedupMetricsReport>>,
97}
98
99impl<I, S> FlatDedupReader<I, S> {
100 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 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 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 if let Some(reporter) = &self.metrics_reporter {
155 reporter.report(&mut self.metrics);
156 }
157 }
158}
159
160pub trait RecordBatchDedupStrategy: Send {
162 fn push_batch(
166 &mut self,
167 batch: RecordBatch,
168 metrics: &mut DedupMetrics,
169 ) -> Result<Option<RecordBatch>>;
170
171 fn finish(&mut self, metrics: &mut DedupMetrics) -> Result<Option<RecordBatch>>;
176}
177
178pub struct FlatLastRow {
180 prev_batch: Option<BatchLastRow>,
183 filter_deleted: bool,
185}
186
187impl FlatLastRow {
188 pub fn new(filter_deleted: bool) -> Self {
190 Self {
191 prev_batch: None,
192 filter_deleted,
193 }
194 }
195
196 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 let mask = find_boundaries(timestamps).context(ComputeArrowSnafu)?;
207 if mask.count_set_bits() == num_rows - 1 {
208 return Ok(batch);
210 }
211
212 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 fn dedup_by_partitions(batch: RecordBatch, partitions: &Partitions) -> Result<RecordBatch> {
229 let ranges = partitions.ranges();
230 let num_duplications: usize = ranges.iter().map(|r| r.end - r.start - 1).sum();
232 if num_duplications == 0 {
233 return Ok(batch);
235 }
236
237 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 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 prev_batch.is_last_row_duplicated(&batch) {
262 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 metrics.dedup_cost += start.elapsed();
273 return Ok(None);
274 };
275
276 self.prev_batch = Some(batch_last_row);
282
283 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
296pub struct FlatLastNonNull {
298 field_column_start: usize,
300 filter_deleted: bool,
302 buffer: Option<BatchLastRow>,
306 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 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 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 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 let schema = batch.schema();
369 let merged = concat_batches(&schema, &[last_row, batch]).context(ComputeArrowSnafu)?;
370 let merged_row_count = merged.num_rows();
371 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 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 fn dedup_one_batch(
413 batch: RecordBatch,
414 field_column_start: usize,
415 prev_batch_contains_delete: bool,
416 ) -> Result<(RecordBatch, bool)> {
417 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 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 return Ok((batch, contains_delete));
443 }
444
445 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 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 let num_duplications: usize = ranges.iter().map(|r| r.end - r.start - 1).sum();
480 if num_duplications == 0 {
481 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 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 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 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 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
564struct BatchLastRow {
566 last_batch: RecordBatch,
569 primary_key: PrimaryKeyArray,
571 timestamp: i64,
573}
574
575impl BatchLastRow {
576 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 fn is_last_row_duplicated(&self, batch: &RecordBatch) -> bool {
601 if batch.num_rows() == 0 {
602 return false;
603 }
604
605 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 let batch_key = primary_key_at(primary_key, 0);
622
623 last_key == batch_key
624 }
625}
626
627fn 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 let cmp = make_comparator(&v1, &v2, SortOptions::default())?;
642 Ok((0..slice_len).map(|i| !cmp(i, i).is_eq()).collect())
643}
644
645fn 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 if batch.num_rows() == 0 {
657 return Ok(None);
658 }
659 Ok(Some(batch))
660}
661
662fn 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 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
699fn 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
710pub(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 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 build_test_pk_string_dict_array(primary_keys),
750 Arc::new(Int64Array::from_iter(
752 fields.iter().map(|v| Some(*v as i64)),
753 )),
754 Arc::new(TimestampMillisecondArray::from_iter_values(
756 timestamps.iter().copied(),
757 )),
758 build_test_pk_array(primary_keys),
760 Arc::new(UInt64Array::from_iter_values(sequences.iter().copied())),
762 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 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 build_test_pk_string_dict_array(primary_keys),
788 Arc::new(Int64Array::from_iter(
790 fields.iter().map(|field| field.0.map(|v| v as i64)),
791 )),
792 Arc::new(Int64Array::from_iter(
794 fields.iter().map(|field| field.1.map(|v| v as i64)),
795 )),
796 Arc::new(TimestampMillisecondArray::from_iter_values(
798 timestamps.iter().copied(),
799 )),
800 build_test_pk_array(primary_keys),
802 Arc::new(UInt64Array::from_iter_values(sequences.iter().copied())),
804 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 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 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 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 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 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 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 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 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 new_record_batch(&[], &[], &[], &[], &[]),
952 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 new_record_batch(&[b"k3"], &[2], &[19], &[OpType::Delete], &[0]),
972 ];
973
974 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 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 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 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 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 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 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 new_record_batch_multi_fields(&[], &[], &[], &[], &[]),
1157 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 new_record_batch_multi_fields(
1196 &[b"k3"],
1197 &[2],
1198 &[19],
1199 &[OpType::Delete],
1200 &[(None, None)],
1201 ),
1202 ];
1203
1204 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 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 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}