1use 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
39pub struct SstableStreamIterator {
41 block_stream: SstableBlockStream,
42 block_iter: Option<BlockIterator>,
44 stats_ptr: Arc<AtomicU64>,
46
47 read_table_ids: HashSet<TableId>,
49 task_progress: Arc<TaskProgress>,
50
51 key_range_left: FullKey<Vec<u8>>,
53 key_range_right: FullKey<Vec<u8>>,
54 key_range_right_exclusive: bool,
55}
56
57impl SstableStreamIterator {
58 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 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 (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 pub async fn seek(&mut self, seek_key: Option<FullKey<&[u8]>>) -> HummockResult<()> {
119 self.next_block().await?;
121
122 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 self.next_block().await?;
142 }
143 }
144
145 self.prune_from_valid_block_iter().await?;
146 Ok(())
147 }
148
149 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 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 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 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 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
275pub struct ConcatSstableIterator {
278 key_range: KeyRange,
281
282 sstable_iter: Option<SstableStreamIterator>,
284
285 cur_idx: usize,
287
288 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 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 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 sstable_iter.next().await?;
412 if sstable_iter.is_valid() {
413 Ok(())
414 } else {
415 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 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 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 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 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 return 0..0;
573 }
574
575 ret
576 }
577 .saturating_sub(1);
578
579 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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}