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;
44
45type PinnedSnapshotEpochsFetcher =
46    Box<dyn Fn() -> BoxFuture<'static, Option<HashMap<TableId, HashSet<u64>>>> + Send>;
47
48// Building and compressing a large checkpoint can spend a long time without yielding. Keep that
49// work off the main meta runtime so it cannot monopolize one of its worker threads.
50static 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
59/// Start hummock's asynchronous tasks.
60pub 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    // These critical tasks are put in their own timer loop deliberately, to avoid long-running ones
67    // from blocking others.
68    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
88/// Starts a task to periodically vacuum stale metadata.
89pub 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                // Wait for interval
100                _ = min_trigger_interval.tick() => {},
101                // Shutdown vacuum
102                _ = &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                // Wait for interval
127                _ = min_trigger_interval.tick() => {},
128                // Shutdown vacuum
129                _ = &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                // Wait for interval
164                _ = min_trigger_interval.tick() => {},
165                // Shutdown checkpoint
166                _ = &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}