1use std::cmp::Ordering;
16use std::collections::HashMap;
17use std::collections::vec_deque::VecDeque;
18use std::ops::Bound::{self};
19use std::sync::Arc;
20use std::time::Instant;
21
22use bytes::Bytes;
23use futures::future::try_join_all;
24use itertools::Itertools;
25use parking_lot::RwLock;
26use risingwave_common::array::VectorRef;
27use risingwave_common::bitmap::Bitmap;
28use risingwave_common::catalog::{TableId, TableOption};
29use risingwave_common::hash::VirtualNode;
30use risingwave_common::util::epoch::MAX_SPILL_TIMES;
31use risingwave_hummock_sdk::key::{
32 FullKey, TableKey, TableKeyRange, UserKey, bound_table_key_range,
33};
34use risingwave_hummock_sdk::key_range::KeyRangeCommon;
35use risingwave_hummock_sdk::sstable_info::SstableInfo;
36use risingwave_hummock_sdk::table_watermark::{
37 TableWatermarksIndex, VnodeWatermark, WatermarkDirection, WatermarkSerdeType,
38};
39use risingwave_hummock_sdk::vector_index::VectorIndexImpl;
40use risingwave_hummock_sdk::{EpochWithGap, HummockEpoch, LocalSstableInfo};
41use risingwave_pb::hummock::LevelType;
42use sync_point::sync_point;
43use tracing::warn;
44
45use crate::error::StorageResult;
46use crate::hummock::event_handler::LocalInstanceId;
47use crate::hummock::iterator::change_log::ChangeLogIterator;
48use crate::hummock::iterator::{
49 BackwardUserIterator, HummockIterator, IteratorFactory, MergeIterator, UserIterator,
50};
51use crate::hummock::local_version::pinned_version::PinnedVersion;
52use crate::hummock::sstable::{SstableIteratorReadOptions, SstableIteratorType};
53use crate::hummock::sstable_store::SstableStoreRef;
54use crate::hummock::table_change_log_manager::TableChangeLogManager;
55use crate::hummock::utils::{
56 MemoryTracker, filter_single_sst, prune_nonoverlapping_ssts, prune_overlapping_ssts,
57 range_overlap, search_sst_idx,
58};
59use crate::hummock::vector::file::{FileVectorStore, FileVectorStoreCtx};
60use crate::hummock::vector::monitor::{VectorStoreCacheStats, report_hnsw_stat};
61use crate::hummock::{
62 BackwardIteratorFactory, ForwardIteratorFactory, HummockError, HummockResult,
63 HummockStorageIterator, HummockStorageIteratorInner, HummockStorageRevIteratorInner,
64 ReadVersionTuple, Sstable, SstableIterator, get_from_batch, get_from_sstable_info,
65 hit_sstable_filter,
66};
67use crate::mem_table::{
68 ImmId, ImmutableMemtable, MemTableHummockIterator, MemTableHummockRevIterator,
69};
70use crate::monitor::{
71 GetLocalMetricsGuard, HummockStateStoreMetrics, IterLocalMetricsGuard, StoreLocalStatistic,
72};
73use crate::store::{
74 OnNearestItemFn, ReadLogOptions, ReadOptions, VectorNearestOptions, gen_min_epoch,
75};
76use crate::vector::hnsw::nearest;
77use crate::vector::{MeasureDistanceBuilder, NearestBuilder};
78
79pub type CommittedVersion = PinnedVersion;
80
81#[derive(Clone, Debug, PartialEq)]
87pub struct StagingSstableInfo {
88 sstable_infos: Vec<LocalSstableInfo>,
90 old_value_sstable_infos: Vec<LocalSstableInfo>,
91 epochs: Vec<HummockEpoch>,
94 imm_ids: HashMap<LocalInstanceId, Vec<ImmId>>,
96 imm_size: usize,
97}
98
99impl StagingSstableInfo {
100 pub fn new(
101 sstable_infos: Vec<LocalSstableInfo>,
102 old_value_sstable_infos: Vec<LocalSstableInfo>,
103 epochs: Vec<HummockEpoch>,
104 imm_ids: HashMap<LocalInstanceId, Vec<ImmId>>,
105 imm_size: usize,
106 ) -> Self {
107 assert!(epochs.is_sorted_by(|epoch1, epoch2| epoch2 <= epoch1));
109 Self {
110 sstable_infos,
111 old_value_sstable_infos,
112 epochs,
113 imm_ids,
114 imm_size,
115 }
116 }
117
118 pub fn sstable_infos(&self) -> &Vec<LocalSstableInfo> {
119 &self.sstable_infos
120 }
121
122 pub fn old_value_sstable_infos(&self) -> &Vec<LocalSstableInfo> {
123 &self.old_value_sstable_infos
124 }
125
126 pub fn imm_size(&self) -> usize {
127 self.imm_size
128 }
129
130 pub fn epochs(&self) -> &Vec<HummockEpoch> {
131 &self.epochs
132 }
133
134 pub fn imm_ids(&self) -> &HashMap<LocalInstanceId, Vec<ImmId>> {
135 &self.imm_ids
136 }
137}
138
139pub enum VersionUpdate {
140 Sst(Arc<StagingSstableInfo>),
141 CommittedSnapshot(CommittedVersion),
142 NewTableWatermark {
143 direction: WatermarkDirection,
144 epoch: HummockEpoch,
145 vnode_watermarks: Vec<VnodeWatermark>,
146 watermark_type: WatermarkSerdeType,
147 },
148}
149
150pub struct StagingVersion {
151 pending_imm_size: usize,
152 pub pending_imms: Vec<(ImmutableMemtable, MemoryTracker)>,
159 pub uploading_imms: VecDeque<ImmutableMemtable>,
165
166 pub sst: VecDeque<Arc<StagingSstableInfo>>,
168}
169
170impl StagingVersion {
171 pub fn prune_overlap<'a>(
174 &'a self,
175 max_epoch_inclusive: HummockEpoch,
176 table_id: TableId,
177 table_key_range: &'a TableKeyRange,
178 ) -> (
179 impl Iterator<Item = &'a ImmutableMemtable> + 'a,
180 impl Iterator<Item = &'a SstableInfo> + 'a,
181 ) {
182 let (left, right) = table_key_range;
183 let left = left.as_ref().map(|key| TableKey(key.0.as_ref()));
184 let right = right.as_ref().map(|key| TableKey(key.0.as_ref()));
185 let overlapped_imms = self
186 .pending_imms
187 .iter()
188 .map(|(imm, _)| imm)
189 .rev() .chain(self.uploading_imms.iter())
191 .filter(move |imm| {
192 imm.epoch() <= max_epoch_inclusive
194 && imm.table_id == table_id
195 && range_overlap(
196 &(left, right),
197 &imm.start_table_key(),
198 Bound::Included(&imm.end_table_key()),
199 )
200 });
201
202 let overlapped_ssts = self
204 .sst
205 .iter()
206 .filter(move |staging_sst| {
207 let sst_max_epoch = *staging_sst.epochs.last().expect("epochs not empty");
208 sst_max_epoch <= max_epoch_inclusive
209 })
210 .flat_map(move |staging_sst| {
211 staging_sst
214 .sstable_infos
215 .iter()
216 .map(|sstable| &sstable.sst_info)
217 .filter(move |sstable: &&SstableInfo| {
218 filter_single_sst(sstable, table_id, table_key_range)
219 })
220 });
221 (overlapped_imms, overlapped_ssts)
222 }
223
224 pub fn is_empty(&self) -> bool {
225 self.pending_imms.is_empty() && self.uploading_imms.is_empty() && self.sst.is_empty()
226 }
227}
228
229pub struct HummockReadVersion {
231 table_id: TableId,
232 instance_id: LocalInstanceId,
233
234 is_initialized: bool,
235
236 staging: StagingVersion,
238
239 committed: CommittedVersion,
241
242 is_replicated: bool,
247
248 table_watermarks: Option<TableWatermarksIndex>,
249
250 vnodes: Arc<Bitmap>,
253}
254
255impl HummockReadVersion {
256 pub fn new_with_replication_option(
257 table_id: TableId,
258 instance_id: LocalInstanceId,
259 committed_version: CommittedVersion,
260 is_replicated: bool,
261 vnodes: Arc<Bitmap>,
262 ) -> Self {
263 assert!(committed_version.is_valid());
267 Self {
268 table_id,
269 instance_id,
270 table_watermarks: {
271 match committed_version.table_watermarks.get(&table_id) {
272 Some(table_watermarks) => Some(TableWatermarksIndex::new_committed(
273 table_watermarks.clone(),
274 committed_version
275 .state_table_info
276 .info()
277 .get(&table_id)
278 .expect("should exist")
279 .committed_epoch,
280 table_watermarks.watermark_type,
281 )),
282 None => None,
283 }
284 },
285 staging: StagingVersion {
286 pending_imm_size: 0,
287 pending_imms: Vec::default(),
288 uploading_imms: VecDeque::default(),
289 sst: VecDeque::default(),
290 },
291
292 committed: committed_version,
293
294 is_replicated,
295 vnodes,
296 is_initialized: false,
297 }
298 }
299
300 pub fn new(
301 table_id: TableId,
302 instance_id: LocalInstanceId,
303 committed_version: CommittedVersion,
304 vnodes: Arc<Bitmap>,
305 ) -> Self {
306 Self::new_with_replication_option(table_id, instance_id, committed_version, false, vnodes)
307 }
308
309 pub fn table_id(&self) -> TableId {
310 self.table_id
311 }
312
313 pub fn init(&mut self) {
314 assert!(!self.is_initialized);
315 self.is_initialized = true;
316 }
317
318 pub fn add_pending_imm(&mut self, imm: ImmutableMemtable, tracker: MemoryTracker) {
319 assert!(self.is_initialized);
320 assert!(!self.is_replicated);
321 if let Some(item) = self
322 .staging
323 .pending_imms
324 .last()
325 .map(|(imm, _)| imm)
326 .or_else(|| self.staging.uploading_imms.front())
327 {
328 assert!(item.batch_id() < imm.batch_id());
330 }
331
332 self.staging.pending_imm_size += imm.size();
333 self.staging.pending_imms.push((imm, tracker));
334 }
335
336 pub fn add_replicated_imm(&mut self, imm: ImmutableMemtable) {
337 assert!(self.is_initialized);
338 assert!(self.is_replicated);
339 assert!(self.staging.pending_imms.is_empty());
340 if let Some(item) = self.staging.uploading_imms.front() {
341 assert!(item.batch_id() < imm.batch_id());
343 }
344 self.staging.uploading_imms.push_front(imm);
345 }
346
347 pub fn pending_imm_size(&self) -> usize {
348 self.staging.pending_imm_size
349 }
350
351 pub fn start_upload_pending_imms(&mut self) -> Vec<(ImmutableMemtable, MemoryTracker)> {
352 assert!(self.is_initialized);
353 assert!(!self.is_replicated);
354 let pending_imms = std::mem::take(&mut self.staging.pending_imms);
355 for (imm, _) in &pending_imms {
356 self.staging.uploading_imms.push_front(imm.clone());
357 }
358 self.staging.pending_imm_size = 0;
359 pending_imms
360 }
361
362 pub fn update(&mut self, info: VersionUpdate) {
372 match info {
373 VersionUpdate::Sst(staging_sst_ref) => {
374 {
375 assert!(!self.is_replicated);
376 let Some(imms) = staging_sst_ref.imm_ids.get(&self.instance_id) else {
377 warn!(
378 instance_id = self.instance_id,
379 "no related imm in sst input"
380 );
381 return;
382 };
383
384 for imm_id in imms.iter().rev() {
386 let check_err = match self.staging.uploading_imms.pop_back() {
387 None => Some("empty".to_owned()),
388 Some(prev_imm_id) => {
389 if prev_imm_id.batch_id() == *imm_id {
390 None
391 } else {
392 Some(format!(
393 "miss match id {} {}",
394 prev_imm_id.batch_id(),
395 *imm_id
396 ))
397 }
398 }
399 };
400 assert!(
401 check_err.is_none(),
402 "should be valid staging_sst.size {},
403 staging_sst.imm_ids {:?},
404 staging_sst.epochs {:?},
405 local_pending_imm_ids {:?},
406 local_uploading_imm_ids {:?},
407 instance_id {}
408 check_err {:?}",
409 staging_sst_ref.imm_size,
410 staging_sst_ref.imm_ids,
411 staging_sst_ref.epochs,
412 self.staging
413 .pending_imms
414 .iter()
415 .map(|(imm, _)| imm.batch_id())
416 .collect_vec(),
417 self.staging
418 .uploading_imms
419 .iter()
420 .map(|imm| imm.batch_id())
421 .collect_vec(),
422 self.instance_id,
423 check_err
424 );
425 }
426
427 self.staging.sst.push_front(staging_sst_ref);
428 }
429 }
430
431 VersionUpdate::CommittedSnapshot(committed_version) => {
432 if let Some(info) = committed_version
433 .state_table_info
434 .info()
435 .get(&self.table_id)
436 {
437 let committed_epoch = info.committed_epoch;
438 if self.is_replicated {
439 self.staging
440 .uploading_imms
441 .retain(|imm| imm.epoch() > committed_epoch);
442 self.staging
443 .pending_imms
444 .retain(|(imm, _)| imm.epoch() > committed_epoch);
445 } else {
446 self.staging
447 .pending_imms
448 .iter()
449 .map(|(imm, _)| imm)
450 .chain(self.staging.uploading_imms.iter())
451 .for_each(|imm| {
452 assert!(
453 imm.epoch() > committed_epoch,
454 "imm of table {} min_epoch {} should be greater than committed_epoch {}",
455 imm.table_id,
456 imm.epoch(),
457 committed_epoch
458 )
459 });
460 }
461
462 self.staging.sst.retain(|sst| {
463 sst.epochs.first().expect("epochs not empty") > &committed_epoch
464 });
465
466 assert!(self.staging.sst.iter().all(|sst| {
468 sst.epochs.last().expect("epochs not empty") > &committed_epoch
469 }));
470
471 if let Some(committed_watermarks) =
472 committed_version.table_watermarks.get(&self.table_id)
473 {
474 if let Some(watermark_index) = &mut self.table_watermarks {
475 watermark_index.apply_committed_watermarks(
476 committed_watermarks.clone(),
477 committed_epoch,
478 );
479 } else {
480 self.table_watermarks = Some(TableWatermarksIndex::new_committed(
481 committed_watermarks.clone(),
482 committed_epoch,
483 committed_watermarks.watermark_type,
484 ));
485 }
486 }
487 }
488
489 self.committed = committed_version;
490 }
491 VersionUpdate::NewTableWatermark {
492 direction,
493 epoch,
494 vnode_watermarks,
495 watermark_type,
496 } => {
497 if let Some(watermark_index) = &mut self.table_watermarks {
498 watermark_index.add_epoch_watermark(
499 epoch,
500 Arc::from(vnode_watermarks),
501 direction,
502 );
503 } else {
504 self.table_watermarks = Some(TableWatermarksIndex::new(
505 direction,
506 epoch,
507 vnode_watermarks,
508 self.committed.table_committed_epoch(self.table_id),
509 watermark_type,
510 ));
511 }
512 }
513 }
514 }
515
516 pub fn staging(&self) -> &StagingVersion {
517 &self.staging
518 }
519
520 pub fn committed(&self) -> &CommittedVersion {
521 &self.committed
522 }
523
524 pub fn filter_regress_watermarks(&self, watermarks: &mut Vec<VnodeWatermark>) {
529 if let Some(watermark_index) = &self.table_watermarks {
530 watermark_index.filter_regress_watermarks(watermarks)
531 }
532 }
533
534 pub fn latest_watermark(&self, vnode: VirtualNode) -> Option<Bytes> {
535 self.table_watermarks
536 .as_ref()
537 .and_then(|watermark_index| watermark_index.latest_watermark(vnode))
538 }
539
540 pub fn is_initialized(&self) -> bool {
541 self.is_initialized
542 }
543
544 pub fn is_replicated(&self) -> bool {
545 self.is_replicated
546 }
547
548 pub fn update_vnode_bitmap(&mut self, vnodes: Arc<Bitmap>) -> Arc<Bitmap> {
549 std::mem::replace(&mut self.vnodes, vnodes)
550 }
551
552 pub fn contains(&self, vnode: VirtualNode) -> bool {
553 self.vnodes.is_set(vnode.to_index())
554 }
555
556 pub fn vnodes(&self) -> Arc<Bitmap> {
557 self.vnodes.clone()
558 }
559}
560
561pub fn read_filter_for_version(
562 epoch: HummockEpoch,
563 table_id: TableId,
564 mut table_key_range: TableKeyRange,
565 read_version: &RwLock<HummockReadVersion>,
566) -> StorageResult<(TableKeyRange, ReadVersionTuple)> {
567 let read_version_guard = read_version.read();
568
569 let committed_version = read_version_guard.committed().clone();
570
571 if let Some(watermark) = read_version_guard.table_watermarks.as_ref() {
572 watermark.rewrite_range_with_table_watermark(epoch, &mut table_key_range)
573 }
574
575 let (imm_iter, sst_iter) =
576 read_version_guard
577 .staging()
578 .prune_overlap(epoch, table_id, &table_key_range);
579
580 let imms = imm_iter.cloned().collect();
581 let ssts = sst_iter.cloned().collect();
582
583 Ok((table_key_range, (imms, ssts, committed_version)))
584}
585
586#[derive(Clone)]
587pub struct HummockVersionReader {
588 sstable_store: SstableStoreRef,
589
590 state_store_metrics: Arc<HummockStateStoreMetrics>,
592 preload_retry_times: usize,
593}
594
595impl HummockVersionReader {
598 pub fn new(
599 sstable_store: SstableStoreRef,
600 state_store_metrics: Arc<HummockStateStoreMetrics>,
601 preload_retry_times: usize,
602 ) -> Self {
603 Self {
604 sstable_store,
605 state_store_metrics,
606 preload_retry_times,
607 }
608 }
609
610 pub fn stats(&self) -> &Arc<HummockStateStoreMetrics> {
611 &self.state_store_metrics
612 }
613}
614
615const SLOW_ITER_FETCH_META_DURATION_SECOND: f64 = 5.0;
616
617impl HummockVersionReader {
618 fn skip_get_by_vnode_user_key_range(
619 sstable_info: &SstableInfo,
620 vnode: VirtualNode,
621 user_key: UserKey<&[u8]>,
622 local_stats: &mut StoreLocalStatistic,
623 ) -> bool {
624 if let Some(vnode_statistics) = &sstable_info.vnode_statistics {
625 if let Some((vnode_min, vnode_max)) = vnode_statistics.get_vnode_user_key_range(vnode) {
628 local_stats.vnode_checked_get_count += 1;
629 if user_key < vnode_min.as_ref() || user_key > vnode_max.as_ref() {
630 local_stats.vnode_pruned_get_count += 1;
631 return true;
632 }
633 }
634 }
635 false
636 }
637
638 pub async fn get<'a, O>(
639 &'a self,
640 table_key: TableKey<Bytes>,
641 epoch: u64,
642 table_id: TableId,
643 table_option: TableOption,
644 read_options: ReadOptions,
645 read_version_tuple: ReadVersionTuple,
646 on_key_value_fn: impl crate::store::KeyValueFn<'a, O>,
647 ) -> StorageResult<Option<O>> {
648 let (imms, uncommitted_ssts, committed_version) = read_version_tuple;
649
650 let min_epoch = gen_min_epoch(epoch, table_option.retention_seconds);
651 let mut stats_guard = GetLocalMetricsGuard::new(self.state_store_metrics.clone(), table_id);
652 let local_stats = &mut stats_guard.local_stats;
653 local_stats.found_key = true;
654
655 for imm in &imms {
657 if imm.epoch() < min_epoch {
659 continue;
660 }
661
662 local_stats.staging_imm_get_count += 1;
663
664 if let Some((data, data_epoch)) = get_from_batch(
665 imm,
666 TableKey(table_key.as_ref()),
667 epoch,
668 &read_options,
669 local_stats,
670 ) {
671 return Ok(if data_epoch.pure_epoch() < min_epoch {
672 None
673 } else {
674 data.into_user_value()
675 .map(|v| {
676 on_key_value_fn(
677 FullKey::new_with_gap_epoch(
678 table_id,
679 table_key.to_ref(),
680 data_epoch,
681 ),
682 v.as_ref(),
683 )
684 })
685 .transpose()?
686 });
687 }
688 }
689
690 let dist_key_hash = read_options
692 .prefix_hint
693 .as_ref()
694 .map(|dist_key| Sstable::hash_for_filter(dist_key.as_ref(), table_id.as_raw_id()));
695
696 let full_key = FullKey::new_with_gap_epoch(
699 table_id,
700 TableKey(table_key.clone()),
701 EpochWithGap::new(epoch, MAX_SPILL_TIMES),
702 );
703 let single_table_key_range = table_key.clone()..=table_key.clone();
704
705 let pruned_uncommitted_ssts =
707 prune_overlapping_ssts(&uncommitted_ssts, table_id, &single_table_key_range);
708 for local_sst in pruned_uncommitted_ssts {
709 local_stats.staging_sst_get_count += 1;
710 if let Some(iter) = get_from_sstable_info(
711 self.sstable_store.clone(),
712 local_sst,
713 full_key.to_ref(),
714 &read_options,
715 dist_key_hash,
716 local_stats,
717 )
718 .await?
719 {
720 debug_assert!(iter.is_valid());
721 let data_epoch = iter.key().epoch_with_gap;
722 return Ok(if data_epoch.pure_epoch() < min_epoch {
723 None
724 } else {
725 iter.value()
726 .into_user_value()
727 .map(|v| {
728 on_key_value_fn(
729 FullKey::new_with_gap_epoch(
730 table_id,
731 table_key.to_ref(),
732 data_epoch,
733 ),
734 v,
735 )
736 })
737 .transpose()?
738 });
739 }
740 }
741 assert!(committed_version.is_valid());
745 for level in committed_version.levels(table_id) {
746 if level.table_infos.is_empty() {
747 continue;
748 }
749
750 match level.level_type {
751 LevelType::Overlapping | LevelType::Unspecified => {
752 let sstable_infos = prune_overlapping_ssts(
753 &level.table_infos,
754 table_id,
755 &single_table_key_range,
756 );
757 for sstable_info in sstable_infos {
758 if Self::skip_get_by_vnode_user_key_range(
760 sstable_info,
761 VirtualNode::from_index(full_key.user_key.get_vnode_id()),
762 full_key.user_key.as_ref(),
763 local_stats,
764 ) {
765 continue;
766 }
767
768 local_stats.overlapping_get_count += 1;
769 if let Some(iter) = get_from_sstable_info(
770 self.sstable_store.clone(),
771 sstable_info,
772 full_key.to_ref(),
773 &read_options,
774 dist_key_hash,
775 local_stats,
776 )
777 .await?
778 {
779 debug_assert!(iter.is_valid());
780 let data_epoch = iter.key().epoch_with_gap;
781 return Ok(if data_epoch.pure_epoch() < min_epoch {
782 None
783 } else {
784 iter.value()
785 .into_user_value()
786 .map(|v| {
787 on_key_value_fn(
788 FullKey::new_with_gap_epoch(
789 table_id,
790 table_key.to_ref(),
791 data_epoch,
792 ),
793 v,
794 )
795 })
796 .transpose()?
797 });
798 }
799 }
800 }
801 LevelType::Nonoverlapping => {
802 let mut table_info_idx =
803 search_sst_idx(&level.table_infos, full_key.user_key.as_ref());
804 if table_info_idx == 0 {
805 continue;
806 }
807 table_info_idx = table_info_idx.saturating_sub(1);
808 let sstable_info = &level.table_infos[table_info_idx];
809
810 if sstable_info.table_ids.binary_search(&table_id).is_err() {
811 continue;
812 }
813
814 let ord = sstable_info
816 .key_range
817 .compare_right_with_user_key(full_key.user_key.as_ref());
818 if ord == Ordering::Less {
820 sync_point!("HUMMOCK_V2::GET::SKIP_BY_NO_FILE");
821 continue;
822 }
823
824 if Self::skip_get_by_vnode_user_key_range(
826 sstable_info,
827 VirtualNode::from_index(full_key.user_key.get_vnode_id()),
828 full_key.user_key.as_ref(),
829 local_stats,
830 ) {
831 continue;
832 }
833
834 local_stats.non_overlapping_get_count += 1;
835 if let Some(iter) = get_from_sstable_info(
836 self.sstable_store.clone(),
837 sstable_info,
838 full_key.to_ref(),
839 &read_options,
840 dist_key_hash,
841 local_stats,
842 )
843 .await?
844 {
845 debug_assert!(iter.is_valid());
846 let data_epoch = iter.key().epoch_with_gap;
847 return Ok(if data_epoch.pure_epoch() < min_epoch {
848 None
849 } else {
850 iter.value()
851 .into_user_value()
852 .map(|v| {
853 on_key_value_fn(
854 FullKey::new_with_gap_epoch(
855 table_id,
856 table_key.to_ref(),
857 data_epoch,
858 ),
859 v,
860 )
861 })
862 .transpose()?
863 });
864 }
865 }
866 }
867 }
868 stats_guard.local_stats.found_key = false;
869 Ok(None)
870 }
871
872 pub async fn iter(
873 &self,
874 table_key_range: TableKeyRange,
875 epoch: u64,
876 table_id: TableId,
877 table_option: TableOption,
878 read_options: ReadOptions,
879 read_version_tuple: (Vec<ImmutableMemtable>, Vec<SstableInfo>, CommittedVersion),
880 ) -> StorageResult<HummockStorageIterator> {
881 self.iter_with_memtable(
882 table_key_range,
883 epoch,
884 table_id,
885 table_option,
886 read_options,
887 read_version_tuple,
888 None,
889 )
890 .await
891 }
892
893 pub async fn iter_with_memtable<'b>(
894 &self,
895 table_key_range: TableKeyRange,
896 epoch: u64,
897 table_id: TableId,
898 table_option: TableOption,
899 read_options: ReadOptions,
900 read_version_tuple: (Vec<ImmutableMemtable>, Vec<SstableInfo>, CommittedVersion),
901 memtable_iter: Option<MemTableHummockIterator<'b>>,
902 ) -> StorageResult<HummockStorageIteratorInner<'b>> {
903 let user_key_range_ref = bound_table_key_range(table_id, &table_key_range);
904 let user_key_range = (
905 user_key_range_ref.0.map(|key| key.cloned()),
906 user_key_range_ref.1.map(|key| key.cloned()),
907 );
908 let mut factory = ForwardIteratorFactory::default();
909 let mut local_stats = StoreLocalStatistic::default();
910 let (imms, uncommitted_ssts, committed) = read_version_tuple;
911 let min_epoch = gen_min_epoch(epoch, table_option.retention_seconds);
912 self.iter_inner(
913 table_key_range,
914 epoch,
915 table_id,
916 read_options,
917 imms,
918 uncommitted_ssts,
919 &committed,
920 &mut local_stats,
921 &mut factory,
922 )
923 .await?;
924 let merge_iter = factory.build(memtable_iter);
925 let mut user_iter = UserIterator::new(
927 merge_iter,
928 user_key_range,
929 epoch,
930 min_epoch,
931 Some(committed),
932 );
933 user_iter.rewind().await?;
934 Ok(HummockStorageIteratorInner::new(
935 user_iter,
936 self.state_store_metrics.clone(),
937 table_id,
938 local_stats,
939 ))
940 }
941
942 pub async fn rev_iter<'b>(
943 &self,
944 table_key_range: TableKeyRange,
945 epoch: u64,
946 table_id: TableId,
947 table_option: TableOption,
948 read_options: ReadOptions,
949 read_version_tuple: (Vec<ImmutableMemtable>, Vec<SstableInfo>, CommittedVersion),
950 memtable_iter: Option<MemTableHummockRevIterator<'b>>,
951 ) -> StorageResult<HummockStorageRevIteratorInner<'b>> {
952 let user_key_range_ref = bound_table_key_range(table_id, &table_key_range);
953 let user_key_range = (
954 user_key_range_ref.0.map(|key| key.cloned()),
955 user_key_range_ref.1.map(|key| key.cloned()),
956 );
957 let mut factory = BackwardIteratorFactory::default();
958 let mut local_stats = StoreLocalStatistic::default();
959 let (imms, uncommitted_ssts, committed) = read_version_tuple;
960 let min_epoch = gen_min_epoch(epoch, table_option.retention_seconds);
961 self.iter_inner(
962 table_key_range,
963 epoch,
964 table_id,
965 read_options,
966 imms,
967 uncommitted_ssts,
968 &committed,
969 &mut local_stats,
970 &mut factory,
971 )
972 .await?;
973 let merge_iter = factory.build(memtable_iter);
974 let mut user_iter = BackwardUserIterator::new(
976 merge_iter,
977 user_key_range,
978 epoch,
979 min_epoch,
980 Some(committed),
981 );
982 user_iter.rewind().await?;
983 Ok(HummockStorageRevIteratorInner::new(
984 user_iter,
985 self.state_store_metrics.clone(),
986 table_id,
987 local_stats,
988 ))
989 }
990
991 async fn iter_inner<F: IteratorFactory>(
992 &self,
993 table_key_range: TableKeyRange,
994 epoch: u64,
995 table_id: TableId,
996 read_options: ReadOptions,
997 imms: Vec<ImmutableMemtable>,
998 uncommitted_ssts: Vec<SstableInfo>,
999 committed: &CommittedVersion,
1000 local_stats: &mut StoreLocalStatistic,
1001 factory: &mut F,
1002 ) -> StorageResult<()> {
1003 {
1004 fn bound_inner<T>(bound: &Bound<T>) -> Option<&T> {
1005 match bound {
1006 Bound::Included(bound) | Bound::Excluded(bound) => Some(bound),
1007 Bound::Unbounded => None,
1008 }
1009 }
1010 let (left, right) = &table_key_range;
1011 if let (Some(left), Some(right)) = (bound_inner(left), bound_inner(right))
1012 && right < left
1013 {
1014 if cfg!(debug_assertions) {
1015 panic!("invalid iter key range: {table_id} {left:?} {right:?}")
1016 } else {
1017 return Err(HummockError::other(format!(
1018 "invalid iter key range: {table_id} {left:?} {right:?}"
1019 ))
1020 .into());
1021 }
1022 }
1023 }
1024
1025 local_stats.staging_imm_iter_count = imms.len() as u64;
1026 for imm in imms {
1027 factory.add_batch_iter(imm);
1028 }
1029
1030 let user_key_range = bound_table_key_range(table_id, &table_key_range);
1034 let user_key_range_ref = (
1035 user_key_range.0.as_ref().map(UserKey::as_ref),
1036 user_key_range.1.as_ref().map(UserKey::as_ref),
1037 );
1038 let mut staging_sst_iter_count = 0;
1039 let filter_prefix_hash = read_options
1041 .prefix_hint
1042 .as_ref()
1043 .map(|hint| Sstable::hash_for_filter(hint, table_id.as_raw_id()));
1044 let mut sst_read_options = SstableIteratorReadOptions::from_read_options(&read_options);
1045 sst_read_options.read_table_id = Some(table_id);
1046 sst_read_options.scan_end_user_key = Some(user_key_range.1.map(|key| key.cloned()));
1047 sst_read_options.prefetch = read_options.prefetch_options.prefetch;
1048 if sst_read_options.prefetch {
1049 sst_read_options.max_preload_retry_times = self.preload_retry_times;
1050 }
1051 let sst_read_options = Arc::new(sst_read_options);
1052 for sstable_info in &uncommitted_ssts {
1053 let table_holder = self
1054 .sstable_store
1055 .sstable(sstable_info, local_stats)
1056 .await?;
1057
1058 if let Some(prefix_hash) = filter_prefix_hash.as_ref()
1059 && !hit_sstable_filter(
1060 &table_holder,
1061 &user_key_range_ref,
1062 *prefix_hash,
1063 local_stats,
1064 )
1065 {
1066 continue;
1067 }
1068
1069 staging_sst_iter_count += 1;
1070 factory.add_staging_sst_iter(F::SstableIteratorType::create(
1071 table_holder,
1072 self.sstable_store.clone(),
1073 sst_read_options.clone(),
1074 sstable_info,
1075 ));
1076 }
1077 local_stats.staging_sst_iter_count = staging_sst_iter_count;
1078
1079 let timer = Instant::now();
1080
1081 for level in committed.levels(table_id) {
1082 if level.table_infos.is_empty() {
1083 continue;
1084 }
1085
1086 if level.level_type == LevelType::Nonoverlapping {
1087 let mut table_infos =
1088 prune_nonoverlapping_ssts(&level.table_infos, user_key_range_ref, table_id)
1089 .peekable();
1090
1091 if table_infos.peek().is_none() {
1092 continue;
1093 }
1094 let sstable_infos = table_infos.cloned().collect_vec();
1095 if sstable_infos.len() > 1 {
1096 factory.add_concat_sst_iter(
1097 sstable_infos,
1098 self.sstable_store.clone(),
1099 sst_read_options.clone(),
1100 );
1101 local_stats.non_overlapping_iter_count += 1;
1102 } else {
1103 let sstable_info = &sstable_infos[0];
1104
1105 let sstable = self
1106 .sstable_store
1107 .sstable(sstable_info, local_stats)
1108 .await?;
1109
1110 if let Some(dist_hash) = filter_prefix_hash.as_ref()
1111 && !hit_sstable_filter(
1112 &sstable,
1113 &user_key_range_ref,
1114 *dist_hash,
1115 local_stats,
1116 )
1117 {
1118 continue;
1119 }
1120 factory.add_overlapping_sst_iter(F::SstableIteratorType::create(
1126 sstable,
1127 self.sstable_store.clone(),
1128 sst_read_options.clone(),
1129 sstable_info,
1130 ));
1131 local_stats.non_overlapping_iter_count += 1;
1132 }
1133 } else {
1134 let table_infos =
1135 prune_overlapping_ssts(&level.table_infos, table_id, &table_key_range);
1136 let fetch_meta_req = table_infos.rev().collect_vec();
1138 if fetch_meta_req.is_empty() {
1139 continue;
1140 }
1141 for sstable_info in fetch_meta_req {
1142 let sstable = self
1143 .sstable_store
1144 .sstable(sstable_info, local_stats)
1145 .await?;
1146 assert_eq!(sstable_info.object_id, sstable.id);
1147 if let Some(dist_hash) = filter_prefix_hash.as_ref()
1148 && !hit_sstable_filter(
1149 &sstable,
1150 &user_key_range_ref,
1151 *dist_hash,
1152 local_stats,
1153 )
1154 {
1155 continue;
1156 }
1157 factory.add_overlapping_sst_iter(F::SstableIteratorType::create(
1158 sstable,
1159 self.sstable_store.clone(),
1160 sst_read_options.clone(),
1161 sstable_info,
1162 ));
1163 local_stats.overlapping_iter_count += 1;
1164 }
1165 }
1166 }
1167 let fetch_meta_duration_sec = timer.elapsed().as_secs_f64();
1168 if fetch_meta_duration_sec > SLOW_ITER_FETCH_META_DURATION_SECOND {
1169 let table_id_string = table_id.to_string();
1170 tracing::warn!(
1171 "Fetching meta while creating an iter to read table_id {:?} at epoch {:?} is slow: duration = {:?}s, cache unhits = {:?}.",
1172 table_id_string,
1173 epoch,
1174 fetch_meta_duration_sec,
1175 local_stats.cache_meta_block_miss
1176 );
1177 self.state_store_metrics
1178 .iter_slow_fetch_meta_cache_unhits
1179 .set(local_stats.cache_meta_block_miss as i64);
1180 }
1181 Ok(())
1182 }
1183
1184 pub async fn iter_log(
1185 &self,
1186 epoch_range: (u64, u64),
1187 key_range: TableKeyRange,
1188 options: ReadLogOptions,
1189 table_change_log_manager: Arc<TableChangeLogManager>,
1190 ) -> HummockResult<ChangeLogIterator> {
1191 let change_log: Vec<_> = {
1193 let table_change_logs = table_change_log_manager
1194 .fetch_table_change_logs(options.table_id, epoch_range, false, None)
1195 .await?;
1196 if let Some(change_log) = table_change_logs.get(&options.table_id) {
1197 change_log.filter_epoch(epoch_range).cloned().collect_vec()
1198 } else {
1199 Vec::new()
1200 }
1201 };
1202
1203 if let Some(max_epoch_change_log) = change_log.last() {
1204 let (_, max_epoch) = epoch_range;
1205 if !max_epoch_change_log.epochs().contains(&max_epoch) {
1206 warn!(
1207 max_epoch,
1208 change_log_epochs = ?change_log.iter().flat_map(|epoch_log| epoch_log.epochs()).collect_vec(),
1209 table_id = %options.table_id,
1210 "max_epoch does not exist"
1211 );
1212 }
1213 }
1214 let read_options = Arc::new(SstableIteratorReadOptions {
1215 cache_policy: Default::default(),
1216 read_table_id: Some(options.table_id),
1217 scan_end_user_key: None,
1218 prefetch: false,
1219 max_preload_retry_times: 0,
1220 });
1221
1222 async fn make_iter(
1223 sstable_infos: impl Iterator<Item = &SstableInfo>,
1224 sstable_store: &SstableStoreRef,
1225 read_options: Arc<SstableIteratorReadOptions>,
1226 local_stat: &mut StoreLocalStatistic,
1227 ) -> HummockResult<MergeIterator<SstableIterator>> {
1228 let iters = try_join_all(sstable_infos.map(|sstable_info| {
1229 let sstable_store = sstable_store.clone();
1230 let read_options = read_options.clone();
1231 async move {
1232 let mut local_stat = StoreLocalStatistic::default();
1233 let table_holder = sstable_store.sstable(sstable_info, &mut local_stat).await?;
1234 Ok::<_, HummockError>((
1235 SstableIterator::new(
1236 table_holder,
1237 sstable_store,
1238 read_options,
1239 sstable_info,
1240 ),
1241 local_stat,
1242 ))
1243 }
1244 }))
1245 .await?;
1246 Ok::<_, HummockError>(MergeIterator::new(iters.into_iter().map(
1247 |(iter, stats)| {
1248 local_stat.add(&stats);
1249 iter
1250 },
1251 )))
1252 }
1253
1254 let mut local_stat = StoreLocalStatistic::default();
1255
1256 let new_value_iter = make_iter(
1257 change_log
1258 .iter()
1259 .flat_map(|log| log.new_value.iter())
1260 .filter(|sst| filter_single_sst(sst, options.table_id, &key_range)),
1261 &self.sstable_store,
1262 read_options.clone(),
1263 &mut local_stat,
1264 )
1265 .await?;
1266 let old_value_iter = make_iter(
1267 change_log
1268 .iter()
1269 .flat_map(|log| log.old_value.iter())
1270 .filter(|sst| filter_single_sst(sst, options.table_id, &key_range)),
1271 &self.sstable_store,
1272 read_options.clone(),
1273 &mut local_stat,
1274 )
1275 .await?;
1276 ChangeLogIterator::new(
1277 epoch_range,
1278 key_range,
1279 new_value_iter,
1280 old_value_iter,
1281 options.table_id,
1282 IterLocalMetricsGuard::new(
1283 self.state_store_metrics.clone(),
1284 options.table_id,
1285 local_stat,
1286 ),
1287 )
1288 .await
1289 }
1290
1291 pub async fn nearest<'a, M: MeasureDistanceBuilder, O: Send>(
1292 &'a self,
1293 version: PinnedVersion,
1294 table_id: TableId,
1295 target: VectorRef<'a>,
1296 options: VectorNearestOptions,
1297 on_nearest_item_fn: impl OnNearestItemFn<'a, O>,
1298 ) -> HummockResult<Vec<O>> {
1299 let Some(index) = version.vector_indexes.get(&table_id) else {
1300 return Ok(vec![]);
1301 };
1302 if target.dimension() != index.dimension {
1303 return Err(HummockError::other(format!(
1304 "target dimension {} not match index dimension {}",
1305 target.dimension(),
1306 index.dimension
1307 )));
1308 }
1309 match &index.inner {
1310 VectorIndexImpl::Flat(flat) => {
1311 let mut builder = NearestBuilder::<'_, O, M>::new(target, options.top_n);
1312 let mut cache_stat = VectorStoreCacheStats::default();
1313 for vector_file in &flat.vector_store_info.vector_files {
1314 let meta = self
1315 .sstable_store
1316 .get_vector_file_meta(vector_file, &mut cache_stat)
1317 .await?;
1318 for (i, block_meta) in meta.block_metas.iter().enumerate() {
1319 let block = self
1320 .sstable_store
1321 .get_vector_block(vector_file, i, block_meta, &mut cache_stat)
1322 .await?;
1323 builder.add(&**block, &on_nearest_item_fn);
1324 }
1325 }
1326 cache_stat.report(table_id, "flat", self.stats());
1327 Ok(builder.finish())
1328 }
1329 VectorIndexImpl::HnswFlat(hnsw_flat) => {
1330 let Some(graph_file) = &hnsw_flat.graph_file else {
1331 return Ok(vec![]);
1332 };
1333
1334 let mut ctx = FileVectorStoreCtx::default();
1335
1336 let graph = self
1337 .sstable_store
1338 .get_hnsw_graph(graph_file, &mut ctx.stats)
1339 .await?;
1340
1341 let vector_store =
1342 FileVectorStore::new_for_reader(hnsw_flat, self.sstable_store.clone());
1343 let (items, stats) = nearest::<O, M, _>(
1344 &vector_store,
1345 &mut ctx,
1346 &*graph,
1347 target,
1348 on_nearest_item_fn,
1349 options.hnsw_ef_search,
1350 options.top_n,
1351 )
1352 .await?;
1353 ctx.stats.report(table_id, "hnsw_read", self.stats());
1354 report_hnsw_stat(
1355 self.stats(),
1356 table_id,
1357 "hnsw_read",
1358 options.top_n,
1359 options.hnsw_ef_search,
1360 [stats],
1361 );
1362 Ok(items)
1363 }
1364 }
1365 }
1366}
1367
1368#[cfg(test)]
1369mod tests {
1370 use std::collections::{BTreeMap, HashMap, HashSet};
1371 use std::sync::Arc;
1372
1373 use bytes::Bytes;
1374 use prometheus::Registry;
1375 use risingwave_common::catalog::{TableId, TableOption};
1376 use risingwave_common::config::MetricLevel;
1377 use risingwave_common::hash::VirtualNode;
1378 use risingwave_common::util::epoch::test_epoch;
1379 use risingwave_hummock_sdk::compaction_group::StaticCompactionGroupId;
1380 use risingwave_hummock_sdk::key::{FullKey, TableKey, UserKey, gen_key_from_bytes};
1381 use risingwave_hummock_sdk::key_range::KeyRange;
1382 use risingwave_hummock_sdk::level::{Level, Levels, OverlappingLevel};
1383 use risingwave_hummock_sdk::sstable_info::{SstableInfo, SstableInfoInner, VnodeStatistics};
1384 use risingwave_hummock_sdk::version::HummockVersion;
1385 use risingwave_hummock_sdk::{EpochWithGap, HummockSstableObjectId};
1386 use risingwave_pb::hummock::hummock_version::PbLevels;
1387 use risingwave_pb::hummock::{
1388 LevelType as PbLevelType, PbHummockVersion, PbLevel, PbOverlappingLevel,
1389 PbSstableFilterLayout, PbSstableFilterType, PbStateTableInfo, StateTableInfoDelta,
1390 };
1391 use tokio::sync::mpsc::unbounded_channel;
1392
1393 use crate::hummock::HummockValue;
1394 use crate::hummock::iterator::test_utils::mock_sstable_store;
1395 use crate::hummock::local_version::pinned_version::{PinVersionAction, PinnedVersion};
1396 use crate::hummock::store::version::{CommittedVersion, HummockVersionReader};
1397 use crate::hummock::test_utils::{
1398 default_builder_opt_for_test, gen_test_sstable_with_table_ids,
1399 };
1400 use crate::monitor::{HummockStateStoreMetrics, flush_local_metrics_for_test};
1401 use crate::store::ReadOptions;
1402
1403 #[tokio::test]
1407 async fn test_get_skips_sst_by_table_id_filter() {
1408 let query_table_id = TableId::new(100);
1409 let epoch: u64 = (31 * 1000) << 16;
1410 let compaction_group_id = StaticCompactionGroupId::StateDefault;
1411
1412 let sst_info = SstableInfoInner {
1414 sst_id: 1.into(),
1415 object_id: 1.into(),
1416 key_range: KeyRange {
1417 left: Bytes::from(
1418 FullKey::for_test(TableId::new(50), b"aaa".to_vec(), epoch).encode(),
1419 ),
1420 right: Bytes::from(
1421 FullKey::for_test(TableId::new(150), b"zzz".to_vec(), epoch).encode(),
1422 ),
1423 right_exclusive: false,
1424 },
1425 table_ids: vec![TableId::new(50), TableId::new(150)],
1426 file_size: 1024,
1427 ..Default::default()
1428 }
1429 .into();
1430
1431 let level = Level {
1432 level_idx: 1,
1433 level_type: PbLevelType::Nonoverlapping,
1434 table_infos: vec![sst_info],
1435 total_file_size: 0,
1436 sub_level_id: 0,
1437 uncompressed_file_size: 0,
1438 vnode_partition_count: 0,
1439 };
1440
1441 #[allow(deprecated)]
1442 let levels = Levels {
1443 levels: vec![level],
1444 l0: OverlappingLevel::default(),
1445 group_id: compaction_group_id,
1446 parent_group_id: compaction_group_id,
1447 member_table_ids: vec![],
1448 compaction_group_version_id: 0,
1449 };
1450
1451 let mut version = HummockVersion::from_persisted_protobuf_owned(PbHummockVersion {
1452 id: 1u64.into(),
1453 ..Default::default()
1454 });
1455 version.levels.insert(compaction_group_id, levels);
1456 version.state_table_info.apply_delta(
1457 &HashMap::from([(
1458 query_table_id,
1459 StateTableInfoDelta {
1460 committed_epoch: epoch,
1461 compaction_group_id,
1462 },
1463 )]),
1464 &HashSet::new(),
1465 );
1466
1467 let pinned_version = PinnedVersion::new(version, unbounded_channel().0);
1468 let reader = HummockVersionReader::new(
1469 mock_sstable_store().await,
1470 Arc::new(HummockStateStoreMetrics::unused()),
1471 0,
1472 );
1473
1474 let result = reader
1475 .get(
1476 TableKey(Bytes::from("test_key")),
1477 epoch,
1478 query_table_id,
1479 TableOption::default(),
1480 ReadOptions::default(),
1481 (vec![], vec![], pinned_version),
1482 |_key, _value| Ok(()),
1483 )
1484 .await
1485 .unwrap();
1486
1487 assert!(result.is_none());
1488 }
1489
1490 #[allow(deprecated)]
1492 fn build_version_with_vnode_stats(
1493 table_id: TableId,
1494 vnode_stats: VnodeStatistics,
1495 key_range: (Vec<u8>, Vec<u8>),
1496 level_type: PbLevelType,
1497 ) -> (SstableInfo, CommittedVersion) {
1498 let object_id = HummockSstableObjectId::new(1);
1499 let left_full_key = FullKey::new_with_gap_epoch(
1500 table_id,
1501 TableKey(Bytes::from(key_range.0)),
1502 EpochWithGap::new_from_epoch(test_epoch(0)),
1503 )
1504 .encode();
1505 let right_full_key = FullKey::new_with_gap_epoch(
1506 table_id,
1507 TableKey(Bytes::from(key_range.1)),
1508 EpochWithGap::new_from_epoch(test_epoch(0)),
1509 )
1510 .encode();
1511
1512 let sstable_info: SstableInfo = SstableInfoInner {
1513 object_id,
1514 sst_id: object_id.as_raw_id().into(),
1515 key_range: KeyRange {
1516 left: Bytes::from(left_full_key),
1517 right: Bytes::from(right_full_key),
1518 right_exclusive: false,
1519 },
1520 file_size: 1,
1521 table_ids: vec![table_id],
1522 meta_offset: 0,
1523 stale_key_count: 0,
1524 total_key_count: 0,
1525 min_epoch: 0,
1526 max_epoch: 0,
1527 uncompressed_file_size: 0,
1528 range_tombstone_count: 0,
1529 filter_type: PbSstableFilterType::SstableFilterXor16,
1530 filter_layout: PbSstableFilterLayout::Plain,
1531 sst_size: 1,
1532 vnode_statistics: Some(vnode_stats),
1533 }
1534 .into();
1535 let pb_level = PbLevel {
1536 level_idx: if level_type == PbLevelType::Overlapping {
1537 0
1538 } else {
1539 1
1540 },
1541 level_type: level_type as i32,
1542 table_infos: vec![sstable_info.clone().into()],
1543 total_file_size: 1,
1544 sub_level_id: 0,
1545 uncompressed_file_size: 1,
1546 vnode_partition_count: 0,
1547 };
1548
1549 let (levels, l0) = if level_type == PbLevelType::Overlapping {
1550 (
1551 vec![],
1552 Some(PbOverlappingLevel {
1553 sub_levels: vec![pb_level],
1554 total_file_size: 1,
1555 uncompressed_file_size: 1,
1556 }),
1557 )
1558 } else {
1559 (vec![pb_level], Some(PbOverlappingLevel::default()))
1560 };
1561
1562 let pb_levels = PbLevels {
1563 levels,
1564 l0,
1565 group_id: StaticCompactionGroupId::NewCompactionGroup,
1566 parent_group_id: 0.into(),
1567 member_table_ids: vec![],
1568 compaction_group_version_id: 0,
1569 };
1570
1571 let pb_version = PbHummockVersion {
1572 id: 1.into(),
1573 levels: HashMap::from_iter([(StaticCompactionGroupId::NewCompactionGroup, pb_levels)]),
1574 max_committed_epoch: 0,
1575 table_watermarks: HashMap::new(),
1576 table_change_logs: HashMap::new(),
1577 state_table_info: HashMap::from_iter([(
1578 table_id,
1579 PbStateTableInfo {
1580 committed_epoch: 0,
1581 compaction_group_id: StaticCompactionGroupId::NewCompactionGroup,
1582 },
1583 )]),
1584 vector_indexes: HashMap::new(),
1585 };
1586
1587 let version = HummockVersion::from(&pb_version);
1588 let (tx, _rx) = unbounded_channel::<PinVersionAction>();
1589 let pinned = PinnedVersion::new(version, tx);
1590 (sstable_info, pinned)
1591 }
1592
1593 #[allow(deprecated)]
1595 fn build_version_from_sstables(
1596 table_id: TableId,
1597 sstable_infos: Vec<SstableInfo>,
1598 level_type: PbLevelType,
1599 ) -> CommittedVersion {
1600 let total_file_size = sstable_infos.iter().map(|sst| sst.file_size).sum::<u64>();
1601 let uncompressed_file_size = sstable_infos
1602 .iter()
1603 .map(|sst| sst.uncompressed_file_size)
1604 .sum::<u64>();
1605 let pb_level = PbLevel {
1606 level_idx: if level_type == PbLevelType::Overlapping {
1607 0
1608 } else {
1609 1
1610 },
1611 level_type: level_type as i32,
1612 table_infos: sstable_infos.into_iter().map(Into::into).collect(),
1613 total_file_size,
1614 sub_level_id: 0,
1615 uncompressed_file_size,
1616 vnode_partition_count: 0,
1617 };
1618
1619 let (levels, l0) = if level_type == PbLevelType::Overlapping {
1620 (
1621 vec![],
1622 Some(PbOverlappingLevel {
1623 sub_levels: vec![pb_level],
1624 total_file_size,
1625 uncompressed_file_size,
1626 }),
1627 )
1628 } else {
1629 (vec![pb_level], Some(PbOverlappingLevel::default()))
1630 };
1631
1632 let pb_levels = PbLevels {
1633 levels,
1634 l0,
1635 group_id: StaticCompactionGroupId::NewCompactionGroup,
1636 parent_group_id: 0.into(),
1637 member_table_ids: vec![],
1638 compaction_group_version_id: 0,
1639 };
1640
1641 let pb_version = PbHummockVersion {
1642 id: 1.into(),
1643 levels: HashMap::from_iter([(StaticCompactionGroupId::NewCompactionGroup, pb_levels)]),
1644 max_committed_epoch: 0,
1645 table_watermarks: HashMap::new(),
1646 table_change_logs: HashMap::new(),
1647 state_table_info: HashMap::from_iter([(
1648 table_id,
1649 PbStateTableInfo {
1650 committed_epoch: 0,
1651 compaction_group_id: StaticCompactionGroupId::NewCompactionGroup,
1652 },
1653 )]),
1654 vector_indexes: HashMap::new(),
1655 };
1656
1657 let version = HummockVersion::from(&pb_version);
1658 let (tx, _rx) = unbounded_channel::<PinVersionAction>();
1659 PinnedVersion::new(version, tx)
1660 }
1661
1662 #[allow(deprecated)]
1664 fn build_version_from_sstable(
1665 table_id: TableId,
1666 sstable_info: SstableInfo,
1667 ) -> CommittedVersion {
1668 build_version_from_sstables(table_id, vec![sstable_info], PbLevelType::Nonoverlapping)
1669 }
1670
1671 fn vnode_prune_counts(
1672 metrics: &HummockStateStoreMetrics,
1673 table_id: TableId,
1674 operation: &str,
1675 ) -> (u64, u64) {
1676 let table_label = table_id.to_string();
1677 let checked = metrics
1678 .vnode_pruning_counts
1679 .with_guarded_label_values(&[
1680 table_label.clone(),
1681 operation.to_owned(),
1682 "checked".to_owned(),
1683 ])
1684 .get();
1685 let pruned = metrics
1686 .vnode_pruning_counts
1687 .with_guarded_label_values(&[table_label, operation.to_owned(), "pruned".to_owned()])
1688 .get();
1689 (checked, pruned)
1690 }
1691
1692 async fn assert_vnode_prune_get_skips_out_of_range_key(
1693 table_id: TableId,
1694 epoch: u64,
1695 level_type: PbLevelType,
1696 ) {
1697 let sstable_store = mock_sstable_store().await;
1698 let registry = Registry::new();
1699 let metrics = Arc::new(HummockStateStoreMetrics::new(®istry, MetricLevel::Debug));
1700 let reader = HummockVersionReader::new(sstable_store, metrics.clone(), 0);
1701 let (checked_before, pruned_before) = vnode_prune_counts(&metrics, table_id, "get");
1702
1703 let make_user_key = |vnode: VirtualNode, suffix: &str| {
1704 let mut raw = vnode.to_be_bytes().to_vec();
1705 raw.extend_from_slice(suffix.as_bytes());
1706 UserKey::new(table_id, TableKey(raw.into()))
1707 };
1708
1709 let vnode_stats = VnodeStatistics::from_map(BTreeMap::from_iter([(
1711 VirtualNode::from_index(1),
1712 (
1713 make_user_key(VirtualNode::from_index(1), "aa"),
1714 make_user_key(VirtualNode::from_index(1), "bb"),
1715 ),
1716 )]));
1717
1718 let key_range = {
1720 let mut left = VirtualNode::from_index(0).to_be_bytes().to_vec();
1721 left.extend_from_slice(b"aa");
1722 let mut right = VirtualNode::from_index(1).to_be_bytes().to_vec();
1723 right.extend_from_slice(b"zzzz");
1724 (left, right)
1725 };
1726
1727 let (_sst, committed) =
1728 build_version_with_vnode_stats(table_id, vnode_stats, key_range, level_type);
1729
1730 let mut raw = VirtualNode::from_index(1).to_be_bytes().to_vec();
1732 raw.extend_from_slice(b"zz");
1733 let table_key = TableKey(Bytes::from(raw.clone()));
1734
1735 let result = reader
1736 .get(
1737 table_key,
1738 epoch,
1739 table_id,
1740 TableOption::default(),
1741 ReadOptions::default(),
1742 (vec![], vec![], committed),
1743 |_k, v| Ok(Bytes::copy_from_slice(v)),
1744 )
1745 .await
1746 .unwrap();
1747 flush_local_metrics_for_test();
1748
1749 assert!(
1750 result.is_none(),
1751 "vnode pruning should skip SST without reading data"
1752 );
1753 let (checked_after, pruned_after) = vnode_prune_counts(&metrics, table_id, "get");
1754 assert_eq!(checked_before + 1, checked_after);
1755 assert_eq!(pruned_before + 1, pruned_after);
1756 }
1757
1758 async fn assert_vnode_prune_get_not_pruned_nonoverlapping() {
1759 let table_id = TableId::new(42);
1760 let epoch = test_epoch(3);
1761 let sstable_store = mock_sstable_store().await;
1762 let registry = Registry::new();
1763 let metrics = Arc::new(HummockStateStoreMetrics::new(®istry, MetricLevel::Debug));
1764 let reader = HummockVersionReader::new(sstable_store.clone(), metrics.clone(), 0);
1765 let (checked_before, pruned_before) = vnode_prune_counts(&metrics, table_id, "get");
1766
1767 let mut opts = default_builder_opt_for_test();
1768 opts.max_vnode_key_range_bytes = None;
1769 let mut kvs = vec![
1770 (
1771 FullKey::new_with_gap_epoch(
1772 table_id,
1773 gen_key_from_bytes(VirtualNode::from_index(1), b"aa"),
1774 EpochWithGap::new_from_epoch(epoch),
1775 ),
1776 HummockValue::put(Bytes::from_static(b"v1")),
1777 ),
1778 (
1779 FullKey::new_with_gap_epoch(
1780 table_id,
1781 gen_key_from_bytes(VirtualNode::ZERO, b"cc"),
1782 EpochWithGap::new_from_epoch(epoch),
1783 ),
1784 HummockValue::put(Bytes::from_static(b"v0")),
1785 ),
1786 ];
1787 kvs.sort_by(|(k1, _), (k2, _)| k1.cmp(k2));
1788 let (_, mut sstable_info): (crate::hummock::sstable_store::TableHolder, SstableInfo) =
1789 gen_test_sstable_with_table_ids(
1790 opts,
1791 10,
1792 kvs.into_iter(),
1793 sstable_store.clone(),
1794 vec![table_id.as_raw_id()],
1795 )
1796 .await;
1797 let mut inner = sstable_info.get_inner();
1799 inner.vnode_statistics = Some(VnodeStatistics::from_map(BTreeMap::from_iter([(
1800 VirtualNode::from_index(1),
1801 (
1802 UserKey::new(
1803 table_id,
1804 gen_key_from_bytes(VirtualNode::from_index(1), b"aa"),
1805 ),
1806 UserKey::new(
1807 table_id,
1808 gen_key_from_bytes(VirtualNode::from_index(1), b"zz"),
1809 ),
1810 ),
1811 )])));
1812 sstable_info = inner.into();
1813 let committed = build_version_from_sstable(table_id, sstable_info.clone());
1814
1815 let mut raw = VirtualNode::from_index(1).to_be_bytes().to_vec();
1817 raw.extend_from_slice(b"aa");
1818 let table_key = TableKey(Bytes::from(raw.clone()));
1819
1820 let result = reader
1821 .get(
1822 table_key,
1823 epoch,
1824 table_id,
1825 TableOption::default(),
1826 ReadOptions::default(),
1827 (vec![], vec![], committed),
1828 |_k, v| Ok(Bytes::copy_from_slice(v)),
1829 )
1830 .await
1831 .unwrap();
1832 flush_local_metrics_for_test();
1833 assert!(result.is_some(), "key should be read when not pruned");
1834 let (checked_after, pruned_after) = vnode_prune_counts(&metrics, table_id, "get");
1835 assert_eq!(checked_before + 1, checked_after);
1836 assert_eq!(pruned_before, pruned_after);
1837 }
1838
1839 #[tokio::test]
1840 async fn test_vnode_prune_get_single_sst_cases() {
1841 assert_vnode_prune_get_skips_out_of_range_key(
1842 TableId::default(),
1843 test_epoch(1),
1844 PbLevelType::Nonoverlapping,
1845 )
1846 .await;
1847 assert_vnode_prune_get_skips_out_of_range_key(
1848 TableId::new(7),
1849 test_epoch(2),
1850 PbLevelType::Overlapping,
1851 )
1852 .await;
1853 assert_vnode_prune_get_not_pruned_nonoverlapping().await;
1854 }
1855
1856 #[tokio::test]
1857 async fn test_vnode_prune_get_overlapping_distribution_prunes_only_out_of_range_sst() {
1858 let table_id = TableId::new(77);
1859 let epoch = test_epoch(4);
1860 let vnode = VirtualNode::from_index(1);
1861 let sstable_store = mock_sstable_store().await;
1862 let registry = Registry::new();
1863 let metrics = Arc::new(HummockStateStoreMetrics::new(®istry, MetricLevel::Debug));
1864 let reader = HummockVersionReader::new(sstable_store.clone(), metrics.clone(), 0);
1865 let (checked_before, pruned_before) = vnode_prune_counts(&metrics, table_id, "get");
1866
1867 let mut opts = default_builder_opt_for_test();
1868 opts.max_vnode_key_range_bytes = None;
1869
1870 let mut kvs1 = vec![
1871 (
1872 FullKey::new_with_gap_epoch(
1873 table_id,
1874 gen_key_from_bytes(vnode, b"aa"),
1875 EpochWithGap::new_from_epoch(epoch),
1876 ),
1877 HummockValue::put(Bytes::from_static(b"s1_aa")),
1878 ),
1879 (
1880 FullKey::new_with_gap_epoch(
1881 table_id,
1882 gen_key_from_bytes(vnode, b"zz"),
1883 EpochWithGap::new_from_epoch(epoch),
1884 ),
1885 HummockValue::put(Bytes::from_static(b"s1_zz")),
1886 ),
1887 ];
1888 kvs1.sort_by(|(k1, _), (k2, _)| k1.cmp(k2));
1889 let (_, mut sst1): (crate::hummock::sstable_store::TableHolder, SstableInfo) =
1890 gen_test_sstable_with_table_ids(
1891 opts.clone(),
1892 11,
1893 kvs1.into_iter(),
1894 sstable_store.clone(),
1895 vec![table_id.as_raw_id()],
1896 )
1897 .await;
1898 let mut sst1_inner = sst1.get_inner();
1899 sst1_inner.vnode_statistics = Some(VnodeStatistics::from_map(BTreeMap::from_iter([(
1900 vnode,
1901 (
1902 UserKey::new(table_id, gen_key_from_bytes(vnode, b"aa")),
1903 UserKey::new(table_id, gen_key_from_bytes(vnode, b"bb")),
1904 ),
1905 )])));
1906 sst1 = sst1_inner.into();
1907
1908 let mut kvs2 = vec![(
1909 FullKey::new_with_gap_epoch(
1910 table_id,
1911 gen_key_from_bytes(vnode, b"mm"),
1912 EpochWithGap::new_from_epoch(epoch),
1913 ),
1914 HummockValue::put(Bytes::from_static(b"hit")),
1915 )];
1916 kvs2.sort_by(|(k1, _), (k2, _)| k1.cmp(k2));
1917 let (_, mut sst2): (crate::hummock::sstable_store::TableHolder, SstableInfo) =
1918 gen_test_sstable_with_table_ids(
1919 opts,
1920 12,
1921 kvs2.into_iter(),
1922 sstable_store.clone(),
1923 vec![table_id.as_raw_id()],
1924 )
1925 .await;
1926 let mut sst2_inner = sst2.get_inner();
1927 sst2_inner.vnode_statistics = Some(VnodeStatistics::from_map(BTreeMap::from_iter([(
1928 vnode,
1929 (
1930 UserKey::new(table_id, gen_key_from_bytes(vnode, b"aa")),
1931 UserKey::new(table_id, gen_key_from_bytes(vnode, b"zz")),
1932 ),
1933 )])));
1934 sst2 = sst2_inner.into();
1935
1936 let committed =
1937 build_version_from_sstables(table_id, vec![sst1, sst2], PbLevelType::Overlapping);
1938
1939 let table_key = gen_key_from_bytes(vnode, b"mm");
1940 let result = reader
1941 .get(
1942 table_key,
1943 epoch,
1944 table_id,
1945 TableOption::default(),
1946 ReadOptions::default(),
1947 (vec![], vec![], committed),
1948 |_k, v| Ok(Bytes::copy_from_slice(v)),
1949 )
1950 .await
1951 .unwrap();
1952 flush_local_metrics_for_test();
1953
1954 assert_eq!(result, Some(Bytes::from_static(b"hit")));
1955 let (checked_after, pruned_after) = vnode_prune_counts(&metrics, table_id, "get");
1956 assert_eq!(checked_before + 2, checked_after);
1958 assert_eq!(pruned_before + 1, pruned_after);
1959 }
1960}