Skip to main content

mito2/sst/index/
primary_key.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
15use 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
23/// Iterates consecutive dictionary-key runs, preserving their row order without decoding.
24pub(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(|&&current| 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}