mito2/sst/index/
primary_key.rs1use datatypes::arrow::array::{Array, BinaryArray};
16use datatypes::arrow::record_batch::RecordBatch;
17use snafu::{OptionExt, ensure};
18
19use crate::error::{InvalidRecordBatchSnafu, Result};
20use crate::sst::parquet::flat_format::primary_key_column_index;
21use crate::sst::parquet::format::PrimaryKeyArray;
22
23pub(crate) struct PrimaryKeyRuns<'a> {
25 keys: &'a [u32],
26 values: &'a BinaryArray,
27}
28
29impl<'a> PrimaryKeyRuns<'a> {
30 pub(crate) fn try_new(batch: &'a RecordBatch) -> Result<Self> {
31 let pk = batch
32 .column(primary_key_column_index(batch.num_columns()))
33 .as_any()
34 .downcast_ref::<PrimaryKeyArray>()
35 .context(InvalidRecordBatchSnafu {
36 reason: "Primary key column is not a dictionary array",
37 })?;
38 let values = pk.values().as_any().downcast_ref::<BinaryArray>().context(
39 InvalidRecordBatchSnafu {
40 reason: "Primary key values are not binary array",
41 },
42 )?;
43 ensure!(
44 pk.null_count() == 0 && values.null_count() == 0,
45 InvalidRecordBatchSnafu {
46 reason: "Primary keys must not be null"
47 }
48 );
49 Ok(Self {
50 keys: pk.keys().values(),
51 values,
52 })
53 }
54}
55
56impl<'a> Iterator for PrimaryKeyRuns<'a> {
57 type Item = (&'a [u8], usize);
58
59 fn next(&mut self) -> Option<Self::Item> {
60 let &key = self.keys.first()?;
61 let count = self
62 .keys
63 .iter()
64 .take_while(|&¤t| current == key)
65 .count();
66 self.keys = &self.keys[count..];
67 Some((self.values.value(key as usize), count))
68 }
69}
70
71#[cfg(test)]
72mod tests {
73 use std::sync::Arc;
74
75 use datatypes::arrow::array::{ArrayRef, BinaryDictionaryBuilder, UInt8Array};
76 use datatypes::arrow::datatypes::UInt32Type;
77
78 use super::*;
79
80 #[test]
81 fn sliced_runs_preserve_order_and_nonconsecutive_keys() {
82 let mut keys = BinaryDictionaryBuilder::<UInt32Type>::new();
83 for key in ["a", "a", "b", "b", "b", "a", "c"] {
84 keys.append(key).unwrap();
85 }
86 let batch = RecordBatch::try_from_iter([
87 ("pk", Arc::new(keys.finish()) as ArrayRef),
88 ("seq", Arc::new(UInt8Array::from(vec![0; 7])) as ArrayRef),
89 ("op", Arc::new(UInt8Array::from(vec![0; 7])) as ArrayRef),
90 ])
91 .unwrap();
92 let slice = batch.slice(1, 5);
93 assert_eq!(
94 PrimaryKeyRuns::try_new(&slice).unwrap().collect::<Vec<_>>(),
95 vec![
96 (b"a".as_slice(), 1),
97 (b"b".as_slice(), 3),
98 (b"a".as_slice(), 1)
99 ]
100 );
101 assert!(
102 PrimaryKeyRuns::try_new(&batch.slice(0, 0))
103 .unwrap()
104 .next()
105 .is_none()
106 );
107 }
108}