Skip to main content

risingwave_meta/hummock/
mod.rs

1// Copyright 2022 RisingWave Labs
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15pub 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
51// Building and compressing a large checkpoint can spend a long time without yielding. Keep that
52// work off the main meta runtime so it cannot monopolize one of its worker threads.
53static 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
62/// Start hummock's asynchronous tasks.
63pub 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    // These critical tasks are put in their own timer loop deliberately, to avoid long-running ones
71    // from blocking others.
72    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
131/// Starts a task to periodically vacuum stale metadata.
132pub 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                // Wait for interval
143                _ = min_trigger_interval.tick() => {},
144                // Shutdown vacuum
145                _ = &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                // Wait for interval
170                _ = min_trigger_interval.tick() => {},
171                // Shutdown vacuum
172                _ = &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                // Wait for interval
207                _ = min_trigger_interval.tick() => {},
208                // Shutdown checkpoint
209                _ = &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}