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