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