risingwave_storage/hummock/
validator.rs1use std::borrow::BorrowMut;
16use std::cmp;
17use std::collections::HashMap;
18use std::sync::Arc;
19
20use risingwave_hummock_sdk::compact_task::ValidationTask;
21use risingwave_hummock_sdk::key::FullKey;
22
23use crate::hummock::iterator::HummockIterator;
24use crate::hummock::sstable::SstableIteratorReadOptions;
25use crate::hummock::sstable_store::SstableStoreRef;
26use crate::hummock::{CachePolicy, SstableIterator};
27use crate::monitor::StoreLocalStatistic;
28
29pub async fn validate_ssts(task: ValidationTask, sstable_store: SstableStoreRef) {
33 let mut visited_keys = HashMap::new();
34 let mut unused = StoreLocalStatistic::default();
35 for sstable_info in task.sst_infos {
36 let mut key_counts = 0;
37 let worker_id = *task
38 .sst_id_to_worker_id
39 .get(&sstable_info.object_id)
40 .expect("valid worker_id");
41 tracing::debug!(
42 "Validating SST sst_id {} object_id {} from worker {}",
43 sstable_info.sst_id,
44 sstable_info.object_id,
45 worker_id,
46 );
47 let holder = match sstable_store
48 .sstable(&sstable_info, unused.borrow_mut())
49 .await
50 {
51 Ok(holder) => holder,
52 Err(_err) => {
53 tracing::info!(
55 "Skip sanity check for SST sst_id {} object_id {} .",
56 sstable_info.sst_id,
57 sstable_info.object_id,
58 );
59 continue;
60 }
61 };
62
63 let mut iter = SstableIterator::new(
64 holder,
65 sstable_store.clone(),
66 Arc::new(SstableIteratorReadOptions {
67 cache_policy: CachePolicy::NotFill,
68 read_table_id: None,
69 scan_end_user_key: None,
70 prefetch: false,
71 max_preload_retry_times: 0,
72 }),
73 &sstable_info,
74 );
75 let mut previous_key: Option<FullKey<Vec<u8>>> = None;
76 if let Err(_err) = iter.rewind().await {
77 tracing::info!(
78 "Skip sanity check for SST sst_id {} object_id {}.",
79 sstable_info.sst_id,
80 sstable_info.object_id
81 );
82 }
83 while iter.is_valid() {
84 key_counts += 1;
85 let current_key = iter.key().to_vec();
86 if let Some((duplicate_sst_object_id, duplicate_worker_id)) =
88 visited_keys.get(¤t_key).cloned()
89 {
90 panic!(
91 "SST sanity check failed: Duplicate key {:x?} in SST object {} from worker {} and SST object {} from worker {}",
92 current_key,
93 sstable_info.object_id,
94 worker_id,
95 duplicate_sst_object_id,
96 duplicate_worker_id
97 )
98 }
99 visited_keys.insert(current_key.clone(), (sstable_info.object_id, worker_id));
100 if let Some(previous_key) = previous_key.take() {
102 let cmp = previous_key.cmp(¤t_key);
103 if cmp != cmp::Ordering::Less {
104 panic!(
105 "SST sanity check failed: For SST sst_id {} object_id {}, expect {:x?} < {:x?}, got {:#?}",
106 sstable_info.sst_id, sstable_info.object_id, previous_key, current_key, cmp
107 )
108 }
109 }
110 previous_key = Some(current_key);
111 if let Err(_err) = iter.next().await {
112 tracing::info!(
113 "Skip remaining sanity check for SST {}",
114 sstable_info.object_id
115 );
116 break;
117 }
118 }
119 tracing::debug!(
120 "Validated {} keys for SST sst_id {} object_id {}",
121 key_counts,
122 sstable_info.sst_id,
123 sstable_info.object_id
124 );
125 iter.collect_local_statistic(&mut unused);
126 unused.ignore();
127 }
128}