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::time::Duration;
31
32pub use compactor_manager::*;
33use futures::future::BoxFuture;
34#[cfg(any(test, feature = "test"))]
35pub use mock_hummock_meta_client::MockHummockMetaClient;
36use risingwave_common::catalog::TableId;
37use tokio::sync::oneshot::Sender;
38use tokio::task::JoinHandle;
39
40use crate::MetaOpts;
41use crate::backup_restore::BackupManagerRef;
42
43type PinnedSnapshotEpochsFetcher =
44 Box<dyn Fn() -> BoxFuture<'static, Option<HashMap<TableId, HashSet<u64>>>> + Send>;
45
46pub fn start_hummock_workers(
48 hummock_manager: HummockManagerRef,
49 backup_manager: BackupManagerRef,
50 meta_opts: &MetaOpts,
51 get_pinned_snapshot_epochs: PinnedSnapshotEpochsFetcher,
52) -> Vec<(JoinHandle<()>, Sender<()>)> {
53 let workers = vec![
56 start_checkpoint_loop(
57 hummock_manager.clone(),
58 backup_manager,
59 Duration::from_secs(meta_opts.hummock_version_checkpoint_interval_sec),
60 meta_opts.min_delta_log_num_for_hummock_version_checkpoint,
61 ),
62 start_vacuum_metadata_loop(
63 hummock_manager.clone(),
64 Duration::from_secs(meta_opts.vacuum_interval_sec),
65 ),
66 start_vacuum_time_travel_metadata_loop(
67 hummock_manager,
68 Duration::from_secs(meta_opts.time_travel_vacuum_interval_sec),
69 get_pinned_snapshot_epochs,
70 ),
71 ];
72 workers
73}
74
75pub fn start_vacuum_metadata_loop(
77 hummock_manager: HummockManagerRef,
78 interval: Duration,
79) -> (JoinHandle<()>, Sender<()>) {
80 let (shutdown_tx, mut shutdown_rx) = tokio::sync::oneshot::channel();
81 let join_handle = tokio::spawn(async move {
82 let mut min_trigger_interval = tokio::time::interval(interval);
83 min_trigger_interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
84 loop {
85 tokio::select! {
86 _ = min_trigger_interval.tick() => {},
88 _ = &mut shutdown_rx => {
90 tracing::info!("Vacuum metadata loop is stopped");
91 return;
92 }
93 }
94 if let Err(err) = hummock_manager.delete_version_deltas().await {
95 tracing::warn!(error = %err.as_report(), "Vacuum metadata error");
96 }
97 }
98 });
99 (join_handle, shutdown_tx)
100}
101
102pub fn start_vacuum_time_travel_metadata_loop(
103 hummock_manager: HummockManagerRef,
104 interval: Duration,
105 get_pinned_snapshot_epochs: PinnedSnapshotEpochsFetcher,
106) -> (JoinHandle<()>, Sender<()>) {
107 let (shutdown_tx, mut shutdown_rx) = tokio::sync::oneshot::channel();
108 let join_handle = tokio::spawn(async move {
109 let mut min_trigger_interval = tokio::time::interval(interval);
110 min_trigger_interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
111 loop {
112 tokio::select! {
113 _ = min_trigger_interval.tick() => {},
115 _ = &mut shutdown_rx => {
117 tracing::info!("Vacuum time travel metadata loop is stopped");
118 return;
119 }
120 }
121 let Some(pinned_snapshot_epochs) = get_pinned_snapshot_epochs().await else {
122 tracing::warn!(
123 "time travel vacuum paused because pinned snapshots are unavailable"
124 );
125 continue;
126 };
127 if let Err(err) = hummock_manager
128 .delete_time_travel_metadata(pinned_snapshot_epochs)
129 .await
130 {
131 tracing::warn!(error = %err.as_report(), "Vacuum time travel metadata error");
132 }
133 }
134 });
135 (join_handle, shutdown_tx)
136}
137
138pub fn start_checkpoint_loop(
139 hummock_manager: HummockManagerRef,
140 backup_manager: BackupManagerRef,
141 interval: Duration,
142 min_delta_log_num: u64,
143) -> (JoinHandle<()>, Sender<()>) {
144 let (shutdown_tx, mut shutdown_rx) = tokio::sync::oneshot::channel();
145 let join_handle = tokio::spawn(async move {
146 let mut min_trigger_interval = tokio::time::interval(interval);
147 min_trigger_interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
148 loop {
149 tokio::select! {
150 _ = min_trigger_interval.tick() => {},
152 _ = &mut shutdown_rx => {
154 tracing::info!("Hummock version checkpoint is stopped");
155 return;
156 }
157 }
158 if hummock_manager.is_version_checkpoint_paused()
159 || hummock_manager.env.opts.compaction_deterministic_test
160 {
161 continue;
162 }
163 match hummock_manager
164 .create_version_checkpoint(min_delta_log_num)
165 .await
166 {
167 Err(err) => {
168 tracing::warn!(error = %err.as_report(), "Hummock version checkpoint error.");
169 }
170 _ => {
171 let backup_manager_2 = backup_manager.clone();
172 let hummock_manager_2 = hummock_manager.clone();
173 tokio::task::spawn(async move {
174 let _ = hummock_manager_2
175 .try_start_minor_gc(backup_manager_2)
176 .await
177 .inspect_err(|err| {
178 tracing::warn!(error = %err.as_report(), "Hummock minor GC error.");
179 });
180 });
181 }
182 }
183 }
184 });
185 (join_handle, shutdown_tx)
186}