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 scan_end_user_key: None,
69 prefetch: false,
70 max_preload_retry_times: 0,
71 }),
72 &sstable_info,
73 );
74 let mut previous_key: Option<FullKey<Vec<u8>>> = None;
75 if let Err(_err) = iter.rewind().await {
76 tracing::info!(
77 "Skip sanity check for SST sst_id {} object_id {}.",
78 sstable_info.sst_id,
79 sstable_info.object_id
80 );
81 }
82 while iter.is_valid() {
83 key_counts += 1;
84 let current_key = iter.key().to_vec();
85 if let Some((duplicate_sst_object_id, duplicate_worker_id)) =
87 visited_keys.get(¤t_key).cloned()
88 {
89 panic!(
90 "SST sanity check failed: Duplicate key {:x?} in SST object {} from worker {} and SST object {} from worker {}",
91 current_key,
92 sstable_info.object_id,
93 worker_id,
94 duplicate_sst_object_id,
95 duplicate_worker_id
96 )
97 }
98 visited_keys.insert(current_key.clone(), (sstable_info.object_id, worker_id));
99 if let Some(previous_key) = previous_key.take() {
101 let cmp = previous_key.cmp(¤t_key);
102 if cmp != cmp::Ordering::Less {
103 panic!(
104 "SST sanity check failed: For SST sst_id {} object_id {}, expect {:x?} < {:x?}, got {:#?}",
105 sstable_info.sst_id, sstable_info.object_id, previous_key, current_key, cmp
106 )
107 }
108 }
109 previous_key = Some(current_key);
110 if let Err(_err) = iter.next().await {
111 tracing::info!(
112 "Skip remaining sanity check for SST {}",
113 sstable_info.object_id
114 );
115 break;
116 }
117 }
118 tracing::debug!(
119 "Validated {} keys for SST sst_id {} object_id {}",
120 key_counts,
121 sstable_info.sst_id,
122 sstable_info.object_id
123 );
124 iter.collect_local_statistic(&mut unused);
125 unused.ignore();
126 }
127}