Skip to main content

risingwave_storage/hummock/compactor/
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::cmp::Ordering;
16use std::collections::HashSet;
17use std::ops::Range;
18use std::sync::atomic::AtomicU64;
19use std::sync::{Arc, atomic};
20use std::time::Instant;
21
22use await_tree::{InstrumentAwait, SpanExt};
23use risingwave_common::catalog::TableId;
24use risingwave_hummock_sdk::KeyComparator;
25use risingwave_hummock_sdk::key::FullKey;
26use risingwave_hummock_sdk::key_range::KeyRange;
27use risingwave_hummock_sdk::sstable_info::SstableInfo;
28
29use crate::hummock::compactor::block_stream::SstableBlockStream;
30use crate::hummock::compactor::task_progress::TaskProgress;
31use crate::hummock::iterator::{Forward, HummockIterator, ValueMeta};
32use crate::hummock::sstable_store::SstableStoreRef;
33use crate::hummock::value::HummockValue;
34use crate::hummock::{Block, BlockHolder, BlockIterator, BlockMeta, HummockResult, TableHolder};
35use crate::monitor::StoreLocalStatistic;
36
37const PROGRESS_KEY_INTERVAL: usize = 100;
38
39/// Iterates over the KV-pairs in a selected range of an SST while downloading its blocks.
40pub struct SstableStreamIterator {
41    block_stream: SstableBlockStream,
42    /// Iterates over the KV-pairs of the current block.
43    block_iter: Option<BlockIterator>,
44    /// Counts the time used for IO.
45    stats_ptr: Arc<AtomicU64>,
46
47    /// Table ids used to decide which blocks should be read from this SST.
48    read_table_ids: HashSet<TableId>,
49    task_progress: Arc<TaskProgress>,
50
51    // key range cache
52    key_range_left: FullKey<Vec<u8>>,
53    key_range_right: FullKey<Vec<u8>>,
54    key_range_right_exclusive: bool,
55}
56
57impl SstableStreamIterator {
58    /// Restricts the selected blocks to the virtual SST's key range before streaming them.
59    pub fn new(
60        sstable: TableHolder,
61        block_metas_range: Range<usize>,
62        sstable_info: SstableInfo,
63        stats: &StoreLocalStatistic,
64        task_progress: Arc<TaskProgress>,
65        sstable_store: SstableStoreRef,
66        max_io_retry_times: usize,
67    ) -> Self {
68        let read_table_ids = HashSet::from_iter(sstable_info.table_ids.iter().copied());
69        // Further filter the block_metas_range with sstable_info.key_range
70        // This is necessary when the SST is split into multiple SstableInfo with different key ranges
71        let block_metas_range = {
72            let block_metas = &sstable.meta.block_metas[block_metas_range.clone()];
73            let inner_range =
74                filter_block_metas(block_metas, &read_table_ids, sstable_info.key_range.clone());
75            // Adjust the range to be relative to the original block_metas
76            (block_metas_range.start + inner_range.start)
77                ..(block_metas_range.start + inner_range.end)
78        };
79
80        let key_range_left = FullKey::decode(&sstable_info.key_range.left).to_vec();
81        let key_range_right = FullKey::decode(&sstable_info.key_range.right).to_vec();
82        let key_range_right_exclusive = sstable_info.key_range.right_exclusive;
83
84        task_progress.inc_num_pending_read_io();
85        Self {
86            block_stream: SstableBlockStream::new(
87                sstable,
88                block_metas_range,
89                sstable_info,
90                sstable_store,
91                max_io_retry_times,
92            ),
93            block_iter: None,
94            stats_ptr: stats.remote_io_time.clone(),
95            read_table_ids,
96            task_progress,
97            key_range_left,
98            key_range_right,
99            key_range_right_exclusive,
100        }
101    }
102
103    async fn prune_from_valid_block_iter(&mut self) -> HummockResult<()> {
104        while let Some(block_iter) = self.block_iter.as_mut() {
105            if self.read_table_ids.contains(&block_iter.table_id()) {
106                return Ok(());
107            } else {
108                self.next_block().await?;
109            }
110        }
111        Ok(())
112    }
113
114    /// Initialises the iterator by moving it to the first KV-pair in the stream's first block where
115    /// key >= `seek_key`. If that block does not contain such a KV-pair, the iterator continues to
116    /// the first KV-pair of the next block. If `seek_key` is not given, the iterator will move to
117    /// the very first KV-pair of the stream's first block.
118    pub async fn seek(&mut self, seek_key: Option<FullKey<&[u8]>>) -> HummockResult<()> {
119        // Load first block.
120        self.next_block().await?;
121
122        // We assume that a block always contains at least one KV pair. Subsequently, if
123        // `next_block()` loads a new block (i.e., `block_iter` is not `None`), then `block_iter` is
124        // also valid and pointing on the block's first KV-pair.
125
126        let seek_key = if let Some(seek_key) = seek_key {
127            if seek_key.cmp(&self.key_range_left.to_ref()).is_lt() {
128                Some(self.key_range_left.to_ref())
129            } else {
130                Some(seek_key)
131            }
132        } else {
133            Some(self.key_range_left.to_ref())
134        };
135
136        if let (Some(block_iter), Some(seek_key)) = (self.block_iter.as_mut(), seek_key) {
137            block_iter.seek(seek_key);
138
139            if !block_iter.is_valid() {
140                // `seek_key` is larger than everything in the first block.
141                self.next_block().await?;
142            }
143        }
144
145        self.prune_from_valid_block_iter().await?;
146        Ok(())
147    }
148
149    /// Loads and decodes the next block, or clears the iterator at the end of the stream.
150    async fn next_block(&mut self) -> HummockResult<()> {
151        let now = Instant::now();
152        let _time_stat = scopeguard::guard(self.stats_ptr.clone(), |stats_ptr: Arc<AtomicU64>| {
153            let add = (now.elapsed().as_secs_f64() * 1000.0).ceil();
154            stats_ptr.fetch_add(add as u64, atomic::Ordering::Relaxed);
155        });
156        self.block_iter = match self.block_stream.next_block().await? {
157            Some((buf, uncompressed_size)) => {
158                // Decode errors are terminal and must not recreate the I/O stream.
159                let block = Box::new(Block::decode(buf, uncompressed_size)?);
160                let mut iter = BlockIterator::new(BlockHolder::from_owned_block(block));
161                iter.seek_to_first();
162                Some(iter)
163            }
164            None => None,
165        };
166        Ok(())
167    }
168
169    /// Moves to the next KV-pair in the table. Assumes that the current position is valid. Even if
170    /// the next position is invalid, the function returns `Ok(())`.
171    ///
172    /// Do not use `next()` to initialise the iterator (i.e. do not use it to find the first
173    /// KV-pair). Instead, use `seek()`. Afterwards, use `next()` to reach the second KV-pair and
174    /// onwards.
175    pub async fn next(&mut self) -> HummockResult<()> {
176        if !self.is_valid() {
177            return Ok(());
178        }
179
180        let block_iter = self.block_iter.as_mut().expect("no block iter");
181        block_iter.next();
182        if !block_iter.is_valid() {
183            self.next_block().await?;
184            self.prune_from_valid_block_iter().await?;
185        }
186
187        if !self.is_valid() {
188            return Ok(());
189        }
190
191        // Check if we need to skip the block.
192        let key = self
193            .block_iter
194            .as_ref()
195            .unwrap_or_else(|| panic!("no block iter sstinfo={}", self.sst_debug_info()))
196            .key();
197
198        if self.exceed_key_range_right(key) {
199            self.block_iter = None;
200        }
201
202        Ok(())
203    }
204
205    pub fn key(&self) -> FullKey<&[u8]> {
206        let key = self
207            .block_iter
208            .as_ref()
209            .unwrap_or_else(|| panic!("no block iter sstinfo={}", self.sst_debug_info()))
210            .key();
211
212        assert!(
213            !self.exceed_key_range_left(key),
214            "key {:?} key_range_left {:?}",
215            key,
216            self.key_range_left.to_ref()
217        );
218
219        assert!(
220            !self.exceed_key_range_right(key),
221            "key {:?} key_range_right {:?} key_range_right_exclusive {}",
222            key,
223            self.key_range_right.to_ref(),
224            self.key_range_right_exclusive
225        );
226
227        key
228    }
229
230    pub fn value(&self) -> HummockValue<&[u8]> {
231        let raw_value = self
232            .block_iter
233            .as_ref()
234            .unwrap_or_else(|| panic!("no block iter sstinfo={}", self.sst_debug_info()))
235            .value();
236        HummockValue::from_slice(raw_value)
237            .unwrap_or_else(|_| panic!("decode error sstinfo={}", self.sst_debug_info()))
238    }
239
240    pub fn is_valid(&self) -> bool {
241        // True iff block_iter exists and is valid.
242        self.block_iter.as_ref().is_some_and(|i| i.is_valid())
243    }
244
245    fn sst_debug_info(&self) -> String {
246        let sstable_info = &self.block_stream.sstable_info;
247        format!(
248            "object_id={}, sst_id={}, meta_offset={}, table_ids={:?}",
249            sstable_info.object_id,
250            sstable_info.sst_id,
251            sstable_info.meta_offset,
252            sstable_info.table_ids
253        )
254    }
255
256    fn exceed_key_range_left(&self, key: FullKey<&[u8]>) -> bool {
257        key.cmp(&self.key_range_left.to_ref()).is_lt()
258    }
259
260    fn exceed_key_range_right(&self, key: FullKey<&[u8]>) -> bool {
261        if self.key_range_right_exclusive {
262            key.cmp(&self.key_range_right.to_ref()).is_ge()
263        } else {
264            key.cmp(&self.key_range_right.to_ref()).is_gt()
265        }
266    }
267}
268
269impl Drop for SstableStreamIterator {
270    fn drop(&mut self) {
271        self.task_progress.dec_num_pending_read_io()
272    }
273}
274
275/// Iterates over the KV-pairs of a given list of SSTs. The key-ranges of these SSTs are assumed to
276/// be consecutive and non-overlapping.
277pub struct ConcatSstableIterator {
278    /// **CAUTION:** `key_range` is used for optimization. It doesn't guarantee value returned by
279    /// the iterator is in this range.
280    key_range: KeyRange,
281
282    /// The iterator of the current table.
283    sstable_iter: Option<SstableStreamIterator>,
284
285    /// Current table index.
286    cur_idx: usize,
287
288    /// All non-overlapping tables.
289    sstables: Vec<SstableInfo>,
290
291    sstable_store: SstableStoreRef,
292
293    stats: StoreLocalStatistic,
294    task_progress: Arc<TaskProgress>,
295    max_io_retry_times: usize,
296}
297
298impl ConcatSstableIterator {
299    /// Caller should make sure that `tables` are non-overlapping,
300    /// arranged in ascending order when it serves as a forward iterator,
301    /// and arranged in descending order when it serves as a backward iterator.
302    pub fn new(
303        sst_infos: Vec<SstableInfo>,
304        key_range: KeyRange,
305        sstable_store: SstableStoreRef,
306        task_progress: Arc<TaskProgress>,
307        max_io_retry_times: usize,
308    ) -> Self {
309        Self {
310            key_range,
311            sstable_iter: None,
312            cur_idx: 0,
313            sstables: sst_infos,
314            sstable_store,
315            task_progress,
316            stats: StoreLocalStatistic::default(),
317            max_io_retry_times,
318        }
319    }
320
321    #[cfg(test)]
322    pub fn for_test(
323        sst_infos: Vec<SstableInfo>,
324        key_range: KeyRange,
325        sstable_store: SstableStoreRef,
326    ) -> Self {
327        Self::new(
328            sst_infos,
329            key_range,
330            sstable_store,
331            Arc::new(TaskProgress::default()),
332            0,
333        )
334    }
335
336    /// Resets the iterator, loads the specified SST, and seeks in that SST to `seek_key` if given.
337    async fn seek_idx(
338        &mut self,
339        idx: usize,
340        seek_key: Option<FullKey<&[u8]>>,
341    ) -> HummockResult<()> {
342        self.sstable_iter.take();
343        let mut seek_key: Option<FullKey<&[u8]>> = match (seek_key, self.key_range.left.is_empty())
344        {
345            (Some(seek_key), false) => match seek_key.cmp(&FullKey::decode(&self.key_range.left)) {
346                Ordering::Less | Ordering::Equal => Some(FullKey::decode(&self.key_range.left)),
347                Ordering::Greater => Some(seek_key),
348            },
349            (Some(seek_key), true) => Some(seek_key),
350            (None, true) => None,
351            (None, false) => Some(FullKey::decode(&self.key_range.left)),
352        };
353
354        self.cur_idx = idx;
355        while self.cur_idx < self.sstables.len() {
356            let table_info = &self.sstables[self.cur_idx];
357            let read_table_ids = HashSet::from_iter(table_info.table_ids.iter().copied());
358            if read_table_ids.is_empty() {
359                self.cur_idx += 1;
360                seek_key = None;
361                continue;
362            }
363            let sstable = self
364                .sstable_store
365                .sstable(table_info, &mut self.stats)
366                .instrument_await("stream_iter_sstable".verbose())
367                .await?;
368
369            let filter_key_range = match seek_key {
370                Some(seek_key) => {
371                    KeyRange::new(seek_key.encode().into(), self.key_range.right.clone())
372                }
373                None => self.key_range.clone(),
374            };
375
376            let block_metas_range =
377                filter_block_metas(&sstable.meta.block_metas, &read_table_ids, filter_key_range);
378
379            if !block_metas_range.is_empty() {
380                let mut sstable_iter = SstableStreamIterator::new(
381                    sstable,
382                    block_metas_range,
383                    table_info.clone(),
384                    &self.stats,
385                    self.task_progress.clone(),
386                    self.sstable_store.clone(),
387                    self.max_io_retry_times,
388                );
389                sstable_iter.seek(seek_key).await?;
390
391                if sstable_iter.is_valid() {
392                    self.sstable_iter = Some(sstable_iter);
393                    return Ok(());
394                }
395            }
396
397            self.cur_idx += 1;
398            seek_key = None;
399        }
400        Ok(())
401    }
402}
403
404impl HummockIterator for ConcatSstableIterator {
405    type Direction = Forward;
406
407    async fn next(&mut self) -> HummockResult<()> {
408        let sstable_iter = self.sstable_iter.as_mut().expect("no table iter");
409
410        // Does just calling `next()` suffice?
411        sstable_iter.next().await?;
412        if sstable_iter.is_valid() {
413            Ok(())
414        } else {
415            // No, seek to next table.
416            self.seek_idx(self.cur_idx + 1, None).await?;
417            Ok(())
418        }
419    }
420
421    fn key(&self) -> FullKey<&[u8]> {
422        self.sstable_iter.as_ref().expect("no table iter").key()
423    }
424
425    fn value(&self) -> HummockValue<&[u8]> {
426        self.sstable_iter.as_ref().expect("no table iter").value()
427    }
428
429    fn is_valid(&self) -> bool {
430        self.sstable_iter.as_ref().is_some_and(|i| i.is_valid())
431    }
432
433    async fn rewind(&mut self) -> HummockResult<()> {
434        self.seek_idx(0, None).await
435    }
436
437    /// Resets the iterator and seeks to the first position where the stored key >= `key`.
438    async fn seek<'a>(&'a mut self, key: FullKey<&'a [u8]>) -> HummockResult<()> {
439        let seek_key = if self.key_range.left.is_empty() {
440            key
441        } else {
442            match key.cmp(&FullKey::decode(&self.key_range.left)) {
443                Ordering::Less | Ordering::Equal => FullKey::decode(&self.key_range.left),
444                Ordering::Greater => key,
445            }
446        };
447        let table_idx = self.sstables.partition_point(|table| {
448            // We use the maximum key of an SST for the search. That way, we guarantee that the
449            // resulting SST contains either that key or the next-larger KV-pair. Subsequently,
450            // we avoid calling `seek_idx()` twice if the determined SST does not contain `key`.
451
452            // Note that we need to use `<` instead of `<=` to ensure that all keys in an SST
453            // (including its max. key) produce the same search result.
454            let max_sst_key = &table.key_range.right;
455            FullKey::decode(max_sst_key).cmp(&seek_key) == Ordering::Less
456        });
457
458        self.seek_idx(table_idx, Some(key)).await
459    }
460
461    fn collect_local_statistic(&self, stats: &mut StoreLocalStatistic) {
462        stats.add(&self.stats)
463    }
464
465    fn value_meta(&self) -> ValueMeta {
466        let iter = self.sstable_iter.as_ref().expect("no table iter");
467        // The stream cursor is absolute and points past the currently decoded block.
468        assert!(iter.block_iter.is_some());
469        ValueMeta {
470            object_id: Some(iter.block_stream.sstable_info.object_id),
471            block_id: Some((iter.block_stream.next_block_index() - 1) as u64),
472        }
473    }
474}
475
476pub struct MonitoredCompactorIterator<I> {
477    inner: I,
478    task_progress: Arc<TaskProgress>,
479
480    processed_key_num: usize,
481}
482
483impl<I: HummockIterator<Direction = Forward>> MonitoredCompactorIterator<I> {
484    pub fn new(inner: I, task_progress: Arc<TaskProgress>) -> Self {
485        Self {
486            inner,
487            task_progress,
488            processed_key_num: 0,
489        }
490    }
491}
492
493impl<I: HummockIterator<Direction = Forward>> HummockIterator for MonitoredCompactorIterator<I> {
494    type Direction = Forward;
495
496    async fn next(&mut self) -> HummockResult<()> {
497        self.inner.next().await?;
498        self.processed_key_num += 1;
499
500        if self.processed_key_num.is_multiple_of(PROGRESS_KEY_INTERVAL) {
501            self.task_progress
502                .inc_progress_key(PROGRESS_KEY_INTERVAL as _);
503        }
504
505        Ok(())
506    }
507
508    fn key(&self) -> FullKey<&[u8]> {
509        self.inner.key()
510    }
511
512    fn value(&self) -> HummockValue<&[u8]> {
513        self.inner.value()
514    }
515
516    fn is_valid(&self) -> bool {
517        self.inner.is_valid()
518    }
519
520    async fn rewind(&mut self) -> HummockResult<()> {
521        self.processed_key_num = 0;
522        self.inner.rewind().await?;
523        Ok(())
524    }
525
526    async fn seek<'a>(&'a mut self, key: FullKey<&'a [u8]>) -> HummockResult<()> {
527        self.processed_key_num = 0;
528        self.inner.seek(key).await?;
529        Ok(())
530    }
531
532    fn collect_local_statistic(&self, stats: &mut StoreLocalStatistic) {
533        self.inner.collect_local_statistic(stats)
534    }
535
536    fn value_meta(&self) -> ValueMeta {
537        self.inner.value_meta()
538    }
539}
540
541pub(crate) fn filter_block_metas(
542    block_metas: &[BlockMeta],
543    read_table_ids: &HashSet<TableId>,
544    key_range: KeyRange,
545) -> Range<usize> {
546    if block_metas.is_empty() {
547        return 0..0;
548    }
549
550    let mut start_index = if key_range.left.is_empty() {
551        0
552    } else {
553        // start_index points to the greatest block whose smallest_key <= seek_key.
554        block_metas
555            .partition_point(|block| {
556                KeyComparator::compare_encoded_full_key(&key_range.left, &block.smallest_key)
557                    != Ordering::Less
558            })
559            .saturating_sub(1)
560    };
561
562    let mut end_index = if key_range.right.is_empty() {
563        block_metas.len()
564    } else {
565        let ret = block_metas.partition_point(|block| {
566            KeyComparator::compare_encoded_full_key(&block.smallest_key, &key_range.right)
567                != Ordering::Greater
568        });
569
570        if ret == 0 {
571            // not found
572            return 0..0;
573        }
574
575        ret
576    }
577    .saturating_sub(1);
578
579    // Skip blocks that are not in the SST read table ids.
580    while start_index <= end_index {
581        let start_block_table_id = block_metas[start_index].table_id();
582        if read_table_ids.contains(&start_block_table_id) {
583            break;
584        }
585
586        // skip this table_id
587        let old_start_index = start_index;
588        let block_metas_to_search = &block_metas[start_index..=end_index];
589
590        start_index += block_metas_to_search
591            .partition_point(|block_meta| block_meta.table_id() == start_block_table_id);
592
593        if old_start_index == start_index {
594            // no more blocks with the same table_id
595            break;
596        }
597    }
598
599    while start_index <= end_index {
600        let end_block_table_id = block_metas[end_index].table_id();
601        if read_table_ids.contains(&end_block_table_id) {
602            break;
603        }
604
605        let old_end_index = end_index;
606        let block_metas_to_search = &block_metas[start_index..=end_index];
607
608        end_index = start_index
609            + block_metas_to_search
610                .partition_point(|block_meta| block_meta.table_id() < end_block_table_id)
611                .saturating_sub(1);
612
613        if end_index == old_end_index {
614            // no more blocks with the same table_id
615            break;
616        }
617    }
618
619    if start_index > end_index {
620        return 0..0;
621    }
622
623    start_index..(end_index + 1)
624}
625
626#[cfg(test)]
627mod tests {
628    use std::cmp::Ordering;
629    use std::collections::HashSet;
630
631    use risingwave_common::catalog::TableId;
632    use risingwave_common::util::epoch::test_epoch;
633    use risingwave_hummock_sdk::key::{FullKey, FullKeyTracker, next_full_key, prev_full_key};
634    use risingwave_hummock_sdk::key_range::KeyRange;
635    use risingwave_hummock_sdk::sstable_info::{SstableInfo, SstableInfoInner};
636
637    use crate::hummock::BlockMeta;
638    use crate::hummock::compactor::ConcatSstableIterator;
639    use crate::hummock::iterator::test_utils::mock_sstable_store;
640    use crate::hummock::iterator::{HummockIterator, MergeIterator};
641    use crate::hummock::test_utils::{
642        TEST_KEYS_COUNT, default_builder_opt_for_test, gen_test_sstable_info,
643        gen_test_sstable_with_table_ids, test_key_of, test_value_of,
644    };
645    use crate::hummock::value::HummockValue;
646
647    #[tokio::test]
648    async fn test_concat_iterator() {
649        let sstable_store = mock_sstable_store().await;
650        let mut table_infos = vec![];
651        for object_id in 0..3 {
652            let start_index = object_id * TEST_KEYS_COUNT;
653            let end_index = (object_id + 1) * TEST_KEYS_COUNT;
654            let table_info = gen_test_sstable_info(
655                default_builder_opt_for_test(),
656                object_id as u64,
657                (start_index..end_index)
658                    .map(|i| (test_key_of(i), HummockValue::put(test_value_of(i)))),
659                sstable_store.clone(),
660            )
661            .await;
662            table_infos.push(table_info);
663        }
664        let start_index = 5000;
665        let end_index = 25000;
666
667        let kr = KeyRange::new(
668            test_key_of(start_index).encode().into(),
669            test_key_of(end_index).encode().into(),
670        );
671        let mut iter =
672            ConcatSstableIterator::for_test(table_infos.clone(), kr.clone(), sstable_store.clone());
673        iter.seek(FullKey::decode(&kr.left)).await.unwrap();
674
675        for idx in start_index..end_index {
676            let key = iter.key();
677            let val = iter.value();
678            assert_eq!(key, test_key_of(idx).to_ref(), "failed at {}", idx);
679            assert_eq!(
680                val.into_user_value().unwrap(),
681                test_value_of(idx).as_slice()
682            );
683            iter.next().await.unwrap();
684        }
685
686        // seek non-overlap range
687        let kr = KeyRange::new(
688            test_key_of(30000).encode().into(),
689            test_key_of(40000).encode().into(),
690        );
691        let mut iter =
692            ConcatSstableIterator::for_test(table_infos.clone(), kr.clone(), sstable_store.clone());
693        iter.seek(FullKey::decode(&kr.left)).await.unwrap();
694        assert!(!iter.is_valid());
695        let kr = KeyRange::new(
696            test_key_of(start_index).encode().into(),
697            test_key_of(40000).encode().into(),
698        );
699        let mut iter =
700            ConcatSstableIterator::for_test(table_infos.clone(), kr.clone(), sstable_store.clone());
701        iter.seek(FullKey::decode(&kr.left)).await.unwrap();
702        for idx in start_index..30000 {
703            let key = iter.key();
704            let val = iter.value();
705            assert_eq!(key, test_key_of(idx).to_ref(), "failed at {}", idx);
706            assert_eq!(
707                val.into_user_value().unwrap(),
708                test_value_of(idx).as_slice()
709            );
710            iter.next().await.unwrap();
711        }
712        assert!(!iter.is_valid());
713
714        // Test seek. Result is dominated by given seek key rather than key range.
715        let kr = KeyRange::new(
716            test_key_of(0).encode().into(),
717            test_key_of(40000).encode().into(),
718        );
719        let mut iter =
720            ConcatSstableIterator::for_test(table_infos.clone(), kr.clone(), sstable_store.clone());
721        iter.seek(test_key_of(10000).to_ref()).await.unwrap();
722        assert!(iter.is_valid() && iter.cur_idx == 1 && iter.key() == test_key_of(10000).to_ref());
723        iter.seek(test_key_of(10001).to_ref()).await.unwrap();
724        assert!(iter.is_valid() && iter.cur_idx == 1 && iter.key() == test_key_of(10001).to_ref());
725        iter.seek(test_key_of(9999).to_ref()).await.unwrap();
726        assert!(iter.is_valid() && iter.cur_idx == 0 && iter.key() == test_key_of(9999).to_ref());
727        iter.seek(test_key_of(1).to_ref()).await.unwrap();
728        assert!(iter.is_valid() && iter.cur_idx == 0 && iter.key() == test_key_of(1).to_ref());
729        iter.seek(test_key_of(29999).to_ref()).await.unwrap();
730        assert!(iter.is_valid() && iter.cur_idx == 2 && iter.key() == test_key_of(29999).to_ref());
731        iter.seek(test_key_of(30000).to_ref()).await.unwrap();
732        assert!(!iter.is_valid());
733
734        // Test seek. Result is dominated by key range rather than given seek key.
735        let kr = KeyRange::new(
736            test_key_of(6000).encode().into(),
737            test_key_of(16000).encode().into(),
738        );
739        let mut iter =
740            ConcatSstableIterator::for_test(table_infos.clone(), kr.clone(), sstable_store.clone());
741        iter.seek(test_key_of(17000).to_ref()).await.unwrap();
742        assert!(!iter.is_valid());
743        iter.seek(test_key_of(1).to_ref()).await.unwrap();
744        assert!(iter.is_valid() && iter.cur_idx == 0 && iter.key() == FullKey::decode(&kr.left));
745    }
746
747    #[tokio::test]
748    async fn test_concat_iterator_seek_idx() {
749        let sstable_store = mock_sstable_store().await;
750        let mut table_infos = vec![];
751        for object_id in 0..3 {
752            let start_index = object_id * TEST_KEYS_COUNT + TEST_KEYS_COUNT / 2;
753            let end_index = (object_id + 1) * TEST_KEYS_COUNT;
754            let table_info = gen_test_sstable_info(
755                default_builder_opt_for_test(),
756                object_id as u64,
757                (start_index..end_index)
758                    .map(|i| (test_key_of(i), HummockValue::put(test_value_of(i)))),
759                sstable_store.clone(),
760            )
761            .await;
762            table_infos.push(table_info);
763        }
764
765        // Test seek_idx. Result is dominated by given seek key rather than key range.
766        let kr = KeyRange::new(
767            test_key_of(0).encode().into(),
768            test_key_of(40000).encode().into(),
769        );
770        let mut iter =
771            ConcatSstableIterator::for_test(table_infos.clone(), kr.clone(), sstable_store.clone());
772        let sst = sstable_store
773            .sstable(&iter.sstables[0], &mut iter.stats)
774            .await
775            .unwrap();
776        let block_metas = &sst.meta.block_metas;
777        let block_1_smallest_key = block_metas[1].smallest_key.clone();
778        let block_2_smallest_key = block_metas[2].smallest_key.clone();
779        // Use block_1_smallest_key as seek key and result in the first KV of block 1.
780        let seek_key = block_1_smallest_key.clone();
781        iter.seek_idx(0, Some(FullKey::decode(&seek_key)))
782            .await
783            .unwrap();
784        assert!(iter.is_valid() && iter.key() == FullKey::decode(block_1_smallest_key.as_slice()));
785        // Use prev_full_key(block_1_smallest_key) as seek key and result in the first KV of block
786        // 1.
787        let seek_key = prev_full_key(block_1_smallest_key.as_slice());
788        iter.seek_idx(0, Some(FullKey::decode(&seek_key)))
789            .await
790            .unwrap();
791        assert!(iter.is_valid() && iter.key() == FullKey::decode(block_1_smallest_key.as_slice()));
792        iter.next().await.unwrap();
793        let block_1_second_key = iter.key().to_vec();
794        // Use a big enough seek key and result in invalid iterator.
795        let seek_key = test_key_of(30001);
796        iter.seek_idx(table_infos.len() - 1, Some(seek_key.to_ref()))
797            .await
798            .unwrap();
799        assert!(!iter.is_valid());
800
801        // Test seek_idx. Result is dominated by key range rather than given seek key.
802        let kr = KeyRange::new(
803            next_full_key(&block_1_smallest_key).into(),
804            prev_full_key(&block_2_smallest_key).into(),
805        );
806        let mut iter =
807            ConcatSstableIterator::for_test(table_infos.clone(), kr.clone(), sstable_store.clone());
808        // Use block_2_smallest_key as seek key and result in invalid iterator.
809        let seek_key = FullKey::decode(&block_2_smallest_key);
810        assert!(seek_key.cmp(&FullKey::decode(&kr.right)) == Ordering::Greater);
811        iter.seek_idx(0, Some(seek_key)).await.unwrap();
812        assert!(!iter.is_valid());
813        // Use a small enough seek key and result in the second KV of block 1.
814        let seek_key = test_key_of(0).encode();
815        iter.seek_idx(0, Some(FullKey::decode(&seek_key)))
816            .await
817            .unwrap();
818        assert!(iter.is_valid());
819        assert_eq!(iter.key(), block_1_second_key.to_ref());
820
821        // Use None seek key and result in the second KV of block 1.
822        iter.seek_idx(0, None).await.unwrap();
823        assert!(iter.is_valid());
824        assert_eq!(iter.key(), block_1_second_key.to_ref());
825    }
826
827    #[tokio::test]
828    async fn test_filter_block_metas() {
829        use crate::hummock::compactor::iterator::filter_block_metas;
830
831        {
832            let block_metas = Vec::default();
833
834            let ret = filter_block_metas(&block_metas, &HashSet::default(), KeyRange::default());
835
836            assert!(ret.is_empty());
837        }
838
839        {
840            let block_metas = vec![
841                BlockMeta {
842                    smallest_key: FullKey::for_test(TableId::new(1), Vec::default(), 0).encode(),
843                    ..Default::default()
844                },
845                BlockMeta {
846                    smallest_key: FullKey::for_test(TableId::new(2), Vec::default(), 0).encode(),
847                    ..Default::default()
848                },
849                BlockMeta {
850                    smallest_key: FullKey::for_test(TableId::new(3), Vec::default(), 0).encode(),
851                    ..Default::default()
852                },
853            ];
854
855            let ret = filter_block_metas(
856                &block_metas,
857                &HashSet::from_iter(vec![1_u32.into(), 2.into(), 3.into()].into_iter()),
858                KeyRange::default(),
859            );
860            let ret = &block_metas[ret];
861
862            assert_eq!(3, ret.len());
863            assert_eq!(
864                1,
865                FullKey::decode(&ret[0].smallest_key)
866                    .user_key
867                    .table_id
868                    .as_raw_id()
869            );
870            assert_eq!(
871                3,
872                FullKey::decode(&ret[2].smallest_key)
873                    .user_key
874                    .table_id
875                    .as_raw_id()
876            );
877        }
878
879        {
880            let block_metas = vec![
881                BlockMeta {
882                    smallest_key: FullKey::for_test(TableId::new(1), Vec::default(), 0).encode(),
883                    ..Default::default()
884                },
885                BlockMeta {
886                    smallest_key: FullKey::for_test(TableId::new(2), Vec::default(), 0).encode(),
887                    ..Default::default()
888                },
889                BlockMeta {
890                    smallest_key: FullKey::for_test(TableId::new(3), Vec::default(), 0).encode(),
891                    ..Default::default()
892                },
893            ];
894
895            let ret = filter_block_metas(
896                &block_metas,
897                &HashSet::from_iter(vec![2_u32.into(), 3.into()].into_iter()),
898                KeyRange::default(),
899            );
900            let ret = &block_metas[ret];
901
902            assert_eq!(2, ret.len());
903            assert_eq!(
904                2,
905                FullKey::decode(&ret[0].smallest_key)
906                    .user_key
907                    .table_id
908                    .as_raw_id()
909            );
910            assert_eq!(
911                3,
912                FullKey::decode(&ret[1].smallest_key)
913                    .user_key
914                    .table_id
915                    .as_raw_id()
916            );
917        }
918
919        {
920            let block_metas = vec![
921                BlockMeta {
922                    smallest_key: FullKey::for_test(TableId::new(1), Vec::default(), 0).encode(),
923                    ..Default::default()
924                },
925                BlockMeta {
926                    smallest_key: FullKey::for_test(TableId::new(2), Vec::default(), 0).encode(),
927                    ..Default::default()
928                },
929                BlockMeta {
930                    smallest_key: FullKey::for_test(TableId::new(3), Vec::default(), 0).encode(),
931                    ..Default::default()
932                },
933            ];
934
935            let ret = filter_block_metas(
936                &block_metas,
937                &HashSet::from_iter(vec![1_u32.into(), 2_u32.into()].into_iter()),
938                KeyRange::default(),
939            );
940            let ret = &block_metas[ret];
941
942            assert_eq!(2, ret.len());
943            assert_eq!(
944                1,
945                FullKey::decode(&ret[0].smallest_key)
946                    .user_key
947                    .table_id
948                    .as_raw_id()
949            );
950            assert_eq!(
951                2,
952                FullKey::decode(&ret[1].smallest_key)
953                    .user_key
954                    .table_id
955                    .as_raw_id()
956            );
957        }
958
959        {
960            let block_metas = vec![
961                BlockMeta {
962                    smallest_key: FullKey::for_test(TableId::new(1), Vec::default(), 0).encode(),
963                    ..Default::default()
964                },
965                BlockMeta {
966                    smallest_key: FullKey::for_test(TableId::new(2), Vec::default(), 0).encode(),
967                    ..Default::default()
968                },
969                BlockMeta {
970                    smallest_key: FullKey::for_test(TableId::new(3), Vec::default(), 0).encode(),
971                    ..Default::default()
972                },
973            ];
974            let ret = filter_block_metas(
975                &block_metas,
976                &HashSet::from_iter(vec![2_u32.into()].into_iter()),
977                KeyRange::default(),
978            );
979            let ret = &block_metas[ret];
980
981            assert_eq!(1, ret.len());
982            assert_eq!(
983                2,
984                FullKey::decode(&ret[0].smallest_key)
985                    .user_key
986                    .table_id
987                    .as_raw_id()
988            );
989        }
990
991        {
992            let block_metas = vec![
993                BlockMeta {
994                    smallest_key: FullKey::for_test(TableId::new(1), Vec::default(), 0).encode(),
995                    ..Default::default()
996                },
997                BlockMeta {
998                    smallest_key: FullKey::for_test(TableId::new(1), Vec::default(), 0).encode(),
999                    ..Default::default()
1000                },
1001                BlockMeta {
1002                    smallest_key: FullKey::for_test(TableId::new(1), Vec::default(), 0).encode(),
1003                    ..Default::default()
1004                },
1005                BlockMeta {
1006                    smallest_key: FullKey::for_test(TableId::new(2), Vec::default(), 0).encode(),
1007                    ..Default::default()
1008                },
1009                BlockMeta {
1010                    smallest_key: FullKey::for_test(TableId::new(3), Vec::default(), 0).encode(),
1011                    ..Default::default()
1012                },
1013            ];
1014            let ret = filter_block_metas(
1015                &block_metas,
1016                &HashSet::from_iter(vec![2_u32.into()].into_iter()),
1017                KeyRange::default(),
1018            );
1019            let ret = &block_metas[ret];
1020
1021            assert_eq!(1, ret.len());
1022            assert_eq!(
1023                2,
1024                FullKey::decode(&ret[0].smallest_key)
1025                    .user_key
1026                    .table_id
1027                    .as_raw_id()
1028            );
1029        }
1030
1031        {
1032            let block_metas = vec![
1033                BlockMeta {
1034                    smallest_key: FullKey::for_test(TableId::new(1), Vec::default(), 0).encode(),
1035                    ..Default::default()
1036                },
1037                BlockMeta {
1038                    smallest_key: FullKey::for_test(TableId::new(2), Vec::default(), 0).encode(),
1039                    ..Default::default()
1040                },
1041                BlockMeta {
1042                    smallest_key: FullKey::for_test(TableId::new(3), Vec::default(), 0).encode(),
1043                    ..Default::default()
1044                },
1045                BlockMeta {
1046                    smallest_key: FullKey::for_test(TableId::new(3), Vec::default(), 0).encode(),
1047                    ..Default::default()
1048                },
1049                BlockMeta {
1050                    smallest_key: FullKey::for_test(TableId::new(3), Vec::default(), 0).encode(),
1051                    ..Default::default()
1052                },
1053            ];
1054
1055            let ret = filter_block_metas(
1056                &block_metas,
1057                &HashSet::from_iter(vec![2_u32.into()].into_iter()),
1058                KeyRange::default(),
1059            );
1060            let ret = &block_metas[ret];
1061
1062            assert_eq!(1, ret.len());
1063            assert_eq!(
1064                2,
1065                FullKey::decode(&ret[0].smallest_key)
1066                    .user_key
1067                    .table_id
1068                    .as_raw_id()
1069            );
1070        }
1071
1072        {
1073            let block_metas = vec![
1074                BlockMeta {
1075                    smallest_key: FullKey::for_test(TableId::new(1), Vec::default(), 0).encode(),
1076                    ..Default::default()
1077                },
1078                BlockMeta {
1079                    smallest_key: FullKey::for_test(TableId::new(2), Vec::default(), 0).encode(),
1080                    ..Default::default()
1081                },
1082                BlockMeta {
1083                    smallest_key: FullKey::for_test(TableId::new(3), Vec::default(), 0).encode(),
1084                    ..Default::default()
1085                },
1086            ];
1087
1088            let ret = filter_block_metas(
1089                &block_metas,
1090                &HashSet::from_iter(vec![1_u32.into(), 3_u32.into()].into_iter()),
1091                KeyRange::default(),
1092            );
1093            let ret = &block_metas[ret];
1094
1095            assert_eq!(3, ret.len());
1096            assert_eq!(
1097                1,
1098                FullKey::decode(&ret[0].smallest_key)
1099                    .user_key
1100                    .table_id
1101                    .as_raw_id()
1102            );
1103            assert_eq!(
1104                2,
1105                FullKey::decode(&ret[1].smallest_key)
1106                    .user_key
1107                    .table_id
1108                    .as_raw_id()
1109            );
1110            assert_eq!(
1111                3,
1112                FullKey::decode(&ret[2].smallest_key)
1113                    .user_key
1114                    .table_id
1115                    .as_raw_id()
1116            );
1117        }
1118    }
1119
1120    #[tokio::test]
1121    async fn test_iterator_same_obj() {
1122        let sstable_store = mock_sstable_store().await;
1123
1124        let table_info = gen_test_sstable_info(
1125            default_builder_opt_for_test(),
1126            1_u64,
1127            (1..10000).map(|i| (test_key_of(i), HummockValue::put(test_value_of(i)))),
1128            sstable_store.clone(),
1129        )
1130        .await;
1131
1132        let split_key = test_key_of(5000).encode();
1133        let sst_1: SstableInfo = SstableInfoInner {
1134            key_range: KeyRange {
1135                left: table_info.key_range.left.clone(),
1136                right: split_key.clone().into(),
1137                right_exclusive: true,
1138            },
1139            ..table_info.get_inner()
1140        }
1141        .into();
1142
1143        let total_key_count = sst_1.total_key_count;
1144        let sst_2: SstableInfo = SstableInfoInner {
1145            sst_id: sst_1.sst_id + 1,
1146            key_range: KeyRange {
1147                left: split_key.clone().into(),
1148                right: table_info.key_range.right.clone(),
1149                right_exclusive: table_info.key_range.right_exclusive,
1150            },
1151            ..table_info.get_inner()
1152        }
1153        .into();
1154
1155        {
1156            // test concate
1157            let mut full_key_tracker = FullKeyTracker::<Vec<u8>>::new(FullKey::default());
1158
1159            let mut iter = ConcatSstableIterator::for_test(
1160                vec![sst_1.clone(), sst_2.clone()],
1161                KeyRange::default(),
1162                sstable_store.clone(),
1163            );
1164
1165            iter.rewind().await.unwrap();
1166
1167            let mut key_count = 0;
1168            while iter.is_valid() {
1169                let is_new_user_key = full_key_tracker.observe(iter.key());
1170                assert!(is_new_user_key);
1171                key_count += 1;
1172                iter.next().await.unwrap();
1173            }
1174
1175            assert_eq!(total_key_count, key_count);
1176        }
1177
1178        {
1179            let mut full_key_tracker = FullKeyTracker::<Vec<u8>>::new(FullKey::default());
1180            let concat_1 = ConcatSstableIterator::for_test(
1181                vec![sst_1.clone()],
1182                KeyRange::default(),
1183                sstable_store.clone(),
1184            );
1185
1186            let concat_2 = ConcatSstableIterator::for_test(
1187                vec![sst_2.clone()],
1188                KeyRange::default(),
1189                sstable_store.clone(),
1190            );
1191
1192            let mut key_count = 0;
1193            let mut iter = MergeIterator::for_compactor(vec![concat_1, concat_2]);
1194            iter.rewind().await.unwrap();
1195            while iter.is_valid() {
1196                full_key_tracker.observe(iter.key());
1197                key_count += 1;
1198                iter.next().await.unwrap();
1199            }
1200            assert_eq!(total_key_count, key_count);
1201        }
1202    }
1203
1204    #[tokio::test]
1205    async fn test_concat_iterator_skips_hole_table_blocks() {
1206        let sstable_store = mock_sstable_store().await;
1207
1208        let key_1 = FullKey::for_test(TableId::new(1), b"a".to_vec(), test_epoch(1));
1209        let key_2 = FullKey::for_test(TableId::new(2), b"b".to_vec(), test_epoch(1));
1210        let key_3 = FullKey::for_test(TableId::new(3), b"c".to_vec(), test_epoch(1));
1211        let kv_pairs = vec![
1212            (key_1.clone(), HummockValue::put(b"value-1".to_vec())),
1213            (key_2.clone(), HummockValue::put(b"value-2".to_vec())),
1214            (key_3.clone(), HummockValue::put(b"value-3".to_vec())),
1215        ];
1216
1217        let (_sstable, table_info) = gen_test_sstable_with_table_ids(
1218            default_builder_opt_for_test(),
1219            10,
1220            kv_pairs.into_iter(),
1221            sstable_store.clone(),
1222            vec![1, 2, 3],
1223        )
1224        .await;
1225
1226        let table_info: SstableInfo = SstableInfoInner {
1227            table_ids: vec![1.into(), 3.into()],
1228            ..table_info.get_inner()
1229        }
1230        .into();
1231
1232        let mut iter = ConcatSstableIterator::for_test(
1233            vec![table_info.clone()],
1234            KeyRange::default(),
1235            sstable_store.clone(),
1236        );
1237        iter.rewind().await.unwrap();
1238        assert!(iter.is_valid());
1239        assert_eq!(iter.key(), key_1.to_ref());
1240
1241        iter.next().await.unwrap();
1242        assert!(iter.is_valid());
1243        assert_eq!(iter.key(), key_3.to_ref());
1244
1245        iter.next().await.unwrap();
1246        assert!(!iter.is_valid());
1247
1248        let mut iter =
1249            ConcatSstableIterator::for_test(vec![table_info], KeyRange::default(), sstable_store);
1250        iter.seek(key_2.to_ref()).await.unwrap();
1251        assert!(iter.is_valid());
1252        assert_eq!(iter.key(), key_3.to_ref());
1253    }
1254}