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::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
46/// Start hummock's asynchronous tasks.
47pub 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    // These critical tasks are put in their own timer loop deliberately, to avoid long-running ones
54    // from blocking others.
55    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
75/// Starts a task to periodically vacuum stale metadata.
76pub 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                // Wait for interval
87                _ = min_trigger_interval.tick() => {},
88                // Shutdown vacuum
89                _ = &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                // Wait for interval
114                _ = min_trigger_interval.tick() => {},
115                // Shutdown vacuum
116                _ = &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                // Wait for interval
151                _ = min_trigger_interval.tick() => {},
152                // Shutdown checkpoint
153                _ = &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}