Skip to main content

index/bloom_filter/
reader.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::ops::{Range, Rem};
16use std::sync::Arc;
17use std::time::{Duration, Instant};
18
19use async_trait::async_trait;
20use bytemuck::try_cast_slice;
21use bytes::Bytes;
22use common_base::range_read::RangeReader;
23use fastbloom::BloomFilter;
24use greptime_proto::v1::index::{BloomFilterLoc, BloomFilterMeta};
25use prost::Message;
26use snafu::{ResultExt, ensure};
27
28use crate::bloom_filter::error::{
29    DecodeProtoSnafu, FileSizeTooSmallSnafu, IoSnafu, Result, UnexpectedMetaSizeSnafu,
30};
31use crate::bloom_filter::{PrehashedBloomFilter, PrehashedBuildHasher, SEED};
32
33/// Minimum size of the bloom filter, which is the size of the length of the bloom filter.
34const BLOOM_META_LEN_SIZE: u64 = 4;
35
36/// Default prefetch size of bloom filter meta.
37pub const DEFAULT_PREFETCH_SIZE: u64 = 8192; // 8KiB
38
39/// Metrics for bloom filter read operations.
40#[derive(Default, Clone)]
41pub struct BloomFilterReadMetrics {
42    /// Total byte size to read.
43    pub total_bytes: u64,
44    /// Total number of ranges to read.
45    pub total_ranges: usize,
46    /// Elapsed time to fetch data.
47    pub fetch_elapsed: Duration,
48    /// Number of cache hits.
49    pub cache_hit: usize,
50    /// Number of cache misses.
51    pub cache_miss: usize,
52}
53
54impl std::fmt::Debug for BloomFilterReadMetrics {
55    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
56        let Self {
57            total_bytes,
58            total_ranges,
59            fetch_elapsed,
60            cache_hit,
61            cache_miss,
62        } = self;
63
64        // If both total_bytes and cache_hit are 0, we didn't read anything.
65        if *total_bytes == 0 && *cache_hit == 0 {
66            return write!(f, "{{}}");
67        }
68        write!(f, "{{")?;
69
70        if *total_bytes > 0 {
71            write!(f, "\"total_bytes\":{}", total_bytes)?;
72        }
73        if *cache_hit > 0 {
74            if *total_bytes > 0 {
75                write!(f, ", ")?;
76            }
77            write!(f, "\"cache_hit\":{}", cache_hit)?;
78        }
79
80        if *total_ranges > 0 {
81            write!(f, ", \"total_ranges\":{}", total_ranges)?;
82        }
83        if !fetch_elapsed.is_zero() {
84            write!(f, ", \"fetch_elapsed\":\"{:?}\"", fetch_elapsed)?;
85        }
86        if *cache_miss > 0 {
87            write!(f, ", \"cache_miss\":{}", cache_miss)?;
88        }
89
90        write!(f, "}}")
91    }
92}
93
94impl BloomFilterReadMetrics {
95    /// Merges another metrics into this one.
96    pub fn merge_from(&mut self, other: &Self) {
97        self.total_bytes += other.total_bytes;
98        self.total_ranges += other.total_ranges;
99        self.fetch_elapsed += other.fetch_elapsed;
100        self.cache_hit += other.cache_hit;
101        self.cache_miss += other.cache_miss;
102    }
103}
104
105/// Safely converts bytes to Vec<u64> using bytemuck for optimal performance.
106/// Faster than chunking and converting each piece individually.
107///
108/// The input bytes are a sequence of little-endian u64s.
109pub fn bytes_to_u64_vec(bytes: &Bytes) -> Vec<u64> {
110    // drop tailing things, this keeps the same behavior with `chunks_exact`.
111    let aligned_length = bytes.len() - bytes.len().rem(std::mem::size_of::<u64>());
112    let byte_slice = &bytes[..aligned_length];
113
114    // Try fast path first: direct cast if aligned
115    let u64_vec = if let Ok(u64_slice) = try_cast_slice::<u8, u64>(byte_slice) {
116        u64_slice.to_vec()
117    } else {
118        // Slow path: create aligned Vec<u64> and copy data
119        let u64_count = byte_slice.len() / std::mem::size_of::<u64>();
120        let mut u64_vec = Vec::<u64>::with_capacity(u64_count);
121
122        // SAFETY: We're creating a properly sized slice from uninitialized but allocated memory
123        // to copy bytes into. The slice has exactly the right size for the byte data.
124        let dest_slice = unsafe {
125            std::slice::from_raw_parts_mut(u64_vec.as_mut_ptr() as *mut u8, byte_slice.len())
126        };
127        dest_slice.copy_from_slice(byte_slice);
128
129        // SAFETY: We've just initialized exactly u64_count elements worth of bytes
130        unsafe { u64_vec.set_len(u64_count) };
131        u64_vec
132    };
133
134    // Convert from platform endianness to little endian if needed
135    // Just in case.
136    #[cfg(target_endian = "little")]
137    {
138        u64_vec
139    }
140    #[cfg(target_endian = "big")]
141    {
142        u64_vec.into_iter().map(|x| x.swap_bytes()).collect()
143    }
144}
145
146/// `BloomFilterReader` reads the bloom filter from the file.
147#[async_trait]
148pub trait BloomFilterReader: Sync {
149    /// Reads range of bytes from the file.
150    async fn range_read(
151        &self,
152        offset: u64,
153        size: u32,
154        metrics: Option<&mut BloomFilterReadMetrics>,
155    ) -> Result<Bytes>;
156
157    /// Reads bunch of ranges from the file.
158    async fn read_vec(
159        &self,
160        ranges: &[Range<u64>],
161        metrics: Option<&mut BloomFilterReadMetrics>,
162    ) -> Result<Vec<Bytes>>;
163
164    /// Reads the meta information of the bloom filter.
165    async fn metadata(
166        &self,
167        metrics: Option<&mut BloomFilterReadMetrics>,
168    ) -> Result<Arc<BloomFilterMeta>>;
169
170    /// Reads a bloom filter with the given location.
171    async fn bloom_filter(
172        &self,
173        loc: &BloomFilterLoc,
174        metrics: Option<&mut BloomFilterReadMetrics>,
175    ) -> Result<BloomFilter> {
176        let bytes = self.range_read(loc.offset, loc.size as _, metrics).await?;
177        let vec = bytes_to_u64_vec(&bytes);
178        let bm = BloomFilter::from_vec(vec)
179            .seed(&SEED)
180            .expected_items(loc.element_count as _);
181        Ok(bm)
182    }
183
184    /// Reads multiple bloom filters; probe them with [`crate::bloom_filter::element_hash`].
185    async fn bloom_filter_vec(
186        &self,
187        locs: &[BloomFilterLoc],
188        metrics: Option<&mut BloomFilterReadMetrics>,
189    ) -> Result<Vec<PrehashedBloomFilter>> {
190        let ranges = locs
191            .iter()
192            .map(|l| l.offset..l.offset + l.size)
193            .collect::<Vec<_>>();
194        let bss = self.read_vec(&ranges, metrics).await?;
195
196        let mut result = Vec::with_capacity(bss.len());
197        for (bs, loc) in bss.into_iter().zip(locs.iter()) {
198            let vec = bytes_to_u64_vec(&bs);
199            let bm = BloomFilter::from_vec(vec)
200                .hasher(PrehashedBuildHasher::default())
201                .expected_items(loc.element_count as _);
202            result.push(bm);
203        }
204
205        Ok(result)
206    }
207}
208
209/// `BloomFilterReaderImpl` reads the bloom filter from the file.
210pub struct BloomFilterReaderImpl<R: RangeReader> {
211    /// The underlying reader.
212    reader: R,
213}
214
215impl<R: RangeReader> BloomFilterReaderImpl<R> {
216    /// Creates a new `BloomFilterReaderImpl` with the given reader.
217    pub fn new(reader: R) -> Self {
218        Self { reader }
219    }
220}
221
222#[async_trait]
223impl<R: RangeReader> BloomFilterReader for BloomFilterReaderImpl<R> {
224    async fn range_read(
225        &self,
226        offset: u64,
227        size: u32,
228        metrics: Option<&mut BloomFilterReadMetrics>,
229    ) -> Result<Bytes> {
230        let start = metrics.as_ref().map(|_| Instant::now());
231        let result = self
232            .reader
233            .read(offset..offset + size as u64)
234            .await
235            .context(IoSnafu)?;
236
237        if let Some(m) = metrics {
238            m.total_ranges += 1;
239            m.total_bytes += size as u64;
240            if let Some(start) = start {
241                m.fetch_elapsed += start.elapsed();
242            }
243        }
244
245        Ok(result)
246    }
247
248    async fn read_vec(
249        &self,
250        ranges: &[Range<u64>],
251        metrics: Option<&mut BloomFilterReadMetrics>,
252    ) -> Result<Vec<Bytes>> {
253        let start = metrics.as_ref().map(|_| Instant::now());
254        let result = self.reader.read_vec(ranges).await.context(IoSnafu)?;
255
256        if let Some(m) = metrics {
257            m.total_ranges += ranges.len();
258            m.total_bytes += ranges.iter().map(|r| r.end - r.start).sum::<u64>();
259            if let Some(start) = start {
260                m.fetch_elapsed += start.elapsed();
261            }
262        }
263
264        Ok(result)
265    }
266
267    async fn metadata(
268        &self,
269        metrics: Option<&mut BloomFilterReadMetrics>,
270    ) -> Result<Arc<BloomFilterMeta>> {
271        let metadata = self.reader.metadata().await.context(IoSnafu)?;
272        let file_size = metadata.content_length;
273
274        let mut meta_reader =
275            BloomFilterMetaReader::new(&self.reader, file_size, Some(DEFAULT_PREFETCH_SIZE));
276        meta_reader.metadata(metrics).await.map(Arc::new)
277    }
278}
279
280/// `BloomFilterMetaReader` reads the metadata of the bloom filter.
281struct BloomFilterMetaReader<R: RangeReader> {
282    reader: R,
283    file_size: u64,
284    prefetch_size: u64,
285}
286
287impl<R: RangeReader> BloomFilterMetaReader<R> {
288    pub fn new(reader: R, file_size: u64, prefetch_size: Option<u64>) -> Self {
289        Self {
290            reader,
291            file_size,
292            prefetch_size: prefetch_size
293                .unwrap_or(BLOOM_META_LEN_SIZE)
294                .max(BLOOM_META_LEN_SIZE),
295        }
296    }
297
298    /// Reads the metadata of the bloom filter.
299    ///
300    /// It will first prefetch some bytes from the end of the file,
301    /// then parse the metadata from the prefetch bytes.
302    pub async fn metadata(
303        &mut self,
304        metrics: Option<&mut BloomFilterReadMetrics>,
305    ) -> Result<BloomFilterMeta> {
306        ensure!(
307            self.file_size >= BLOOM_META_LEN_SIZE,
308            FileSizeTooSmallSnafu {
309                size: self.file_size,
310            }
311        );
312
313        let start = metrics.as_ref().map(|_| Instant::now());
314        let meta_start = self.file_size.saturating_sub(self.prefetch_size);
315        let suffix = self
316            .reader
317            .read(meta_start..self.file_size)
318            .await
319            .context(IoSnafu)?;
320        let suffix_len = suffix.len();
321        let length = u32::from_le_bytes(Self::read_tailing_four_bytes(&suffix)?) as u64;
322        self.validate_meta_size(length)?;
323
324        if length > suffix_len as u64 - BLOOM_META_LEN_SIZE {
325            let metadata_start = self.file_size - length - BLOOM_META_LEN_SIZE;
326            let meta = self
327                .reader
328                .read(metadata_start..self.file_size - BLOOM_META_LEN_SIZE)
329                .await
330                .context(IoSnafu)?;
331
332            if let Some(m) = metrics {
333                // suffix read + meta read
334                m.total_ranges += 2;
335                // Ignores the meta length size to simplify the calculation.
336                m.total_bytes += self.file_size.min(self.prefetch_size) + length;
337                if let Some(start) = start {
338                    m.fetch_elapsed += start.elapsed();
339                }
340            }
341
342            BloomFilterMeta::decode(meta).context(DecodeProtoSnafu)
343        } else {
344            if let Some(m) = metrics {
345                // suffix read only
346                m.total_ranges += 1;
347                m.total_bytes += self.file_size.min(self.prefetch_size);
348                if let Some(start) = start {
349                    m.fetch_elapsed += start.elapsed();
350                }
351            }
352
353            let metadata_start = self.file_size - length - BLOOM_META_LEN_SIZE - meta_start;
354            let meta = &suffix[metadata_start as usize..suffix_len - BLOOM_META_LEN_SIZE as usize];
355            BloomFilterMeta::decode(meta).context(DecodeProtoSnafu)
356        }
357    }
358
359    fn read_tailing_four_bytes(suffix: &[u8]) -> Result<[u8; 4]> {
360        let suffix_len = suffix.len();
361        ensure!(
362            suffix_len >= 4,
363            FileSizeTooSmallSnafu {
364                size: suffix_len as u64
365            }
366        );
367        let mut bytes = [0; 4];
368        bytes.copy_from_slice(&suffix[suffix_len - 4..suffix_len]);
369
370        Ok(bytes)
371    }
372
373    fn validate_meta_size(&self, length: u64) -> Result<()> {
374        let max_meta_size = self.file_size - BLOOM_META_LEN_SIZE;
375        ensure!(
376            length <= max_meta_size,
377            UnexpectedMetaSizeSnafu {
378                max_meta_size,
379                actual_meta_size: length,
380            }
381        );
382        Ok(())
383    }
384}
385
386#[cfg(test)]
387mod tests {
388    use std::sync::Arc;
389    use std::sync::atomic::AtomicUsize;
390
391    use futures::io::Cursor;
392
393    use super::*;
394    use crate::bloom_filter::creator::BloomFilterCreator;
395    use crate::external_provider::MockExternalTempFileProvider;
396
397    async fn mock_bloom_filter_bytes() -> Vec<u8> {
398        let mut writer = Cursor::new(vec![]);
399        let mut creator = BloomFilterCreator::new(
400            2,
401            0.01,
402            Arc::new(MockExternalTempFileProvider::new()),
403            Arc::new(AtomicUsize::new(0)),
404            None,
405        );
406
407        creator
408            .push_row_elems(vec![b"a".to_vec(), b"b".to_vec()])
409            .await
410            .unwrap();
411        creator
412            .push_row_elems(vec![b"c".to_vec(), b"d".to_vec()])
413            .await
414            .unwrap();
415        creator
416            .push_row_elems(vec![b"e".to_vec(), b"f".to_vec()])
417            .await
418            .unwrap();
419
420        creator.finish(&mut writer).await.unwrap();
421
422        writer.into_inner()
423    }
424
425    #[tokio::test]
426    async fn test_bloom_filter_meta_reader() {
427        let bytes = mock_bloom_filter_bytes().await;
428        let file_size = bytes.len() as u64;
429
430        for prefetch in [0u64, file_size / 2, file_size, file_size + 10] {
431            let mut reader =
432                BloomFilterMetaReader::new(bytes.clone(), file_size as _, Some(prefetch));
433            let meta = reader.metadata(None).await.unwrap();
434
435            assert_eq!(meta.rows_per_segment, 2);
436            assert_eq!(meta.segment_count, 2);
437            assert_eq!(meta.row_count, 3);
438            assert_eq!(meta.bloom_filter_locs.len(), 2);
439
440            assert_eq!(meta.bloom_filter_locs[0].offset, 0);
441            assert_eq!(meta.bloom_filter_locs[0].element_count, 4);
442            assert_eq!(
443                meta.bloom_filter_locs[1].offset,
444                meta.bloom_filter_locs[0].size
445            );
446            assert_eq!(meta.bloom_filter_locs[1].element_count, 2);
447        }
448    }
449
450    #[tokio::test]
451    async fn test_bloom_filter_reader() {
452        let bytes = mock_bloom_filter_bytes().await;
453
454        let reader = BloomFilterReaderImpl::new(bytes);
455        let meta = reader.metadata(None).await.unwrap();
456
457        assert_eq!(meta.bloom_filter_locs.len(), 2);
458        let bf = reader
459            .bloom_filter(&meta.bloom_filter_locs[0], None)
460            .await
461            .unwrap();
462        assert!(bf.contains(&b"a"));
463        assert!(bf.contains(&b"b"));
464        assert!(bf.contains(&b"c"));
465        assert!(bf.contains(&b"d"));
466
467        let bf = reader
468            .bloom_filter(&meta.bloom_filter_locs[1], None)
469            .await
470            .unwrap();
471        assert!(bf.contains(&b"e"));
472        assert!(bf.contains(&b"f"));
473    }
474}