1use std::clone::Clone;
16use std::collections::VecDeque;
17use std::ops::Deref;
18use std::pin::Pin;
19use std::sync::Arc;
20use std::sync::atomic::{AtomicUsize, Ordering};
21
22use await_tree::{InstrumentAwait, SpanExt};
23use bytes::Bytes;
24use fail::fail_point;
25use foyer::{
26 Cache, CacheBuilder, CacheEntry, EventListener, Hint, HybridCache, HybridCacheBuilder,
27 HybridCacheEntry, HybridCacheProperties,
28};
29use futures::{FutureExt, StreamExt, future};
30use prost::Message;
31use risingwave_hummock_sdk::sstable_info::SstableInfo;
32use risingwave_hummock_sdk::vector_index::{HnswGraphFileInfo, VectorFileInfo};
33use risingwave_hummock_sdk::{
34 HummockHnswGraphFileId, HummockObjectId, HummockRawObjectId, HummockSstableObjectId,
35 HummockVectorFileId, SST_OBJECT_SUFFIX,
36};
37use risingwave_hummock_trace::TracedCachePolicy;
38use risingwave_object_store::object::{
39 ObjectError, ObjectMetadataIter, ObjectResult, ObjectStoreRef, ObjectStreamingUploader,
40};
41use risingwave_pb::hummock::PbHnswGraph;
42use serde::{Deserialize, Serialize};
43use thiserror_ext::AsReport;
44use tokio::time::Instant;
45
46use super::{
47 BatchUploadWriter, Block, BlockMeta, BlockResponse, RecentFilter, Sstable, SstableMeta,
48 SstableWriterOptions,
49};
50use crate::hummock::block_stream::{BlockDataStream, MemoryUsageTracker, PrefetchBlockStream};
51use crate::hummock::none::NoneRecentFilter;
52use crate::hummock::vector::file::{VectorBlock, VectorBlockMeta, VectorFileMeta};
53use crate::hummock::vector::monitor::VectorStoreCacheStats;
54use crate::hummock::{BlockEntry, BlockHolder, HummockError, HummockResult, RecentFilterTrait};
55use crate::monitor::{HummockStateStoreMetrics, StoreLocalStatistic};
56
57macro_rules! impl_vector_index_meta_file {
58 ($($type_name:ident),+) => {
59 pub enum HummockVectorIndexMetaFile {
60 $(
61 $type_name(Pin<Box<$type_name>>),
62 )+
63 }
64
65 $(
66 impl From<$type_name> for HummockVectorIndexMetaFile {
67 fn from(v: $type_name) -> Self {
68 Self::$type_name(Box::pin(v))
69 }
70 }
71
72 unsafe impl Send for VectorMetaFileHolder<$type_name> {}
73
74 impl VectorMetaFileHolder<$type_name> {
75 fn try_from_entry(
76 entry: CacheEntry<HummockRawObjectId, HummockVectorIndexMetaFile>,
77 object_id: HummockRawObjectId
78 ) -> HummockResult<Self> {
79 let HummockVectorIndexMetaFile::$type_name(file_meta) = &*entry else {
80 return Err(HummockError::decode_error(format!(
81 "expect {} for object {}",
82 stringify!($type_name),
83 object_id
84 )));
85 };
86 let ptr = file_meta.as_ref().get_ref() as *const _;
87 Ok(VectorMetaFileHolder {
88 _cache_entry: entry,
89 ptr,
90 })
91 }
92 }
93 )+
94 };
95}
96
97impl_vector_index_meta_file!(VectorFileMeta, PbHnswGraph);
98
99pub struct VectorMetaFileHolder<T> {
100 _cache_entry: CacheEntry<HummockRawObjectId, HummockVectorIndexMetaFile>,
101 ptr: *const T,
102}
103
104impl<T> Deref for VectorMetaFileHolder<T> {
105 type Target = T;
106
107 fn deref(&self) -> &Self::Target {
108 unsafe { &*self.ptr }
110 }
111}
112
113pub type TableHolder = HybridCacheEntry<HummockSstableObjectId, Box<Sstable>>;
114
115pub type VectorBlockHolder = CacheEntry<(HummockVectorFileId, usize), Box<VectorBlock>>;
116
117pub type VectorFileHolder = VectorMetaFileHolder<VectorFileMeta>;
118pub type HnswGraphFileHolder = VectorMetaFileHolder<PbHnswGraph>;
119
120#[derive(Debug, Clone, Copy, PartialEq, PartialOrd, Eq, Ord, Hash, Serialize, Deserialize)]
121pub struct SstableBlockIndex {
122 pub sst_id: HummockSstableObjectId,
123 pub block_idx: u64,
124}
125
126pub struct BlockCacheEventListener {
127 metrics: Arc<HummockStateStoreMetrics>,
128}
129
130impl BlockCacheEventListener {
131 pub fn new(metrics: Arc<HummockStateStoreMetrics>) -> Self {
132 Self { metrics }
133 }
134}
135
136impl EventListener for BlockCacheEventListener {
137 type Key = SstableBlockIndex;
138 type Value = Box<Block>;
139
140 fn on_leave(&self, _reason: foyer::Event, _key: &Self::Key, value: &Self::Value)
141 where
142 Self::Key: foyer::Key,
143 Self::Value: foyer::Value,
144 {
145 self.metrics
146 .block_efficiency_histogram
147 .observe(value.efficiency());
148 }
149}
150
151#[derive(Clone, Copy, Eq, PartialEq)]
153pub enum CachePolicy {
154 Disable,
156 Fill(Hint),
158 NotFill,
160}
161
162impl Default for CachePolicy {
163 fn default() -> Self {
164 CachePolicy::Fill(Hint::Normal)
165 }
166}
167
168impl From<TracedCachePolicy> for CachePolicy {
169 fn from(policy: TracedCachePolicy) -> Self {
170 match policy {
171 TracedCachePolicy::Disable => Self::Disable,
172 TracedCachePolicy::Fill(priority) => Self::Fill(priority.into()),
173 TracedCachePolicy::NotFill => Self::NotFill,
174 }
175 }
176}
177
178impl From<CachePolicy> for TracedCachePolicy {
179 fn from(policy: CachePolicy) -> Self {
180 match policy {
181 CachePolicy::Disable => Self::Disable,
182 CachePolicy::Fill(priority) => Self::Fill(priority.into()),
183 CachePolicy::NotFill => Self::NotFill,
184 }
185 }
186}
187
188pub struct SstableStoreConfig {
189 pub store: ObjectStoreRef,
190 pub path: String,
191
192 pub prefetch_buffer_capacity: usize,
193 pub max_prefetch_block_number: usize,
194 pub recent_filter: Arc<RecentFilter<(HummockSstableObjectId, usize)>>,
195 pub state_store_metrics: Arc<HummockStateStoreMetrics>,
196 pub use_new_object_prefix_strategy: bool,
197 pub skip_bloom_filter_in_serde: bool,
198
199 pub meta_cache: HybridCache<HummockSstableObjectId, Box<Sstable>>,
200 pub block_cache: HybridCache<SstableBlockIndex, Box<Block>>,
201
202 pub vector_meta_cache: Cache<HummockRawObjectId, HummockVectorIndexMetaFile>,
203 pub vector_block_cache: Cache<(HummockVectorFileId, usize), Box<VectorBlock>>,
204}
205
206pub struct SstableStore {
207 path: String,
208 store: ObjectStoreRef,
209
210 meta_cache: HybridCache<HummockSstableObjectId, Box<Sstable>>,
211 block_cache: HybridCache<SstableBlockIndex, Box<Block>>,
212 pub vector_meta_cache: Cache<HummockRawObjectId, HummockVectorIndexMetaFile>,
213 pub vector_block_cache: Cache<(HummockVectorFileId, usize), Box<VectorBlock>>,
214
215 recent_filter: Arc<RecentFilter<(HummockSstableObjectId, usize)>>,
219 prefetch_buffer_usage: Arc<AtomicUsize>,
220 prefetch_buffer_capacity: usize,
221 max_prefetch_block_number: usize,
222 use_new_object_prefix_strategy: bool,
231
232 skip_bloom_filter_in_serde: bool,
236}
237
238impl SstableStore {
239 pub fn new(config: SstableStoreConfig) -> Self {
240 Self {
244 path: config.path,
245 store: config.store,
246
247 meta_cache: config.meta_cache,
248 block_cache: config.block_cache,
249 vector_meta_cache: config.vector_meta_cache,
250 vector_block_cache: config.vector_block_cache,
251
252 recent_filter: config.recent_filter,
253 prefetch_buffer_usage: Arc::new(AtomicUsize::new(0)),
254 prefetch_buffer_capacity: config.prefetch_buffer_capacity,
255 max_prefetch_block_number: config.max_prefetch_block_number,
256 use_new_object_prefix_strategy: config.use_new_object_prefix_strategy,
257 skip_bloom_filter_in_serde: config.skip_bloom_filter_in_serde,
258 }
259 }
260
261 #[expect(clippy::borrowed_box)]
264 pub async fn for_compactor(
265 store: ObjectStoreRef,
266 path: String,
267 block_cache_capacity: usize,
268 meta_cache_capacity: usize,
269 use_new_object_prefix_strategy: bool,
270 ) -> HummockResult<Self> {
271 let meta_cache = HybridCacheBuilder::new()
272 .memory(meta_cache_capacity)
273 .with_shards(1)
274 .with_weighter(|_: &HummockSstableObjectId, value: &Box<Sstable>| {
275 std::mem::size_of::<HummockSstableObjectId>()
276 + value.estimated_meta_cache_memory_weight()
277 })
278 .storage()
279 .build()
280 .await
281 .map_err(HummockError::foyer_error)?;
282
283 let block_cache = HybridCacheBuilder::new()
284 .memory(block_cache_capacity)
285 .with_shards(1)
286 .with_weighter(|_: &SstableBlockIndex, value: &Box<Block>| {
287 std::mem::size_of::<SstableBlockIndex>() + value.estimated_memory_weight()
288 })
289 .storage()
290 .build()
291 .await
292 .map_err(HummockError::foyer_error)?;
293
294 Ok(Self {
295 path,
296 store,
297
298 prefetch_buffer_usage: Arc::new(AtomicUsize::new(0)),
299 prefetch_buffer_capacity: block_cache_capacity,
300 max_prefetch_block_number: 16, recent_filter: Arc::new(NoneRecentFilter::default().into()),
302 use_new_object_prefix_strategy,
303 skip_bloom_filter_in_serde: false,
304
305 meta_cache,
306 block_cache,
307 vector_meta_cache: CacheBuilder::new(1 << 10).build(),
308 vector_block_cache: CacheBuilder::new(1 << 10).build(),
309 })
310 }
311
312 pub async fn delete(&self, object_id: HummockSstableObjectId) -> HummockResult<()> {
313 self.store
314 .delete(self.get_sst_data_path(object_id).as_str())
315 .await?;
316 self.meta_cache.remove(&object_id);
317 Ok(())
319 }
320
321 pub fn delete_cache(&self, object_id: HummockSstableObjectId) -> HummockResult<()> {
322 self.meta_cache.remove(&object_id);
323 Ok(())
324 }
325
326 pub(crate) async fn put_sst_data(
327 &self,
328 object_id: HummockSstableObjectId,
329 data: Bytes,
330 ) -> HummockResult<()> {
331 let data_path = self.get_sst_data_path(object_id);
332 self.store
333 .upload(&data_path, data)
334 .await
335 .map_err(Into::into)
336 }
337
338 pub async fn prefetch_blocks(
342 &self,
343 sst: &Sstable,
344 block_index: usize,
345 end_index: usize,
346 policy: CachePolicy,
347 stats: &mut StoreLocalStatistic,
348 ) -> HummockResult<Box<PrefetchBlockStream>> {
349 let object_id = sst.id;
350 if self.prefetch_buffer_usage.load(Ordering::Acquire) > self.prefetch_buffer_capacity {
351 let block = self.get(sst, block_index, policy, stats).await?;
352 return Ok(Box::new(PrefetchBlockStream::new(
353 VecDeque::from([block]),
354 block_index,
355 None,
356 )));
357 }
358 if policy != CachePolicy::Disable
359 && let Some(entry) = self
360 .block_cache
361 .get(&SstableBlockIndex {
362 sst_id: object_id,
363 block_idx: block_index as _,
364 })
365 .await
366 .map_err(HummockError::foyer_error)?
367 {
368 stats.cache_data_block_total += 1;
369 if entry.source() == foyer::Source::Outer {
370 stats.cache_data_block_miss += 1;
371 }
372 let block = BlockHolder::from_hybrid_cache_entry(entry);
373 return Ok(Box::new(PrefetchBlockStream::new(
374 VecDeque::from([block]),
375 block_index,
376 None,
377 )));
378 }
379 let end_index = std::cmp::min(
380 end_index,
381 block_index.saturating_add(self.max_prefetch_block_number),
382 );
383 let mut end_index = std::cmp::min(end_index, sst.meta.block_metas.len());
384 let start_offset = sst.meta.block_metas[block_index].offset as usize;
385 if policy != CachePolicy::Disable {
386 let mut min_hit_index = end_index;
387 let mut hit_count = 0;
388 for idx in block_index..end_index {
389 if self.block_cache.contains(&SstableBlockIndex {
390 sst_id: object_id,
391 block_idx: idx as _,
392 }) {
393 if min_hit_index > idx && idx > block_index {
394 min_hit_index = idx;
395 }
396 hit_count += 1;
397 }
398 }
399
400 if hit_count * 3 >= (end_index - block_index)
401 || min_hit_index * 2 > block_index + end_index
402 {
403 end_index = min_hit_index;
404 }
405 }
406 stats.cache_data_prefetch_count += 1;
407 stats.cache_data_prefetch_block_count += (end_index - block_index) as u64;
408 let end_offset = start_offset
409 + sst.meta.block_metas[block_index..end_index]
410 .iter()
411 .map(|meta| meta.len as usize)
412 .sum::<usize>();
413 let data_path = self.get_sst_data_path(object_id);
414 let memory_usage = end_offset - start_offset;
415 let tracker = MemoryUsageTracker::new(self.prefetch_buffer_usage.clone(), memory_usage);
416 let span = await_tree::span!("Prefetch SST-{}", object_id).verbose();
417 let store = self.store.clone();
418 let join_handle = tokio::spawn(async move {
419 store
420 .read(&data_path, start_offset..end_offset)
421 .instrument_await(span)
422 .await
423 });
424 let buf = match join_handle.await {
425 Ok(Ok(data)) => data,
426 Ok(Err(e)) => {
427 tracing::error!(
428 "prefetch meet error when read {}..{} from sst-{} ({})",
429 start_offset,
430 end_offset,
431 object_id,
432 sst.meta.estimated_size,
433 );
434 return Err(e.into());
435 }
436 Err(_) => {
437 return Err(HummockError::other("cancel by other thread"));
438 }
439 };
440 let mut offset = 0;
441 let mut blocks = VecDeque::default();
442 for idx in block_index..end_index {
443 let end = offset + sst.meta.block_metas[idx].len as usize;
444 if end > buf.len() {
445 return Err(ObjectError::internal("read unexpected EOF").into());
446 }
447 let block = Block::decode_with_copy(
449 buf.slice(offset..end),
450 sst.meta.block_metas[idx].uncompressed_size as usize,
451 true,
452 )?;
453 let holder = if let CachePolicy::Fill(hint) = policy {
454 let hint = if idx == block_index { hint } else { Hint::Low };
455 let entry = self.block_cache.insert_with_properties(
456 SstableBlockIndex {
457 sst_id: object_id,
458 block_idx: idx as _,
459 },
460 Box::new(block),
461 HybridCacheProperties::default().with_hint(hint),
462 );
463 BlockHolder::from_hybrid_cache_entry(entry)
464 } else {
465 BlockHolder::from_owned_block(Box::new(block))
466 };
467
468 blocks.push_back(holder);
469 offset = end;
470 }
471 Ok(Box::new(PrefetchBlockStream::new(
472 blocks,
473 block_index,
474 Some(tracker),
475 )))
476 }
477
478 pub async fn get_block_response(
479 &self,
480 sst: &Sstable,
481 block_index: usize,
482 policy: CachePolicy,
483 ) -> HummockResult<BlockResponse> {
484 let object_id = sst.id;
485 let (range, uncompressed_capacity) = sst.calculate_block_info(block_index);
486 let store = self.store.clone();
487
488 let file_size = sst.meta.estimated_size;
489 let data_path = Arc::new(self.get_sst_data_path(object_id));
490
491 let disable_cache: fn() -> bool = || {
492 fail_point!("disable_block_cache", |_| true);
493 false
494 };
495
496 let policy = if disable_cache() {
497 CachePolicy::Disable
498 } else {
499 policy
500 };
501
502 let idx = SstableBlockIndex {
503 sst_id: object_id,
504 block_idx: block_index as _,
505 };
506
507 self.recent_filter
508 .extend([(object_id, usize::MAX), (object_id, block_index)]);
509
510 let fetch_block = async move {
512 let block_data = match store
513 .read(&data_path, range.clone())
514 .instrument_await("get_block_response".verbose())
515 .await
516 {
517 Ok(data) => data,
518 Err(e) => {
519 tracing::error!(
520 "get_block_response meet error when read {:?} from sst-{}, total length: {}",
521 range,
522 object_id,
523 file_size
524 );
525 return Err(HummockError::from(e));
526 }
527 };
528 let block = Box::new(Block::decode(block_data, uncompressed_capacity)?);
529 Ok(block)
530 };
531
532 match policy {
533 CachePolicy::Fill(hint) => {
534 let properties = HybridCacheProperties::default().with_hint(hint);
535 let fetch = self.block_cache.get_or_fetch(&idx, || {
536 fetch_block.map(|res| res.map(|block| (block, properties)))
537 });
538 Ok(BlockResponse::Fetch(fetch))
539 }
540 CachePolicy::NotFill => {
541 match self
542 .block_cache
543 .get(&idx)
544 .await
545 .map_err(HummockError::foyer_error)?
546 {
547 Some(entry) => Ok(BlockResponse::Block(BlockHolder::from_hybrid_cache_entry(
548 entry,
549 ))),
550 _ => {
551 let block = fetch_block.await?;
552 Ok(BlockResponse::Block(BlockHolder::from_owned_block(block)))
553 }
554 }
555 }
556 CachePolicy::Disable => {
557 let block = fetch_block.await?;
558 Ok(BlockResponse::Block(BlockHolder::from_owned_block(block)))
559 }
560 }
561 }
562
563 pub async fn get(
564 &self,
565 sst: &Sstable,
566 block_index: usize,
567 policy: CachePolicy,
568 stats: &mut StoreLocalStatistic,
569 ) -> HummockResult<BlockHolder> {
570 let block_response = self.get_block_response(sst, block_index, policy).await?;
571 let block_holder = block_response.wait().await?;
572 stats.cache_data_block_total += 1;
573 if let BlockEntry::HybridCache(entry) = block_holder.entry()
574 && entry.source() == foyer::Source::Outer
575 {
576 stats.cache_data_block_miss += 1;
577 }
578 Ok(block_holder)
579 }
580
581 pub async fn get_vector_file_meta(
582 &self,
583 vector_file: &VectorFileInfo,
584 stats: &mut VectorStoreCacheStats,
585 ) -> HummockResult<VectorFileHolder> {
586 let store = self.store.clone();
587 let path = self.get_object_data_path(HummockObjectId::VectorFile(vector_file.object_id));
588 let meta_offset = vector_file.meta_offset;
589 let entry = self
590 .vector_meta_cache
591 .get_or_fetch(&vector_file.object_id.as_raw(), || async move {
592 let encoded_footer = store.read(&path, meta_offset..).await?;
593 let meta = VectorFileMeta::decode_footer(&encoded_footer)?;
594 Ok::<_, anyhow::Error>(HummockVectorIndexMetaFile::from(meta))
595 })
596 .await?;
597 stats.file_meta_total += 1;
598 if entry.source() == foyer::Source::Outer {
599 stats.file_meta_miss += 1;
600 }
601 VectorFileHolder::try_from_entry(entry, vector_file.object_id.as_raw())
602 }
603
604 pub async fn get_vector_block(
605 &self,
606 vector_file: &VectorFileInfo,
607 block_idx: usize,
608 block_meta: &VectorBlockMeta,
609 stats: &mut VectorStoreCacheStats,
610 ) -> HummockResult<VectorBlockHolder> {
611 let store = self.store.clone();
612 let path = self.get_object_data_path(HummockObjectId::VectorFile(vector_file.object_id));
613 let start_offset = block_meta.offset;
614 let end_offset = start_offset + block_meta.block_size;
615 let entry = self
616 .vector_block_cache
617 .get_or_fetch(&(vector_file.object_id, block_idx), || async move {
618 let encoded_block = store.read(&path, start_offset..end_offset).await?;
619 let block = VectorBlock::decode(&encoded_block)?;
620 Ok::<_, anyhow::Error>(Box::new(block))
621 })
622 .await
623 .map_err(HummockError::foyer_error)?;
624
625 stats.file_block_total += 1;
626 if entry.source() == foyer::Source::Outer {
627 stats.file_block_miss += 1;
628 }
629 Ok(entry)
630 }
631
632 pub fn insert_vector_cache(
633 &self,
634 object_id: HummockVectorFileId,
635 meta: VectorFileMeta,
636 blocks: Vec<VectorBlock>,
637 ) {
638 self.vector_meta_cache
639 .insert(object_id.as_raw(), meta.into());
640 for (idx, block) in blocks.into_iter().enumerate() {
641 self.vector_block_cache
642 .insert((object_id, idx), Box::new(block));
643 }
644 }
645
646 pub fn insert_hnsw_graph_cache(&self, object_id: HummockHnswGraphFileId, graph: PbHnswGraph) {
647 self.vector_meta_cache
648 .insert(object_id.as_raw(), graph.into());
649 }
650
651 pub async fn get_hnsw_graph(
652 &self,
653 graph_file: &HnswGraphFileInfo,
654 stats: &mut VectorStoreCacheStats,
655 ) -> HummockResult<HnswGraphFileHolder> {
656 let store = self.store.clone();
657 let graph_file_path =
658 self.get_object_data_path(HummockObjectId::HnswGraphFile(graph_file.object_id));
659 let entry = self
660 .vector_meta_cache
661 .get_or_fetch(&graph_file.object_id.as_raw(), || async move {
662 let encoded_graph = store.read(&graph_file_path, ..).await?;
663 let graph = PbHnswGraph::decode(encoded_graph.as_ref())?;
664 Ok::<_, anyhow::Error>(HummockVectorIndexMetaFile::from(graph))
665 })
666 .await
667 .map_err(HummockError::foyer_error)?;
668 stats.hnsw_graph_total += 1;
669 if entry.source() == foyer::Source::Outer {
670 stats.hnsw_graph_miss += 1;
671 }
672 HnswGraphFileHolder::try_from_entry(entry, graph_file.object_id.as_raw())
673 }
674
675 pub fn get_sst_data_path(&self, object_id: impl Into<HummockSstableObjectId>) -> String {
676 self.get_object_data_path(HummockObjectId::Sstable(object_id.into()))
677 }
678
679 pub fn get_object_data_path(&self, object_id: HummockObjectId) -> String {
680 let obj_prefix = self.store.get_object_prefix(
681 object_id.as_raw().as_raw_id(),
682 self.use_new_object_prefix_strategy,
683 );
684 risingwave_hummock_sdk::get_object_data_path(&obj_prefix, &self.path, object_id)
685 }
686
687 pub fn get_object_id_from_path(path: &str) -> HummockObjectId {
688 risingwave_hummock_sdk::get_object_id_from_path(path)
689 }
690
691 pub fn store(&self) -> ObjectStoreRef {
692 self.store.clone()
693 }
694
695 #[cfg(any(test, feature = "test"))]
696 pub async fn clear_block_cache(&self) -> HummockResult<()> {
697 self.block_cache
698 .clear()
699 .await
700 .map_err(HummockError::foyer_error)
701 }
702
703 #[cfg(any(test, feature = "test"))]
704 pub async fn clear_meta_cache(&self) -> HummockResult<()> {
705 self.meta_cache
706 .clear()
707 .await
708 .map_err(HummockError::foyer_error)
709 }
710
711 pub async fn sstable_cached(
712 &self,
713 sst_obj_id: HummockSstableObjectId,
714 ) -> HummockResult<Option<HybridCacheEntry<HummockSstableObjectId, Box<Sstable>>>> {
715 self.meta_cache
716 .get(&sst_obj_id)
717 .await
718 .map_err(HummockError::foyer_error)
719 }
720
721 pub async fn sstable(
723 &self,
724 sstable_info_ref: &SstableInfo,
725 stats: &mut StoreLocalStatistic,
726 ) -> HummockResult<TableHolder> {
727 let object_id = sstable_info_ref.object_id;
728 let store = self.store.clone();
729 let meta_path = self.get_sst_data_path(object_id);
730 let stats_ptr = stats.remote_io_time.clone();
731 let range = sstable_info_ref.meta_offset as usize..;
732 let skip_bloom_filter_in_serde = self.skip_bloom_filter_in_serde;
733
734 let fetch = self.meta_cache.get_or_fetch(&object_id, || async move {
735 let now = Instant::now();
736 let buf = store
737 .read(&meta_path, range)
738 .instrument_await("get_meta_response".verbose())
739 .await?;
740 let meta = SstableMeta::decode(&buf[..])?;
741
742 let sst = Sstable::new(object_id, meta, skip_bloom_filter_in_serde);
743 let add = (now.elapsed().as_secs_f64() * 1000.0).ceil();
744 stats_ptr.fetch_add(add as u64, Ordering::Relaxed);
745 Ok::<_, anyhow::Error>(Box::new(sst))
746 });
747
748 stats.cache_meta_block_total += 1;
749 let entry = fetch
750 .instrument_await("fetch_meta".verbose())
751 .await
752 .map_err(HummockError::foyer_error);
753 if let Ok(ref entry) = entry
754 && entry.source() == foyer::Source::Outer
755 {
756 stats.cache_meta_block_miss += 1;
757 }
758 entry
759 }
760
761 pub async fn list_sst_object_metadata_from_object_store(
762 &self,
763 prefix: Option<String>,
764 start_after: Option<String>,
765 limit: Option<usize>,
766 ) -> HummockResult<ObjectMetadataIter> {
767 let list_path = format!("{}/{}", self.path, prefix.unwrap_or("".into()));
768 let raw_iter = self.store.list(&list_path, start_after, limit).await?;
769 let iter = raw_iter.filter(|r| match r {
770 Ok(i) => future::ready(i.key.ends_with(&format!(".{}", SST_OBJECT_SUFFIX))),
771 Err(_) => future::ready(true),
772 });
773 Ok(Box::pin(iter))
774 }
775
776 pub fn create_sst_writer(
777 self: Arc<Self>,
778 object_id: impl Into<HummockSstableObjectId>,
779 options: SstableWriterOptions,
780 ) -> BatchUploadWriter {
781 BatchUploadWriter::new(object_id, self, options)
782 }
783
784 pub fn insert_meta_cache(&self, object_id: HummockSstableObjectId, meta: SstableMeta) {
785 let sst = Sstable::new(object_id, meta, self.skip_bloom_filter_in_serde);
786 self.meta_cache.insert(object_id, Box::new(sst));
787 }
788
789 pub fn insert_block_cache(
790 &self,
791 object_id: HummockSstableObjectId,
792 block_index: u64,
793 block: Box<Block>,
794 ) {
795 self.block_cache.insert(
796 SstableBlockIndex {
797 sst_id: object_id,
798 block_idx: block_index,
799 },
800 block,
801 );
802 }
803
804 pub fn get_prefetch_memory_usage(&self) -> usize {
805 self.prefetch_buffer_usage.load(Ordering::Acquire)
806 }
807
808 pub async fn get_stream_for_blocks(
809 &self,
810 object_id: HummockSstableObjectId,
811 metas: &[BlockMeta],
812 ) -> HummockResult<BlockDataStream> {
813 fail_point!("get_stream_err");
814 let data_path = self.get_sst_data_path(object_id);
815 let store = self.store();
816 let block_meta = &metas[0];
817 let start_pos = block_meta.offset as usize;
818 let end_pos = metas.iter().map(|meta| meta.len as usize).sum::<usize>() + start_pos;
819 let range = start_pos..end_pos;
820 let ret = tokio::spawn(async move { store.streaming_read(&data_path, range).await }).await;
822
823 let reader = match ret {
824 Ok(Ok(reader)) => reader,
825 Ok(Err(e)) => return Err(HummockError::from(e)),
826 Err(e) => {
827 return Err(HummockError::other(format!(
828 "failed to get result, this read request may be canceled: {}",
829 e.as_report()
830 )));
831 }
832 };
833 Ok(BlockDataStream::new(reader, metas))
834 }
835
836 pub fn meta_cache(&self) -> &HybridCache<HummockSstableObjectId, Box<Sstable>> {
837 &self.meta_cache
838 }
839
840 pub fn block_cache(&self) -> &HybridCache<SstableBlockIndex, Box<Block>> {
841 &self.block_cache
842 }
843
844 pub fn recent_filter(&self) -> &Arc<RecentFilter<(HummockSstableObjectId, usize)>> {
845 &self.recent_filter
846 }
847
848 pub async fn create_streaming_uploader(
849 &self,
850 path: &str,
851 ) -> ObjectResult<ObjectStreamingUploader> {
852 self.store.streaming_upload(path).await
853 }
854}
855
856pub type SstableStoreRef = Arc<SstableStore>;
857#[cfg(test)]
858mod tests {
859 use std::ops::Range;
860 use std::sync::Arc;
861
862 use risingwave_hummock_sdk::HummockObjectId;
863 use risingwave_hummock_sdk::sstable_info::SstableInfo;
864
865 use super::{SstableStoreRef, SstableWriterOptions};
866 use crate::hummock::iterator::HummockIterator;
867 use crate::hummock::iterator::test_utils::{iterator_test_key_of, mock_sstable_store};
868 use crate::hummock::sstable::SstableIteratorReadOptions;
869 use crate::hummock::test_utils::{
870 default_builder_opt_for_test, gen_default_test_sstable, gen_test_sstable_data, put_sst,
871 test_key_of,
872 };
873 use crate::hummock::value::HummockValue;
874 use crate::hummock::{CachePolicy, SstableIterator, SstableMeta, SstableStore};
875 use crate::monitor::StoreLocalStatistic;
876
877 const SST_ID: u64 = 1;
878
879 fn get_hummock_value(x: usize) -> HummockValue<Vec<u8>> {
880 HummockValue::put(format!("overlapped_new_{}", x).as_bytes().to_vec())
881 }
882
883 async fn validate_sst(
884 sstable_store: SstableStoreRef,
885 info: &SstableInfo,
886 mut meta: SstableMeta,
887 x_range: Range<usize>,
888 ) {
889 let mut stats = StoreLocalStatistic::default();
890 let holder = sstable_store.sstable(info, &mut stats).await.unwrap();
891 std::mem::take(&mut meta.bloom_filter);
892 assert_eq!(holder.meta, meta);
893 let holder = sstable_store.sstable(info, &mut stats).await.unwrap();
894 assert_eq!(holder.meta, meta);
895 let mut iter = SstableIterator::new(
896 holder,
897 sstable_store,
898 Arc::new(SstableIteratorReadOptions::default()),
899 info,
900 );
901 iter.rewind().await.unwrap();
902 for i in x_range {
903 let key = iter.key();
904 let value = iter.value();
905 assert_eq!(key, iterator_test_key_of(i).to_ref());
906 assert_eq!(value, get_hummock_value(i).as_slice());
907 iter.next().await.unwrap();
908 }
909 }
910
911 #[tokio::test]
912 async fn test_batch_upload() {
913 let sstable_store = mock_sstable_store().await;
914 let x_range = 0..100;
915 let (data, meta) = gen_test_sstable_data(
916 default_builder_opt_for_test(),
917 x_range
918 .clone()
919 .map(|x| (iterator_test_key_of(x), get_hummock_value(x))),
920 )
921 .await;
922 let writer_opts = SstableWriterOptions {
923 capacity_hint: None,
924 tracker: None,
925 policy: CachePolicy::Disable,
926 };
927 let info = put_sst(
928 SST_ID,
929 data.clone(),
930 meta.clone(),
931 sstable_store.clone(),
932 writer_opts,
933 vec![0],
934 )
935 .await
936 .unwrap();
937
938 validate_sst(sstable_store, &info, meta, x_range).await;
939 }
940
941 #[tokio::test]
942 async fn test_streaming_upload() {
943 let sstable_store = mock_sstable_store().await;
945 let x_range = 0..100;
946 let (data, meta) = gen_test_sstable_data(
947 default_builder_opt_for_test(),
948 x_range
949 .clone()
950 .map(|x| (iterator_test_key_of(x), get_hummock_value(x))),
951 )
952 .await;
953 let writer_opts = SstableWriterOptions {
954 capacity_hint: None,
955 tracker: None,
956 policy: CachePolicy::Disable,
957 };
958 let info = put_sst(
959 SST_ID,
960 data.clone(),
961 meta.clone(),
962 sstable_store.clone(),
963 writer_opts,
964 vec![0],
965 )
966 .await
967 .unwrap();
968
969 validate_sst(sstable_store, &info, meta, x_range).await;
970 }
971
972 #[tokio::test]
973 async fn test_empty_prefetch_falls_back_to_block_get() {
974 let mut sstable_store = mock_sstable_store().await;
975 Arc::get_mut(&mut sstable_store)
977 .unwrap()
978 .max_prefetch_block_number = 0;
979 let (sstable, info) =
980 gen_default_test_sstable(default_builder_opt_for_test(), 0, sstable_store.clone())
981 .await;
982 sstable_store.clear_block_cache().await.unwrap();
983 let mut iter = SstableIterator::new(
984 sstable,
985 sstable_store.clone(),
986 Arc::new(SstableIteratorReadOptions {
987 prefetch: true,
988 cache_policy: CachePolicy::Disable,
989 ..Default::default()
990 }),
991 &info,
992 );
993 tokio::time::timeout(std::time::Duration::from_secs(5), iter.rewind())
994 .await
995 .expect("empty prefetch must not cause an unbounded refill loop")
996 .unwrap();
997 assert_eq!(iter.key(), test_key_of(0).to_ref());
998 let mut stats = StoreLocalStatistic::default();
999 iter.collect_local_statistic(&mut stats);
1000 assert_eq!(stats.cache_data_prefetch_count, 1);
1001 assert_eq!(stats.cache_data_block_total, 1);
1002 assert_eq!(sstable_store.get_prefetch_memory_usage(), 0);
1003 }
1004
1005 #[tokio::test]
1006 async fn test_prefetch_failure_falls_back_and_releases_tracker() {
1007 for truncate in [true, false] {
1008 let sstable_store = mock_sstable_store().await;
1009 let (sstable, info) =
1010 gen_default_test_sstable(default_builder_opt_for_test(), 0, sstable_store.clone())
1011 .await;
1012 assert!(sstable.meta.block_metas.len() > 1);
1013 sstable_store.clear_block_cache().await.unwrap();
1014 let path = sstable_store.get_sst_data_path(sstable.id);
1015 let mut data = sstable_store.store.read(&path, ..).await.unwrap().to_vec();
1016 let second_offset = sstable.meta.block_metas[1].offset as usize;
1017 if truncate {
1018 data.truncate(second_offset);
1020 } else {
1021 data[second_offset] ^= 1;
1023 }
1024 sstable_store
1025 .store
1026 .upload(&path, data.into())
1027 .await
1028 .unwrap();
1029 let mut iter = SstableIterator::new(
1030 sstable.clone(),
1031 sstable_store.clone(),
1032 Arc::new(SstableIteratorReadOptions {
1033 prefetch: true,
1034 cache_policy: CachePolicy::Disable,
1035 ..Default::default()
1036 }),
1037 &info,
1038 );
1039 iter.rewind().await.unwrap();
1040 assert_eq!(iter.key(), test_key_of(0).to_ref());
1041 assert_eq!(sstable_store.get_prefetch_memory_usage(), 0);
1042 let mut stats = StoreLocalStatistic::default();
1043 iter.collect_local_statistic(&mut stats);
1044 assert_eq!(stats.cache_data_prefetch_count, 1);
1045 assert_eq!(stats.cache_data_block_total, 1);
1046 let second_key = risingwave_hummock_sdk::key::FullKey::decode(
1048 &sstable.meta.block_metas[1].smallest_key,
1049 );
1050 assert!(iter.seek(second_key).await.is_err());
1051 assert_eq!(sstable_store.get_prefetch_memory_usage(), 0);
1052 }
1053 }
1054
1055 #[tokio::test]
1056 async fn test_prefetch_releases_old_budget_before_refill() {
1057 let mut sstable_store = mock_sstable_store().await;
1058 Arc::get_mut(&mut sstable_store)
1061 .unwrap()
1062 .prefetch_buffer_capacity = 0;
1063 let (sstable, info) =
1064 gen_default_test_sstable(default_builder_opt_for_test(), 0, sstable_store.clone())
1065 .await;
1066 sstable_store.clear_block_cache().await.unwrap();
1067 let next_batch = sstable_store.max_prefetch_block_number;
1068 assert!(sstable.meta.block_metas.len() > next_batch + 1);
1069 let mut iter = SstableIterator::new(
1070 sstable.clone(),
1071 sstable_store.clone(),
1072 Arc::new(SstableIteratorReadOptions {
1073 prefetch: true,
1074 cache_policy: CachePolicy::Disable,
1075 ..Default::default()
1076 }),
1077 &info,
1078 );
1079 for idx in [0, next_batch] {
1080 let key = risingwave_hummock_sdk::key::FullKey::decode(
1081 &sstable.meta.block_metas[idx].smallest_key,
1082 );
1083 iter.seek(key).await.unwrap();
1084 assert_eq!(iter.key(), key);
1085 assert!(sstable_store.get_prefetch_memory_usage() > 0);
1086 }
1087 let mut stats = StoreLocalStatistic::default();
1088 iter.collect_local_statistic(&mut stats);
1089 assert_eq!(stats.cache_data_prefetch_count, 2);
1090 assert_eq!(stats.cache_data_block_total, 0);
1091 drop(iter);
1092 assert_eq!(sstable_store.get_prefetch_memory_usage(), 0);
1093 }
1094
1095 #[tokio::test]
1096 async fn test_prefetch_producer_large_limit() {
1097 let mut sstable_store = mock_sstable_store().await;
1098 Arc::get_mut(&mut sstable_store)
1099 .unwrap()
1100 .max_prefetch_block_number = usize::MAX;
1101 let (sstable, _) =
1102 gen_default_test_sstable(default_builder_opt_for_test(), 0, sstable_store.clone())
1103 .await;
1104 sstable_store.clear_block_cache().await.unwrap();
1105 let mut stats = StoreLocalStatistic::default();
1106 let mut stream = sstable_store
1107 .prefetch_blocks(&sstable, 1, 3, CachePolicy::Disable, &mut stats)
1108 .await
1109 .unwrap();
1110 assert_eq!(stats.cache_data_prefetch_block_count, 2);
1111 for idx in [1, 2] {
1112 assert!(matches!(
1113 stream.take_block(idx),
1114 crate::hummock::block_stream::PrefetchLookup::Hit(_)
1115 ));
1116 }
1117 assert!(matches!(
1118 stream.take_block(3),
1119 crate::hummock::block_stream::PrefetchLookup::Exhausted
1120 ));
1121 drop(stream);
1122 assert_eq!(sstable_store.get_prefetch_memory_usage(), 0);
1123 }
1124
1125 #[tokio::test]
1126 async fn test_prefetch_producer_respects_disabled_cache() {
1127 let sstable_store = mock_sstable_store().await;
1128 let (sstable, _) =
1129 gen_default_test_sstable(default_builder_opt_for_test(), 0, sstable_store.clone())
1130 .await;
1131 sstable_store.clear_block_cache().await.unwrap();
1132 sstable_store
1134 .get(
1135 &sstable,
1136 1,
1137 CachePolicy::Fill(foyer::Hint::Normal),
1138 &mut StoreLocalStatistic::default(),
1139 )
1140 .await
1141 .unwrap();
1142 let mut stats = StoreLocalStatistic::default();
1143 let stream = sstable_store
1144 .prefetch_blocks(&sstable, 0, 3, CachePolicy::Disable, &mut stats)
1145 .await
1146 .unwrap();
1147 assert_eq!(stats.cache_data_prefetch_block_count, 3);
1148 drop(stream);
1149 assert!(
1150 !sstable_store
1151 .block_cache
1152 .contains(&super::SstableBlockIndex {
1153 sst_id: sstable.id,
1154 block_idx: 0,
1155 })
1156 );
1157
1158 sstable_store
1160 .get(
1161 &sstable,
1162 0,
1163 CachePolicy::Fill(foyer::Hint::Normal),
1164 &mut StoreLocalStatistic::default(),
1165 )
1166 .await
1167 .unwrap();
1168 sstable_store.delete(sstable.id).await.unwrap();
1169 assert!(
1170 sstable_store
1171 .prefetch_blocks(
1172 &sstable,
1173 0,
1174 3,
1175 CachePolicy::NotFill,
1176 &mut StoreLocalStatistic::default()
1177 )
1178 .await
1179 .is_ok()
1180 );
1181 assert!(
1182 sstable_store
1183 .prefetch_blocks(
1184 &sstable,
1185 0,
1186 3,
1187 CachePolicy::Disable,
1188 &mut StoreLocalStatistic::default()
1189 )
1190 .await
1191 .is_err()
1192 );
1193 assert_eq!(sstable_store.get_prefetch_memory_usage(), 0);
1194 }
1195
1196 #[tokio::test]
1197 async fn test_basic() {
1198 let sstable_store = mock_sstable_store().await;
1199 let object_id = 123;
1200 let data_path = sstable_store.get_sst_data_path(object_id);
1201 assert_eq!(data_path, "test/123.data");
1202 assert_eq!(
1203 SstableStore::get_object_id_from_path(&data_path),
1204 HummockObjectId::Sstable(object_id.into())
1205 );
1206 }
1207
1208 #[tokio::test]
1209 async fn test_clear_file_cache_and_refetch() {
1210 let sstable_store = mock_sstable_store().await;
1211 let x_range = 0..10;
1212 let (data, meta) = gen_test_sstable_data(
1213 default_builder_opt_for_test(),
1214 x_range.map(|x| (iterator_test_key_of(x), get_hummock_value(x))),
1215 )
1216 .await;
1217 let info = put_sst(
1218 SST_ID,
1219 data,
1220 meta,
1221 sstable_store.clone(),
1222 SstableWriterOptions {
1223 capacity_hint: None,
1224 tracker: None,
1225 policy: CachePolicy::Fill(foyer::Hint::Normal),
1226 },
1227 vec![0],
1228 )
1229 .await
1230 .unwrap();
1231
1232 let mut stats = StoreLocalStatistic::default();
1233 assert!(
1234 sstable_store
1235 .sstable_cached(info.object_id)
1236 .await
1237 .unwrap()
1238 .is_some()
1239 );
1240 let sstable = sstable_store.sstable(&info, &mut stats).await.unwrap();
1241 sstable_store
1242 .get(
1243 &sstable,
1244 0,
1245 CachePolicy::Fill(foyer::Hint::Normal),
1246 &mut stats,
1247 )
1248 .await
1249 .unwrap();
1250 assert!(
1251 sstable_store
1252 .block_cache()
1253 .get(&super::SstableBlockIndex {
1254 sst_id: info.object_id,
1255 block_idx: 0,
1256 })
1257 .await
1258 .unwrap()
1259 .is_some()
1260 );
1261
1262 sstable_store.clear_meta_cache().await.unwrap();
1263 assert!(
1264 sstable_store
1265 .sstable_cached(info.object_id)
1266 .await
1267 .unwrap()
1268 .is_none()
1269 );
1270 sstable_store.sstable(&info, &mut stats).await.unwrap();
1271 assert!(
1272 sstable_store
1273 .sstable_cached(info.object_id)
1274 .await
1275 .unwrap()
1276 .is_some()
1277 );
1278
1279 sstable_store.clear_block_cache().await.unwrap();
1280 assert!(
1281 sstable_store
1282 .block_cache()
1283 .get(&super::SstableBlockIndex {
1284 sst_id: info.object_id,
1285 block_idx: 0,
1286 })
1287 .await
1288 .unwrap()
1289 .is_none()
1290 );
1291 sstable_store
1292 .get(
1293 &sstable,
1294 0,
1295 CachePolicy::Fill(foyer::Hint::Normal),
1296 &mut stats,
1297 )
1298 .await
1299 .unwrap();
1300 assert!(
1301 sstable_store
1302 .block_cache()
1303 .get(&super::SstableBlockIndex {
1304 sst_id: info.object_id,
1305 block_idx: 0,
1306 })
1307 .await
1308 .unwrap()
1309 .is_some()
1310 );
1311 }
1312}