1use std::ops::Bound::*;
16use std::sync::Arc;
17
18use await_tree::{InstrumentAwait, SpanExt};
19use risingwave_hummock_sdk::key::FullKey;
20use risingwave_hummock_sdk::sstable_info::SstableInfo;
21use sync_point::sync_point;
22use thiserror_ext::AsReport;
23
24use super::super::{HummockResult, HummockValue};
25use crate::hummock::block_stream::BlockStream;
26use crate::hummock::iterator::{Forward, HummockIterator, ValueMeta};
27use crate::hummock::sstable::SstableIteratorReadOptions;
28use crate::hummock::{BlockIterator, SstableStoreRef, TableHolder};
29use crate::monitor::StoreLocalStatistic;
30
31pub trait SstableIteratorType: HummockIterator + 'static {
32 fn create(
33 sstable: TableHolder,
34 sstable_store: SstableStoreRef,
35 read_options: Arc<SstableIteratorReadOptions>,
36 sstable_info_ref: &SstableInfo,
37 ) -> Self;
38}
39
40pub struct SstableIterator {
42 block_iter: Option<BlockIterator>,
44
45 cur_idx: usize,
47
48 preload_stream: Option<Box<dyn BlockStream>>,
49 pub sst: TableHolder,
51 preload_end_block_idx: usize,
52 preload_retry_times: usize,
53
54 sstable_store: SstableStoreRef,
55 stats: StoreLocalStatistic,
56 options: Arc<SstableIteratorReadOptions>,
57
58 block_start_idx_inclusive: usize,
62 block_end_idx_exclusive: usize,
63}
64
65impl SstableIterator {
66 pub fn new(
67 sstable: TableHolder,
68 sstable_store: SstableStoreRef,
69 options: Arc<SstableIteratorReadOptions>,
70 sstable_info_ref: &SstableInfo,
71 ) -> Self {
72 let mut block_start_idx_inclusive = 0;
73 let mut block_end_idx_exclusive = sstable.meta.block_metas.len();
74 assert!(
75 !sstable_info_ref.table_ids.is_empty(),
76 "SstableIterator: SST {} (object {}) has empty table_ids",
77 sstable_info_ref.sst_id,
78 sstable_info_ref.object_id,
79 );
80 let read_table_id_range = (
81 *sstable_info_ref.table_ids.first().unwrap(),
82 *sstable_info_ref.table_ids.last().unwrap(),
83 );
84 assert!(
85 read_table_id_range.0 <= read_table_id_range.1,
86 "invalid table id range {} - {}",
87 read_table_id_range.0,
88 read_table_id_range.1
89 );
90 let block_meta_count = sstable.meta.block_metas.len();
91 assert!(block_meta_count > 0);
92 assert!(
93 sstable.meta.block_metas[0].table_id() <= read_table_id_range.0,
94 "table id {} not found table_ids in block_meta {:?}",
95 read_table_id_range.0,
96 sstable
97 .meta
98 .block_metas
99 .iter()
100 .map(|meta| meta.table_id())
101 .collect::<Vec<_>>()
102 );
103 assert!(
104 sstable.meta.block_metas[block_meta_count - 1].table_id() >= read_table_id_range.1,
105 "table id {} not found table_ids in block_meta {:?}",
106 read_table_id_range.1,
107 sstable
108 .meta
109 .block_metas
110 .iter()
111 .map(|meta| meta.table_id())
112 .collect::<Vec<_>>()
113 );
114
115 while block_start_idx_inclusive < block_meta_count
116 && sstable.meta.block_metas[block_start_idx_inclusive].table_id()
117 < read_table_id_range.0
118 {
119 block_start_idx_inclusive += 1;
120 }
121 assert!(
123 block_start_idx_inclusive < block_meta_count,
124 "table id {} not found table_ids in block_meta {:?}",
125 read_table_id_range.0,
126 sstable
127 .meta
128 .block_metas
129 .iter()
130 .map(|meta| meta.table_id())
131 .collect::<Vec<_>>()
132 );
133
134 while block_end_idx_exclusive > block_start_idx_inclusive
135 && sstable.meta.block_metas[block_end_idx_exclusive - 1].table_id()
136 > read_table_id_range.1
137 {
138 block_end_idx_exclusive -= 1;
139 }
140 assert!(
141 block_end_idx_exclusive > block_start_idx_inclusive,
142 "block_end_idx_exclusive {} <= block_start_idx_inclusive {} block_meta_count {}",
143 block_end_idx_exclusive,
144 block_start_idx_inclusive,
145 block_meta_count
146 );
147
148 if let Some(end_bound) = options.scan_end_user_key.as_ref() {
149 let block_metas =
150 &sstable.meta.block_metas[block_start_idx_inclusive..block_end_idx_exclusive];
151 let range_end_idx_exclusive = match end_bound {
152 Unbounded => block_metas.len(),
153 Included(end_key) => block_metas.partition_point(|block_meta| {
154 FullKey::decode(&block_meta.smallest_key).user_key <= end_key.as_ref()
155 }),
156 Excluded(end_key) => block_metas.partition_point(|block_meta| {
157 FullKey::decode(&block_meta.smallest_key).user_key < end_key.as_ref()
158 }),
159 };
160 block_end_idx_exclusive = block_start_idx_inclusive + range_end_idx_exclusive;
161 }
162
163 Self {
164 block_iter: None,
165 cur_idx: 0,
166 preload_stream: None,
167 sst: sstable,
168 sstable_store,
169 stats: StoreLocalStatistic::default(),
170 options,
171 preload_end_block_idx: 0,
172 preload_retry_times: 0,
173 block_start_idx_inclusive,
174 block_end_idx_exclusive,
175 }
176 }
177
178 fn init_block_prefetch_range(&mut self, start_idx: usize) {
179 assert!(
180 start_idx >= self.block_start_idx_inclusive && start_idx < self.block_end_idx_exclusive
181 );
182
183 self.preload_end_block_idx = 0;
184 if !self.options.prefetch {
185 return;
186 }
187
188 if start_idx + 1 < self.block_end_idx_exclusive {
190 self.preload_end_block_idx = self.block_end_idx_exclusive;
191 }
192 }
193
194 async fn seek_idx(
196 &mut self,
197 idx: usize,
198 seek_key: Option<FullKey<&[u8]>>,
199 ) -> HummockResult<()> {
200 tracing::debug!(
201 target: "events::storage::sstable::block_seek",
202 "table iterator seek: sstable_object_id = {}, block_id = {}",
203 self.sst.id,
204 idx,
205 );
206
207 tokio::task::consume_budget().await;
211
212 let mut hit_cache = false;
213 if idx >= self.block_end_idx_exclusive {
214 self.block_iter = None;
215 return Ok(());
216 }
217 if self.preload_stream.is_none() && idx + 1 < self.preload_end_block_idx {
219 match self
220 .sstable_store
221 .prefetch_blocks(
222 &self.sst,
223 idx,
224 self.preload_end_block_idx,
225 self.options.cache_policy,
226 &mut self.stats,
227 )
228 .instrument_await("prefetch_blocks".verbose())
229 .await
230 {
231 Ok(preload_stream) => self.preload_stream = Some(preload_stream),
232 Err(e) => {
233 tracing::warn!(error = %e.as_report(), "failed to create stream for prefetch data, fall back to block get")
234 }
235 }
236 }
237
238 if self
239 .preload_stream
240 .as_ref()
241 .map(|preload_stream| preload_stream.next_block_index() <= idx)
242 .unwrap_or(false)
243 {
244 while let Some(preload_stream) = self.preload_stream.as_mut() {
245 let mut ret = Ok(());
246 while preload_stream.next_block_index() < idx {
247 if let Err(e) = preload_stream.next_block().await {
248 ret = Err(e);
249 break;
250 }
251 }
252 assert_eq!(preload_stream.next_block_index(), idx);
253 if ret.is_ok() {
254 match preload_stream.next_block().await {
255 Ok(Some(block)) => {
256 hit_cache = true;
257 self.block_iter = Some(BlockIterator::new(block));
258 break;
259 }
260 Ok(None) => {
261 self.preload_stream.take();
262 }
263 Err(e) => {
264 self.preload_stream.take();
265 ret = Err(e);
266 }
267 }
268 } else {
269 self.preload_stream.take();
270 }
271 if self.preload_stream.is_none() && idx + 1 < self.preload_end_block_idx {
272 if let Err(e) = ret {
273 tracing::warn!(error = %e.as_report(), "recreate stream because the connection to remote storage has closed");
274 if self.preload_retry_times >= self.options.max_preload_retry_times {
275 break;
276 }
277 self.preload_retry_times += 1;
278 }
279
280 match self
281 .sstable_store
282 .prefetch_blocks(
283 &self.sst,
284 idx,
285 self.preload_end_block_idx,
286 self.options.cache_policy,
287 &mut self.stats,
288 )
289 .instrument_await("prefetch_blocks".verbose())
290 .await
291 {
292 Ok(stream) => {
293 self.preload_stream = Some(stream);
294 }
295 Err(e) => {
296 tracing::warn!(error = %e.as_report(), "failed to recreate stream meet IO error");
297 break;
298 }
299 }
300 }
301 }
302 }
303 if !hit_cache {
304 let block = self
305 .sstable_store
306 .get(&self.sst, idx, self.options.cache_policy, &mut self.stats)
307 .await?;
308 self.block_iter = Some(BlockIterator::new(block));
309 };
310 let block_iter = self.block_iter.as_mut().unwrap();
311 if let Some(key) = seek_key {
312 block_iter.seek(key);
313 } else {
314 block_iter.seek_to_first();
315 }
316
317 self.cur_idx = idx;
318
319 Ok(())
320 }
321
322 fn calculate_block_idx_by_key(&self, key: FullKey<&[u8]>) -> usize {
323 self.block_start_idx_inclusive
324 + self.sst.meta.block_metas
325 [self.block_start_idx_inclusive..self.block_end_idx_exclusive]
326 .partition_point(|block_meta| {
327 FullKey::decode(&block_meta.smallest_key).le(&key)
331 })
332 .saturating_sub(1) }
334}
335
336impl HummockIterator for SstableIterator {
337 type Direction = Forward;
338
339 async fn next(&mut self) -> HummockResult<()> {
340 self.stats.total_key_count += 1;
341 let block_iter = self.block_iter.as_mut().expect("no block iter");
342 if !block_iter.try_next() {
343 self.seek_idx(self.cur_idx + 1, None).await?;
345 }
346
347 Ok(())
348 }
349
350 fn key(&self) -> FullKey<&[u8]> {
351 self.block_iter.as_ref().expect("no block iter").key()
352 }
353
354 fn value(&self) -> HummockValue<&[u8]> {
355 let raw_value = self.block_iter.as_ref().expect("no block iter").value();
356
357 HummockValue::from_slice(raw_value).expect("decode error")
358 }
359
360 fn is_valid(&self) -> bool {
361 self.block_iter.as_ref().is_some_and(|i| i.is_valid())
362 }
363
364 async fn rewind(&mut self) -> HummockResult<()> {
365 if self.block_start_idx_inclusive >= self.block_end_idx_exclusive {
366 self.block_iter = None;
367 return Ok(());
368 }
369 self.init_block_prefetch_range(self.block_start_idx_inclusive);
370 self.seek_idx(self.block_start_idx_inclusive, None).await?;
372 Ok(())
373 }
374
375 async fn seek<'a>(&'a mut self, key: FullKey<&'a [u8]>) -> HummockResult<()> {
376 if self.block_start_idx_inclusive >= self.block_end_idx_exclusive {
377 self.block_iter = None;
378 return Ok(());
379 }
380 let block_idx = self.calculate_block_idx_by_key(key);
381 self.init_block_prefetch_range(block_idx);
382
383 self.seek_idx(block_idx, Some(key)).await?;
384 if !self.is_valid() {
385 sync_point!("SSTABLE_ITERATOR::SEEK::BEFORE_NEXT_BLOCK");
387 self.seek_idx(block_idx + 1, None).await?;
388 }
389 Ok(())
390 }
391
392 fn collect_local_statistic(&self, stats: &mut StoreLocalStatistic) {
393 stats.add(&self.stats);
394 }
395
396 fn value_meta(&self) -> ValueMeta {
397 ValueMeta {
398 object_id: Some(self.sst.id),
399 block_id: Some(self.cur_idx as _),
400 }
401 }
402}
403
404impl SstableIteratorType for SstableIterator {
405 fn create(
406 sstable: TableHolder,
407 sstable_store: SstableStoreRef,
408 options: Arc<SstableIteratorReadOptions>,
409 sstable_info_ref: &SstableInfo,
410 ) -> Self {
411 SstableIterator::new(sstable, sstable_store, options, sstable_info_ref)
412 }
413}
414
415#[cfg(test)]
416mod tests {
417 use std::collections::Bound;
418
419 use bytes::Bytes;
420 use foyer::Hint;
421 use itertools::Itertools;
422 use rand::prelude::*;
423 use rand::rng as thread_rng;
424 use risingwave_common::catalog::TableId;
425 use risingwave_common::hash::VirtualNode;
426 use risingwave_common::util::epoch::test_epoch;
427 use risingwave_hummock_sdk::EpochWithGap;
428 use risingwave_hummock_sdk::key::{TableKey, UserKey};
429 use risingwave_hummock_sdk::sstable_info::{SstableInfo, SstableInfoInner};
430
431 use super::*;
432 use crate::assert_bytes_eq;
433 use crate::hummock::CachePolicy;
434 use crate::hummock::iterator::test_utils::mock_sstable_store;
435 use crate::hummock::test_utils::{
436 TEST_KEYS_COUNT, default_builder_opt_for_test, gen_default_test_sstable,
437 gen_test_sstable_info, gen_test_sstable_with_table_ids, test_key_of, test_value_of,
438 };
439
440 async fn inner_test_forward_iterator(
441 sstable_store: SstableStoreRef,
442 handle: TableHolder,
443 sstable_info: SstableInfo,
444 ) {
445 let mut sstable_iter = SstableIterator::create(
448 handle,
449 sstable_store,
450 Arc::new(SstableIteratorReadOptions::default()),
451 &sstable_info,
452 );
453 let mut cnt = 0;
454 sstable_iter.rewind().await.unwrap();
455
456 while sstable_iter.is_valid() {
457 let key = sstable_iter.key();
458 let value = sstable_iter.value();
459 assert_eq!(key, test_key_of(cnt).to_ref());
460 assert_bytes_eq!(value.into_user_value().unwrap(), test_value_of(cnt));
461 cnt += 1;
462 sstable_iter.next().await.unwrap();
463 }
464
465 assert_eq!(cnt, TEST_KEYS_COUNT);
466 }
467
468 #[tokio::test]
469 async fn test_table_iterator() {
470 let sstable_store = mock_sstable_store().await;
472 let (sstable, sstable_info) =
473 gen_default_test_sstable(default_builder_opt_for_test(), 0, sstable_store.clone())
474 .await;
475 assert!(sstable.meta.block_metas.len() > 10);
478
479 inner_test_forward_iterator(sstable_store.clone(), sstable, sstable_info).await;
480 }
481
482 #[tokio::test]
483 async fn test_table_seek() {
484 let sstable_store = mock_sstable_store().await;
485 let (sstable, sstable_info) =
486 gen_default_test_sstable(default_builder_opt_for_test(), 0, sstable_store.clone())
487 .await;
488 assert!(sstable.meta.block_metas.len() > 10);
491 let mut sstable_iter = SstableIterator::create(
492 sstable,
493 sstable_store,
494 Arc::new(SstableIteratorReadOptions::default()),
495 &sstable_info,
496 );
497 let mut all_key_to_test = (0..TEST_KEYS_COUNT).collect_vec();
498 let mut rng = thread_rng();
499 all_key_to_test.shuffle(&mut rng);
500
501 for i in all_key_to_test {
503 sstable_iter.seek(test_key_of(i).to_ref()).await.unwrap();
504 let key = sstable_iter.key();
506 assert_eq!(key, test_key_of(i).to_ref());
507 }
508
509 sstable_iter.seek(test_key_of(500).to_ref()).await.unwrap();
511 for i in 500..TEST_KEYS_COUNT {
512 let key = sstable_iter.key();
513 assert_eq!(key, test_key_of(i).to_ref());
514 sstable_iter.next().await.unwrap();
515 }
516 assert!(!sstable_iter.is_valid());
517
518 let smallest_key = FullKey::for_test(
520 TableId::default(),
521 [
522 VirtualNode::ZERO.to_be_bytes().as_slice(),
523 format!("key_aaaa_{:05}", 0).as_bytes(),
524 ]
525 .concat(),
526 test_epoch(233),
527 );
528 sstable_iter.seek(smallest_key.to_ref()).await.unwrap();
529 let key = sstable_iter.key();
530 assert_eq!(key, test_key_of(0).to_ref());
531
532 let largest_key = FullKey::for_test(
534 TableId::default(),
535 [
536 VirtualNode::ZERO.to_be_bytes().as_slice(),
537 format!("key_zzzz_{:05}", 0).as_bytes(),
538 ]
539 .concat(),
540 test_epoch(233),
541 );
542 sstable_iter.seek(largest_key.to_ref()).await.unwrap();
543 assert!(!sstable_iter.is_valid());
544
545 for idx in 1..TEST_KEYS_COUNT {
547 sstable_iter
552 .seek(
553 FullKey::for_test(
554 TableId::default(),
555 [
556 VirtualNode::ZERO.to_be_bytes().as_slice(),
557 format!("key_test_{:05}", idx * 2 - 1).as_bytes(),
558 ]
559 .concat(),
560 0,
561 )
562 .to_ref(),
563 )
564 .await
565 .unwrap();
566
567 let key = sstable_iter.key();
568 assert_eq!(key, test_key_of(idx).to_ref());
569 sstable_iter.next().await.unwrap();
570 }
571 assert!(!sstable_iter.is_valid());
572 }
573
574 #[tokio::test]
575 async fn test_prefetch_table_read() {
576 let sstable_store = mock_sstable_store().await;
577 let kv_iter =
579 (0..TEST_KEYS_COUNT).map(|i| (test_key_of(i), HummockValue::put(test_value_of(i))));
580 let sst_info = gen_test_sstable_info(
581 default_builder_opt_for_test(),
582 0,
583 kv_iter,
584 sstable_store.clone(),
585 )
586 .await;
587
588 let end_key = test_key_of(TEST_KEYS_COUNT);
589 let uk = UserKey::new(
590 end_key.user_key.table_id,
591 TableKey(Bytes::from(end_key.user_key.table_key.0)),
592 );
593 let options = Arc::new(SstableIteratorReadOptions {
594 cache_policy: CachePolicy::Fill(Hint::Normal),
595 scan_end_user_key: Some(Bound::Included(uk.clone())),
596 prefetch: true,
597 max_preload_retry_times: 0,
598 });
599 let mut stats = StoreLocalStatistic::default();
600 let mut sstable_iter = SstableIterator::create(
601 sstable_store.sstable(&sst_info, &mut stats).await.unwrap(),
602 sstable_store.clone(),
603 options.clone(),
604 &sst_info,
605 );
606 let mut cnt = 1000;
607 sstable_iter.seek(test_key_of(cnt).to_ref()).await.unwrap();
608 while sstable_iter.is_valid() {
609 let key = sstable_iter.key();
610 let value = sstable_iter.value();
611 assert_eq!(
612 key,
613 test_key_of(cnt).to_ref(),
614 "fail at {}, get key :{:?}",
615 cnt,
616 String::from_utf8(key.user_key.table_key.key_part().to_vec()).unwrap()
617 );
618 assert_bytes_eq!(value.into_user_value().unwrap(), test_value_of(cnt));
619 cnt += 1;
620 sstable_iter.next().await.unwrap();
621 }
622 assert_eq!(cnt, TEST_KEYS_COUNT);
623 let mut sstable_iter = SstableIterator::create(
624 sstable_store.sstable(&sst_info, &mut stats).await.unwrap(),
625 sstable_store,
626 options.clone(),
627 &sst_info,
628 );
629 let mut cnt = 1000;
630 sstable_iter.seek(test_key_of(cnt).to_ref()).await.unwrap();
631 while sstable_iter.is_valid() {
632 let key = sstable_iter.key();
633 let value = sstable_iter.value();
634 assert_eq!(key, test_key_of(cnt).to_ref());
635 assert_bytes_eq!(value.into_user_value().unwrap(), test_value_of(cnt));
636 cnt += 1;
637 sstable_iter.next().await.unwrap();
638 }
639 assert_eq!(cnt, TEST_KEYS_COUNT);
640 }
641
642 #[tokio::test]
643 async fn test_scan_end_with_prefetch_on_or_off() {
644 let sstable_store = mock_sstable_store().await;
645 let mut builder_options = default_builder_opt_for_test();
646 builder_options.block_capacity = 128;
647
648 let test_user_key = |table_id, idx| {
649 UserKey::new(
650 TableId::new(table_id),
651 TableKey(Bytes::from(
652 [
653 VirtualNode::ZERO.to_be_bytes().as_slice(),
654 format!("key_{idx:05}").as_bytes(),
655 ]
656 .concat(),
657 )),
658 )
659 };
660 let test_key = |table_id, idx| FullKey {
661 user_key: test_user_key(table_id, idx),
662 epoch_with_gap: EpochWithGap::new_from_epoch(test_epoch(1)),
663 };
664 let kv_pairs = (1..=3).flat_map(|table_id| {
665 (0..8).map(move |idx| {
666 (
667 test_key(table_id, idx),
668 HummockValue::put(Bytes::from(test_value_of(idx))),
669 )
670 })
671 });
672 let (sstable, sstable_info) = gen_test_sstable_with_table_ids(
673 builder_options,
674 10,
675 kv_pairs,
676 sstable_store.clone(),
677 vec![1, 2, 3],
678 )
679 .await;
680
681 let table_2_block_start = sstable
682 .meta
683 .block_metas
684 .partition_point(|block_meta| block_meta.table_id() < TableId::new(2));
685 let table_3_block_start = sstable
686 .meta
687 .block_metas
688 .partition_point(|block_meta| block_meta.table_id() < TableId::new(3));
689 assert!(table_2_block_start + 2 < table_3_block_start);
690
691 let table_3_start = UserKey::new(TableId::new(3), TableKey(Bytes::from_static(b"")));
692 for (case, prefetch) in [("prefetch off", false), ("prefetch on", true)] {
693 let options = Arc::new(SstableIteratorReadOptions {
694 cache_policy: CachePolicy::Disable,
695 scan_end_user_key: Some(Bound::Excluded(table_3_start.clone())),
696 prefetch,
697 max_preload_retry_times: 0,
698 });
699 let mut sstable_iter = SstableIterator::create(
700 sstable.clone(),
701 sstable_store.clone(),
702 options,
703 &sstable_info,
704 );
705
706 assert_eq!(sstable_iter.block_end_idx_exclusive, table_3_block_start);
707 sstable_iter.seek(test_key(2, 0).to_ref()).await.unwrap();
708 if prefetch {
709 assert_eq!(
710 sstable_iter.preload_end_block_idx, sstable_iter.block_end_idx_exclusive,
711 "{case}"
712 );
713 } else {
714 assert_eq!(sstable_iter.preload_end_block_idx, 0, "{case}");
715 }
716 let mut key_count = 0;
717 while sstable_iter.is_valid() {
718 assert_eq!(
719 sstable_iter.key().user_key.table_id,
720 TableId::new(2),
721 "{case}"
722 );
723 key_count += 1;
724 sstable_iter.next().await.unwrap();
725 }
726 assert_eq!(key_count, 8, "{case}");
727
728 let mut stats = StoreLocalStatistic::default();
729 sstable_iter.collect_local_statistic(&mut stats);
730 assert_eq!(
731 stats.cache_data_block_total + stats.cache_data_prefetch_block_count,
732 (table_3_block_start - table_2_block_start) as u64,
733 "{case}"
734 );
735 if prefetch {
736 assert!(stats.cache_data_prefetch_block_count > 0, "{case}");
737 } else {
738 assert_eq!(stats.cache_data_prefetch_block_count, 0, "{case}");
739 }
740 }
741
742 let mut table_2_sstable_info = sstable_info.get_inner();
743 table_2_sstable_info.table_ids = vec![TableId::new(2)];
744 let table_2_sstable_info = SstableInfo::from(table_2_sstable_info);
745 let table_2_start: UserKey<Bytes> =
746 FullKey::decode(&sstable.meta.block_metas[table_2_block_start].smallest_key)
747 .user_key
748 .copy_into();
749 let options = Arc::new(SstableIteratorReadOptions {
750 cache_policy: CachePolicy::Disable,
751 scan_end_user_key: Some(Bound::Excluded(table_2_start)),
752 prefetch: false,
753 max_preload_retry_times: 0,
754 });
755 let mut sstable_iter =
756 SstableIterator::create(sstable, sstable_store, options, &table_2_sstable_info);
757
758 assert_eq!(
759 sstable_iter.block_start_idx_inclusive,
760 sstable_iter.block_end_idx_exclusive
761 );
762 sstable_iter.seek(test_key(2, 0).to_ref()).await.unwrap();
763 assert!(!sstable_iter.is_valid());
764 let mut stats = StoreLocalStatistic::default();
765 sstable_iter.collect_local_statistic(&mut stats);
766 assert_eq!(stats.cache_data_block_total, 0);
767 assert_eq!(stats.cache_data_prefetch_block_count, 0);
768 }
769
770 #[tokio::test]
771 async fn test_read_table_id_range() {
772 {
773 let sstable_store = mock_sstable_store().await;
774 let (sstable, sstable_info) =
775 gen_default_test_sstable(default_builder_opt_for_test(), 0, sstable_store.clone())
776 .await;
777 let mut sstable_iter = SstableIterator::create(
778 sstable,
779 sstable_store.clone(),
780 Arc::new(SstableIteratorReadOptions::default()),
781 &sstable_info,
782 );
783 sstable_iter.rewind().await.unwrap();
784 assert!(sstable_iter.is_valid());
785 assert_eq!(sstable_iter.key(), test_key_of(0).to_ref());
786 }
787
788 {
789 let sstable_store = mock_sstable_store().await;
790 let k1 = {
792 let mut table_key = VirtualNode::ZERO.to_be_bytes().to_vec();
793 table_key.extend_from_slice(format!("key_test_{:05}", 1).as_bytes());
794 let uk = UserKey::for_test(TableId::from(1), table_key);
795 FullKey {
796 user_key: uk,
797 epoch_with_gap: EpochWithGap::new_from_epoch(test_epoch(1)),
798 }
799 };
800
801 let k2 = {
802 let mut table_key = VirtualNode::ZERO.to_be_bytes().to_vec();
803 table_key.extend_from_slice(format!("key_test_{:05}", 2).as_bytes());
804 let uk = UserKey::for_test(TableId::from(2), table_key);
805 FullKey {
806 user_key: uk,
807 epoch_with_gap: EpochWithGap::new_from_epoch(test_epoch(1)),
808 }
809 };
810
811 let k3 = {
812 let mut table_key = VirtualNode::ZERO.to_be_bytes().to_vec();
813 table_key.extend_from_slice(format!("key_test_{:05}", 3).as_bytes());
814 let uk = UserKey::for_test(TableId::from(3), table_key);
815 FullKey {
816 user_key: uk,
817 epoch_with_gap: EpochWithGap::new_from_epoch(test_epoch(1)),
818 }
819 };
820
821 {
822 let kv_pairs = vec![
823 (k1.clone(), HummockValue::put(test_value_of(1))),
824 (k2.clone(), HummockValue::put(test_value_of(2))),
825 (k3.clone(), HummockValue::put(test_value_of(3))),
826 ];
827
828 let (sstable, _sstable_info) = gen_test_sstable_with_table_ids(
829 default_builder_opt_for_test(),
830 10,
831 kv_pairs.into_iter(),
832 sstable_store.clone(),
833 vec![1, 2, 3],
834 )
835 .await;
836 let mut sstable_iter = SstableIterator::create(
837 sstable,
838 sstable_store.clone(),
839 Arc::new(SstableIteratorReadOptions::default()),
840 &SstableInfo::from(SstableInfoInner {
841 table_ids: vec![1.into(), 2.into(), 3.into()],
842 ..Default::default()
843 }),
844 );
845 sstable_iter.rewind().await.unwrap();
846 assert!(sstable_iter.is_valid());
847 assert!(sstable_iter.key().eq(&k1.to_ref()));
848
849 let mut cnt = 0;
850 let mut last_key = k1.clone();
851 while sstable_iter.is_valid() {
852 last_key = sstable_iter.key().to_vec();
853 cnt += 1;
854 sstable_iter.next().await.unwrap();
855 }
856
857 assert_eq!(3, cnt);
858 assert_eq!(last_key, k3.clone());
859 }
860
861 {
862 let kv_pairs = vec![
863 (k1.clone(), HummockValue::put(test_value_of(1))),
864 (k2.clone(), HummockValue::put(test_value_of(2))),
865 (k3.clone(), HummockValue::put(test_value_of(3))),
866 ];
867
868 let (sstable, _sstable_info) = gen_test_sstable_with_table_ids(
869 default_builder_opt_for_test(),
870 10,
871 kv_pairs.into_iter(),
872 sstable_store.clone(),
873 vec![1, 2, 3],
874 )
875 .await;
876
877 let mut sstable_iter = SstableIterator::create(
878 sstable,
879 sstable_store.clone(),
880 Arc::new(SstableIteratorReadOptions::default()),
881 &SstableInfo::from(SstableInfoInner {
882 table_ids: vec![1.into(), 2.into()],
883 ..Default::default()
884 }),
885 );
886 sstable_iter.rewind().await.unwrap();
887 assert!(sstable_iter.is_valid());
888 assert!(sstable_iter.key().eq(&k1.to_ref()));
889
890 let mut cnt = 0;
891 let mut last_key = k1.clone();
892 while sstable_iter.is_valid() {
893 last_key = sstable_iter.key().to_vec();
894 cnt += 1;
895 sstable_iter.next().await.unwrap();
896 }
897
898 assert_eq!(2, cnt);
899 assert_eq!(last_key, k2.clone());
900 }
901
902 {
903 let kv_pairs = vec![
904 (k1.clone(), HummockValue::put(test_value_of(1))),
905 (k2.clone(), HummockValue::put(test_value_of(2))),
906 (k3.clone(), HummockValue::put(test_value_of(3))),
907 ];
908
909 let (sstable, _sstable_info) = gen_test_sstable_with_table_ids(
910 default_builder_opt_for_test(),
911 10,
912 kv_pairs.into_iter(),
913 sstable_store.clone(),
914 vec![1, 2, 3],
915 )
916 .await;
917
918 let mut sstable_iter = SstableIterator::create(
919 sstable,
920 sstable_store.clone(),
921 Arc::new(SstableIteratorReadOptions::default()),
922 &SstableInfo::from(SstableInfoInner {
923 table_ids: vec![2.into(), 3.into()],
924 ..Default::default()
925 }),
926 );
927 sstable_iter.rewind().await.unwrap();
928 assert!(sstable_iter.is_valid());
929 assert!(sstable_iter.key().eq(&k2.to_ref()));
930
931 let mut cnt = 0;
932 let mut last_key = k1.clone();
933 while sstable_iter.is_valid() {
934 last_key = sstable_iter.key().to_vec();
935 cnt += 1;
936 sstable_iter.next().await.unwrap();
937 }
938
939 assert_eq!(2, cnt);
940 assert_eq!(last_key, k3.clone());
941 }
942
943 {
944 let kv_pairs = vec![
945 (k1.clone(), HummockValue::put(test_value_of(1))),
946 (k2.clone(), HummockValue::put(test_value_of(2))),
947 (k3.clone(), HummockValue::put(test_value_of(3))),
948 ];
949
950 let (sstable, _sstable_info) = gen_test_sstable_with_table_ids(
951 default_builder_opt_for_test(),
952 10,
953 kv_pairs.into_iter(),
954 sstable_store.clone(),
955 vec![1, 2, 3],
956 )
957 .await;
958
959 let mut sstable_iter = SstableIterator::create(
960 sstable,
961 sstable_store.clone(),
962 Arc::new(SstableIteratorReadOptions::default()),
963 &SstableInfo::from(SstableInfoInner {
964 table_ids: vec![2.into()],
965 ..Default::default()
966 }),
967 );
968 sstable_iter.rewind().await.unwrap();
969 assert!(sstable_iter.is_valid());
970 assert!(sstable_iter.key().eq(&k2.to_ref()));
971
972 let mut cnt = 0;
973 let mut last_key = k1.clone();
974 while sstable_iter.is_valid() {
975 last_key = sstable_iter.key().to_vec();
976 cnt += 1;
977 sstable_iter.next().await.unwrap();
978 }
979
980 assert_eq!(1, cnt);
981 assert_eq!(last_key, k2.clone());
982 }
983 }
984 }
985}