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;
44
45type PinnedSnapshotEpochsFetcher =
46 Box<dyn Fn() -> BoxFuture<'static, Option<HashMap<TableId, HashSet<u64>>>> + Send>;
47
48static VERSION_CHECKPOINT_RUNTIME: LazyLock<Runtime> = LazyLock::new(|| {
51 tokio::runtime::Builder::new_multi_thread()
52 .worker_threads(1)
53 .thread_name("rw-version-checkpoint")
54 .enable_all()
55 .build()
56 .expect("failed to build hummock version checkpoint runtime")
57});
58
59pub fn start_hummock_workers(
61 hummock_manager: HummockManagerRef,
62 backup_manager: BackupManagerRef,
63 meta_opts: &MetaOpts,
64 get_pinned_snapshot_epochs: PinnedSnapshotEpochsFetcher,
65) -> Vec<(JoinHandle<()>, Sender<()>)> {
66 let workers = vec![
69 start_checkpoint_loop(
70 hummock_manager.clone(),
71 backup_manager,
72 Duration::from_secs(meta_opts.hummock_version_checkpoint_interval_sec),
73 meta_opts.min_delta_log_num_for_hummock_version_checkpoint,
74 ),
75 start_vacuum_metadata_loop(
76 hummock_manager.clone(),
77 Duration::from_secs(meta_opts.vacuum_interval_sec),
78 ),
79 start_vacuum_time_travel_metadata_loop(
80 hummock_manager,
81 Duration::from_secs(meta_opts.time_travel_vacuum_interval_sec),
82 get_pinned_snapshot_epochs,
83 ),
84 ];
85 workers
86}
87
88pub fn start_vacuum_metadata_loop(
90 hummock_manager: HummockManagerRef,
91 interval: Duration,
92) -> (JoinHandle<()>, Sender<()>) {
93 let (shutdown_tx, mut shutdown_rx) = tokio::sync::oneshot::channel();
94 let join_handle = tokio::spawn(async move {
95 let mut min_trigger_interval = tokio::time::interval(interval);
96 min_trigger_interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
97 loop {
98 tokio::select! {
99 _ = min_trigger_interval.tick() => {},
101 _ = &mut shutdown_rx => {
103 tracing::info!("Vacuum metadata loop is stopped");
104 return;
105 }
106 }
107 if let Err(err) = hummock_manager.delete_version_deltas().await {
108 tracing::warn!(error = %err.as_report(), "Vacuum metadata error");
109 }
110 }
111 });
112 (join_handle, shutdown_tx)
113}
114
115pub fn start_vacuum_time_travel_metadata_loop(
116 hummock_manager: HummockManagerRef,
117 interval: Duration,
118 get_pinned_snapshot_epochs: PinnedSnapshotEpochsFetcher,
119) -> (JoinHandle<()>, Sender<()>) {
120 let (shutdown_tx, mut shutdown_rx) = tokio::sync::oneshot::channel();
121 let join_handle = tokio::spawn(async move {
122 let mut min_trigger_interval = tokio::time::interval(interval);
123 min_trigger_interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
124 loop {
125 tokio::select! {
126 _ = min_trigger_interval.tick() => {},
128 _ = &mut shutdown_rx => {
130 tracing::info!("Vacuum time travel metadata loop is stopped");
131 return;
132 }
133 }
134 let Some(pinned_snapshot_epochs) = get_pinned_snapshot_epochs().await else {
135 tracing::warn!(
136 "time travel vacuum paused because pinned snapshots are unavailable"
137 );
138 continue;
139 };
140 if let Err(err) = hummock_manager
141 .delete_time_travel_metadata(pinned_snapshot_epochs)
142 .await
143 {
144 tracing::warn!(error = %err.as_report(), "Vacuum time travel metadata error");
145 }
146 }
147 });
148 (join_handle, shutdown_tx)
149}
150
151pub fn start_checkpoint_loop(
152 hummock_manager: HummockManagerRef,
153 backup_manager: BackupManagerRef,
154 interval: Duration,
155 min_delta_log_num: u64,
156) -> (JoinHandle<()>, Sender<()>) {
157 let (shutdown_tx, mut shutdown_rx) = tokio::sync::oneshot::channel();
158 let join_handle = tokio::spawn(async move {
159 let mut min_trigger_interval = tokio::time::interval(interval);
160 min_trigger_interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
161 loop {
162 tokio::select! {
163 _ = min_trigger_interval.tick() => {},
165 _ = &mut shutdown_rx => {
167 tracing::info!("Hummock version checkpoint is stopped");
168 return;
169 }
170 }
171 if hummock_manager.is_version_checkpoint_paused()
172 || hummock_manager.env.opts.compaction_deterministic_test
173 {
174 continue;
175 }
176 let checkpoint_manager = hummock_manager.clone();
177 match VERSION_CHECKPOINT_RUNTIME
178 .spawn(async move {
179 checkpoint_manager
180 .create_version_checkpoint(min_delta_log_num)
181 .await
182 })
183 .await
184 {
185 Ok(Err(err)) => {
186 tracing::warn!(error = %err.as_report(), "Hummock version checkpoint error.");
187 }
188 Err(err) => {
189 tracing::warn!(error = %err.as_report(), "Hummock version checkpoint task failed.");
190 }
191 Ok(Ok(_)) => {
192 let backup_manager_2 = backup_manager.clone();
193 let hummock_manager_2 = hummock_manager.clone();
194 tokio::task::spawn(async move {
195 let _ = hummock_manager_2
196 .try_start_minor_gc(backup_manager_2)
197 .await
198 .inspect_err(|err| {
199 tracing::warn!(error = %err.as_report(), "Hummock minor GC error.");
200 });
201 });
202 }
203 }
204 }
205 });
206 (join_handle, shutdown_tx)
207}