1pub(crate) mod cache_size;
18
19pub(crate) mod file_cache;
20pub(crate) mod index;
21pub(crate) mod manifest_cache;
22#[cfg(test)]
23pub(crate) mod test_util;
24pub(crate) mod write_cache;
25
26use std::collections::{BTreeMap, HashMap};
27use std::mem;
28use std::ops::Range;
29use std::sync::{Arc, RwLock, Weak};
30
31use bytes::Bytes;
32use common_base::readable_size::ReadableSize;
33use common_datasource::compression::CompressionType;
34use common_telemetry::warn;
35use datatypes::arrow::buffer::BooleanBuffer;
36use datatypes::arrow::record_batch::RecordBatch;
37use datatypes::value::Value;
38use datatypes::vectors::VectorRef;
39use index::bloom_filter_index::{BloomFilterIndexCache, BloomFilterIndexCacheRef};
40use index::result_cache::IndexResultCache;
41use moka::notification::RemovalCause;
42use moka::sync::Cache;
43use object_store::ObjectStore;
44use parquet::arrow::arrow_reader::{RowSelection, RowSelector};
45use parquet::file::metadata::{
46 FileMetaData, PageIndexPolicy, ParquetMetaData, ParquetMetaDataReader, ParquetMetaDataWriter,
47};
48use puffin::puffin_manager::cache::{PuffinMetadataCache, PuffinMetadataCacheRef};
49use smallvec::SmallVec;
50use snafu::{OptionExt, ResultExt};
51use store_api::metadata::{RegionMetadata, RegionMetadataRef};
52use store_api::storage::{ConcreteDataType, FileId, RegionId, TimeSeriesRowSelector};
53
54use crate::cache::cache_size::parquet_meta_size;
55use crate::cache::file_cache::{FileType, IndexKey};
56use crate::cache::index::inverted_index::{InvertedIndexCache, InvertedIndexCacheRef};
57#[cfg(feature = "vector_index")]
58use crate::cache::index::vector_index::{VectorIndexCache, VectorIndexCacheRef};
59use crate::cache::write_cache::WriteCacheRef;
60use crate::error::{
61 CompressObjectSnafu, DecompressObjectSnafu, InvalidMetadataSnafu, InvalidParquetSnafu,
62 JoinSnafu, ReadParquetSnafu, Result, UnexpectedSnafu, WriteParquetSnafu,
63};
64use crate::memtable::record_batch_estimated_size;
65use crate::metrics::{CACHE_BYTES, CACHE_EVICTION, CACHE_HIT, CACHE_MISS};
66use crate::read::Batch;
67use crate::read::range_cache::{RangeScanCacheKey, RangeScanCacheValue};
68use crate::sst::file::{RegionFileId, RegionIndexId};
69use crate::sst::parquet::PARQUET_METADATA_KEY;
70use crate::sst::parquet::read_columns::ParquetReadColumns;
71use crate::sst::parquet::reader::MetadataCacheMetrics;
72
73const SST_META_TYPE: &str = "sst_meta";
75const SST_META_DECODED_TYPE: &str = "sst_meta_decoded";
77const VECTOR_TYPE: &str = "vector";
79const PAGE_TYPE: &str = "page";
81const FILE_TYPE: &str = "file";
83const INDEX_TYPE: &str = "index";
85const SELECTOR_RESULT_TYPE: &str = "selector_result";
87const RANGE_RESULT_TYPE: &str = "range_result";
89const PREFILTER_RESULT_TYPE: &str = "prefilter_result";
91const RANGE_RESULT_CONCAT_MEMORY_LIMIT: ReadableSize = ReadableSize::mb(512);
92const RANGE_RESULT_CONCAT_MEMORY_PERMIT: ReadableSize = ReadableSize::kb(1);
93
94#[derive(Debug)]
95pub(crate) struct RangeResultMemoryLimiter {
96 semaphore: Arc<tokio::sync::Semaphore>,
97 permit_bytes: usize,
98 total_permits: usize,
99}
100
101impl Default for RangeResultMemoryLimiter {
102 fn default() -> Self {
103 Self::new(
104 RANGE_RESULT_CONCAT_MEMORY_LIMIT.as_bytes() as usize,
105 RANGE_RESULT_CONCAT_MEMORY_PERMIT.as_bytes() as usize,
106 )
107 }
108}
109
110impl RangeResultMemoryLimiter {
111 pub(crate) fn new(limit_bytes: usize, permit_bytes: usize) -> Self {
112 let permit_bytes = permit_bytes.max(1);
113 let total_permits = limit_bytes
114 .div_ceil(permit_bytes)
115 .clamp(1, tokio::sync::Semaphore::MAX_PERMITS);
116 Self {
117 semaphore: Arc::new(tokio::sync::Semaphore::new(total_permits)),
118 permit_bytes,
119 total_permits,
120 }
121 }
122
123 #[cfg(test)]
124 pub(crate) fn permit_bytes(&self) -> usize {
125 self.permit_bytes
126 }
127
128 #[cfg(test)]
129 pub(crate) fn available_permits(&self) -> usize {
130 self.semaphore.available_permits()
131 }
132
133 pub(crate) async fn acquire(&self, bytes: usize) -> Result<tokio::sync::SemaphorePermit<'_>> {
134 let permits = bytes.div_ceil(self.permit_bytes).max(1);
135 if permits > self.total_permits {
136 return UnexpectedSnafu {
137 reason: format!(
138 "range result memory request of {bytes} bytes exceeds limiter capacity of {} bytes",
139 self.total_permits.saturating_mul(self.permit_bytes)
140 ),
141 }
142 .fail();
143 }
144 self.semaphore
145 .acquire_many(permits as u32)
146 .await
147 .map_err(|_| {
148 UnexpectedSnafu {
149 reason: "range result memory limiter is unexpectedly closed",
150 }
151 .build()
152 })
153 }
154}
155
156#[derive(Debug)]
161pub(crate) struct CachedSstMeta {
162 parquet_metadata: Arc<ParquetMetaData>,
163 parquet_metadata_size: usize,
164 region_metadata: RegionMetadataRef,
165 region_metadata_weight: usize,
166 page_index_policy: PageIndexPolicy,
167}
168
169#[derive(Debug)]
175pub(crate) struct CompactSstMeta {
176 encoded_metadata: Bytes,
177 decoded_size: usize,
178 region_metadata: Weak<RegionMetadata>,
179 page_index_policy: PageIndexPolicy,
180 decoded: tokio::sync::Mutex<Weak<CachedSstMeta>>,
181}
182
183#[derive(Debug)]
185pub(crate) struct PreparedSstMeta {
186 compact: Arc<CompactSstMeta>,
187 decoded: Arc<CachedSstMeta>,
188}
189
190impl PreparedSstMeta {
191 pub(crate) fn decoded(&self) -> Arc<CachedSstMeta> {
192 self.decoded.clone()
193 }
194}
195
196#[derive(Debug)]
198pub(crate) enum SstMetaPreparation {
199 Prepared(PreparedSstMeta),
201 DecodedOnly {
203 decoded: Arc<CachedSstMeta>,
204 encoding_error: crate::error::Error,
205 },
206}
207
208impl SstMetaPreparation {
209 pub(crate) fn decoded(&self) -> Arc<CachedSstMeta> {
210 match self {
211 Self::Prepared(metadata) => metadata.decoded(),
212 Self::DecodedOnly { decoded, .. } => decoded.clone(),
213 }
214 }
215}
216
217impl CompactSstMeta {
218 async fn decode(&self) -> Result<Arc<CachedSstMeta>> {
219 let mut decoded_guard = self.decoded.lock().await;
220 if let Some(decoded) = decoded_guard.upgrade() {
221 return Ok(decoded);
222 }
223
224 let encoded_metadata = self.encoded_metadata.clone();
225 let decoded_size = self.decoded_size;
226 let region_metadata = self.region_metadata.upgrade();
227 let page_index_policy = self.page_index_policy;
228 let decoded = common_runtime::spawn_blocking_global(move || {
229 let bytes = zstd::bulk::decompress(&encoded_metadata, decoded_size).context(
230 DecompressObjectSnafu {
231 compress_type: CompressionType::Zstd,
232 path: "cached SST metadata",
233 },
234 )?;
235 let bytes = Bytes::from(bytes);
236 let mut reader = ParquetMetaDataReader::new()
237 .with_column_index_policy(PageIndexPolicy::Skip)
238 .with_offset_index_policy(page_index_policy);
239 reader.try_parse(&bytes).context(ReadParquetSnafu {
240 path: "cached SST metadata",
241 })?;
242 let metadata = reader.finish().context(ReadParquetSnafu {
243 path: "cached SST metadata",
244 })?;
245 CachedSstMeta::try_new_with_page_index_policy(
246 "cached SST metadata",
247 metadata,
248 region_metadata,
249 page_index_policy,
250 )
251 .map(Arc::new)
252 })
253 .await
254 .context(JoinSnafu)??;
255
256 *decoded_guard = Arc::downgrade(&decoded);
257 Ok(decoded)
258 }
259
260 fn satisfies_page_index_policy(&self, requested: PageIndexPolicy) -> bool {
261 satisfies_page_index_policy(self.page_index_policy, requested)
262 }
263}
264
265pub(crate) async fn decode_sst_meta(
267 file_path: &str,
268 parquet_metadata: ParquetMetaData,
269 region_metadata: Option<RegionMetadataRef>,
270 page_index_policy: PageIndexPolicy,
271) -> Result<Arc<CachedSstMeta>> {
272 let file_path = file_path.to_string();
273 common_runtime::spawn_blocking_global(move || {
274 let parquet_metadata = strip_column_indexes(parquet_metadata);
275 CachedSstMeta::try_new_with_page_index_policy(
276 &file_path,
277 parquet_metadata,
278 region_metadata,
279 page_index_policy,
280 )
281 .map(Arc::new)
282 })
283 .await
284 .context(JoinSnafu)?
285}
286
287pub(crate) async fn prepare_sst_meta(
289 file_path: &str,
290 parquet_metadata: ParquetMetaData,
291 region_metadata: Option<RegionMetadataRef>,
292 page_index_policy: PageIndexPolicy,
293) -> Result<SstMetaPreparation> {
294 let file_path = file_path.to_string();
295 common_runtime::spawn_blocking_global(move || {
296 prepare_sst_meta_sync(
297 &file_path,
298 parquet_metadata,
299 region_metadata,
300 page_index_policy,
301 )
302 })
303 .await
304 .context(JoinSnafu)?
305}
306
307pub(crate) fn prepare_sst_meta_sync(
309 file_path: &str,
310 parquet_metadata: ParquetMetaData,
311 region_metadata: Option<RegionMetadataRef>,
312 page_index_policy: PageIndexPolicy,
313) -> Result<SstMetaPreparation> {
314 let parquet_metadata = strip_column_indexes(parquet_metadata);
315 let cache_encoding = encode_compact_sst_meta(file_path, &parquet_metadata);
316 finish_sst_meta_preparation(
317 file_path,
318 parquet_metadata,
319 region_metadata,
320 page_index_policy,
321 cache_encoding,
322 )
323}
324
325fn strip_column_indexes(parquet_metadata: ParquetMetaData) -> ParquetMetaData {
326 let mut builder = parquet_metadata.into_builder();
329 builder.take_column_index();
330 builder.build()
331}
332
333fn encode_compact_sst_meta(
334 file_path: &str,
335 parquet_metadata: &ParquetMetaData,
336) -> Result<(Bytes, usize)> {
337 let mut encoded = Vec::new();
338 ParquetMetaDataWriter::new(&mut encoded, parquet_metadata)
339 .finish()
340 .context(WriteParquetSnafu)?;
341 let decoded_size = encoded.len();
342 let encoded_metadata = zstd::bulk::compress(&encoded, 3).context(CompressObjectSnafu {
343 compress_type: CompressionType::Zstd,
344 path: file_path,
345 })?;
346
347 Ok((
351 Bytes::from(encoded_metadata.into_boxed_slice()),
352 decoded_size,
353 ))
354}
355
356fn finish_sst_meta_preparation(
357 file_path: &str,
358 parquet_metadata: ParquetMetaData,
359 region_metadata: Option<RegionMetadataRef>,
360 page_index_policy: PageIndexPolicy,
361 cache_encoding: Result<(Bytes, usize)>,
362) -> Result<SstMetaPreparation> {
363 let decoded = Arc::new(CachedSstMeta::try_new_with_page_index_policy(
364 file_path,
365 parquet_metadata,
366 region_metadata,
367 page_index_policy,
368 )?);
369 let (encoded_metadata, decoded_size) = match cache_encoding {
370 Ok(encoded) => encoded,
371 Err(encoding_error) => {
372 return Ok(SstMetaPreparation::DecodedOnly {
373 decoded,
374 encoding_error,
375 });
376 }
377 };
378 let compact = Arc::new(CompactSstMeta {
379 encoded_metadata,
380 decoded_size,
381 region_metadata: Arc::downgrade(&decoded.region_metadata),
382 page_index_policy,
383 decoded: tokio::sync::Mutex::new(Arc::downgrade(&decoded)),
384 });
385
386 Ok(SstMetaPreparation::Prepared(PreparedSstMeta {
387 compact,
388 decoded,
389 }))
390}
391
392impl CachedSstMeta {
393 #[cfg(test)]
394 pub(crate) fn try_new(file_path: &str, parquet_metadata: ParquetMetaData) -> Result<Self> {
395 let page_index_policy = infer_loaded_page_index_policy(&parquet_metadata);
396 Self::try_new_with_page_index_policy(file_path, parquet_metadata, None, page_index_policy)
397 }
398
399 pub(crate) fn try_new_with_region_metadata(
400 file_path: &str,
401 parquet_metadata: ParquetMetaData,
402 region_metadata: Option<RegionMetadataRef>,
403 ) -> Result<Self> {
404 let page_index_policy = infer_loaded_page_index_policy(&parquet_metadata);
405 Self::try_new_with_page_index_policy(
406 file_path,
407 parquet_metadata,
408 region_metadata,
409 page_index_policy,
410 )
411 }
412
413 pub(crate) fn try_new_with_page_index_policy(
414 file_path: &str,
415 parquet_metadata: ParquetMetaData,
416 region_metadata: Option<RegionMetadataRef>,
417 page_index_policy: PageIndexPolicy,
418 ) -> Result<Self> {
419 let file_metadata = parquet_metadata.file_metadata();
420 let key_values = file_metadata
421 .key_value_metadata()
422 .context(InvalidParquetSnafu {
423 file: file_path,
424 reason: "missing key value meta",
425 })?;
426 let meta_value = key_values
427 .iter()
428 .find(|kv| kv.key == PARQUET_METADATA_KEY)
429 .with_context(|| InvalidParquetSnafu {
430 file: file_path,
431 reason: format!("key {} not found", PARQUET_METADATA_KEY),
432 })?;
433 let json = meta_value
434 .value
435 .as_ref()
436 .with_context(|| InvalidParquetSnafu {
437 file: file_path,
438 reason: format!("No value for key {}", PARQUET_METADATA_KEY),
439 })?;
440 let region_metadata = match region_metadata {
441 Some(region_metadata) => region_metadata,
442 None => Arc::new(
443 store_api::metadata::RegionMetadata::from_json(json)
444 .context(InvalidMetadataSnafu)?,
445 ),
446 };
447 let region_metadata_weight = region_metadata.estimated_size().max(json.len());
449 let parquet_metadata = Arc::new(strip_region_metadata_from_parquet(parquet_metadata));
450 let parquet_metadata_size = parquet_meta_size(&parquet_metadata);
451
452 Ok(Self {
453 parquet_metadata,
454 parquet_metadata_size,
455 region_metadata,
456 region_metadata_weight,
457 page_index_policy,
458 })
459 }
460
461 pub(crate) fn parquet_metadata(&self) -> Arc<ParquetMetaData> {
462 self.parquet_metadata.clone()
463 }
464
465 pub(crate) fn parquet_metadata_size(&self) -> usize {
467 self.parquet_metadata_size
468 }
469
470 pub(crate) fn region_metadata(&self) -> RegionMetadataRef {
471 self.region_metadata.clone()
472 }
473
474 fn satisfies_page_index_policy(&self, requested: PageIndexPolicy) -> bool {
475 satisfies_page_index_policy(self.page_index_policy, requested)
476 }
477}
478
479fn satisfies_page_index_policy(cached: PageIndexPolicy, requested: PageIndexPolicy) -> bool {
480 match requested {
481 PageIndexPolicy::Skip => true,
482 PageIndexPolicy::Optional => cached != PageIndexPolicy::Skip,
483 PageIndexPolicy::Required => cached == PageIndexPolicy::Required,
484 }
485}
486
487fn infer_loaded_page_index_policy(parquet_metadata: &ParquetMetaData) -> PageIndexPolicy {
488 if parquet_metadata.offset_index().is_some() {
489 PageIndexPolicy::Optional
490 } else {
491 PageIndexPolicy::Skip
492 }
493}
494
495fn strip_region_metadata_from_parquet(parquet_metadata: ParquetMetaData) -> ParquetMetaData {
496 let file_metadata = parquet_metadata.file_metadata();
497 let filtered_key_values = file_metadata.key_value_metadata().and_then(|key_values| {
498 let filtered = key_values
499 .iter()
500 .filter(|kv| kv.key != PARQUET_METADATA_KEY)
501 .cloned()
502 .collect::<Vec<_>>();
503 (!filtered.is_empty()).then_some(filtered)
504 });
505 let stripped_file_metadata = FileMetaData::new(
506 file_metadata.version(),
507 file_metadata.num_rows(),
508 file_metadata.created_by().map(ToString::to_string),
509 filtered_key_values,
510 file_metadata.schema_descr_ptr(),
511 file_metadata.column_orders().cloned(),
512 );
513
514 let mut builder = parquet_metadata.into_builder();
515 let row_groups = builder.take_row_groups();
516 let offset_index = builder.take_offset_index();
517
518 parquet::file::metadata::ParquetMetaDataBuilder::new(stripped_file_metadata)
519 .set_row_groups(row_groups)
520 .set_offset_index(offset_index)
521 .build()
522}
523
524fn removal_cause_str(cause: RemovalCause) -> &'static str {
525 match cause {
526 RemovalCause::Expired => "expired",
527 RemovalCause::Explicit => "explicit",
528 RemovalCause::Replaced => "replaced",
529 RemovalCause::Size => "size",
530 }
531}
532
533#[derive(Debug, Clone, PartialEq, Eq, Hash)]
534pub(crate) struct PrefilterRowSelector {
535 row_count: usize,
536 skip: bool,
537}
538
539impl From<&RowSelector> for PrefilterRowSelector {
544 fn from(selector: &RowSelector) -> Self {
545 Self {
546 row_count: selector.row_count,
547 skip: selector.skip,
548 }
549 }
550}
551
552#[derive(Debug, Clone, PartialEq, Eq, Hash)]
554pub(crate) struct PrefilterKey {
555 file_id: FileId,
556 row_group_idx: u32,
557 row_selection: Option<Arc<Vec<PrefilterRowSelector>>>,
558 schema_version: u64,
559 filter_exprs: SmallVec<[String; 1]>,
560 mem_usage: usize,
561}
562
563impl PrefilterKey {
564 pub(crate) fn row_selection_snapshot(
565 row_selection: Option<&RowSelection>,
566 ) -> Option<Arc<Vec<PrefilterRowSelector>>> {
567 row_selection.map(|selection| {
568 Arc::new(
569 selection
570 .iter()
571 .map(PrefilterRowSelector::from)
572 .collect::<Vec<_>>(),
573 )
574 })
575 }
576
577 pub(crate) fn new(
578 file_id: FileId,
579 row_group_idx: u32,
580 row_selection: Option<Arc<Vec<PrefilterRowSelector>>>,
581 schema_version: u64,
582 filter_exprs: SmallVec<[String; 1]>,
583 ) -> Self {
584 let row_selection_bytes = row_selection
585 .as_ref()
586 .map(|selection| selection.len() * mem::size_of::<PrefilterRowSelector>())
587 .unwrap_or(0);
588 let spilled_expr_bytes = if filter_exprs.spilled() {
589 filter_exprs.capacity() * mem::size_of::<String>()
590 } else {
591 0
592 };
593 let expr_bytes = filter_exprs.iter().map(|s| s.capacity()).sum::<usize>();
594
595 Self {
596 file_id,
597 row_group_idx,
598 row_selection,
599 schema_version,
600 filter_exprs,
601 mem_usage: mem::size_of::<Self>()
602 + row_selection_bytes
603 + spilled_expr_bytes
604 + expr_bytes,
605 }
606 }
607
608 fn mem_usage(&self) -> usize {
609 self.mem_usage
610 }
611}
612
613type PrefilterResultCache = Cache<PrefilterKey, Arc<BooleanBuffer>>;
614
615fn new_prefilter_result_cache(capacity: u64) -> PrefilterResultCache {
616 Cache::builder()
617 .max_capacity(capacity)
618 .weigher(prefilter_result_cache_weight)
619 .eviction_listener(|k, v, cause| {
620 let size = prefilter_result_cache_weight(&k, &v);
621 CACHE_BYTES
622 .with_label_values(&[PREFILTER_RESULT_TYPE])
623 .sub(size.into());
624 CACHE_EVICTION
625 .with_label_values(&[PREFILTER_RESULT_TYPE, removal_cause_str(cause)])
626 .inc();
627 })
628 .build()
629}
630
631fn prefilter_result_cache_weight(k: &PrefilterKey, v: &Arc<BooleanBuffer>) -> u32 {
632 (k.mem_usage() + mem::size_of::<BooleanBuffer>() + v.values().len()) as u32
633}
634
635#[derive(Clone)]
637pub enum CacheStrategy {
638 EnableAll(CacheManagerRef),
641 Compaction(CacheManagerRef),
646 Disabled,
648}
649
650impl CacheStrategy {
651 pub(crate) fn sst_meta_cache_enabled(&self) -> bool {
653 match self {
654 CacheStrategy::EnableAll(cache_manager) | CacheStrategy::Compaction(cache_manager) => {
655 cache_manager.sst_meta_cache_enabled()
656 }
657 CacheStrategy::Disabled => false,
658 }
659 }
660
661 pub(crate) async fn get_sst_meta_data(
663 &self,
664 file_id: RegionFileId,
665 metrics: &mut MetadataCacheMetrics,
666 page_index_policy: PageIndexPolicy,
667 ) -> Option<Arc<CachedSstMeta>> {
668 match self {
669 CacheStrategy::EnableAll(cache_manager) | CacheStrategy::Compaction(cache_manager) => {
670 cache_manager
671 .get_sst_meta_data(file_id, metrics, page_index_policy)
672 .await
673 }
674 CacheStrategy::Disabled => {
675 metrics.cache_miss += 1;
676 None
677 }
678 }
679 }
680
681 pub(crate) fn get_sst_meta_data_from_mem_cache(
683 &self,
684 file_id: RegionFileId,
685 page_index_policy: PageIndexPolicy,
686 ) -> Option<Arc<CachedSstMeta>> {
687 match self {
688 CacheStrategy::EnableAll(cache_manager) | CacheStrategy::Compaction(cache_manager) => {
689 cache_manager.get_sst_meta_data_from_mem_cache(file_id, page_index_policy)
690 }
691 CacheStrategy::Disabled => None,
692 }
693 }
694
695 pub fn get_parquet_meta_data_from_mem_cache(
697 &self,
698 file_id: RegionFileId,
699 ) -> Option<Arc<ParquetMetaData>> {
700 self.get_sst_meta_data_from_mem_cache(file_id, PageIndexPolicy::Skip)
701 .map(|metadata| metadata.parquet_metadata())
702 }
703
704 pub(crate) fn put_prepared_sst_meta(
706 &self,
707 file_id: RegionFileId,
708 metadata: PreparedSstMeta,
709 retain_decoded: bool,
710 ) {
711 match self {
712 CacheStrategy::EnableAll(cache_manager) | CacheStrategy::Compaction(cache_manager) => {
713 cache_manager.put_prepared_sst_meta(file_id, metadata, retain_decoded);
714 }
715 CacheStrategy::Disabled => {}
716 }
717 }
718
719 pub fn put_parquet_meta_data(
721 &self,
722 file_id: RegionFileId,
723 metadata: Arc<ParquetMetaData>,
724 region_metadata: Option<RegionMetadataRef>,
725 ) {
726 match self {
727 CacheStrategy::EnableAll(cache_manager) | CacheStrategy::Compaction(cache_manager) => {
728 cache_manager.put_parquet_meta_data(file_id, metadata, region_metadata);
729 }
730 CacheStrategy::Disabled => {}
731 }
732 }
733
734 pub(crate) fn get_prefilter_result(&self, key: &PrefilterKey) -> Option<Arc<BooleanBuffer>> {
737 match self {
738 CacheStrategy::EnableAll(cache_manager) => cache_manager.get_prefilter_result(key),
739 CacheStrategy::Compaction(_) | CacheStrategy::Disabled => None,
740 }
741 }
742
743 pub(crate) fn put_prefilter_result(&self, key: PrefilterKey, result: Arc<BooleanBuffer>) {
746 if let CacheStrategy::EnableAll(cache_manager) = self {
747 cache_manager.put_prefilter_result(key, result);
748 }
749 }
750
751 pub fn remove_parquet_meta_data(&self, file_id: RegionFileId) {
753 match self {
754 CacheStrategy::EnableAll(cache_manager) => {
755 cache_manager.remove_parquet_meta_data(file_id);
756 }
757 CacheStrategy::Compaction(cache_manager) => {
758 cache_manager.remove_parquet_meta_data(file_id);
759 }
760 CacheStrategy::Disabled => {}
761 }
762 }
763
764 pub fn get_repeated_vector(
767 &self,
768 data_type: &ConcreteDataType,
769 value: &Value,
770 ) -> Option<VectorRef> {
771 match self {
772 CacheStrategy::EnableAll(cache_manager) => {
773 cache_manager.get_repeated_vector(data_type, value)
774 }
775 CacheStrategy::Compaction(_) | CacheStrategy::Disabled => None,
776 }
777 }
778
779 pub fn put_repeated_vector(&self, value: Value, vector: VectorRef) {
782 if let CacheStrategy::EnableAll(cache_manager) = self {
783 cache_manager.put_repeated_vector(value, vector);
784 }
785 }
786
787 pub fn get_page_ranges(
790 &self,
791 file_id: FileId,
792 row_group_idx: usize,
793 ranges: &[Range<u64>],
794 ) -> Option<PageRangeLookup> {
795 match self {
796 CacheStrategy::EnableAll(cache_manager) => {
797 cache_manager.get_page_ranges(file_id, row_group_idx, ranges)
798 }
799 CacheStrategy::Compaction(_) | CacheStrategy::Disabled => None,
800 }
801 }
802
803 pub fn put_page_ranges(
806 &self,
807 file_id: FileId,
808 row_group_idx: usize,
809 ranges: &[Range<u64>],
810 pages: &[Bytes],
811 ) {
812 if let CacheStrategy::EnableAll(cache_manager) = self {
813 cache_manager.put_page_ranges(file_id, row_group_idx, ranges, pages);
814 }
815 }
816
817 pub async fn evict_puffin_cache(&self, file_id: RegionIndexId) {
819 match self {
820 CacheStrategy::EnableAll(cache_manager) => {
821 cache_manager.evict_puffin_cache(file_id).await
822 }
823 CacheStrategy::Compaction(cache_manager) => {
824 cache_manager.evict_puffin_cache(file_id).await
825 }
826 CacheStrategy::Disabled => {}
827 }
828 }
829
830 pub fn get_selector_result(
833 &self,
834 selector_key: &SelectorResultKey,
835 ) -> Option<Arc<SelectorResultValue>> {
836 match self {
837 CacheStrategy::EnableAll(cache_manager) => {
838 cache_manager.get_selector_result(selector_key)
839 }
840 CacheStrategy::Compaction(_) | CacheStrategy::Disabled => None,
841 }
842 }
843
844 pub fn put_selector_result(
847 &self,
848 selector_key: SelectorResultKey,
849 result: Arc<SelectorResultValue>,
850 ) {
851 if let CacheStrategy::EnableAll(cache_manager) = self {
852 cache_manager.put_selector_result(selector_key, result);
853 }
854 }
855
856 #[allow(dead_code)]
859 pub(crate) fn get_range_result(
860 &self,
861 key: &RangeScanCacheKey,
862 ) -> Option<Arc<RangeScanCacheValue>> {
863 match self {
864 CacheStrategy::EnableAll(cache_manager) => cache_manager.get_range_result(key),
865 CacheStrategy::Compaction(_) | CacheStrategy::Disabled => None,
866 }
867 }
868
869 pub(crate) fn put_range_result(
872 &self,
873 key: RangeScanCacheKey,
874 result: Arc<RangeScanCacheValue>,
875 ) {
876 if let CacheStrategy::EnableAll(cache_manager) = self {
877 cache_manager.put_range_result(key, result);
878 }
879 }
880
881 pub(crate) fn has_range_result_cache(&self) -> bool {
883 match self {
884 CacheStrategy::EnableAll(cache_manager) => cache_manager.has_range_result_cache(),
885 CacheStrategy::Compaction(_) | CacheStrategy::Disabled => false,
886 }
887 }
888
889 pub(crate) fn range_result_memory_limiter(&self) -> Option<&Arc<RangeResultMemoryLimiter>> {
890 match self {
891 CacheStrategy::EnableAll(cache_manager) => {
892 Some(cache_manager.range_result_memory_limiter())
893 }
894 CacheStrategy::Compaction(_) | CacheStrategy::Disabled => None,
895 }
896 }
897
898 pub(crate) fn range_result_cache_size(&self) -> Option<usize> {
899 match self {
900 CacheStrategy::EnableAll(cache_manager) => {
901 Some(cache_manager.range_result_cache_size())
902 }
903 CacheStrategy::Compaction(_) | CacheStrategy::Disabled => None,
904 }
905 }
906
907 pub fn write_cache(&self) -> Option<&WriteCacheRef> {
910 match self {
911 CacheStrategy::EnableAll(cache_manager) => cache_manager.write_cache(),
912 CacheStrategy::Compaction(cache_manager) => cache_manager.write_cache(),
913 CacheStrategy::Disabled => None,
914 }
915 }
916
917 pub fn inverted_index_cache(&self) -> Option<&InvertedIndexCacheRef> {
920 match self {
921 CacheStrategy::EnableAll(cache_manager) => cache_manager.inverted_index_cache(),
922 CacheStrategy::Compaction(_) | CacheStrategy::Disabled => None,
923 }
924 }
925
926 pub fn bloom_filter_index_cache(&self) -> Option<&BloomFilterIndexCacheRef> {
929 match self {
930 CacheStrategy::EnableAll(cache_manager) => cache_manager.bloom_filter_index_cache(),
931 CacheStrategy::Compaction(_) | CacheStrategy::Disabled => None,
932 }
933 }
934
935 #[cfg(feature = "vector_index")]
938 pub fn vector_index_cache(&self) -> Option<&VectorIndexCacheRef> {
939 match self {
940 CacheStrategy::EnableAll(cache_manager) => cache_manager.vector_index_cache(),
941 CacheStrategy::Compaction(_) | CacheStrategy::Disabled => None,
942 }
943 }
944
945 pub fn puffin_metadata_cache(&self) -> Option<&PuffinMetadataCacheRef> {
948 match self {
949 CacheStrategy::EnableAll(cache_manager) => cache_manager.puffin_metadata_cache(),
950 CacheStrategy::Compaction(_) | CacheStrategy::Disabled => None,
951 }
952 }
953
954 pub fn index_result_cache(&self) -> Option<&IndexResultCache> {
957 match self {
958 CacheStrategy::EnableAll(cache_manager) => cache_manager.index_result_cache(),
959 CacheStrategy::Compaction(_) | CacheStrategy::Disabled => None,
960 }
961 }
962
963 pub fn maybe_download_background(
965 &self,
966 index_key: IndexKey,
967 remote_path: String,
968 remote_store: ObjectStore,
969 file_size: u64,
970 ) {
971 if let CacheStrategy::EnableAll(cache_manager) = self
972 && let Some(write_cache) = cache_manager.write_cache()
973 {
974 write_cache.file_cache().maybe_download_background(
975 index_key,
976 remote_path,
977 remote_store,
978 file_size,
979 );
980 }
981 }
982}
983
984#[derive(Default)]
988pub struct CacheManager {
989 sst_meta_cache: Option<SstMetaCache>,
991 sst_decoded_meta_cache: Option<SstDecodedMetaCache>,
993 vector_cache: Option<VectorCache>,
995 page_cache: Option<Arc<PageRangeCache>>,
997 write_cache: Option<WriteCacheRef>,
999 inverted_index_cache: Option<InvertedIndexCacheRef>,
1001 bloom_filter_index_cache: Option<BloomFilterIndexCacheRef>,
1003 #[cfg(feature = "vector_index")]
1005 vector_index_cache: Option<VectorIndexCacheRef>,
1006 puffin_metadata_cache: Option<PuffinMetadataCacheRef>,
1008 selector_result_cache: Option<SelectorResultCache>,
1010 range_result_cache: Option<RangeResultCache>,
1012 range_result_cache_size: u64,
1014 range_result_memory_limiter: Arc<RangeResultMemoryLimiter>,
1016 index_result_cache: Option<IndexResultCache>,
1018 prefilter_result_cache: Option<PrefilterResultCache>,
1020}
1021
1022pub type CacheManagerRef = Arc<CacheManager>;
1023
1024impl CacheManager {
1025 pub fn builder() -> CacheManagerBuilder {
1027 CacheManagerBuilder::default()
1028 }
1029
1030 pub(crate) async fn get_sst_meta_data(
1033 &self,
1034 file_id: RegionFileId,
1035 metrics: &mut MetadataCacheMetrics,
1036 page_index_policy: PageIndexPolicy,
1037 ) -> Option<Arc<CachedSstMeta>> {
1038 let cache_key = SstMetaKey(file_id.region_id(), file_id.file_id());
1039 let compact = self
1040 .get_compact_sst_meta(&cache_key)
1041 .filter(|metadata| metadata.satisfies_page_index_policy(page_index_policy));
1042 let decoded = self
1043 .get_decoded_sst_meta(&cache_key)
1044 .filter(|metadata| metadata.satisfies_page_index_policy(page_index_policy));
1045
1046 if let Some(compact) = compact {
1047 CACHE_HIT.with_label_values(&[SST_META_TYPE]).inc();
1048 metrics.mem_cache_hit += 1;
1049 if let Some(decoded) = decoded {
1050 CACHE_HIT.with_label_values(&[SST_META_DECODED_TYPE]).inc();
1051 return Some(decoded);
1052 }
1053
1054 CACHE_MISS.with_label_values(&[SST_META_DECODED_TYPE]).inc();
1055 match compact.decode().await {
1056 Ok(decoded) => {
1057 self.put_sst_meta_data(file_id, decoded.clone());
1058 return Some(decoded);
1059 }
1060 Err(err) => {
1061 warn!(err; "Failed to decode compact SST metadata, region_id: {}, file_id: {}", file_id.region_id(), file_id.file_id());
1062 self.remove_parquet_meta_data(file_id);
1063 }
1064 }
1065 } else if let Some(decoded) = decoded {
1066 CACHE_HIT.with_label_values(&[SST_META_TYPE]).inc();
1069 CACHE_HIT.with_label_values(&[SST_META_DECODED_TYPE]).inc();
1070 metrics.mem_cache_hit += 1;
1071 return Some(decoded);
1072 } else {
1073 CACHE_MISS.with_label_values(&[SST_META_TYPE]).inc();
1074 }
1075
1076 let key = IndexKey::new(file_id.region_id(), file_id.file_id(), FileType::Parquet);
1077 if let Some(write_cache) = &self.write_cache {
1078 let file_cache = write_cache.file_cache();
1079 if self.sst_meta_cache_enabled() {
1080 if let Some(metadata) = file_cache
1081 .get_sst_meta_data(key, metrics, page_index_policy)
1082 .await
1083 {
1084 metrics.file_cache_hit += 1;
1085 let decoded = metadata.decoded();
1086 match metadata {
1087 SstMetaPreparation::Prepared(metadata) => {
1088 self.put_prepared_sst_meta(file_id, metadata, true);
1089 }
1090 SstMetaPreparation::DecodedOnly { encoding_error, .. } => {
1091 warn!(
1092 encoding_error;
1093 "Failed to encode file-cached SST metadata for memory cache, region_id: {}, file_id: {}",
1094 file_id.region_id(),
1095 file_id.file_id()
1096 );
1097 }
1098 }
1099 return Some(decoded);
1100 }
1101 } else if let Some(decoded) = file_cache
1102 .get_decoded_sst_meta_data(key, metrics, page_index_policy)
1103 .await
1104 {
1105 metrics.file_cache_hit += 1;
1106 return Some(decoded);
1107 }
1108 }
1109
1110 metrics.cache_miss += 1;
1111 None
1112 }
1113
1114 pub(crate) fn get_sst_meta_data_from_mem_cache(
1117 &self,
1118 file_id: RegionFileId,
1119 page_index_policy: PageIndexPolicy,
1120 ) -> Option<Arc<CachedSstMeta>> {
1121 let key = SstMetaKey(file_id.region_id(), file_id.file_id());
1122 let value = self
1123 .get_decoded_sst_meta(&key)
1124 .filter(|metadata| metadata.satisfies_page_index_policy(page_index_policy));
1125 update_hit_miss(value, SST_META_DECODED_TYPE)
1126 }
1127
1128 pub fn get_parquet_meta_data_from_mem_cache(
1131 &self,
1132 file_id: RegionFileId,
1133 ) -> Option<Arc<ParquetMetaData>> {
1134 self.get_sst_meta_data_from_mem_cache(file_id, PageIndexPolicy::Skip)
1135 .map(|metadata| metadata.parquet_metadata())
1136 }
1137
1138 pub(crate) fn put_sst_meta_data(&self, file_id: RegionFileId, metadata: Arc<CachedSstMeta>) {
1140 if let Some(cache) = &self.sst_decoded_meta_cache {
1141 let key = SstMetaKey(file_id.region_id(), file_id.file_id());
1142 CACHE_BYTES
1143 .with_label_values(&[SST_META_DECODED_TYPE])
1144 .add(decoded_meta_cache_weight(&key, &metadata).into());
1145 cache.insert(key, metadata);
1146 }
1147 }
1148
1149 pub(crate) fn put_prepared_sst_meta(
1151 &self,
1152 file_id: RegionFileId,
1153 metadata: PreparedSstMeta,
1154 retain_decoded: bool,
1155 ) {
1156 let key = SstMetaKey(file_id.region_id(), file_id.file_id());
1157 if let Some(cache) = &self.sst_meta_cache {
1158 CACHE_BYTES
1159 .with_label_values(&[SST_META_TYPE])
1160 .add(meta_cache_weight(&key, &metadata.compact).into());
1161 cache.insert(key.clone(), metadata.compact);
1162 }
1163 if retain_decoded {
1164 self.put_sst_meta_data(file_id, metadata.decoded);
1165 }
1166 }
1167
1168 fn get_compact_sst_meta(&self, key: &SstMetaKey) -> Option<Arc<CompactSstMeta>> {
1169 self.sst_meta_cache.as_ref()?.get(key)
1170 }
1171
1172 fn get_decoded_sst_meta(&self, key: &SstMetaKey) -> Option<Arc<CachedSstMeta>> {
1173 self.sst_decoded_meta_cache.as_ref()?.get(key)
1174 }
1175
1176 pub(crate) async fn get_compact_sst_meta_data(
1181 &self,
1182 file_id: RegionFileId,
1183 page_index_policy: PageIndexPolicy,
1184 ) -> Option<Arc<CachedSstMeta>> {
1185 let key = SstMetaKey(file_id.region_id(), file_id.file_id());
1186 let compact = self
1187 .get_compact_sst_meta(&key)
1188 .filter(|metadata| metadata.satisfies_page_index_policy(page_index_policy))?;
1189 match compact.decode().await {
1190 Ok(metadata) => Some(metadata),
1191 Err(err) => {
1192 warn!(err; "Failed to decode compact SST metadata, region_id: {}, file_id: {}", file_id.region_id(), file_id.file_id());
1193 self.remove_parquet_meta_data(file_id);
1194 None
1195 }
1196 }
1197 }
1198
1199 pub fn put_parquet_meta_data(
1201 &self,
1202 file_id: RegionFileId,
1203 metadata: Arc<ParquetMetaData>,
1204 region_metadata: Option<RegionMetadataRef>,
1205 ) {
1206 if self.sst_decoded_meta_cache.is_some() {
1207 let file_path = format!(
1208 "region_id={}, file_id={}",
1209 file_id.region_id(),
1210 file_id.file_id()
1211 );
1212 match CachedSstMeta::try_new_with_region_metadata(
1213 &file_path,
1214 Arc::unwrap_or_clone(metadata),
1215 region_metadata,
1216 ) {
1217 Ok(metadata) => self.put_sst_meta_data(file_id, Arc::new(metadata)),
1218 Err(err) => warn!(
1219 err; "Failed to decode region metadata while caching parquet metadata, region_id: {}, file_id: {}",
1220 file_id.region_id(),
1221 file_id.file_id()
1222 ),
1223 }
1224 }
1225 }
1226
1227 pub fn remove_parquet_meta_data(&self, file_id: RegionFileId) {
1229 let key = SstMetaKey(file_id.region_id(), file_id.file_id());
1230 if let Some(cache) = &self.sst_meta_cache {
1231 cache.remove(&key);
1232 }
1233 if let Some(cache) = &self.sst_decoded_meta_cache {
1234 cache.remove(&key);
1235 }
1236 }
1237
1238 pub(crate) fn sst_meta_cache_is_full(&self) -> bool {
1240 let Some(cache) = &self.sst_meta_cache else {
1241 return true;
1242 };
1243 cache
1244 .policy()
1245 .max_capacity()
1246 .is_some_and(|capacity| cache.weighted_size() >= capacity)
1247 }
1248
1249 pub(crate) fn sst_meta_cache_enabled(&self) -> bool {
1251 self.sst_meta_cache.is_some()
1252 }
1253
1254 pub fn get_repeated_vector(
1256 &self,
1257 data_type: &ConcreteDataType,
1258 value: &Value,
1259 ) -> Option<VectorRef> {
1260 self.vector_cache.as_ref().and_then(|vector_cache| {
1261 let value = vector_cache.get(&(data_type.clone(), value.clone()));
1262 update_hit_miss(value, VECTOR_TYPE)
1263 })
1264 }
1265
1266 pub fn put_repeated_vector(&self, value: Value, vector: VectorRef) {
1268 if let Some(cache) = &self.vector_cache {
1269 let key = (vector.data_type(), value);
1270 CACHE_BYTES
1271 .with_label_values(&[VECTOR_TYPE])
1272 .add(vector_cache_weight(&key, &vector).into());
1273 cache.insert(key, vector);
1274 }
1275 }
1276
1277 pub fn get_page_ranges(
1279 &self,
1280 file_id: FileId,
1281 row_group_idx: usize,
1282 ranges: &[Range<u64>],
1283 ) -> Option<PageRangeLookup> {
1284 self.page_cache.as_ref().map(|page_cache| {
1285 let lookup = page_cache.lookup(file_id, row_group_idx, ranges);
1286 if lookup.cached_bytes > 0 {
1287 CACHE_HIT.with_label_values(&[PAGE_TYPE]).inc();
1288 }
1289 if !lookup.missing_ranges.is_empty() {
1290 CACHE_MISS.with_label_values(&[PAGE_TYPE]).inc();
1291 }
1292 lookup
1293 })
1294 }
1295
1296 pub fn put_page_ranges(
1298 &self,
1299 file_id: FileId,
1300 row_group_idx: usize,
1301 ranges: &[Range<u64>],
1302 pages: &[Bytes],
1303 ) {
1304 if let Some(cache) = &self.page_cache {
1305 cache.insert_ranges(file_id, row_group_idx, ranges, pages);
1306 }
1307 }
1308
1309 pub async fn evict_puffin_cache(&self, file_id: RegionIndexId) {
1311 if let Some(cache) = &self.bloom_filter_index_cache {
1312 cache.invalidate_file(file_id.file_id());
1313 }
1314
1315 if let Some(cache) = &self.inverted_index_cache {
1316 cache.invalidate_file(file_id.file_id());
1317 }
1318
1319 if let Some(cache) = &self.index_result_cache {
1320 cache.invalidate_file(file_id.file_id());
1321 }
1322
1323 #[cfg(feature = "vector_index")]
1324 if let Some(cache) = &self.vector_index_cache {
1325 cache.invalidate_file(file_id.file_id());
1326 }
1327
1328 if let Some(cache) = &self.puffin_metadata_cache {
1329 cache.remove(&file_id.to_string());
1330 }
1331
1332 if let Some(write_cache) = &self.write_cache {
1333 write_cache
1334 .remove(IndexKey::new(
1335 file_id.region_id(),
1336 file_id.file_id(),
1337 FileType::Puffin(file_id.version),
1338 ))
1339 .await;
1340 }
1341 }
1342
1343 pub fn get_selector_result(
1345 &self,
1346 selector_key: &SelectorResultKey,
1347 ) -> Option<Arc<SelectorResultValue>> {
1348 self.selector_result_cache
1349 .as_ref()
1350 .and_then(|selector_result_cache| selector_result_cache.get(selector_key))
1351 }
1352
1353 pub fn put_selector_result(
1355 &self,
1356 selector_key: SelectorResultKey,
1357 result: Arc<SelectorResultValue>,
1358 ) {
1359 if let Some(cache) = &self.selector_result_cache {
1360 CACHE_BYTES
1361 .with_label_values(&[SELECTOR_RESULT_TYPE])
1362 .add(selector_result_cache_weight(&selector_key, &result).into());
1363 cache.insert(selector_key, result);
1364 }
1365 }
1366
1367 #[allow(dead_code)]
1369 pub(crate) fn get_range_result(
1370 &self,
1371 key: &RangeScanCacheKey,
1372 ) -> Option<Arc<RangeScanCacheValue>> {
1373 self.range_result_cache
1374 .as_ref()
1375 .and_then(|cache| update_hit_miss(cache.get(key), RANGE_RESULT_TYPE))
1376 }
1377
1378 pub(crate) fn put_range_result(
1380 &self,
1381 key: RangeScanCacheKey,
1382 result: Arc<RangeScanCacheValue>,
1383 ) {
1384 if let Some(cache) = &self.range_result_cache {
1385 CACHE_BYTES
1386 .with_label_values(&[RANGE_RESULT_TYPE])
1387 .add(range_result_cache_weight(&key, &result).into());
1388 cache.insert(key, result);
1389 }
1390 }
1391
1392 pub(crate) fn has_range_result_cache(&self) -> bool {
1394 self.range_result_cache.is_some()
1395 }
1396
1397 pub(crate) fn range_result_memory_limiter(&self) -> &Arc<RangeResultMemoryLimiter> {
1398 &self.range_result_memory_limiter
1399 }
1400
1401 pub(crate) fn range_result_cache_size(&self) -> usize {
1402 self.range_result_cache_size as usize
1403 }
1404
1405 pub(crate) fn write_cache(&self) -> Option<&WriteCacheRef> {
1407 self.write_cache.as_ref()
1408 }
1409
1410 pub(crate) fn inverted_index_cache(&self) -> Option<&InvertedIndexCacheRef> {
1411 self.inverted_index_cache.as_ref()
1412 }
1413
1414 pub(crate) fn bloom_filter_index_cache(&self) -> Option<&BloomFilterIndexCacheRef> {
1415 self.bloom_filter_index_cache.as_ref()
1416 }
1417
1418 #[cfg(feature = "vector_index")]
1419 pub(crate) fn vector_index_cache(&self) -> Option<&VectorIndexCacheRef> {
1420 self.vector_index_cache.as_ref()
1421 }
1422
1423 pub(crate) fn puffin_metadata_cache(&self) -> Option<&PuffinMetadataCacheRef> {
1424 self.puffin_metadata_cache.as_ref()
1425 }
1426
1427 pub(crate) fn index_result_cache(&self) -> Option<&IndexResultCache> {
1428 self.index_result_cache.as_ref()
1429 }
1430
1431 pub(crate) fn get_prefilter_result(&self, key: &PrefilterKey) -> Option<Arc<BooleanBuffer>> {
1432 self.prefilter_result_cache
1433 .as_ref()
1434 .and_then(|cache| update_hit_miss(cache.get(key), PREFILTER_RESULT_TYPE))
1435 }
1436
1437 pub(crate) fn put_prefilter_result(&self, key: PrefilterKey, result: Arc<BooleanBuffer>) {
1438 if let Some(cache) = &self.prefilter_result_cache {
1439 CACHE_BYTES
1440 .with_label_values(&[PREFILTER_RESULT_TYPE])
1441 .add(prefilter_result_cache_weight(&key, &result).into());
1442 cache.insert(key, result);
1443 }
1444 }
1445}
1446
1447pub fn selector_result_cache_miss() {
1449 CACHE_MISS.with_label_values(&[SELECTOR_RESULT_TYPE]).inc()
1450}
1451
1452pub fn selector_result_cache_hit() {
1454 CACHE_HIT.with_label_values(&[SELECTOR_RESULT_TYPE]).inc()
1455}
1456
1457#[derive(Default)]
1459pub struct CacheManagerBuilder {
1460 sst_meta_cache_size: u64,
1461 vector_cache_size: u64,
1462 page_cache_size: u64,
1463 index_metadata_size: u64,
1464 index_content_size: u64,
1465 index_content_page_size: u64,
1466 index_result_cache_size: u64,
1467 prefilter_result_cache_size: u64,
1468 puffin_metadata_size: u64,
1469 write_cache: Option<WriteCacheRef>,
1470 selector_result_cache_size: u64,
1471 range_result_cache_size: u64,
1472}
1473
1474impl CacheManagerBuilder {
1475 pub fn sst_meta_cache_size(mut self, bytes: u64) -> Self {
1477 self.sst_meta_cache_size = bytes;
1478 self
1479 }
1480
1481 pub fn vector_cache_size(mut self, bytes: u64) -> Self {
1483 self.vector_cache_size = bytes;
1484 self
1485 }
1486
1487 pub fn page_cache_size(mut self, bytes: u64) -> Self {
1489 self.page_cache_size = bytes;
1490 self
1491 }
1492
1493 pub fn write_cache(mut self, cache: Option<WriteCacheRef>) -> Self {
1495 self.write_cache = cache;
1496 self
1497 }
1498
1499 pub fn index_metadata_size(mut self, bytes: u64) -> Self {
1501 self.index_metadata_size = bytes;
1502 self
1503 }
1504
1505 pub fn index_content_size(mut self, bytes: u64) -> Self {
1507 self.index_content_size = bytes;
1508 self
1509 }
1510
1511 pub fn index_content_page_size(mut self, bytes: u64) -> Self {
1513 self.index_content_page_size = bytes;
1514 self
1515 }
1516
1517 pub fn index_result_cache_size(mut self, bytes: u64) -> Self {
1519 self.index_result_cache_size = bytes;
1520 self
1521 }
1522
1523 pub fn prefilter_result_cache_size(mut self, bytes: u64) -> Self {
1525 self.prefilter_result_cache_size = bytes;
1526 self
1527 }
1528
1529 pub fn puffin_metadata_size(mut self, bytes: u64) -> Self {
1531 self.puffin_metadata_size = bytes;
1532 self
1533 }
1534
1535 pub fn selector_result_cache_size(mut self, bytes: u64) -> Self {
1537 self.selector_result_cache_size = bytes;
1538 self
1539 }
1540
1541 pub fn range_result_cache_size(mut self, bytes: u64) -> Self {
1543 self.range_result_cache_size = bytes;
1544 self
1545 }
1546
1547 pub fn build(self) -> CacheManager {
1549 let compact_meta_capacity = self.sst_meta_cache_size.div_ceil(2);
1553 let decoded_meta_capacity = self.sst_meta_cache_size / 2;
1554 let sst_meta_cache = (compact_meta_capacity != 0).then(|| {
1555 Cache::builder()
1556 .max_capacity(compact_meta_capacity)
1557 .weigher(meta_cache_weight)
1558 .eviction_listener(|k, v, cause| {
1559 let size = meta_cache_weight(&k, &v);
1560 CACHE_BYTES
1561 .with_label_values(&[SST_META_TYPE])
1562 .sub(size.into());
1563 CACHE_EVICTION
1564 .with_label_values(&[SST_META_TYPE, removal_cause_str(cause)])
1565 .inc();
1566 })
1567 .build()
1568 });
1569 let sst_decoded_meta_cache = (decoded_meta_capacity != 0).then(|| {
1570 Cache::builder()
1571 .max_capacity(decoded_meta_capacity)
1572 .weigher(decoded_meta_cache_weight)
1573 .eviction_listener(|k, v, cause| {
1574 let size = decoded_meta_cache_weight(&k, &v);
1575 CACHE_BYTES
1576 .with_label_values(&[SST_META_DECODED_TYPE])
1577 .sub(size.into());
1578 CACHE_EVICTION
1579 .with_label_values(&[SST_META_DECODED_TYPE, removal_cause_str(cause)])
1580 .inc();
1581 })
1582 .build()
1583 });
1584 let vector_cache = (self.vector_cache_size != 0).then(|| {
1585 Cache::builder()
1586 .max_capacity(self.vector_cache_size)
1587 .weigher(vector_cache_weight)
1588 .eviction_listener(|k, v, cause| {
1589 let size = vector_cache_weight(&k, &v);
1590 CACHE_BYTES
1591 .with_label_values(&[VECTOR_TYPE])
1592 .sub(size.into());
1593 CACHE_EVICTION
1594 .with_label_values(&[VECTOR_TYPE, removal_cause_str(cause)])
1595 .inc();
1596 })
1597 .build()
1598 });
1599 let page_cache =
1600 (self.page_cache_size != 0).then(|| PageRangeCache::new(self.page_cache_size));
1601 let inverted_index_cache = InvertedIndexCache::new(
1602 self.index_metadata_size,
1603 self.index_content_size,
1604 self.index_content_page_size,
1605 );
1606 let bloom_filter_index_cache = BloomFilterIndexCache::new(
1608 self.index_metadata_size,
1609 self.index_content_size,
1610 self.index_content_page_size,
1611 );
1612 #[cfg(feature = "vector_index")]
1613 let vector_index_cache = (self.index_content_size != 0)
1614 .then(|| Arc::new(VectorIndexCache::new(self.index_content_size)));
1615 let index_result_cache = (self.index_result_cache_size != 0)
1616 .then(|| IndexResultCache::new(self.index_result_cache_size));
1617 let prefilter_result_cache = (self.prefilter_result_cache_size != 0)
1618 .then(|| new_prefilter_result_cache(self.prefilter_result_cache_size));
1619 let puffin_metadata_cache =
1620 PuffinMetadataCache::new(self.puffin_metadata_size, &CACHE_BYTES);
1621 let selector_result_cache = (self.selector_result_cache_size != 0).then(|| {
1622 Cache::builder()
1623 .max_capacity(self.selector_result_cache_size)
1624 .weigher(selector_result_cache_weight)
1625 .eviction_listener(|k, v, cause| {
1626 let size = selector_result_cache_weight(&k, &v);
1627 CACHE_BYTES
1628 .with_label_values(&[SELECTOR_RESULT_TYPE])
1629 .sub(size.into());
1630 CACHE_EVICTION
1631 .with_label_values(&[SELECTOR_RESULT_TYPE, removal_cause_str(cause)])
1632 .inc();
1633 })
1634 .build()
1635 });
1636 let range_result_cache = (self.range_result_cache_size != 0).then(|| {
1637 Cache::builder()
1638 .max_capacity(self.range_result_cache_size)
1639 .weigher(range_result_cache_weight)
1640 .eviction_listener(move |k, v, cause| {
1641 let size = range_result_cache_weight(&k, &v);
1642 CACHE_BYTES
1643 .with_label_values(&[RANGE_RESULT_TYPE])
1644 .sub(size.into());
1645 CACHE_EVICTION
1646 .with_label_values(&[RANGE_RESULT_TYPE, removal_cause_str(cause)])
1647 .inc();
1648 })
1649 .build()
1650 });
1651 CacheManager {
1652 sst_meta_cache,
1653 sst_decoded_meta_cache,
1654 vector_cache,
1655 page_cache,
1656 write_cache: self.write_cache,
1657 inverted_index_cache: Some(Arc::new(inverted_index_cache)),
1658 bloom_filter_index_cache: Some(Arc::new(bloom_filter_index_cache)),
1659 #[cfg(feature = "vector_index")]
1660 vector_index_cache,
1661 puffin_metadata_cache: Some(Arc::new(puffin_metadata_cache)),
1662 selector_result_cache,
1663 range_result_cache,
1664 range_result_cache_size: self.range_result_cache_size,
1665 range_result_memory_limiter: Arc::new(RangeResultMemoryLimiter::new(
1666 self.range_result_cache_size as usize,
1667 RANGE_RESULT_CONCAT_MEMORY_PERMIT.as_bytes() as usize,
1668 )),
1669 index_result_cache,
1670 prefilter_result_cache,
1671 }
1672 }
1673}
1674
1675fn meta_cache_weight(k: &SstMetaKey, v: &Arc<CompactSstMeta>) -> u32 {
1676 let size = k.estimated_size() + mem::size_of::<CompactSstMeta>() + v.encoded_metadata.len();
1679 u32::try_from(size).unwrap_or(u32::MAX)
1680}
1681
1682fn decoded_meta_cache_weight(k: &SstMetaKey, v: &Arc<CachedSstMeta>) -> u32 {
1683 let size = k.estimated_size() + v.parquet_metadata_size + v.region_metadata_weight;
1684 u32::try_from(size).unwrap_or(u32::MAX)
1685}
1686
1687fn vector_cache_weight(_k: &(ConcreteDataType, Value), v: &VectorRef) -> u32 {
1688 (mem::size_of::<ConcreteDataType>() + mem::size_of::<Value>() + v.memory_size()) as u32
1690}
1691
1692fn page_cache_weight(k: &PageFragmentKey, v: &Bytes) -> u32 {
1693 (k.estimated_size() + mem::size_of::<Bytes>() + v.len()) as u32
1694}
1695
1696fn selector_result_cache_weight(k: &SelectorResultKey, v: &Arc<SelectorResultValue>) -> u32 {
1697 (mem::size_of_val(k) + v.estimated_size()) as u32
1698}
1699
1700fn range_result_cache_weight(k: &RangeScanCacheKey, v: &Arc<RangeScanCacheValue>) -> u32 {
1701 (k.estimated_size() + v.estimated_size()) as u32
1702}
1703
1704fn update_hit_miss<T>(value: Option<T>, cache_type: &str) -> Option<T> {
1706 if value.is_some() {
1707 CACHE_HIT.with_label_values(&[cache_type]).inc();
1708 } else {
1709 CACHE_MISS.with_label_values(&[cache_type]).inc();
1710 }
1711 value
1712}
1713
1714#[derive(Debug, Clone, PartialEq, Eq, Hash)]
1716struct SstMetaKey(RegionId, FileId);
1717
1718impl SstMetaKey {
1719 fn estimated_size(&self) -> usize {
1721 mem::size_of::<Self>()
1722 }
1723}
1724
1725#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
1726struct PageFragmentGroupKey {
1727 file_id: FileId,
1728 row_group_idx: usize,
1729}
1730
1731#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
1733pub struct PageFragmentKey {
1734 file_id: FileId,
1736 row_group_idx: usize,
1738 start: u64,
1740 end: u64,
1742}
1743
1744impl PageFragmentKey {
1745 fn new(file_id: FileId, row_group_idx: usize, range: &Range<u64>) -> PageFragmentKey {
1746 PageFragmentKey {
1747 file_id,
1748 row_group_idx,
1749 start: range.start,
1750 end: range.end,
1751 }
1752 }
1753
1754 fn group_key(&self) -> PageFragmentGroupKey {
1755 PageFragmentGroupKey {
1756 file_id: self.file_id,
1757 row_group_idx: self.row_group_idx,
1758 }
1759 }
1760
1761 fn estimated_size(&self) -> usize {
1763 mem::size_of::<Self>()
1764 }
1765}
1766
1767#[derive(Clone)]
1769pub struct PageRangePart {
1770 pub range: Range<u64>,
1772 pub bytes: Bytes,
1774}
1775
1776pub struct PageRangeLookup {
1778 pub cached_parts: Vec<Vec<PageRangePart>>,
1780 pub missing_ranges: Vec<Range<u64>>,
1782 pub cached_range_count: usize,
1784 pub cached_bytes: u64,
1786}
1787
1788impl PageRangeLookup {
1789 pub fn is_fully_cached(&self) -> bool {
1790 self.missing_ranges.is_empty()
1791 }
1792}
1793
1794type PageFragmentRangeIndex = BTreeMap<(u64, u64), PageFragmentKey>;
1795type PageFragmentIndex = HashMap<PageFragmentGroupKey, PageFragmentRangeIndex>;
1796
1797pub struct PageRangeCache {
1799 cache: Cache<PageFragmentKey, Bytes>,
1800 index: RwLock<PageFragmentIndex>,
1801}
1802
1803impl PageRangeCache {
1804 fn new(capacity: u64) -> Arc<PageRangeCache> {
1805 Arc::new_cyclic(|weak_cache: &std::sync::Weak<PageRangeCache>| {
1806 let cache = Cache::builder()
1807 .max_capacity(capacity)
1808 .weigher(page_cache_weight)
1809 .eviction_listener({
1810 let weak_cache = weak_cache.clone();
1811 move |k, v, cause| {
1812 let size = page_cache_weight(&k, &v);
1813 CACHE_BYTES.with_label_values(&[PAGE_TYPE]).sub(size.into());
1814 CACHE_EVICTION
1815 .with_label_values(&[PAGE_TYPE, removal_cause_str(cause)])
1816 .inc();
1817
1818 if let Some(cache) = weak_cache.upgrade()
1819 && !matches!(cause, RemovalCause::Replaced)
1820 {
1821 cache.remove_index_entry(*k);
1822 }
1823 }
1824 })
1825 .build();
1826
1827 PageRangeCache {
1828 cache,
1829 index: RwLock::new(HashMap::new()),
1830 }
1831 })
1832 }
1833
1834 fn lookup(
1835 &self,
1836 file_id: FileId,
1837 row_group_idx: usize,
1838 ranges: &[Range<u64>],
1839 ) -> PageRangeLookup {
1840 let mut cached_parts = Vec::with_capacity(ranges.len());
1841 let mut missing_ranges = Vec::new();
1842 let mut cached_range_count = 0;
1843 let mut cached_bytes = 0;
1844
1845 for range in ranges {
1846 if range.start >= range.end {
1847 cached_parts.push(Vec::new());
1848 continue;
1849 }
1850
1851 let mut parts = Vec::new();
1852 let candidates = self.find_index_candidates(file_id, row_group_idx, range);
1853 let mut stale_keys = Vec::new();
1854
1855 for fragment_key in candidates {
1856 if let Some(bytes) = self.cache.get(&fragment_key) {
1857 let part_start = range.start.max(fragment_key.start);
1858 let part_end = range.end.min(fragment_key.end);
1859 let slice_start = (part_start - fragment_key.start) as usize;
1860 let slice_end = (part_end - fragment_key.start) as usize;
1861 parts.push(PageRangePart {
1862 range: part_start..part_end,
1863 bytes: bytes.slice(slice_start..slice_end),
1864 });
1865 } else {
1866 stale_keys.push(fragment_key);
1867 }
1868 }
1869 for key in stale_keys {
1870 self.remove_uncached_index_entry(key);
1871 }
1872
1873 let mut cursor = range.start;
1874 let mut compacted_parts: Vec<PageRangePart> = Vec::with_capacity(parts.len());
1875 for part in parts {
1876 if part.range.end <= cursor {
1877 continue;
1878 }
1879
1880 let part = if part.range.start < cursor {
1881 let offset = (cursor - part.range.start) as usize;
1882 PageRangePart {
1883 range: cursor..part.range.end,
1884 bytes: part.bytes.slice(offset..),
1885 }
1886 } else {
1887 part
1888 };
1889
1890 if cursor < part.range.start {
1891 missing_ranges.push(cursor..part.range.start);
1892 }
1893 cached_bytes += part.range.end - part.range.start;
1894 cached_range_count += 1;
1895 cursor = part.range.end;
1896 compacted_parts.push(part);
1897
1898 if cursor >= range.end {
1899 break;
1900 }
1901 }
1902
1903 if cursor < range.end {
1904 missing_ranges.push(cursor..range.end);
1905 }
1906 cached_parts.push(compacted_parts);
1907 }
1908
1909 PageRangeLookup {
1910 cached_parts,
1911 missing_ranges,
1912 cached_range_count,
1913 cached_bytes,
1914 }
1915 }
1916
1917 fn insert_ranges(
1918 &self,
1919 file_id: FileId,
1920 row_group_idx: usize,
1921 ranges: &[Range<u64>],
1922 pages: &[Bytes],
1923 ) {
1924 for (range, bytes) in ranges.iter().zip(pages) {
1925 if range.start >= range.end || bytes.len() as u64 != range.end - range.start {
1926 continue;
1927 }
1928
1929 let key = PageFragmentKey::new(file_id, row_group_idx, range);
1930 let bytes = Bytes::copy_from_slice(bytes.as_ref());
1931 let size = page_cache_weight(&key, &bytes);
1932 CACHE_BYTES.with_label_values(&[PAGE_TYPE]).add(size.into());
1933 self.cache.insert(key, bytes);
1934 let mut index = self.index.write().unwrap();
1935 index
1936 .entry(key.group_key())
1937 .or_default()
1938 .insert((key.start, key.end), key);
1939 }
1940 }
1941
1942 fn find_index_candidates(
1943 &self,
1944 file_id: FileId,
1945 row_group_idx: usize,
1946 range: &Range<u64>,
1947 ) -> Vec<PageFragmentKey> {
1948 let group_key = PageFragmentGroupKey {
1949 file_id,
1950 row_group_idx,
1951 };
1952 let index = self.index.read().unwrap();
1953 index
1954 .get(&group_key)
1955 .map(|ranges| {
1956 ranges
1957 .range(..(range.end, 0))
1958 .filter_map(|(_, fragment_key)| {
1959 (fragment_key.end > range.start).then_some(*fragment_key)
1960 })
1961 .collect()
1962 })
1963 .unwrap_or_default()
1964 }
1965
1966 fn remove_uncached_index_entry(&self, key: PageFragmentKey) {
1967 let group_key = key.group_key();
1968 let mut index = self.index.write().unwrap();
1969 if self.cache.contains_key(&key) {
1970 return;
1971 }
1972
1973 Self::remove_index_entry_locked(&mut index, group_key, key);
1974 }
1975
1976 fn remove_index_entry(&self, key: PageFragmentKey) {
1977 let group_key = key.group_key();
1978 let mut index = self.index.write().unwrap();
1979 Self::remove_index_entry_locked(&mut index, group_key, key);
1980 }
1981
1982 fn remove_index_entry_locked(
1983 index: &mut PageFragmentIndex,
1984 group_key: PageFragmentGroupKey,
1985 key: PageFragmentKey,
1986 ) {
1987 let Some(ranges) = index.get_mut(&group_key) else {
1988 return;
1989 };
1990
1991 let removed = ranges
1992 .get(&(key.start, key.end))
1993 .is_some_and(|current| current == &key);
1994 if removed {
1995 ranges.remove(&(key.start, key.end));
1996 }
1997 if ranges.is_empty() {
1998 index.remove(&group_key);
1999 }
2000 }
2001}
2002
2003#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
2005pub struct SelectorResultKey {
2006 pub file_id: FileId,
2008 pub row_group_idx: usize,
2010 pub selector: TimeSeriesRowSelector,
2012}
2013
2014pub enum SelectorResult {
2016 PrimaryKey(Vec<Batch>),
2018 Flat(Vec<RecordBatch>),
2020}
2021
2022pub struct SelectorResultValue {
2024 pub result: SelectorResult,
2026 pub read_cols: ParquetReadColumns,
2028}
2029
2030impl SelectorResultValue {
2031 pub fn new(result: Vec<Batch>, read_cols: ParquetReadColumns) -> SelectorResultValue {
2033 SelectorResultValue {
2034 result: SelectorResult::PrimaryKey(result),
2035 read_cols,
2036 }
2037 }
2038
2039 pub fn new_flat(
2041 result: Vec<RecordBatch>,
2042 read_cols: ParquetReadColumns,
2043 ) -> SelectorResultValue {
2044 SelectorResultValue {
2045 result: SelectorResult::Flat(result),
2046 read_cols,
2047 }
2048 }
2049
2050 fn estimated_size(&self) -> usize {
2052 match &self.result {
2053 SelectorResult::PrimaryKey(batches) => {
2054 batches.iter().map(|batch| batch.memory_size()).sum()
2055 }
2056 SelectorResult::Flat(batches) => batches.iter().map(record_batch_estimated_size).sum(),
2057 }
2058 }
2059}
2060
2061type SstMetaCache = Cache<SstMetaKey, Arc<CompactSstMeta>>;
2063type SstDecodedMetaCache = Cache<SstMetaKey, Arc<CachedSstMeta>>;
2065type VectorCache = Cache<(ConcreteDataType, Value), VectorRef>;
2069type SelectorResultCache = Cache<SelectorResultKey, Arc<SelectorResultValue>>;
2071type RangeResultCache = Cache<RangeScanCacheKey, Arc<RangeScanCacheValue>>;
2073
2074#[cfg(test)]
2075mod tests {
2076 use std::sync::Arc;
2077
2078 use api::v1::SemanticType;
2079 use api::v1::index::{BloomFilterMeta, InvertedIndexMetas};
2080 use datatypes::schema::ColumnSchema;
2081 use datatypes::vectors::Int64Vector;
2082 use parquet::file::page_index::offset_index::{OffsetIndexMetaData, PageLocation};
2083 use puffin::file_metadata::FileMetadata;
2084 use store_api::metadata::{ColumnMetadata, RegionMetadata, RegionMetadataBuilder};
2085 use store_api::storage::ColumnId;
2086
2087 use super::*;
2088 use crate::cache::index::bloom_filter_index::Tag;
2089 use crate::cache::index::result_cache::PredicateKey;
2090 use crate::cache::test_util::{
2091 parquet_meta, sst_parquet_meta, sst_parquet_meta_with_region_metadata,
2092 };
2093 use crate::read::range_cache::{
2094 RangeScanCacheKey, RangeScanCacheValue, ScanRequestFingerprintBuilder,
2095 };
2096 use crate::read::read_columns::ReadColumns;
2097 use crate::sst::parquet::row_selection::RowGroupSelection;
2098
2099 #[tokio::test]
2100 async fn test_disable_cache() {
2101 let cache = CacheManager::default();
2102 assert!(cache.sst_meta_cache.is_none());
2103 assert!(cache.sst_decoded_meta_cache.is_none());
2104 assert!(cache.vector_cache.is_none());
2105 assert!(cache.page_cache.is_none());
2106
2107 let region_id = RegionId::new(1, 1);
2108 let file_id = RegionFileId::new(region_id, FileId::random());
2109 let metadata = parquet_meta();
2110 let mut metrics = MetadataCacheMetrics::default();
2111 cache.put_parquet_meta_data(file_id, metadata, None);
2112 assert!(
2113 cache
2114 .get_sst_meta_data(file_id, &mut metrics, Default::default())
2115 .await
2116 .is_none()
2117 );
2118
2119 let value = Value::Int64(10);
2120 let vector: VectorRef = Arc::new(Int64Vector::from_slice([10, 10, 10, 10]));
2121 cache.put_repeated_vector(value.clone(), vector.clone());
2122 assert!(
2123 cache
2124 .get_repeated_vector(&ConcreteDataType::int64_datatype(), &value)
2125 .is_none()
2126 );
2127
2128 cache.put_page_ranges(
2129 file_id.file_id(),
2130 1,
2131 &[Range { start: 0, end: 5 }],
2132 &[Bytes::from_static(b"abcde")],
2133 );
2134 assert!(
2135 cache
2136 .get_page_ranges(file_id.file_id(), 1, &[Range { start: 0, end: 5 }])
2137 .is_none()
2138 );
2139
2140 assert!(cache.write_cache().is_none());
2141 }
2142
2143 #[test]
2144 fn test_sst_meta_cache_splits_capacity() {
2145 let cache = CacheManager::builder().sst_meta_cache_size(101).build();
2146
2147 assert_eq!(
2148 Some(51),
2149 cache
2150 .sst_meta_cache
2151 .as_ref()
2152 .unwrap()
2153 .policy()
2154 .max_capacity()
2155 );
2156 assert_eq!(
2157 Some(50),
2158 cache
2159 .sst_decoded_meta_cache
2160 .as_ref()
2161 .unwrap()
2162 .policy()
2163 .max_capacity()
2164 );
2165 }
2166
2167 #[tokio::test]
2168 async fn test_parquet_meta_cache() {
2169 let cache = CacheManager::builder().sst_meta_cache_size(2000).build();
2170 let mut metrics = MetadataCacheMetrics::default();
2171 let region_id = RegionId::new(1, 1);
2172 let file_id = RegionFileId::new(region_id, FileId::random());
2173 assert!(
2174 cache
2175 .get_sst_meta_data(file_id, &mut metrics, Default::default())
2176 .await
2177 .is_none()
2178 );
2179 let (metadata, region_metadata) = sst_parquet_meta();
2180 cache.put_parquet_meta_data(file_id, metadata, None);
2181 let cached = cache
2182 .get_sst_meta_data(file_id, &mut metrics, Default::default())
2183 .await
2184 .unwrap();
2185 assert_eq!(region_metadata, cached.region_metadata());
2186 assert!(
2187 cached
2188 .parquet_metadata()
2189 .file_metadata()
2190 .key_value_metadata()
2191 .is_none_or(|key_values| {
2192 key_values
2193 .iter()
2194 .all(|key_value| key_value.key != PARQUET_METADATA_KEY)
2195 })
2196 );
2197 cache.remove_parquet_meta_data(file_id);
2198 assert!(
2199 cache
2200 .get_sst_meta_data(file_id, &mut metrics, Default::default())
2201 .await
2202 .is_none()
2203 );
2204 }
2205
2206 #[tokio::test]
2207 async fn test_parquet_meta_cache_with_provided_region_metadata() {
2208 let cache = CacheManager::builder().sst_meta_cache_size(2000).build();
2209 let mut metrics = MetadataCacheMetrics::default();
2210 let region_id = RegionId::new(1, 1);
2211 let file_id = RegionFileId::new(region_id, FileId::random());
2212 let (metadata, region_metadata) = sst_parquet_meta();
2213
2214 cache.put_parquet_meta_data(file_id, metadata, Some(region_metadata.clone()));
2215
2216 let cached = cache
2217 .get_sst_meta_data(file_id, &mut metrics, Default::default())
2218 .await
2219 .unwrap();
2220 assert!(Arc::ptr_eq(®ion_metadata, &cached.region_metadata()));
2221 }
2222
2223 #[tokio::test]
2224 async fn test_compact_sst_meta_round_trip() {
2225 let cache = CacheManager::builder()
2226 .sst_meta_cache_size(1024 * 1024)
2227 .build();
2228 let region_id = RegionId::new(1, 1);
2229 let file_id = RegionFileId::new(region_id, FileId::random());
2230 let (metadata, expected_region_metadata) = sst_parquet_meta();
2231 let metadata = Arc::unwrap_or_clone(metadata);
2232 let offset_indexes = metadata
2233 .row_groups()
2234 .iter()
2235 .map(|row_group| {
2236 row_group
2237 .columns()
2238 .iter()
2239 .map(|_| OffsetIndexMetaData {
2240 page_locations: vec![PageLocation {
2241 offset: 42,
2242 compressed_page_size: 128,
2243 first_row_index: 0,
2244 }],
2245 unencoded_byte_array_data_bytes: None,
2246 })
2247 .collect()
2248 })
2249 .collect();
2250 let metadata = metadata
2251 .into_builder()
2252 .set_offset_index(Some(offset_indexes))
2253 .build();
2254 let expected_rows = metadata.file_metadata().num_rows();
2255 let prepared = prepare_sst_meta("test.parquet", metadata, None, PageIndexPolicy::Required)
2256 .await
2257 .unwrap();
2258 let SstMetaPreparation::Prepared(prepared) = prepared else {
2259 panic!("valid metadata should produce a compact cache entry");
2260 };
2261
2262 cache.put_prepared_sst_meta(file_id, prepared, false);
2264 let mut metrics = MetadataCacheMetrics::default();
2265 let cached = cache
2266 .get_sst_meta_data(file_id, &mut metrics, PageIndexPolicy::Required)
2267 .await
2268 .unwrap();
2269
2270 assert_eq!(1, metrics.mem_cache_hit);
2271 assert_eq!(0, metrics.cache_miss);
2272 assert_eq!(
2273 expected_rows,
2274 cached.parquet_metadata.file_metadata().num_rows()
2275 );
2276 assert!(cached.parquet_metadata.offset_index().is_some());
2277 assert_eq!(expected_region_metadata, cached.region_metadata());
2278 assert!(
2279 cached
2280 .parquet_metadata
2281 .file_metadata()
2282 .key_value_metadata()
2283 .is_none_or(|key_values| key_values
2284 .iter()
2285 .all(|key_value| key_value.key != PARQUET_METADATA_KEY))
2286 );
2287 }
2288
2289 #[test]
2290 fn test_cache_encoding_failure_preserves_decoded_metadata() {
2291 let (metadata, expected_region_metadata) = sst_parquet_meta();
2292 let expected_rows = metadata.file_metadata().num_rows();
2293 let cache_encoding: Result<(Bytes, usize)> = Err(UnexpectedSnafu {
2294 reason: "injected compact metadata encoding failure",
2295 }
2296 .build());
2297
2298 let preparation = finish_sst_meta_preparation(
2299 "test.parquet",
2300 Arc::unwrap_or_clone(metadata),
2301 None,
2302 PageIndexPolicy::Skip,
2303 cache_encoding,
2304 )
2305 .unwrap();
2306 let SstMetaPreparation::DecodedOnly {
2307 decoded,
2308 encoding_error,
2309 } = preparation
2310 else {
2311 panic!("cache encoding failure should return decoded-only metadata");
2312 };
2313
2314 assert_eq!(
2315 expected_rows,
2316 decoded.parquet_metadata.file_metadata().num_rows()
2317 );
2318 assert_eq!(expected_region_metadata, decoded.region_metadata());
2319 assert!(
2320 encoding_error
2321 .to_string()
2322 .contains("injected compact metadata encoding failure")
2323 );
2324 }
2325
2326 #[tokio::test]
2327 async fn test_parquet_meta_cache_respects_page_index_policy() {
2328 let cache = CacheManager::builder().sst_meta_cache_size(2000).build();
2329 let region_id = RegionId::new(1, 1);
2330 let file_id = RegionFileId::new(region_id, FileId::random());
2331 let (metadata, _) = sst_parquet_meta();
2332
2333 let skip_metadata = Arc::new(
2334 CachedSstMeta::try_new_with_page_index_policy(
2335 "test.parquet",
2336 Arc::unwrap_or_clone(metadata.clone()),
2337 None,
2338 PageIndexPolicy::Skip,
2339 )
2340 .unwrap(),
2341 );
2342 cache.put_sst_meta_data(file_id, skip_metadata);
2343
2344 let mut metrics = MetadataCacheMetrics::default();
2345 assert!(
2346 cache
2347 .get_sst_meta_data(file_id, &mut metrics, PageIndexPolicy::Optional)
2348 .await
2349 .is_none()
2350 );
2351 assert_eq!(1, metrics.cache_miss);
2352
2353 let optional_metadata = Arc::new(
2354 CachedSstMeta::try_new_with_page_index_policy(
2355 "test.parquet",
2356 Arc::unwrap_or_clone(metadata),
2357 None,
2358 PageIndexPolicy::Optional,
2359 )
2360 .unwrap(),
2361 );
2362 cache.put_sst_meta_data(file_id, optional_metadata);
2363
2364 let mut metrics = MetadataCacheMetrics::default();
2365 assert!(
2366 cache
2367 .get_sst_meta_data(file_id, &mut metrics, PageIndexPolicy::Optional)
2368 .await
2369 .is_some()
2370 );
2371 assert_eq!(1, metrics.mem_cache_hit);
2372
2373 let mut metrics = MetadataCacheMetrics::default();
2374 assert!(
2375 cache
2376 .get_sst_meta_data(file_id, &mut metrics, PageIndexPolicy::Skip)
2377 .await
2378 .is_some()
2379 );
2380 assert_eq!(1, metrics.mem_cache_hit);
2381 }
2382
2383 #[test]
2384 fn test_decoded_meta_cache_weight_accounts_for_region_metadata() {
2385 let region_metadata = Arc::new(wide_region_metadata(128));
2386 let json_len = region_metadata.to_json().unwrap().len();
2387 let metadata = sst_parquet_meta_with_region_metadata(region_metadata.clone());
2388 let cached = Arc::new(
2389 CachedSstMeta::try_new("test.parquet", Arc::unwrap_or_clone(metadata)).unwrap(),
2390 );
2391 let key = SstMetaKey(region_metadata.region_id, FileId::random());
2392
2393 assert!(cached.region_metadata_weight > json_len);
2394 assert_eq!(
2395 decoded_meta_cache_weight(&key, &cached) as usize,
2396 key.estimated_size() + cached.parquet_metadata_size + cached.region_metadata_weight
2397 );
2398 assert_eq!(
2399 cached.parquet_metadata_size,
2400 parquet_meta_size(&cached.parquet_metadata)
2401 );
2402 }
2403
2404 #[test]
2405 fn test_decoded_meta_cache_weight_saturates_on_overflow() {
2406 let region_metadata = Arc::new(wide_region_metadata(1));
2407 let metadata = sst_parquet_meta_with_region_metadata(region_metadata.clone());
2408 let mut cached =
2409 CachedSstMeta::try_new("test.parquet", Arc::unwrap_or_clone(metadata)).unwrap();
2410 cached.region_metadata_weight = u32::MAX as usize + 1;
2411 let cached = Arc::new(cached);
2412 let key = SstMetaKey(region_metadata.region_id, FileId::random());
2413
2414 assert_eq!(u32::MAX, decoded_meta_cache_weight(&key, &cached));
2415 }
2416
2417 #[test]
2418 fn test_repeated_vector_cache() {
2419 let cache = CacheManager::builder().vector_cache_size(4096).build();
2420 let value = Value::Int64(10);
2421 assert!(
2422 cache
2423 .get_repeated_vector(&ConcreteDataType::int64_datatype(), &value)
2424 .is_none()
2425 );
2426 let vector: VectorRef = Arc::new(Int64Vector::from_slice([10, 10, 10, 10]));
2427 cache.put_repeated_vector(value.clone(), vector.clone());
2428 let cached = cache
2429 .get_repeated_vector(&ConcreteDataType::int64_datatype(), &value)
2430 .unwrap();
2431 assert_eq!(vector, cached);
2432 }
2433
2434 #[test]
2435 fn test_page_cache() {
2436 let cache = CacheManager::builder().page_cache_size(1000).build();
2437 let file_id = FileId::random();
2438 let uncached = 0..10;
2439 assert_eq!(
2440 vec![0..10],
2441 cache
2442 .get_page_ranges(file_id, 0, std::slice::from_ref(&uncached))
2443 .unwrap()
2444 .missing_ranges
2445 );
2446
2447 let cached = 100..500;
2448 cache.put_page_ranges(
2449 file_id,
2450 0,
2451 std::slice::from_ref(&cached),
2452 &[Bytes::from(vec![7; 400])],
2453 );
2454
2455 let subrange = 200..300;
2456 let lookup = cache
2457 .get_page_ranges(file_id, 0, std::slice::from_ref(&subrange))
2458 .unwrap();
2459 assert!(lookup.is_fully_cached());
2460 assert_eq!(100, lookup.cached_bytes);
2461 assert_eq!(1, lookup.cached_parts.len());
2462 assert_eq!(200..300, lookup.cached_parts[0][0].range);
2463 assert_eq!(100, lookup.cached_parts[0][0].bytes.len());
2464
2465 let overlapping = 400..600;
2466 let lookup = cache
2467 .get_page_ranges(file_id, 0, std::slice::from_ref(&overlapping))
2468 .unwrap();
2469 assert!(!lookup.is_fully_cached());
2470 assert_eq!(100, lookup.cached_bytes);
2471 assert_eq!(vec![500..600], lookup.missing_ranges);
2472 assert_eq!(400..500, lookup.cached_parts[0][0].range);
2473 }
2474
2475 #[test]
2476 fn test_page_cache_detaches_fragment_bytes() {
2477 let cache = PageRangeCache::new(1000);
2478 let file_id = FileId::random();
2479 let backing = Bytes::from(vec![1; 1024]);
2480 let page = backing.slice(512..522);
2481 let page_ptr = page.as_ptr();
2482 let range = 0..10;
2483
2484 cache.insert_ranges(
2485 file_id,
2486 0,
2487 std::slice::from_ref(&range),
2488 std::slice::from_ref(&page),
2489 );
2490
2491 let lookup = cache.lookup(file_id, 0, std::slice::from_ref(&range));
2492 assert!(lookup.is_fully_cached());
2493 assert_eq!(1, lookup.cached_parts[0].len());
2494 assert_eq!(&page[..], &lookup.cached_parts[0][0].bytes[..]);
2495 assert_ne!(page_ptr, lookup.cached_parts[0][0].bytes.as_ptr());
2496 }
2497
2498 #[test]
2499 fn test_page_cache_replaces_fragment() {
2500 let cache = PageRangeCache::new(1000);
2501 let file_id = FileId::random();
2502 let range = 0..10;
2503
2504 cache.insert_ranges(
2505 file_id,
2506 0,
2507 std::slice::from_ref(&range),
2508 &[Bytes::from(vec![1; 10])],
2509 );
2510 cache.insert_ranges(
2511 file_id,
2512 0,
2513 std::slice::from_ref(&range),
2514 &[Bytes::from(vec![2; 10])],
2515 );
2516 cache.cache.run_pending_tasks();
2517 assert_eq!(
2518 vec![PageFragmentKey::new(file_id, 0, &range)],
2519 cache.find_index_candidates(file_id, 0, &range)
2520 );
2521
2522 let lookup = cache.lookup(file_id, 0, std::slice::from_ref(&range));
2523 assert!(lookup.is_fully_cached());
2524 assert_eq!(&vec![2; 10][..], &lookup.cached_parts[0][0].bytes[..]);
2525 }
2526
2527 #[test]
2528 fn test_page_cache_retains_disjoint_inserts_for_same_row_group() {
2529 let cache = PageRangeCache::new(1000);
2530 let file_id = FileId::random();
2531 let range1 = 0..10;
2532 let range2 = 20..30;
2533
2534 cache.insert_ranges(
2535 file_id,
2536 0,
2537 std::slice::from_ref(&range1),
2538 &[Bytes::from(vec![1; 10])],
2539 );
2540 cache.insert_ranges(
2541 file_id,
2542 0,
2543 std::slice::from_ref(&range2),
2544 &[Bytes::from(vec![2; 10])],
2545 );
2546
2547 let lookup = cache.lookup(file_id, 0, &[range1, range2]);
2548 assert!(lookup.is_fully_cached());
2549 assert_eq!(2, lookup.cached_range_count);
2550 assert_eq!(&vec![1; 10][..], &lookup.cached_parts[0][0].bytes[..]);
2551 assert_eq!(&vec![2; 10][..], &lookup.cached_parts[1][0].bytes[..]);
2552 }
2553
2554 #[test]
2555 fn test_page_cache_fragment_eviction() {
2556 let file_id = FileId::random();
2557 let range = 0..10;
2558 let key = PageFragmentKey::new(file_id, 0, &range);
2559 let page = Bytes::from(vec![1; 10]);
2560 let cache = PageRangeCache::new(page_cache_weight(&key, &page) as u64);
2561
2562 cache.insert_ranges(
2563 file_id,
2564 0,
2565 std::slice::from_ref(&range),
2566 &[Bytes::from(vec![1; 10])],
2567 );
2568 assert!(
2569 cache
2570 .lookup(file_id, 0, std::slice::from_ref(&range))
2571 .is_fully_cached()
2572 );
2573
2574 cache.cache.invalidate(&key);
2575 cache.cache.run_pending_tasks();
2576 assert!(cache.find_index_candidates(file_id, 0, &range).is_empty());
2577
2578 let lookup = cache.lookup(file_id, 0, std::slice::from_ref(&range));
2579 assert!(!lookup.is_fully_cached());
2580 assert_eq!(vec![0..10], lookup.missing_ranges);
2581 }
2582
2583 #[test]
2584 fn test_page_cache_rejects_oversized_fragment() {
2585 let cache = PageRangeCache::new(1);
2586 let file_id = FileId::random();
2587 let range = 0..10;
2588
2589 cache.insert_ranges(
2590 file_id,
2591 0,
2592 std::slice::from_ref(&range),
2593 &[Bytes::from(vec![1; 10])],
2594 );
2595 cache.cache.run_pending_tasks();
2596
2597 let lookup = cache.lookup(file_id, 0, std::slice::from_ref(&range));
2598 assert!(!lookup.is_fully_cached());
2599 assert_eq!(vec![0..10], lookup.missing_ranges);
2600 }
2601
2602 #[test]
2603 fn test_selector_result_cache() {
2604 let cache = CacheManager::builder()
2605 .selector_result_cache_size(1000)
2606 .build();
2607 let file_id = FileId::random();
2608 let key = SelectorResultKey {
2609 file_id,
2610 row_group_idx: 0,
2611 selector: TimeSeriesRowSelector::LastRow,
2612 };
2613 assert!(cache.get_selector_result(&key).is_none());
2614 let result = Arc::new(SelectorResultValue::new(
2615 Vec::new(),
2616 ParquetReadColumns::from_deduped(Vec::new()),
2617 ));
2618 cache.put_selector_result(key, result);
2619 assert!(cache.get_selector_result(&key).is_some());
2620 }
2621
2622 #[test]
2623 fn test_prefilter_result_cache() {
2624 let disabled = CacheManager::builder().build();
2625 let file_id = FileId::random();
2626 let key = PrefilterKey::new(
2627 file_id,
2628 0,
2629 None,
2630 1,
2631 SmallVec::from_vec(vec!["tag_0 IN ([a])".to_string()]),
2632 );
2633 let selection = Arc::new(BooleanBuffer::new_set(3));
2634
2635 disabled.put_prefilter_result(key.clone(), selection.clone());
2636 assert!(disabled.get_prefilter_result(&key).is_none());
2637
2638 let cache = Arc::new(
2639 CacheManager::builder()
2640 .prefilter_result_cache_size(1000)
2641 .build(),
2642 );
2643 assert!(cache.get_prefilter_result(&key).is_none());
2644 cache.put_prefilter_result(key.clone(), selection.clone());
2645 assert_eq!(
2646 cache.get_prefilter_result(&key).unwrap().as_ref(),
2647 selection.as_ref()
2648 );
2649
2650 let enable_all = CacheStrategy::EnableAll(cache.clone());
2651 assert!(enable_all.get_prefilter_result(&key).is_some());
2652
2653 let compaction = CacheStrategy::Compaction(cache.clone());
2654 assert!(compaction.get_prefilter_result(&key).is_none());
2655 compaction.put_prefilter_result(key.clone(), selection.clone());
2656 assert!(cache.get_prefilter_result(&key).is_some());
2657
2658 let disabled_strategy = CacheStrategy::Disabled;
2659 assert!(disabled_strategy.get_prefilter_result(&key).is_none());
2660 disabled_strategy.put_prefilter_result(key.clone(), selection);
2661 assert!(cache.get_prefilter_result(&key).is_some());
2662 }
2663
2664 #[test]
2665 fn test_prefilter_key_distinguishes_dimensions() {
2666 let file_id = FileId::random();
2667 let row_selection = RowSelection::from(vec![RowSelector::skip(1), RowSelector::select(3)]);
2668 let other_row_selection =
2669 RowSelection::from(vec![RowSelector::skip(2), RowSelector::select(2)]);
2670 let row_selection = PrefilterKey::row_selection_snapshot(Some(&row_selection));
2671 let other_row_selection = PrefilterKey::row_selection_snapshot(Some(&other_row_selection));
2672 let base = PrefilterKey::new(
2673 file_id,
2674 0,
2675 row_selection.clone(),
2676 1,
2677 SmallVec::from_vec(vec!["tag_0 IN ([a])".to_string()]),
2678 );
2679
2680 assert_ne!(
2681 base,
2682 PrefilterKey::new(
2683 FileId::random(),
2684 0,
2685 row_selection.clone(),
2686 1,
2687 SmallVec::from_vec(vec!["tag_0 IN ([a])".to_string()])
2688 )
2689 );
2690 assert_ne!(
2691 base,
2692 PrefilterKey::new(
2693 file_id,
2694 1,
2695 row_selection.clone(),
2696 1,
2697 SmallVec::from_vec(vec!["tag_0 IN ([a])".to_string()])
2698 )
2699 );
2700 assert_ne!(
2701 base,
2702 PrefilterKey::new(
2703 file_id,
2704 0,
2705 other_row_selection,
2706 1,
2707 SmallVec::from_vec(vec!["tag_0 IN ([a])".to_string()])
2708 )
2709 );
2710 assert_ne!(
2711 base,
2712 PrefilterKey::new(
2713 file_id,
2714 0,
2715 row_selection.clone(),
2716 1,
2717 SmallVec::from_vec(vec!["tag_0 IN ([b])".to_string()])
2718 )
2719 );
2720 assert_ne!(
2721 base,
2722 PrefilterKey::new(
2723 file_id,
2724 0,
2725 row_selection.clone(),
2726 2,
2727 SmallVec::from_vec(vec!["tag_0 IN ([a])".to_string()])
2728 )
2729 );
2730 let pk_group = PrefilterKey::new(
2731 file_id,
2732 0,
2733 row_selection,
2734 1,
2735 SmallVec::from_vec(vec![
2736 "tag_0 IN ([a])".to_string(),
2737 "tag_1 IN ([x])".to_string(),
2738 ]),
2739 );
2740 assert_ne!(base, pk_group);
2741 }
2742
2743 #[test]
2744 fn test_range_result_cache() {
2745 let cache = Arc::new(
2746 CacheManager::builder()
2747 .range_result_cache_size(1024 * 1024)
2748 .build(),
2749 );
2750
2751 let key = RangeScanCacheKey {
2752 region_id: RegionId::new(1, 1),
2753 row_groups: vec![(FileId::random(), 0)],
2754 scan: ScanRequestFingerprintBuilder {
2755 read_columns: ReadColumns::from_deduped_column_ids(std::iter::empty()),
2756 read_column_types: vec![],
2757 filters: vec!["tag_0 = 1".to_string()],
2758 time_filters: vec![],
2759 series_row_selector: None,
2760 append_mode: false,
2761 filter_deleted: true,
2762 merge_mode: crate::region::options::MergeMode::LastRow,
2763 partition_expr_version: 0,
2764 }
2765 .build(),
2766 };
2767 let value = Arc::new(RangeScanCacheValue::new(Vec::new(), 0));
2768
2769 assert!(cache.get_range_result(&key).is_none());
2770 cache.put_range_result(key.clone(), value.clone());
2771 assert!(cache.get_range_result(&key).is_some());
2772
2773 let enable_all = CacheStrategy::EnableAll(cache.clone());
2774 assert!(enable_all.get_range_result(&key).is_some());
2775
2776 let compaction = CacheStrategy::Compaction(cache.clone());
2777 assert!(compaction.get_range_result(&key).is_none());
2778 compaction.put_range_result(key.clone(), value.clone());
2779 assert!(cache.get_range_result(&key).is_some());
2780
2781 let disabled = CacheStrategy::Disabled;
2782 assert!(disabled.get_range_result(&key).is_none());
2783 disabled.put_range_result(key.clone(), value);
2784 assert!(cache.get_range_result(&key).is_some());
2785 }
2786
2787 #[test]
2788 fn test_range_result_cache_size_configures_limiter() {
2789 let cache_size = 3 * 1024_u64;
2790 let cache = CacheManager::builder()
2791 .range_result_cache_size(cache_size)
2792 .build();
2793
2794 assert_eq!(cache.range_result_cache_size(), cache_size as usize);
2795 assert_eq!(
2796 cache.range_result_memory_limiter().permit_bytes(),
2797 RANGE_RESULT_CONCAT_MEMORY_PERMIT.as_bytes() as usize
2798 );
2799 assert_eq!(
2800 cache.range_result_memory_limiter().available_permits(),
2801 (cache_size as usize).div_ceil(RANGE_RESULT_CONCAT_MEMORY_PERMIT.as_bytes() as usize)
2802 );
2803 }
2804
2805 #[tokio::test]
2806 async fn range_result_memory_limiter_rejects_oversized_request() {
2807 let limiter = RangeResultMemoryLimiter::new(2 * 1024, 1024);
2808 assert_eq!(limiter.available_permits(), 2);
2809
2810 let err = limiter.acquire(10 * 1024).await.unwrap_err();
2811 assert!(
2812 err.to_string().contains("exceeds limiter capacity"),
2813 "unexpected error: {err}"
2814 );
2815 assert_eq!(limiter.available_permits(), 2);
2816 }
2817
2818 #[tokio::test]
2819 async fn range_result_memory_limiter_allows_request_up_to_capacity() {
2820 let limiter = RangeResultMemoryLimiter::new(2 * 1024, 1024);
2821 let permit = limiter.acquire(2 * 1024).await.unwrap();
2822 assert_eq!(limiter.available_permits(), 0);
2823 drop(permit);
2824 assert_eq!(limiter.available_permits(), 2);
2825 }
2826
2827 #[tokio::test]
2828 async fn test_evict_puffin_cache_clears_all_entries() {
2829 use std::collections::{BTreeMap, HashMap};
2830
2831 let cache = CacheManager::builder()
2832 .index_metadata_size(128)
2833 .index_content_size(128)
2834 .index_content_page_size(64)
2835 .index_result_cache_size(128)
2836 .puffin_metadata_size(128)
2837 .build();
2838 let cache = Arc::new(cache);
2839
2840 let region_id = RegionId::new(1, 1);
2841 let index_id = RegionIndexId::new(RegionFileId::new(region_id, FileId::random()), 0);
2842 let column_id: ColumnId = 1;
2843
2844 let bloom_cache = cache.bloom_filter_index_cache().unwrap().clone();
2845 let inverted_cache = cache.inverted_index_cache().unwrap().clone();
2846 let result_cache = cache.index_result_cache().unwrap();
2847 let puffin_metadata_cache = cache.puffin_metadata_cache().unwrap().clone();
2848
2849 let bloom_key = (
2850 index_id.file_id(),
2851 index_id.version,
2852 column_id,
2853 Tag::Skipping,
2854 );
2855 bloom_cache.put_metadata(bloom_key, Arc::new(BloomFilterMeta::default()));
2856 inverted_cache.put_metadata(
2857 (index_id.file_id(), index_id.version),
2858 Arc::new(InvertedIndexMetas::default()),
2859 );
2860 let predicate = PredicateKey::new_bloom(Arc::new(BTreeMap::new()));
2861 let selection = Arc::new(RowGroupSelection::default());
2862 result_cache.put(predicate.clone(), index_id.file_id(), selection);
2863 let file_id_str = index_id.to_string();
2864 let metadata = Arc::new(FileMetadata {
2865 blobs: Vec::new(),
2866 properties: HashMap::new(),
2867 });
2868 puffin_metadata_cache.put_metadata(file_id_str.clone(), metadata);
2869
2870 assert!(bloom_cache.get_metadata(bloom_key).is_some());
2871 assert!(
2872 inverted_cache
2873 .get_metadata((index_id.file_id(), index_id.version))
2874 .is_some()
2875 );
2876 assert!(result_cache.get(&predicate, index_id.file_id()).is_some());
2877 assert!(puffin_metadata_cache.get_metadata(&file_id_str).is_some());
2878
2879 cache.evict_puffin_cache(index_id).await;
2880
2881 assert!(bloom_cache.get_metadata(bloom_key).is_none());
2882 assert!(
2883 inverted_cache
2884 .get_metadata((index_id.file_id(), index_id.version))
2885 .is_none()
2886 );
2887 assert!(result_cache.get(&predicate, index_id.file_id()).is_none());
2888 assert!(puffin_metadata_cache.get_metadata(&file_id_str).is_none());
2889
2890 bloom_cache.put_metadata(bloom_key, Arc::new(BloomFilterMeta::default()));
2892 inverted_cache.put_metadata(
2893 (index_id.file_id(), index_id.version),
2894 Arc::new(InvertedIndexMetas::default()),
2895 );
2896 result_cache.put(
2897 predicate.clone(),
2898 index_id.file_id(),
2899 Arc::new(RowGroupSelection::default()),
2900 );
2901 puffin_metadata_cache.put_metadata(
2902 file_id_str.clone(),
2903 Arc::new(FileMetadata {
2904 blobs: Vec::new(),
2905 properties: HashMap::new(),
2906 }),
2907 );
2908
2909 let strategy = CacheStrategy::EnableAll(cache.clone());
2910 strategy.evict_puffin_cache(index_id).await;
2911
2912 assert!(bloom_cache.get_metadata(bloom_key).is_none());
2913 assert!(
2914 inverted_cache
2915 .get_metadata((index_id.file_id(), index_id.version))
2916 .is_none()
2917 );
2918 assert!(result_cache.get(&predicate, index_id.file_id()).is_none());
2919 assert!(puffin_metadata_cache.get_metadata(&file_id_str).is_none());
2920 }
2921
2922 fn wide_region_metadata(column_count: u32) -> RegionMetadata {
2923 let region_id = RegionId::new(1024, 7);
2924 let mut builder = RegionMetadataBuilder::new(region_id);
2925 let mut primary_key = Vec::new();
2926
2927 for column_id in 0..column_count {
2928 let semantic_type = if column_id < 32 {
2929 primary_key.push(column_id);
2930 SemanticType::Tag
2931 } else {
2932 SemanticType::Field
2933 };
2934 let mut column_schema = ColumnSchema::new(
2935 format!("wide_column_{column_id}"),
2936 ConcreteDataType::string_datatype(),
2937 true,
2938 );
2939 column_schema
2940 .mut_metadata()
2941 .insert(format!("cache_key_{column_id}"), "cache_value".repeat(4));
2942 builder.push_column_metadata(ColumnMetadata {
2943 column_schema,
2944 semantic_type,
2945 column_id,
2946 });
2947 }
2948
2949 builder.push_column_metadata(ColumnMetadata {
2950 column_schema: ColumnSchema::new(
2951 "ts",
2952 ConcreteDataType::timestamp_millisecond_datatype(),
2953 false,
2954 ),
2955 semantic_type: SemanticType::Timestamp,
2956 column_id: column_count,
2957 });
2958 builder.primary_key(primary_key);
2959
2960 builder.build().unwrap()
2961 }
2962}