Skip to main content

index/inverted_index/search/index_apply/
predicates_apply.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::mem::size_of;
16
17use async_trait::async_trait;
18use greptime_proto::v1::index::InvertedIndexMetas;
19
20use crate::bitmap::Bitmap;
21use crate::inverted_index::error::{IndexNotFoundSnafu, Result};
22use crate::inverted_index::format::reader::{InvertedIndexReadMetrics, InvertedIndexReader};
23use crate::inverted_index::search::fst_apply::{
24    FstApplier, IntersectionFstApplier, KeysFstApplier,
25};
26use crate::inverted_index::search::fst_values_mapper::ParallelFstValuesMapper;
27use crate::inverted_index::search::index_apply::{
28    ApplyOutput, IndexApplier, IndexNotFoundStrategy, SearchContext,
29};
30use crate::inverted_index::search::predicate::Predicate;
31
32type IndexName = String;
33
34/// `PredicatesIndexApplier` contains a collection of `FstApplier`s, each associated with an index name,
35/// to process and filter index data based on compiled predicates.
36pub struct PredicatesIndexApplier {
37    /// A list of `FstApplier`s, each associated with a specific index name
38    /// (e.g. a tag field uses its column name as index name)
39    fst_appliers: Vec<(IndexName, Box<dyn FstApplier>)>,
40}
41
42#[async_trait]
43impl IndexApplier for PredicatesIndexApplier {
44    /// Applies all `FstApplier`s to the data in the inverted index reader, intersecting the individual
45    /// bitmaps obtained for each index to result in a final set of indices.
46    async fn apply<'a, 'b>(
47        &self,
48        context: SearchContext,
49        reader: &mut (dyn InvertedIndexReader + 'a),
50        metrics: Option<&'b mut InvertedIndexReadMetrics>,
51    ) -> Result<ApplyOutput> {
52        let mut metrics = metrics;
53        let metadata = reader.metadata(metrics.as_deref_mut()).await?;
54        let mut output = ApplyOutput {
55            matched_segment_ids: Bitmap::new_bitvec(),
56            total_row_count: metadata.total_row_count as _,
57            segment_row_count: metadata.segment_row_count as _,
58        };
59
60        // TODO(zhongzc): optimize the order of applying to make it quicker to return empty.
61        let mut appliers = Vec::with_capacity(self.fst_appliers.len());
62        let mut fst_ranges = Vec::with_capacity(self.fst_appliers.len());
63
64        for (name, fst_applier) in &self.fst_appliers {
65            let Some(meta) = metadata.metas.get(name) else {
66                match context.index_not_found_strategy {
67                    IndexNotFoundStrategy::ReturnEmpty => {
68                        return Ok(output);
69                    }
70                    IndexNotFoundStrategy::Ignore => {
71                        continue;
72                    }
73                    IndexNotFoundStrategy::ThrowError => {
74                        return IndexNotFoundSnafu { name }.fail();
75                    }
76                }
77            };
78            let fst_offset = meta.base_offset + meta.relative_fst_offset as u64;
79            let fst_size = meta.fst_size as u64;
80            appliers.push((fst_applier, meta));
81            fst_ranges.push(fst_offset..fst_offset + fst_size);
82        }
83
84        if fst_ranges.is_empty() {
85            output.matched_segment_ids = Self::bitmap_full_range(&metadata);
86            return Ok(output);
87        }
88
89        let fsts = reader.fst_vec(&fst_ranges, metrics.as_deref_mut()).await?;
90        let value_and_meta_vec = fsts
91            .into_iter()
92            .zip(appliers)
93            .map(|(fst, (fst_applier, meta))| (fst_applier.apply(&fst), meta))
94            .collect::<Vec<_>>();
95
96        let mut mapper = ParallelFstValuesMapper::new(reader);
97        let bm_vec = mapper.map_values_vec(&value_and_meta_vec, metrics).await?;
98
99        let mut iter = bm_vec.into_iter();
100        let mut bitmap = iter.next().unwrap(); // SAFETY: `fst_ranges` is not empty
101        for bm in iter {
102            bitmap.intersect(bm);
103            if bitmap.count_ones() == 0 {
104                break;
105            }
106        }
107
108        output.matched_segment_ids = bitmap;
109        Ok(output)
110    }
111
112    /// Returns the memory usage of the applier.
113    fn memory_usage(&self) -> usize {
114        let mut size = self.fst_appliers.capacity() * size_of::<(IndexName, Box<dyn FstApplier>)>();
115        for (name, fst_applier) in &self.fst_appliers {
116            size += name.capacity();
117            size += fst_applier.memory_usage();
118        }
119        size
120    }
121}
122
123impl PredicatesIndexApplier {
124    /// Constructs an instance of `PredicatesIndexApplier` based on a list of tag predicates.
125    /// Chooses an appropriate `FstApplier` for each index name based on the nature of its predicates.
126    pub fn try_from(mut predicates: Vec<(IndexName, Vec<Predicate>)>) -> Result<Self> {
127        let mut fst_appliers = Vec::with_capacity(predicates.len());
128
129        // InList predicates are applied first to benefit from higher selectivity.
130        let in_list_index =
131            crate::inverted_index::search::partition_in_place(&mut predicates, |(_, ps)| {
132                ps.iter().any(|p| matches!(p, Predicate::InList(_)))
133            });
134        let mut iter = predicates.into_iter();
135        for _ in 0..in_list_index {
136            let (column_name, predicates) = iter.next().unwrap();
137            let fst_applier = Box::new(KeysFstApplier::try_from(predicates)?) as _;
138            fst_appliers.push((column_name, fst_applier));
139        }
140
141        for (column_name, predicates) in iter {
142            if predicates.is_empty() {
143                continue;
144            }
145            let fst_applier = Box::new(IntersectionFstApplier::try_from(predicates)?) as _;
146            fst_appliers.push((column_name, fst_applier));
147        }
148
149        Ok(PredicatesIndexApplier { fst_appliers })
150    }
151
152    /// Creates a `Bitmap` representing the full range of data in the index for initial scanning.
153    fn bitmap_full_range(metadata: &InvertedIndexMetas) -> Bitmap {
154        let total_count = metadata.total_row_count;
155        let segment_count = metadata.segment_row_count;
156        let len = total_count.div_ceil(segment_count);
157        Bitmap::full_bitvec(len as _)
158    }
159}
160
161impl TryFrom<Vec<(String, Vec<Predicate>)>> for PredicatesIndexApplier {
162    type Error = crate::inverted_index::error::Error;
163    fn try_from(predicates: Vec<(String, Vec<Predicate>)>) -> Result<Self> {
164        Self::try_from(predicates)
165    }
166}
167
168#[cfg(test)]
169mod tests {
170    use std::collections::VecDeque;
171    use std::sync::Arc;
172
173    use greptime_proto::v1::index::{BitmapType, InvertedIndexMeta};
174
175    use super::*;
176    use crate::bitmap::Bitmap;
177    use crate::inverted_index::FstMap;
178    use crate::inverted_index::error::Error;
179    use crate::inverted_index::format::reader::MockInvertedIndexReader;
180    use crate::inverted_index::search::fst_apply::MockFstApplier;
181
182    fn s(s: &'static str) -> String {
183        s.to_owned()
184    }
185
186    fn mock_metas(tags: impl IntoIterator<Item = (&'static str, u32)>) -> Arc<InvertedIndexMetas> {
187        let mut metas = InvertedIndexMetas {
188            total_row_count: 8,
189            segment_row_count: 1,
190            ..Default::default()
191        };
192        for (tag, idx) in tags.into_iter() {
193            let meta = InvertedIndexMeta {
194                name: s(tag),
195                relative_fst_offset: idx,
196                bitmap_type: BitmapType::Roaring.into(),
197                ..Default::default()
198            };
199            metas.metas.insert(s(tag), meta);
200        }
201        Arc::new(metas)
202    }
203
204    fn key_fst_applier(value: &'static str) -> Box<dyn FstApplier> {
205        let mut mock_fst_applier = MockFstApplier::new();
206        mock_fst_applier
207            .expect_apply()
208            .returning(move |fst| fst.get(value).into_iter().collect());
209        Box::new(mock_fst_applier)
210    }
211
212    fn fst_value(offset: u32, size: u32) -> u64 {
213        bytemuck::cast::<_, u64>([offset, size])
214    }
215
216    #[tokio::test]
217    async fn test_index_applier_apply_get_key() {
218        // An index applier that point-gets "tag-0_value-0" on tag "tag-0"
219        let applier = PredicatesIndexApplier {
220            fst_appliers: vec![(s("tag-0"), key_fst_applier("tag-0_value-0"))],
221        };
222
223        // An index reader with a single tag "tag-0" and a corresponding value "tag-0_value-0"
224        let mut mock_reader = MockInvertedIndexReader::new();
225        mock_reader
226            .expect_metadata()
227            .returning(|_| Ok(mock_metas([("tag-0", 0)])));
228        mock_reader.expect_fst_vec().returning(|_ranges, _metrics| {
229            Ok(vec![
230                FstMap::from_iter([(b"tag-0_value-0", fst_value(2, 1))]).unwrap(),
231            ])
232        });
233
234        mock_reader
235            .expect_bitmap_deque()
236            .returning(|arg, _metrics| {
237                assert_eq!(arg.len(), 1);
238                let range = &arg[0].0;
239                let bitmap_type = arg[0].1;
240                assert_eq!(*range, 2..3);
241                assert_eq!(bitmap_type, BitmapType::Roaring);
242                Ok(VecDeque::from([Bitmap::from_lsb0_bytes(
243                    &[0b10101010],
244                    bitmap_type,
245                )]))
246            });
247        let output = applier
248            .apply(SearchContext::default(), &mut mock_reader, None)
249            .await
250            .unwrap();
251        assert_eq!(
252            output.matched_segment_ids,
253            Bitmap::from_lsb0_bytes(&[0b10101010], BitmapType::Roaring)
254        );
255
256        // An index reader with a single tag "tag-0" but without value "tag-0_value-0"
257        let mut mock_reader = MockInvertedIndexReader::new();
258        mock_reader
259            .expect_metadata()
260            .returning(|_| Ok(mock_metas([("tag-0", 0)])));
261        mock_reader.expect_fst_vec().returning(|_range, _metrics| {
262            Ok(vec![
263                FstMap::from_iter([(b"tag-0_value-1", fst_value(2, 1))]).unwrap(),
264            ])
265        });
266        let output = applier
267            .apply(SearchContext::default(), &mut mock_reader, None)
268            .await
269            .unwrap();
270        assert_eq!(output.matched_segment_ids.count_ones(), 0);
271    }
272
273    #[tokio::test]
274    async fn test_index_applier_apply_intersection_with_two_tags() {
275        // An index applier that intersects "tag-0_value-0" on tag "tag-0" and "tag-1_value-a" on tag "tag-1"
276        let applier = PredicatesIndexApplier {
277            fst_appliers: vec![
278                (s("tag-0"), key_fst_applier("tag-0_value-0")),
279                (s("tag-1"), key_fst_applier("tag-1_value-a")),
280            ],
281        };
282
283        // An index reader with two tags "tag-0" and "tag-1" and respective values "tag-0_value-0" and "tag-1_value-a"
284        let mut mock_reader = MockInvertedIndexReader::new();
285        mock_reader
286            .expect_metadata()
287            .returning(|_| Ok(mock_metas([("tag-0", 0), ("tag-1", 1)])));
288        mock_reader.expect_fst_vec().returning(|ranges, _metrics| {
289            let mut output = vec![];
290            for range in ranges {
291                match range.start {
292                    0 => output
293                        .push(FstMap::from_iter([(b"tag-0_value-0", fst_value(1, 1))]).unwrap()),
294                    1 => output
295                        .push(FstMap::from_iter([(b"tag-1_value-a", fst_value(2, 1))]).unwrap()),
296                    _ => unreachable!(),
297                }
298            }
299            Ok(output)
300        });
301        mock_reader
302            .expect_bitmap_deque()
303            .returning(|ranges, _metrics| {
304                let mut output = VecDeque::new();
305                for (range, bitmap_type) in ranges {
306                    let offset = range.start;
307                    let size = range.end - range.start;
308                    match (offset, size, bitmap_type) {
309                        (1, 1, BitmapType::Roaring) => {
310                            output.push_back(Bitmap::from_lsb0_bytes(&[0b10101010], *bitmap_type))
311                        }
312                        (2, 1, BitmapType::Roaring) => {
313                            output.push_back(Bitmap::from_lsb0_bytes(&[0b11011011], *bitmap_type))
314                        }
315                        _ => unreachable!(),
316                    }
317                }
318
319                Ok(output)
320            });
321
322        let output = applier
323            .apply(SearchContext::default(), &mut mock_reader, None)
324            .await
325            .unwrap();
326        assert_eq!(
327            output.matched_segment_ids,
328            Bitmap::from_lsb0_bytes(&[0b10001010], BitmapType::Roaring)
329        );
330    }
331
332    #[tokio::test]
333    async fn test_index_applier_without_predicates() {
334        let applier = PredicatesIndexApplier {
335            fst_appliers: vec![],
336        };
337
338        let mut mock_reader: MockInvertedIndexReader = MockInvertedIndexReader::new();
339        mock_reader
340            .expect_metadata()
341            .returning(|_| Ok(mock_metas([("tag-0", 0)])));
342
343        let output = applier
344            .apply(SearchContext::default(), &mut mock_reader, None)
345            .await
346            .unwrap();
347        assert_eq!(output.matched_segment_ids, Bitmap::full_bitvec(8)); // full range to scan
348    }
349
350    #[tokio::test]
351    async fn test_index_applier_with_empty_index() {
352        let mut mock_reader = MockInvertedIndexReader::new();
353        mock_reader.expect_metadata().returning(move |_| {
354            Ok(Arc::new(InvertedIndexMetas {
355                total_row_count: 0, // No rows
356                segment_row_count: 1,
357                ..Default::default()
358            }))
359        });
360
361        let mut mock_fst_applier = MockFstApplier::new();
362        mock_fst_applier.expect_apply().never();
363
364        let applier = PredicatesIndexApplier {
365            fst_appliers: vec![(s("tag-0"), Box::new(mock_fst_applier))],
366        };
367
368        let output = applier
369            .apply(SearchContext::default(), &mut mock_reader, None)
370            .await
371            .unwrap();
372        assert!(output.matched_segment_ids.is_empty());
373    }
374
375    #[tokio::test]
376    async fn test_index_applier_with_nonexistent_index() {
377        let mut mock_reader = MockInvertedIndexReader::new();
378        mock_reader
379            .expect_metadata()
380            .returning(|_| Ok(mock_metas(vec![])));
381
382        let mut mock_fst_applier = MockFstApplier::new();
383        mock_fst_applier.expect_apply().never();
384
385        let applier = PredicatesIndexApplier {
386            fst_appliers: vec![(s("tag-0"), Box::new(mock_fst_applier))],
387        };
388
389        let result = applier
390            .apply(
391                SearchContext {
392                    index_not_found_strategy: IndexNotFoundStrategy::ThrowError,
393                },
394                &mut mock_reader,
395                None,
396            )
397            .await;
398        assert!(matches!(result, Err(Error::IndexNotFound { .. })));
399
400        let output = applier
401            .apply(
402                SearchContext {
403                    index_not_found_strategy: IndexNotFoundStrategy::ReturnEmpty,
404                },
405                &mut mock_reader,
406                None,
407            )
408            .await
409            .unwrap();
410        assert!(output.matched_segment_ids.is_empty());
411
412        let output = applier
413            .apply(
414                SearchContext {
415                    index_not_found_strategy: IndexNotFoundStrategy::Ignore,
416                },
417                &mut mock_reader,
418                None,
419            )
420            .await
421            .unwrap();
422        assert_eq!(output.matched_segment_ids, Bitmap::full_bitvec(8));
423    }
424
425    #[test]
426    fn test_index_applier_memory_usage() {
427        let mut mock_fst_applier = MockFstApplier::new();
428        mock_fst_applier.expect_memory_usage().returning(|| 100);
429
430        let applier = PredicatesIndexApplier {
431            fst_appliers: vec![(s("tag-0"), Box::new(mock_fst_applier))],
432        };
433
434        assert_eq!(
435            applier.memory_usage(),
436            size_of::<(IndexName, Box<dyn FstApplier>)>() + 5 + 100
437        );
438    }
439}