1use 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
33const BLOOM_META_LEN_SIZE: u64 = 4;
35
36pub const DEFAULT_PREFETCH_SIZE: u64 = 8192; #[derive(Default, Clone)]
41pub struct BloomFilterReadMetrics {
42 pub total_bytes: u64,
44 pub total_ranges: usize,
46 pub fetch_elapsed: Duration,
48 pub cache_hit: usize,
50 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 *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 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
105pub fn bytes_to_u64_vec(bytes: &Bytes) -> Vec<u64> {
110 let aligned_length = bytes.len() - bytes.len().rem(std::mem::size_of::<u64>());
112 let byte_slice = &bytes[..aligned_length];
113
114 let u64_vec = if let Ok(u64_slice) = try_cast_slice::<u8, u64>(byte_slice) {
116 u64_slice.to_vec()
117 } else {
118 let u64_count = byte_slice.len() / std::mem::size_of::<u64>();
120 let mut u64_vec = Vec::<u64>::with_capacity(u64_count);
121
122 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 unsafe { u64_vec.set_len(u64_count) };
131 u64_vec
132 };
133
134 #[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#[async_trait]
148pub trait BloomFilterReader: Sync {
149 async fn range_read(
151 &self,
152 offset: u64,
153 size: u32,
154 metrics: Option<&mut BloomFilterReadMetrics>,
155 ) -> Result<Bytes>;
156
157 async fn read_vec(
159 &self,
160 ranges: &[Range<u64>],
161 metrics: Option<&mut BloomFilterReadMetrics>,
162 ) -> Result<Vec<Bytes>>;
163
164 async fn metadata(
166 &self,
167 metrics: Option<&mut BloomFilterReadMetrics>,
168 ) -> Result<Arc<BloomFilterMeta>>;
169
170 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 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
209pub struct BloomFilterReaderImpl<R: RangeReader> {
211 reader: R,
213}
214
215impl<R: RangeReader> BloomFilterReaderImpl<R> {
216 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
280struct 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 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 m.total_ranges += 2;
335 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 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}