Skip to main content

mito2/sst/range_index/
searcher.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 std::cmp::Ordering;
16use std::ops::Range;
17
18use datafusion_expr::{col, lit};
19use datatypes::arrow::array::{Int64Array, UInt32Array, UInt64Array};
20use datatypes::arrow::datatypes::{DataType, SchemaRef};
21use futures::TryStreamExt;
22use object_store::ObjectStore;
23use snafu::{OptionExt, ensure};
24use table::predicate::Predicate;
25
26use crate::error::{InvalidRecordBatchSnafu, Result, UnexpectedSnafu};
27use crate::series_index::MetricSeriesId;
28use crate::sst::parquet::index_reader::ParquetIndexReader;
29use crate::sst::range_index::{
30    END_COLUMN, ROW_GROUP_ID_COLUMN, START_COLUMN, TABLE_ID_COLUMN, TSID_COLUMN,
31};
32
33/// Searches per-SST range-index files for the rows of candidate metric series.
34pub struct SstRangeIndexSearcher {
35    reader: ParquetIndexReader,
36}
37
38impl SstRangeIndexSearcher {
39    /// Opens the range-index file at `path` and loads its Parquet metadata.
40    pub async fn open(object_store: ObjectStore, path: &str) -> Result<Self> {
41        let reader = ParquetIndexReader::open(object_store, path).await?;
42        validate_index_schema(reader.schema())?;
43        Ok(Self { reader })
44    }
45
46    /// Returns the row ranges for `series` in one source SST row group.
47    ///
48    /// `series` is one batch emitted by a
49    /// [`MetricSeriesIdStream`](crate::series_index::MetricSeriesIdStream). The
50    /// returned half-open ranges are relative to the start of `row_group_id`,
51    /// sorted, non-overlapping, and coalesced when adjacent. The number of
52    /// returned ranges may be less than the number of input series if some
53    /// series don't exist in the row group.
54    pub async fn search(
55        &self,
56        row_group_id: u32,
57        series: &[MetricSeriesId],
58    ) -> Result<Vec<Range<usize>>> {
59        if series.is_empty() {
60            return Ok(Vec::new());
61        }
62
63        validate_sorted_series(series)?;
64        let predicate = search_predicate(row_group_id, series)?;
65        let mut batches = self.reader.read(
66            &predicate,
67            &[
68                ROW_GROUP_ID_COLUMN,
69                TABLE_ID_COLUMN,
70                TSID_COLUMN,
71                START_COLUMN,
72                END_COLUMN,
73            ],
74        )?;
75        let mut merge = RangeMergeState::new(row_group_id, series);
76
77        while let Some(batch) = batches.try_next().await? {
78            if merge.append_batch(&batch)? {
79                break;
80            }
81        }
82
83        Ok(merge.finish())
84    }
85}
86
87fn validate_sorted_series(series: &[MetricSeriesId]) -> Result<()> {
88    if let Some(pair) = series.windows(2).find(|pair| pair[0] > pair[1]) {
89        return InvalidRecordBatchSnafu {
90            reason: format!(
91                "range index search series are not sorted: {:?} appears before {:?}",
92                pair[0], pair[1]
93            ),
94        }
95        .fail();
96    }
97    Ok(())
98}
99
100fn search_predicate(row_group_id: u32, series: &[MetricSeriesId]) -> Result<Predicate> {
101    let min_table_id = series
102        .first()
103        .context(UnexpectedSnafu {
104            reason: "cannot build a range-index predicate for an empty series set",
105        })?
106        .table_id;
107    let max_table_id = series
108        .last()
109        .context(UnexpectedSnafu {
110            reason: "cannot build a range-index predicate for an empty series set",
111        })?
112        .table_id;
113
114    Ok(Predicate::new(vec![
115        col(ROW_GROUP_ID_COLUMN).eq(lit(row_group_id)),
116        col(TABLE_ID_COLUMN).gt_eq(lit(min_table_id)),
117        col(TABLE_ID_COLUMN).lt_eq(lit(max_table_id)),
118    ]))
119}
120
121fn validate_index_schema(schema: &SchemaRef) -> Result<()> {
122    for (name, data_type) in [
123        (ROW_GROUP_ID_COLUMN, DataType::UInt32),
124        (TABLE_ID_COLUMN, DataType::UInt32),
125        (TSID_COLUMN, DataType::UInt64),
126        (START_COLUMN, DataType::Int64),
127        (END_COLUMN, DataType::Int64),
128    ] {
129        let field = schema
130            .field_with_name(name)
131            .ok()
132            .with_context(|| InvalidRecordBatchSnafu {
133                reason: format!("range index is missing column {name}"),
134            })?;
135        ensure!(
136            field.data_type() == &data_type && !field.is_nullable(),
137            InvalidRecordBatchSnafu {
138                reason: format!(
139                    "range index column {name} must be non-nullable {data_type:?}, got {:?}",
140                    field.data_type()
141                ),
142            }
143        );
144    }
145    Ok(())
146}
147
148struct RangeMergeState<'a> {
149    /// Source SST row group whose ranges are being searched.
150    row_group_id: u32,
151    /// Sorted metric series to match against the range index.
152    series: &'a [MetricSeriesId],
153    /// Cursor to the next series to match.
154    series_index: usize,
155    /// Last range-index key read, used to validate ordering across batches.
156    last_index_key: Option<(u32, MetricSeriesId)>,
157    /// Matching row ranges, sorted and coalesced when adjacent.
158    ranges: Vec<Range<usize>>,
159}
160
161impl<'a> RangeMergeState<'a> {
162    fn new(row_group_id: u32, series: &'a [MetricSeriesId]) -> Self {
163        Self {
164            row_group_id,
165            series,
166            series_index: 0,
167            last_index_key: None,
168            ranges: Vec::new(),
169        }
170    }
171
172    /// Appends matches from `batch` and returns whether the merge is complete.
173    fn append_batch(
174        &mut self,
175        batch: &datatypes::arrow::record_batch::RecordBatch,
176    ) -> Result<bool> {
177        let row_group_ids = typed_column::<UInt32Array>(batch, ROW_GROUP_ID_COLUMN, "UInt32")?;
178        let table_ids = typed_column::<UInt32Array>(batch, TABLE_ID_COLUMN, "UInt32")?;
179        let tsids = typed_column::<UInt64Array>(batch, TSID_COLUMN, "UInt64")?;
180        let starts = typed_column::<Int64Array>(batch, START_COLUMN, "Int64")?;
181        let ends = typed_column::<Int64Array>(batch, END_COLUMN, "Int64")?;
182
183        for row in 0..batch.num_rows() {
184            let index_series = MetricSeriesId {
185                table_id: table_ids.value(row),
186                tsid: tsids.value(row),
187            };
188            let index_key = (row_group_ids.value(row), index_series);
189            ensure!(
190                self.last_index_key.is_none_or(|last| last < index_key),
191                InvalidRecordBatchSnafu {
192                    reason: format!(
193                        "range index rows are not strictly sorted: {index_key:?} follows {:?}",
194                        self.last_index_key
195                    ),
196                }
197            );
198            self.last_index_key = Some(index_key);
199
200            match index_key.0.cmp(&self.row_group_id) {
201                Ordering::Less => continue,
202                Ordering::Greater => return Ok(true),
203                Ordering::Equal => {}
204            }
205
206            while self.series_index < self.series.len()
207                && self.series[self.series_index] < index_series
208            {
209                self.advance_series();
210            }
211            if self.series_index == self.series.len() {
212                return Ok(true);
213            }
214
215            match self.series[self.series_index].cmp(&index_series) {
216                Ordering::Less => {
217                    return UnexpectedSnafu {
218                        reason: "range-index merge cursor did not advance past a smaller series",
219                    }
220                    .fail();
221                }
222                Ordering::Greater => continue,
223                Ordering::Equal => {
224                    self.append_range(starts.value(row), ends.value(row), row)?;
225                    self.advance_series();
226                    if self.series_index == self.series.len() {
227                        return Ok(true);
228                    }
229                }
230            }
231        }
232        Ok(false)
233    }
234
235    fn advance_series(&mut self) {
236        let current = self.series[self.series_index];
237        while self.series_index < self.series.len() && self.series[self.series_index] == current {
238            self.series_index += 1;
239        }
240    }
241
242    fn append_range(&mut self, start: i64, end: i64, row: usize) -> Result<()> {
243        let start = usize::try_from(start).map_err(|_| {
244            InvalidRecordBatchSnafu {
245                reason: format!("range index contains negative start offset at row {row}"),
246            }
247            .build()
248        })?;
249        let end = usize::try_from(end).map_err(|_| {
250            InvalidRecordBatchSnafu {
251                reason: format!("range index contains negative end offset at row {row}"),
252            }
253            .build()
254        })?;
255        ensure!(
256            start < end,
257            InvalidRecordBatchSnafu {
258                reason: format!("range index contains invalid range {start}..{end} at row {row}"),
259            }
260        );
261
262        if let Some(last) = self.ranges.last_mut() {
263            ensure!(
264                start >= last.end,
265                InvalidRecordBatchSnafu {
266                    reason: format!(
267                        "range index contains overlapping or unsorted range {start}..{end} after {}..{}",
268                        last.start, last.end
269                    ),
270                }
271            );
272            if start == last.end {
273                last.end = end;
274                return Ok(());
275            }
276        }
277        self.ranges.push(start..end);
278        Ok(())
279    }
280
281    fn finish(self) -> Vec<Range<usize>> {
282        self.ranges
283    }
284}
285
286fn typed_column<'a, T: 'static>(
287    batch: &'a datatypes::arrow::record_batch::RecordBatch,
288    name: &str,
289    data_type: &str,
290) -> Result<&'a T> {
291    let index = batch
292        .schema()
293        .index_of(name)
294        .ok()
295        .with_context(|| InvalidRecordBatchSnafu {
296            reason: format!("range index batch is missing column {name}"),
297        })?;
298    batch
299        .column(index)
300        .as_any()
301        .downcast_ref::<T>()
302        .with_context(|| InvalidRecordBatchSnafu {
303            reason: format!("range index column {name} is not {data_type}"),
304        })
305}
306
307#[cfg(test)]
308mod tests {
309    use std::sync::Arc;
310
311    use datatypes::arrow::array::{ArrayRef, BinaryArray};
312    use datatypes::arrow::datatypes::{Field, Schema};
313    use datatypes::arrow::record_batch::RecordBatch;
314    use object_store::services::Memory;
315    use store_api::codec::PrimaryKeyEncoding;
316    use store_api::metadata::RegionMetadataRef;
317    use store_api::storage::consts::PRIMARY_KEY_COLUMN_NAME;
318
319    use super::*;
320    use crate::sst::range_index::{
321        SstRangeIndexWriter, SstRangeIndexWriterOptions, range_index_schema,
322    };
323    use crate::test_util::sst_util::{new_sparse_primary_key, sst_region_metadata_with_encoding};
324
325    fn object_store() -> ObjectStore {
326        ObjectStore::new(Memory::default()).unwrap()
327    }
328
329    fn series(table_id: u32, tsid: u64) -> MetricSeriesId {
330        MetricSeriesId { table_id, tsid }
331    }
332
333    fn primary_key_batch(metadata: &RegionMetadataRef, ids: &[(u32, u64)]) -> RecordBatch {
334        let primary_keys = ids
335            .iter()
336            .map(|(table_id, tsid)| new_sparse_primary_key(&["a", "x"], metadata, *table_id, *tsid))
337            .collect::<Vec<_>>();
338        let schema = Arc::new(Schema::new(vec![Field::new(
339            PRIMARY_KEY_COLUMN_NAME,
340            DataType::Binary,
341            false,
342        )]));
343        RecordBatch::try_new(
344            schema,
345            vec![Arc::new(BinaryArray::from_iter_values(
346                primary_keys.iter().map(Vec::as_slice),
347            ))],
348        )
349        .unwrap()
350    }
351
352    async fn write_index(store: &ObjectStore, path: &str) {
353        let metadata = Arc::new(sst_region_metadata_with_encoding(
354            PrimaryKeyEncoding::Sparse,
355        ));
356        let mut writer = SstRangeIndexWriter::try_new(
357            metadata.clone(),
358            store.clone(),
359            path,
360            SstRangeIndexWriterOptions {
361                index_row_group_size: 2,
362            },
363        )
364        .await
365        .unwrap();
366        writer
367            .write(
368                0,
369                &primary_key_batch(
370                    &metadata,
371                    &[(1, 10), (1, 10), (1, 20), (2, 10), (2, 20), (2, 20)],
372                ),
373            )
374            .await
375            .unwrap();
376        writer
377            .write(1, &primary_key_batch(&metadata, &[(2, 20), (2, 20)]))
378            .await
379            .unwrap();
380        writer.finish().await.unwrap();
381    }
382
383    #[tokio::test]
384    async fn search_filters_exact_series_pairs_and_coalesces_ranges() {
385        let store = object_store();
386        let path = "range-search.parquet";
387        write_index(&store, path).await;
388        let searcher = SstRangeIndexSearcher::open(store, path).await.unwrap();
389
390        let ranges = searcher
391            .search(0, &[series(1, 10), series(2, 20)])
392            .await
393            .unwrap();
394        assert_eq!(ranges, vec![0..2, 4..6]);
395
396        let ranges = searcher
397            .search(0, &[series(1, 10), series(1, 20)])
398            .await
399            .unwrap();
400        assert_eq!(ranges, vec![0..3]);
401
402        let ranges = searcher
403            .search(0, &[series(1, 15), series(2, 20)])
404            .await
405            .unwrap();
406        assert_eq!(ranges, vec![4..6]);
407
408        let ranges = searcher
409            .search(0, &[series(1, 10), series(2, 30)])
410            .await
411            .unwrap();
412        assert_eq!(ranges, vec![0..2]);
413
414        let ranges = searcher
415            .search(1, &[series(2, 20), series(2, 20)])
416            .await
417            .unwrap();
418        assert_eq!(ranges, vec![0..2]);
419
420        assert!(
421            searcher
422                .search(1, &[series(1, 10)])
423                .await
424                .unwrap()
425                .is_empty()
426        );
427
428        assert!(searcher.search(0, &[]).await.unwrap().is_empty());
429
430        let error = searcher
431            .search(0, &[series(2, 20), series(1, 10)])
432            .await
433            .unwrap_err();
434        assert!(error.to_string().contains("not sorted"), "{error}");
435    }
436
437    #[tokio::test]
438    async fn opening_a_missing_index_fails() {
439        assert!(
440            SstRangeIndexSearcher::open(object_store(), "does-not-exist.parquet")
441                .await
442                .is_err()
443        );
444    }
445
446    #[tokio::test]
447    async fn pruning_uses_the_source_row_group_and_table_id_range() {
448        let store = object_store();
449        let path = "range-pruning.parquet";
450        write_index(&store, path).await;
451        let reader = ParquetIndexReader::open(store, path).await.unwrap();
452        let predicate = search_predicate(0, &[series(1, 999), series(2, 999)]).unwrap();
453
454        assert_eq!(reader.row_groups_to_read(&predicate), vec![0, 1]);
455    }
456
457    #[test]
458    fn validates_schema_and_range_offsets() {
459        let nullable_schema = Arc::new(Schema::new(vec![
460            Field::new(ROW_GROUP_ID_COLUMN, DataType::UInt32, false),
461            Field::new(TABLE_ID_COLUMN, DataType::UInt32, false),
462            Field::new(TSID_COLUMN, DataType::UInt64, false),
463            Field::new(START_COLUMN, DataType::Int64, true),
464            Field::new(END_COLUMN, DataType::Int64, false),
465        ]));
466        assert!(validate_index_schema(&nullable_schema).is_err());
467
468        let batch = RecordBatch::try_new(
469            range_index_schema(),
470            vec![
471                Arc::new(UInt32Array::from(vec![0])) as ArrayRef,
472                Arc::new(UInt32Array::from(vec![1])),
473                Arc::new(UInt64Array::from(vec![10])),
474                Arc::new(Int64Array::from(vec![-1])),
475                Arc::new(Int64Array::from(vec![2])),
476            ],
477        )
478        .unwrap();
479        let selected = [series(1, 10)];
480        let mut merge = RangeMergeState::new(0, &selected);
481        assert!(merge.append_batch(&batch).is_err());
482
483        let unsorted_batch = RecordBatch::try_new(
484            range_index_schema(),
485            vec![
486                Arc::new(UInt32Array::from(vec![0, 0])) as ArrayRef,
487                Arc::new(UInt32Array::from(vec![1, 1])),
488                Arc::new(UInt64Array::from(vec![20, 10])),
489                Arc::new(Int64Array::from(vec![0, 1])),
490                Arc::new(Int64Array::from(vec![1, 2])),
491            ],
492        )
493        .unwrap();
494        let selected = [series(1, 20), series(1, 30)];
495        let mut merge = RangeMergeState::new(0, &selected);
496        assert!(merge.append_batch(&unsorted_batch).is_err());
497
498        let make_batch = |tsid, start, end| {
499            RecordBatch::try_new(
500                range_index_schema(),
501                vec![
502                    Arc::new(UInt32Array::from(vec![0])) as ArrayRef,
503                    Arc::new(UInt32Array::from(vec![1])),
504                    Arc::new(UInt64Array::from(vec![tsid])),
505                    Arc::new(Int64Array::from(vec![start])),
506                    Arc::new(Int64Array::from(vec![end])),
507                ],
508            )
509            .unwrap()
510        };
511        let selected = [series(1, 10), series(1, 20)];
512        let mut merge = RangeMergeState::new(0, &selected);
513        assert!(!merge.append_batch(&make_batch(10, 0, 1)).unwrap());
514        assert!(merge.append_batch(&make_batch(20, 1, 2)).unwrap());
515        assert_eq!(merge.finish(), vec![0..2]);
516    }
517}