Skip to main content

risingwave_storage/hummock/iterator/
forward_concat.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 crate::hummock::SstableIterator;
16use crate::hummock::iterator::concat_inner::ConcatIteratorInner;
17
18/// Iterates on multiple non-overlapping tables.
19pub type ConcatIterator = ConcatIteratorInner<SstableIterator>;
20
21#[cfg(test)]
22mod tests {
23    use std::sync::Arc;
24
25    #[cfg(feature = "failpoints")]
26    use foyer::Hint;
27
28    use super::*;
29    #[cfg(feature = "failpoints")]
30    use crate::hummock::CachePolicy;
31    use crate::hummock::iterator::HummockIterator;
32    use crate::hummock::iterator::test_utils::{
33        TEST_KEYS_COUNT, default_builder_opt_for_test, gen_iterator_test_sstable_info,
34        iterator_test_key_of, iterator_test_value_of, mock_sstable_store,
35    };
36    use crate::hummock::sstable::SstableIteratorReadOptions;
37    #[cfg(feature = "failpoints")]
38    use crate::monitor::StoreLocalStatistic;
39
40    #[tokio::test]
41    async fn test_concat_iterator() {
42        let sstable_store = mock_sstable_store().await;
43        let table0 = gen_iterator_test_sstable_info(
44            0,
45            default_builder_opt_for_test(),
46            |x| x,
47            sstable_store.clone(),
48            TEST_KEYS_COUNT,
49        )
50        .await;
51        let table1 = gen_iterator_test_sstable_info(
52            1,
53            default_builder_opt_for_test(),
54            |x| TEST_KEYS_COUNT + x,
55            sstable_store.clone(),
56            TEST_KEYS_COUNT,
57        )
58        .await;
59        let table2 = gen_iterator_test_sstable_info(
60            2,
61            default_builder_opt_for_test(),
62            |x| TEST_KEYS_COUNT * 2 + x,
63            sstable_store.clone(),
64            TEST_KEYS_COUNT,
65        )
66        .await;
67        let mut iter = ConcatIterator::new(
68            vec![table0, table1, table2],
69            sstable_store,
70            Arc::new(SstableIteratorReadOptions::default()),
71        );
72        let mut i = 0;
73        iter.rewind().await.unwrap();
74
75        while iter.is_valid() {
76            let key = iter.key();
77            let val = iter.value();
78            assert_eq!(key, iterator_test_key_of(i).to_ref());
79            assert_eq!(
80                val.into_user_value().unwrap(),
81                iterator_test_value_of(i).as_slice()
82            );
83            i += 1;
84            iter.next().await.unwrap();
85            if i == TEST_KEYS_COUNT * 3 {
86                assert!(!iter.is_valid());
87                break;
88            }
89        }
90
91        iter.rewind().await.unwrap();
92        let key = iter.key();
93        let val = iter.value();
94        assert_eq!(key, iterator_test_key_of(0).to_ref());
95        assert_eq!(
96            val.into_user_value().unwrap(),
97            iterator_test_value_of(0).as_slice()
98        );
99    }
100
101    #[tokio::test]
102    async fn test_concat_seek() {
103        let sstable_store = mock_sstable_store().await;
104        let table0 = gen_iterator_test_sstable_info(
105            0,
106            default_builder_opt_for_test(),
107            |x| x,
108            sstable_store.clone(),
109            TEST_KEYS_COUNT,
110        )
111        .await;
112        let table1 = gen_iterator_test_sstable_info(
113            1,
114            default_builder_opt_for_test(),
115            |x| TEST_KEYS_COUNT + x,
116            sstable_store.clone(),
117            TEST_KEYS_COUNT,
118        )
119        .await;
120        let table2 = gen_iterator_test_sstable_info(
121            2,
122            default_builder_opt_for_test(),
123            |x| TEST_KEYS_COUNT * 2 + x,
124            sstable_store.clone(),
125            TEST_KEYS_COUNT,
126        )
127        .await;
128        let mut iter = ConcatIterator::new(
129            vec![table0, table1, table2],
130            sstable_store,
131            Arc::new(SstableIteratorReadOptions::default()),
132        );
133
134        iter.seek(iterator_test_key_of(TEST_KEYS_COUNT + 1).to_ref())
135            .await
136            .unwrap();
137
138        let key = iter.key();
139        let val = iter.value();
140        assert_eq!(key, iterator_test_key_of(TEST_KEYS_COUNT + 1).to_ref());
141        assert_eq!(
142            val.into_user_value().unwrap(),
143            iterator_test_value_of(TEST_KEYS_COUNT + 1).as_slice()
144        );
145
146        // Left edge case
147        iter.seek(iterator_test_key_of(0).to_ref()).await.unwrap();
148        let key = iter.key();
149        let val = iter.value();
150        assert_eq!(key, iterator_test_key_of(0).to_ref());
151        assert_eq!(
152            val.into_user_value().unwrap(),
153            iterator_test_value_of(0).as_slice()
154        );
155
156        // Right edge case
157        iter.seek(iterator_test_key_of(3 * TEST_KEYS_COUNT - 1).to_ref())
158            .await
159            .unwrap();
160
161        let key = iter.key();
162        let val = iter.value();
163        assert_eq!(key, iterator_test_key_of(3 * TEST_KEYS_COUNT - 1).to_ref());
164        assert_eq!(
165            val.into_user_value().unwrap(),
166            iterator_test_value_of(3 * TEST_KEYS_COUNT - 1).as_slice()
167        );
168
169        // Right overflow case
170        iter.seek(iterator_test_key_of(3 * TEST_KEYS_COUNT).to_ref())
171            .await
172            .unwrap();
173        assert!(!iter.is_valid());
174    }
175
176    #[tokio::test]
177    async fn test_concat_seek_not_exists() {
178        let sstable_store = mock_sstable_store().await;
179        let table0 = gen_iterator_test_sstable_info(
180            0,
181            default_builder_opt_for_test(),
182            |x| x * 2,
183            sstable_store.clone(),
184            TEST_KEYS_COUNT,
185        )
186        .await;
187        let table1 = gen_iterator_test_sstable_info(
188            1,
189            default_builder_opt_for_test(),
190            |x| (TEST_KEYS_COUNT + x) * 2,
191            sstable_store.clone(),
192            TEST_KEYS_COUNT,
193        )
194        .await;
195        let table2 = gen_iterator_test_sstable_info(
196            2,
197            default_builder_opt_for_test(),
198            |x| (2 * TEST_KEYS_COUNT + x) * 2,
199            sstable_store.clone(),
200            TEST_KEYS_COUNT,
201        )
202        .await;
203        let mut iter = ConcatIterator::new(
204            vec![table0, table1, table2],
205            sstable_store,
206            Arc::new(SstableIteratorReadOptions::default()),
207        );
208
209        iter.seek(iterator_test_key_of(TEST_KEYS_COUNT + 1).to_ref())
210            .await
211            .unwrap();
212
213        let key = iter.key();
214        let val = iter.value();
215        assert_eq!(key, iterator_test_key_of(TEST_KEYS_COUNT + 2).to_ref());
216        assert_eq!(
217            val.into_user_value().unwrap(),
218            iterator_test_value_of(TEST_KEYS_COUNT + 2).as_slice()
219        );
220
221        // seek the last of table1
222        iter.seek(iterator_test_key_of((TEST_KEYS_COUNT + 9) * 2 + 1).to_ref())
223            .await
224            .unwrap();
225
226        let key = iter.key();
227        let val = iter.value();
228        assert_eq!(key, iterator_test_key_of(TEST_KEYS_COUNT * 4).to_ref());
229        assert_eq!(
230            val.into_user_value().unwrap(),
231            iterator_test_value_of(TEST_KEYS_COUNT * 4).as_slice()
232        );
233    }
234
235    #[tokio::test]
236    #[cfg(feature = "failpoints")]
237    async fn test_concat_init_error_retains_current_iterator_and_candidate_stats() {
238        let sstable_store = mock_sstable_store().await;
239        let table0 = gen_iterator_test_sstable_info(
240            0,
241            default_builder_opt_for_test(),
242            |x| x,
243            sstable_store.clone(),
244            TEST_KEYS_COUNT,
245        )
246        .await;
247        let table1 = gen_iterator_test_sstable_info(
248            1,
249            default_builder_opt_for_test(),
250            |x| TEST_KEYS_COUNT + x,
251            sstable_store.clone(),
252            TEST_KEYS_COUNT,
253        )
254        .await;
255
256        // Warm the candidate's first block so its initialization records a cache hit.
257        let mut warmup_stats = StoreLocalStatistic::default();
258        let sstable = sstable_store
259            .sstable(&table1, &mut warmup_stats)
260            .await
261            .unwrap();
262        sstable_store
263            .get(
264                &sstable,
265                0,
266                CachePolicy::Fill(Hint::Normal),
267                &mut warmup_stats,
268            )
269            .await
270            .unwrap();
271        warmup_stats.discard();
272
273        let mut iter = ConcatIterator::new(
274            vec![table0, table1],
275            sstable_store,
276            Arc::new(SstableIteratorReadOptions::default()),
277        );
278        iter.rewind().await.unwrap();
279
280        let mut before = StoreLocalStatistic::default();
281        iter.collect_local_statistic(&mut before);
282
283        fail::cfg("concat_iter_after_seek_before_init_complete", "return").unwrap();
284        let result = iter
285            .seek(iterator_test_key_of(TEST_KEYS_COUNT).to_ref())
286            .await;
287        fail::remove("concat_iter_after_seek_before_init_complete");
288
289        assert!(result.is_err());
290        assert!(iter.is_valid());
291        assert_eq!(iter.key(), iterator_test_key_of(0).to_ref());
292
293        let mut after = StoreLocalStatistic::default();
294        iter.collect_local_statistic(&mut after);
295        assert_eq!(
296            after.cache_data_block_total,
297            before.cache_data_block_total + 1,
298            "failed candidate initialization must settle its cache-hit statistic"
299        );
300    }
301}