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