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 = (
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        // We assume that the table id read must exist in the sstable, otherwise it is a fatal error.
122        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        // Prefetch never extends beyond the block window already clipped by the scan hard bound.
189        if start_idx + 1 < self.block_end_idx_exclusive {
190            self.preload_end_block_idx = self.block_end_idx_exclusive;
191        }
192    }
193
194    /// Seeks to a block, and then seeks to the key if `seek_key` is given.
195    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        // When all data are in block cache, it is highly possible that this iterator will stay on a
208        // worker thread for a full time. Therefore, we use tokio's unstable API consume_budget to
209        // do cooperative scheduling.
210        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        // Maybe the previous preload stream breaks on some cached block, so here we can try to preload some data again
218        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                    // compare by version comparator
328                    // Note: we are comparing against the `smallest_key` of the `block`, thus the
329                    // partition point should be `prev(<=)` instead of `<`.
330                    FullKey::decode(&block_meta.smallest_key).le(&key)
331                })
332                .saturating_sub(1) // considering the boundary of 0
333    }
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            // seek to next block
344            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        // seek_idx will update the current block iter state
371        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            // seek to next block
386            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        // We should have at least 10 blocks, so that sstable iterator test could cover more code
446        // path.
447        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        // Build remote sstable
471        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        // We should have at least 10 blocks, so that sstable iterator test could cover more code
476        // path.
477        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        // We should have at least 10 blocks, so that sstable iterator test could cover more code
489        // path.
490        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        // We seek and access all the keys in random order
502        for i in all_key_to_test {
503            sstable_iter.seek(test_key_of(i).to_ref()).await.unwrap();
504            // sstable_iter.next().await.unwrap();
505            let key = sstable_iter.key();
506            assert_eq!(key, test_key_of(i).to_ref());
507        }
508
509        // Seek to key #500 and start iterating.
510        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        // Seek to < first key
519        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        // Seek to > last key
533        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        // Seek to non-existing key
546        for idx in 1..TEST_KEYS_COUNT {
547            // Seek to the previous key of each existing key. e.g.,
548            // Our key space is `key_test_00000`, `key_test_00002`, `key_test_00004`, ...
549            // And we seek to `key_test_00001` (will produce `key_test_00002`), `key_test_00003`
550            // (will produce `key_test_00004`).
551            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        // when upload data is successful, but upload meta is fail and delete is fail
578        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            // test key_range right
791            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}