Skip to main content

risingwave_storage/hummock/sstable/
forward_sstable_iterator.rs

1// Copyright 2022 RisingWave Labs
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use 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
40/// Iterates on a sstable.
41pub struct SstableIterator {
42    /// The iterator of the current block.
43    block_iter: Option<BlockIterator>,
44
45    /// Current block index.
46    cur_idx: usize,
47
48    preload_stream: Option<Box<dyn BlockStream>>,
49    /// Reference to the sst
50    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    // The readable block window is first restricted by table IDs and then by `scan_end_user_key`.
59    // Included/Excluded bounds set its exclusive block end; outer iterators still filter keys
60    // precisely inside the boundary block.
61    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        // We assume that the table id read must exist in the sstable, otherwise it is a fatal error.
131        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        // Prefetch never extends beyond the block window already clipped by the scan hard bound.
198        if start_idx + 1 < self.block_end_idx_exclusive {
199            self.preload_end_block_idx = self.block_end_idx_exclusive;
200        }
201    }
202
203    /// Seeks to a block, and then seeks to the key if `seek_key` is given.
204    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        // When all data are in block cache, it is highly possible that this iterator will stay on a
217        // worker thread for a full time. Therefore, we use tokio's unstable API consume_budget to
218        // do cooperative scheduling.
219        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        // Maybe the previous preload stream breaks on some cached block, so here we can try to preload some data again
227        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                    // compare by version comparator
337                    // Note: we are comparing against the `smallest_key` of the `block`, thus the
338                    // partition point should be `prev(<=)` instead of `<`.
339                    FullKey::decode(&block_meta.smallest_key).le(&key)
340                })
341                .saturating_sub(1) // considering the boundary of 0
342    }
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            // seek to next block
353            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        // seek_idx will update the current block iter state
380        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            // seek to next block
395            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        // We should have at least 10 blocks, so that sstable iterator test could cover more code
455        // path.
456        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        // Build remote sstable
480        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        // We should have at least 10 blocks, so that sstable iterator test could cover more code
485        // path.
486        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        // We should have at least 10 blocks, so that sstable iterator test could cover more code
498        // path.
499        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        // We seek and access all the keys in random order
511        for i in all_key_to_test {
512            sstable_iter.seek(test_key_of(i).to_ref()).await.unwrap();
513            // sstable_iter.next().await.unwrap();
514            let key = sstable_iter.key();
515            assert_eq!(key, test_key_of(i).to_ref());
516        }
517
518        // Seek to key #500 and start iterating.
519        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        // Seek to < first key
528        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        // Seek to > last key
542        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        // Seek to non-existing key
555        for idx in 1..TEST_KEYS_COUNT {
556            // Seek to the previous key of each existing key. e.g.,
557            // Our key space is `key_test_00000`, `key_test_00002`, `key_test_00004`, ...
558            // And we seek to `key_test_00001` (will produce `key_test_00002`), `key_test_00003`
559            // (will produce `key_test_00004`).
560            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        // when upload data is successful, but upload meta is fail and delete is fail
587        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            // test key_range right
847            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}