risingwave_meta/hummock/
mod.rs1pub mod compaction;
16pub mod compactor_manager;
17pub mod error;
18mod manager;
19pub(crate) use manager::checkpoint::{compress_payload, xxhash64_checksum};
20pub use manager::*;
21use thiserror_ext::AsReport;
22
23mod level_handler;
24mod metrics_utils;
25#[cfg(any(test, feature = "test"))]
26pub mod mock_hummock_meta_client;
27pub mod model;
28pub mod test_utils;
29use std::collections::{HashMap, HashSet};
30use std::sync::LazyLock;
31use std::time::Duration;
32
33pub use compactor_manager::*;
34use futures::future::BoxFuture;
35#[cfg(any(test, feature = "test"))]
36pub use mock_hummock_meta_client::MockHummockMetaClient;
37use risingwave_common::catalog::TableId;
38use tokio::runtime::Runtime;
39use tokio::sync::oneshot::Sender;
40use tokio::task::JoinHandle;
41
42use crate::MetaOpts;
43use crate::backup_restore::BackupManagerRef;
44use crate::controller::streaming_job::TableChangeLogTruncateInfo;
45
46type PinnedSnapshotEpochsFetcher =
47 Box<dyn Fn() -> BoxFuture<'static, Option<HashMap<TableId, HashSet<u64>>>> + Send>;
48type TableChangeLogTruncateInfoFetcher =
49 Box<dyn Fn() -> BoxFuture<'static, Option<TableChangeLogTruncateInfo>> + Send>;
50
51static VERSION_CHECKPOINT_RUNTIME: LazyLock<Runtime> = LazyLock::new(|| {
54 tokio::runtime::Builder::new_multi_thread()
55 .worker_threads(1)
56 .thread_name("rw-version-checkpoint")
57 .enable_all()
58 .build()
59 .expect("failed to build hummock version checkpoint runtime")
60});
61
62pub fn start_hummock_workers(
64 hummock_manager: HummockManagerRef,
65 backup_manager: BackupManagerRef,
66 meta_opts: &MetaOpts,
67 get_table_change_log_truncate_info: TableChangeLogTruncateInfoFetcher,
68 get_pinned_snapshot_epochs: PinnedSnapshotEpochsFetcher,
69) -> Vec<(JoinHandle<()>, Sender<()>)> {
70 let workers = vec![
73 start_checkpoint_loop(
74 hummock_manager.clone(),
75 backup_manager,
76 Duration::from_secs(meta_opts.hummock_version_checkpoint_interval_sec),
77 meta_opts.min_delta_log_num_for_hummock_version_checkpoint,
78 ),
79 start_vacuum_metadata_loop(
80 hummock_manager.clone(),
81 Duration::from_secs(meta_opts.vacuum_interval_sec),
82 ),
83 start_truncate_table_change_log_loop(
84 hummock_manager.clone(),
85 Duration::from_secs(meta_opts.table_change_log_truncate_interval_sec),
86 get_table_change_log_truncate_info,
87 ),
88 start_vacuum_time_travel_metadata_loop(
89 hummock_manager,
90 Duration::from_secs(meta_opts.time_travel_vacuum_interval_sec),
91 get_pinned_snapshot_epochs,
92 ),
93 ];
94 workers
95}
96
97pub fn start_truncate_table_change_log_loop(
98 hummock_manager: HummockManagerRef,
99 interval: Duration,
100 get_truncate_info: TableChangeLogTruncateInfoFetcher,
101) -> (JoinHandle<()>, Sender<()>) {
102 let (shutdown_tx, mut shutdown_rx) = tokio::sync::oneshot::channel();
103 let join_handle = tokio::spawn(async move {
104 let mut min_trigger_interval = tokio::time::interval(interval);
105 min_trigger_interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
106 loop {
107 tokio::select! {
108 _ = min_trigger_interval.tick() => {},
109 _ = &mut shutdown_rx => {
110 tracing::info!("Table change log truncate loop is stopped");
111 return;
112 }
113 }
114 let Some(truncate_info) = get_truncate_info().await else {
115 tracing::warn!(
116 "table change log truncation paused because retention metadata is unavailable"
117 );
118 continue;
119 };
120 if let Err(err) = hummock_manager
121 .truncate_table_change_log(truncate_info)
122 .await
123 {
124 tracing::warn!(error = %err.as_report(), "Table change log truncation error");
125 }
126 }
127 });
128 (join_handle, shutdown_tx)
129}
130
131pub fn start_vacuum_metadata_loop(
133 hummock_manager: HummockManagerRef,
134 interval: Duration,
135) -> (JoinHandle<()>, Sender<()>) {
136 let (shutdown_tx, mut shutdown_rx) = tokio::sync::oneshot::channel();
137 let join_handle = tokio::spawn(async move {
138 let mut min_trigger_interval = tokio::time::interval(interval);
139 min_trigger_interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
140 loop {
141 tokio::select! {
142 _ = min_trigger_interval.tick() => {},
144 _ = &mut shutdown_rx => {
146 tracing::info!("Vacuum metadata loop is stopped");
147 return;
148 }
149 }
150 if let Err(err) = hummock_manager.delete_version_deltas().await {
151 tracing::warn!(error = %err.as_report(), "Vacuum metadata error");
152 }
153 }
154 });
155 (join_handle, shutdown_tx)
156}
157
158pub fn start_vacuum_time_travel_metadata_loop(
159 hummock_manager: HummockManagerRef,
160 interval: Duration,
161 get_pinned_snapshot_epochs: PinnedSnapshotEpochsFetcher,
162) -> (JoinHandle<()>, Sender<()>) {
163 let (shutdown_tx, mut shutdown_rx) = tokio::sync::oneshot::channel();
164 let join_handle = tokio::spawn(async move {
165 let mut min_trigger_interval = tokio::time::interval(interval);
166 min_trigger_interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
167 loop {
168 tokio::select! {
169 _ = min_trigger_interval.tick() => {},
171 _ = &mut shutdown_rx => {
173 tracing::info!("Vacuum time travel metadata loop is stopped");
174 return;
175 }
176 }
177 let Some(pinned_snapshot_epochs) = get_pinned_snapshot_epochs().await else {
178 tracing::warn!(
179 "time travel vacuum paused because pinned snapshots are unavailable"
180 );
181 continue;
182 };
183 if let Err(err) = hummock_manager
184 .delete_time_travel_metadata(pinned_snapshot_epochs)
185 .await
186 {
187 tracing::warn!(error = %err.as_report(), "Vacuum time travel metadata error");
188 }
189 }
190 });
191 (join_handle, shutdown_tx)
192}
193
194pub fn start_checkpoint_loop(
195 hummock_manager: HummockManagerRef,
196 backup_manager: BackupManagerRef,
197 interval: Duration,
198 min_delta_log_num: u64,
199) -> (JoinHandle<()>, Sender<()>) {
200 let (shutdown_tx, mut shutdown_rx) = tokio::sync::oneshot::channel();
201 let join_handle = tokio::spawn(async move {
202 let mut min_trigger_interval = tokio::time::interval(interval);
203 min_trigger_interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
204 loop {
205 tokio::select! {
206 _ = min_trigger_interval.tick() => {},
208 _ = &mut shutdown_rx => {
210 tracing::info!("Hummock version checkpoint is stopped");
211 return;
212 }
213 }
214 if hummock_manager.is_version_checkpoint_paused()
215 || hummock_manager.env.opts.compaction_deterministic_test
216 {
217 continue;
218 }
219 let checkpoint_manager = hummock_manager.clone();
220 match VERSION_CHECKPOINT_RUNTIME
221 .spawn(async move {
222 checkpoint_manager
223 .create_version_checkpoint(min_delta_log_num)
224 .await
225 })
226 .await
227 {
228 Ok(Err(err)) => {
229 tracing::warn!(error = %err.as_report(), "Hummock version checkpoint error.");
230 }
231 Err(err) => {
232 tracing::warn!(error = %err.as_report(), "Hummock version checkpoint task failed.");
233 }
234 Ok(Ok(_)) => {
235 let backup_manager_2 = backup_manager.clone();
236 let hummock_manager_2 = hummock_manager.clone();
237 tokio::task::spawn(async move {
238 let _ = hummock_manager_2
239 .try_start_minor_gc(backup_manager_2)
240 .await
241 .inspect_err(|err| {
242 tracing::warn!(error = %err.as_report(), "Hummock minor GC error.");
243 });
244 });
245 }
246 }
247 }
248 });
249 (join_handle, shutdown_tx)
250}