Skip to main content

risingwave_meta/hummock/manager/
gc.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
15use std::cmp;
16use std::collections::{HashMap, HashSet};
17use std::ops::DerefMut;
18use std::sync::atomic::{AtomicBool, Ordering};
19use std::time::{Duration, Instant, SystemTime};
20
21use chrono::DateTime;
22use futures::future::try_join_all;
23use futures::{StreamExt, TryStreamExt, future};
24use itertools::Itertools;
25use risingwave_common::catalog::TableId;
26use risingwave_common::system_param::reader::SystemParamsRead;
27use risingwave_common::util::epoch::Epoch;
28use risingwave_hummock_sdk::{
29    HummockEpoch, HummockObjectId, HummockRawObjectId, VALID_OBJECT_ID_SUFFIXES,
30    get_object_data_path, get_object_id_from_path,
31};
32use risingwave_meta_model::hummock_sequence::HUMMOCK_NOW;
33use risingwave_meta_model::{hummock_gc_history, hummock_sequence, hummock_version_delta};
34use risingwave_meta_model_migration::OnConflict;
35use risingwave_object_store::object::{ObjectMetadataIter, ObjectStoreRef};
36use risingwave_pb::stream_service::GetMinUncommittedObjectIdRequest;
37use risingwave_rpc_client::StreamClientPool;
38use sea_orm::{ActiveValue, ColumnTrait, EntityTrait, QueryFilter, Set};
39
40use crate::MetaResult;
41use crate::backup_restore::BackupManagerRef;
42use crate::hummock::HummockManager;
43use crate::hummock::error::{Error, Result};
44use crate::manager::MetadataManager;
45
46pub(crate) struct GcManager {
47    store: ObjectStoreRef,
48    path_prefix: String,
49    use_new_object_prefix_strategy: bool,
50    /// These objects may still be used by backup or time travel.
51    may_delete_object_ids: parking_lot::Mutex<HashSet<HummockObjectId>>,
52}
53
54impl GcManager {
55    pub fn new(
56        store: ObjectStoreRef,
57        path_prefix: &str,
58        use_new_object_prefix_strategy: bool,
59    ) -> Self {
60        Self {
61            store,
62            path_prefix: path_prefix.to_owned(),
63            use_new_object_prefix_strategy,
64            may_delete_object_ids: Default::default(),
65        }
66    }
67
68    /// Deletes all objects specified in the given list of IDs from storage.
69    pub async fn delete_objects(
70        &self,
71        object_id_list: impl Iterator<Item = HummockObjectId>,
72    ) -> Result<()> {
73        let mut paths = Vec::with_capacity(1000);
74        for object_id in object_id_list {
75            let obj_prefix = self.store.get_object_prefix(
76                object_id.as_raw().as_raw_id(),
77                self.use_new_object_prefix_strategy,
78            );
79            paths.push(get_object_data_path(
80                &obj_prefix,
81                &self.path_prefix,
82                object_id,
83            ));
84        }
85        self.store.delete_objects(&paths).await?;
86        Ok(())
87    }
88
89    async fn list_object_metadata_from_object_store(
90        &self,
91        prefix: Option<String>,
92        start_after: Option<String>,
93        limit: Option<usize>,
94    ) -> Result<ObjectMetadataIter> {
95        let list_path = format!("{}/{}", self.path_prefix, prefix.unwrap_or("".into()));
96        let raw_iter = self.store.list(&list_path, start_after, limit).await?;
97        let valid_suffixes = VALID_OBJECT_ID_SUFFIXES.map(|suffix| format!(".{}", suffix));
98        let iter = raw_iter.filter(move |r| match r {
99            Ok(i) => future::ready(valid_suffixes.iter().any(|suffix| i.key.ends_with(suffix))),
100            Err(_) => future::ready(true),
101        });
102        Ok(Box::pin(iter))
103    }
104
105    /// Returns **filtered** object ids, and **unfiltered** total object count and size.
106    pub async fn list_objects(
107        &self,
108        object_retention_watermark: u64,
109        prefix: Option<String>,
110        start_after: Option<String>,
111        limit: Option<u64>,
112    ) -> Result<(HashSet<HummockObjectId>, u64, u64, Option<String>)> {
113        tracing::debug!(
114            object_retention_watermark,
115            prefix,
116            start_after,
117            limit,
118            "Try to list objects."
119        );
120        let mut total_object_count = 0;
121        let mut total_object_size = 0;
122        let mut next_start_after: Option<String> = None;
123        let metadata_iter = self
124            .list_object_metadata_from_object_store(prefix, start_after, limit.map(|i| i as usize))
125            .await?;
126        let filtered = metadata_iter
127            .filter_map(|r| {
128                let result = match r {
129                    Ok(o) => {
130                        total_object_count += 1;
131                        total_object_size += o.total_size;
132                        // Determine if the LIST has been truncated.
133                        // A false positives would at most cost one additional LIST later.
134                        if let Some(limit) = limit
135                            && limit == total_object_count
136                        {
137                            next_start_after = Some(o.key.clone());
138                            tracing::debug!(next_start_after, "set next start after");
139                        }
140                        if o.last_modified < object_retention_watermark as f64 {
141                            Some(Ok(get_object_id_from_path(&o.key)))
142                        } else {
143                            None
144                        }
145                    }
146                    Err(e) => Some(Err(Error::ObjectStore(e))),
147                };
148                async move { result }
149            })
150            .try_collect::<HashSet<HummockObjectId>>()
151            .await?;
152        Ok((
153            filtered,
154            total_object_count,
155            total_object_size as u64,
156            next_start_after,
157        ))
158    }
159
160    pub fn add_may_delete_object_ids(
161        &self,
162        may_delete_object_ids: impl Iterator<Item = HummockObjectId>,
163    ) {
164        self.may_delete_object_ids
165            .lock()
166            .extend(may_delete_object_ids);
167    }
168
169    /// Takes if `least_count` elements available.
170    pub fn try_take_may_delete_object_ids(
171        &self,
172        least_count: usize,
173    ) -> Option<HashSet<HummockObjectId>> {
174        let mut guard = self.may_delete_object_ids.lock();
175        if guard.len() < least_count {
176            None
177        } else {
178            Some(std::mem::take(&mut *guard))
179        }
180    }
181}
182
183impl HummockManager {
184    /// Deletes version deltas.
185    ///
186    /// Returns number of deleted deltas
187    pub async fn delete_version_deltas(&self) -> Result<usize> {
188        let mut versioning_guard = self.versioning.write().await;
189        let versioning = versioning_guard.deref_mut();
190        let context_info = self.context_info.read().await;
191        // If there is any safe point, skip this to ensure meta backup has required delta logs to
192        // replay version.
193        if !context_info.version_safe_points.is_empty() {
194            return Ok(0);
195        }
196        // The context_info lock must be held to prevent any potential metadata backup.
197        // The lock order requires version lock to be held as well.
198        let version_id = versioning.checkpoint.version.id;
199        let res = hummock_version_delta::Entity::delete_many()
200            .filter(hummock_version_delta::Column::Id.lte(version_id))
201            .exec(&self.env.meta_store_ref().conn)
202            .await?;
203        tracing::debug!(rows_affected = res.rows_affected, "Deleted version deltas");
204        versioning
205            .hummock_version_deltas
206            .retain(|id, _| *id > version_id);
207        #[cfg(test)]
208        {
209            drop(context_info);
210            drop(versioning_guard);
211            self.check_state_consistency().await;
212        }
213        Ok(res.rows_affected as usize)
214    }
215
216    /// Filters by Hummock version and Writes GC history.
217    pub async fn finalize_objects_to_delete(
218        &self,
219        object_ids: impl Iterator<Item = HummockObjectId>,
220    ) -> Result<Vec<HummockObjectId>> {
221        // This lock ensures `commit_epoch` and `report_compat_task` can see the latest GC history during sanity check.
222        let versioning = self.versioning.read().await;
223        let tracked_object_ids: HashSet<HummockObjectId> = versioning
224            .get_tracked_object_ids(self.context_info.read().await.min_pinned_version_id());
225        let to_delete = object_ids
226            .filter(|object_id| !tracked_object_ids.contains(object_id))
227            .collect_vec();
228        self.write_gc_history(to_delete.iter().copied()).await?;
229        Ok(to_delete)
230    }
231
232    /// LIST object store and DELETE stale objects, in batches.
233    /// GC can be very slow. Spawn a dedicated tokio task for it.
234    pub async fn start_full_gc(
235        &self,
236        object_retention_time: Duration,
237        prefix: Option<String>,
238        backup_manager: Option<BackupManagerRef>,
239    ) -> Result<()> {
240        if !self.full_gc_state.try_start() {
241            return Err(anyhow::anyhow!("failed to start GC due to an ongoing process").into());
242        }
243        let _guard = scopeguard::guard(self.full_gc_state.clone(), |full_gc_state| {
244            full_gc_state.stop()
245        });
246        self.metrics.full_gc_trigger_count.inc();
247        let object_retention_time = cmp::max(
248            object_retention_time,
249            Duration::from_secs(self.env.opts.min_sst_retention_time_sec),
250        );
251        let limit = self.env.opts.full_gc_object_limit;
252        let mut start_after = None;
253        let object_retention_watermark = self
254            .now()
255            .await?
256            .saturating_sub(object_retention_time.as_secs());
257        let mut total_object_count = 0;
258        let mut total_object_size = 0;
259        tracing::info!(
260            retention_sec = object_retention_time.as_secs(),
261            prefix,
262            limit,
263            "Start GC."
264        );
265        loop {
266            tracing::debug!(
267                retention_sec = object_retention_time.as_secs(),
268                prefix,
269                start_after,
270                limit,
271                "Start a GC batch."
272            );
273            let (object_ids, batch_object_count, batch_object_size, next_start_after) = self
274                .gc_manager
275                .list_objects(
276                    object_retention_watermark,
277                    prefix.clone(),
278                    start_after.clone(),
279                    Some(limit),
280                )
281                .await?;
282            total_object_count += batch_object_count;
283            total_object_size += batch_object_size;
284            tracing::debug!(
285                ?object_ids,
286                batch_object_count,
287                batch_object_size,
288                "Finish listing a GC batch."
289            );
290            self.complete_gc_batch(object_ids, backup_manager.clone())
291                .await?;
292            if next_start_after.is_none() {
293                break;
294            }
295            start_after = next_start_after;
296        }
297        tracing::info!(total_object_count, total_object_size, "Finish GC");
298        self.metrics.total_object_size.set(total_object_size as _);
299        self.metrics.total_object_count.set(total_object_count as _);
300        match self.time_travel_pinned_object_count().await {
301            Ok(count) => {
302                self.metrics.time_travel_object_count.set(count as _);
303            }
304            Err(err) => {
305                use thiserror_ext::AsReport;
306                tracing::warn!(error = %err.as_report(), "Failed to count time travel objects.");
307            }
308        }
309        Ok(())
310    }
311
312    /// Given candidate objects to delete, filter out false positive.
313    /// Returns number of objects to delete.
314    pub(crate) async fn complete_gc_batch(
315        &self,
316        object_ids: HashSet<HummockObjectId>,
317        backup_manager: Option<BackupManagerRef>,
318    ) -> Result<usize> {
319        if object_ids.is_empty() {
320            return Ok(0);
321        }
322        let pinned_by_metadata_backup = match backup_manager.as_ref() {
323            Some(b) => b.list_pinned_object_ids().await,
324            None => HashSet::default(),
325        };
326        // It's crucial to collect_min_uncommitted_object_id (i.e. `min_object_id`) only after LIST object store (i.e. `object_ids`).
327        // Because after getting `min_object_id`, new compute nodes may join and generate new uncommitted objects that are not covered by `min_sst_id`.
328        // By getting `min_object_id` after `object_ids`, it's ensured `object_ids` won't include any objects from those new compute nodes.
329        let min_object_id = collect_min_uncommitted_object_id(
330            &self.metadata_manager,
331            self.env.stream_client_pool(),
332        )
333        .await?;
334        let metrics = &self.metrics;
335        let candidate_object_number = object_ids.len();
336        metrics
337            .full_gc_candidate_object_count
338            .observe(candidate_object_number as _);
339        // filter by metadata backup
340        let object_ids = object_ids
341            .into_iter()
342            .filter(|s| !pinned_by_metadata_backup.contains(&s.as_raw()))
343            .collect_vec();
344        let after_metadata_backup = object_ids.len();
345        // filter by time travel archive
346        let filter_by_time_travel_start_time = Instant::now();
347        let object_ids = self
348            .filter_out_objects_by_time_travel(object_ids.into_iter())
349            .await?;
350        tracing::info!(elapsed = ?filter_by_time_travel_start_time.elapsed(), "filter out objects by time travel in full GC");
351        let after_time_travel = object_ids.len();
352        // filter by object id watermark, i.e. minimum id of uncommitted objects reported by compute nodes.
353        let object_ids = object_ids
354            .into_iter()
355            .filter(|id| id.as_raw() < min_object_id)
356            .collect_vec();
357        let after_min_object_id = object_ids.len();
358        // filter by version
359        let after_version = self
360            .finalize_objects_to_delete(object_ids.into_iter())
361            .await?;
362        let after_version_count = after_version.len();
363        metrics
364            .full_gc_selected_object_count
365            .observe(after_version_count as _);
366        tracing::info!(
367            candidate_object_number,
368            after_metadata_backup,
369            after_time_travel,
370            after_min_object_id,
371            after_version_count,
372            "complete gc batch"
373        );
374        self.delete_objects(after_version).await?;
375        Ok(after_version_count)
376    }
377
378    pub async fn now(&self) -> Result<u64> {
379        let mut guard = self.now.lock().await;
380        let new_now = SystemTime::now()
381            .duration_since(SystemTime::UNIX_EPOCH)
382            .expect("Clock may have gone backwards")
383            .as_secs();
384        if new_now < *guard {
385            return Err(anyhow::anyhow!(format!(
386                "unexpected decreasing now, old={}, new={}",
387                *guard, new_now
388            ))
389            .into());
390        }
391        *guard = new_now;
392        drop(guard);
393        // Persist now to maintain non-decreasing even after a meta node reboot.
394        let m = hummock_sequence::ActiveModel {
395            name: ActiveValue::Set(HUMMOCK_NOW.into()),
396            seq: ActiveValue::Set(new_now.try_into().unwrap()),
397        };
398        hummock_sequence::Entity::insert(m)
399            .on_conflict(
400                OnConflict::column(hummock_sequence::Column::Name)
401                    .update_column(hummock_sequence::Column::Seq)
402                    .to_owned(),
403            )
404            .exec(&self.env.meta_store_ref().conn)
405            .await?;
406        Ok(new_now)
407    }
408
409    pub(crate) async fn load_now(&self) -> Result<Option<u64>> {
410        let now = hummock_sequence::Entity::find_by_id(HUMMOCK_NOW.to_owned())
411            .one(&self.env.meta_store_ref().conn)
412            .await?
413            .map(|m| m.seq.try_into().unwrap());
414        Ok(now)
415    }
416
417    async fn write_gc_history(
418        &self,
419        object_ids: impl Iterator<Item = HummockObjectId>,
420    ) -> Result<()> {
421        if self.env.opts.gc_history_retention_time_sec == 0 {
422            return Ok(());
423        }
424        let now = self.now().await?;
425        let dt = DateTime::from_timestamp(now.try_into().unwrap(), 0).unwrap();
426        let mut models = object_ids.map(|o| hummock_gc_history::ActiveModel {
427            object_id: Set(o.as_raw().into()),
428            mark_delete_at: Set(dt.naive_utc()),
429        });
430        let db = &self.meta_store_ref().conn;
431        let gc_history_low_watermark = DateTime::from_timestamp(
432            now.saturating_sub(self.env.opts.gc_history_retention_time_sec)
433                .try_into()
434                .unwrap(),
435            0,
436        )
437        .unwrap();
438        hummock_gc_history::Entity::delete_many()
439            .filter(hummock_gc_history::Column::MarkDeleteAt.lt(gc_history_low_watermark))
440            .exec(db)
441            .await?;
442        let mut is_finished = false;
443        while !is_finished {
444            let mut batch = vec![];
445            let mut count: usize = self.env.opts.hummock_gc_history_insert_batch_size;
446            while count > 0 {
447                let Some(m) = models.next() else {
448                    is_finished = true;
449                    break;
450                };
451                count -= 1;
452                batch.push(m);
453            }
454            if batch.is_empty() {
455                break;
456            }
457            hummock_gc_history::Entity::insert_many(batch)
458                .on_conflict_do_nothing()
459                .exec(db)
460                .await?;
461        }
462        Ok(())
463    }
464
465    pub async fn delete_time_travel_metadata(
466        &self,
467        pinned_snapshot_epochs: HashMap<TableId, HashSet<HummockEpoch>>,
468    ) -> MetaResult<()> {
469        let current_epoch_time = Epoch::now().physical_time();
470        let epoch_watermark = Epoch::from_physical_time(
471            current_epoch_time.saturating_sub(
472                self.env
473                    .system_params_reader()
474                    .await
475                    .time_travel_retention_ms(),
476            ),
477        )
478        .0;
479        self.truncate_time_travel_metadata(epoch_watermark, pinned_snapshot_epochs)
480            .await?;
481        Ok(())
482    }
483
484    /// Deletes stale objects from object store.
485    ///
486    /// Returns the total count of deleted objects.
487    pub async fn delete_objects(
488        &self,
489        mut objects_to_delete: Vec<HummockObjectId>,
490    ) -> Result<usize> {
491        let total = objects_to_delete.len();
492        let mut batch_size = 1000usize;
493        while !objects_to_delete.is_empty() {
494            if self.env.opts.vacuum_spin_interval_ms != 0 {
495                tokio::time::sleep(Duration::from_millis(self.env.opts.vacuum_spin_interval_ms))
496                    .await;
497            }
498            batch_size = cmp::min(objects_to_delete.len(), batch_size);
499            if batch_size == 0 {
500                break;
501            }
502            let delete_batch: HashSet<_> = objects_to_delete.drain(..batch_size).collect();
503            tracing::info!(?delete_batch, "Attempt to delete objects.");
504            let deleted_object_ids = delete_batch.clone();
505            self.gc_manager
506                .delete_objects(delete_batch.into_iter())
507                .await?;
508            tracing::debug!(?deleted_object_ids, "Finish deleting objects.");
509        }
510        Ok(total)
511    }
512
513    /// Minor GC attempts to delete objects that were part of Hummock version but are no longer in use.
514    pub async fn try_start_minor_gc(&self, backup_manager: BackupManagerRef) -> Result<()> {
515        const MIN_MINOR_GC_OBJECT_COUNT: usize = 1000;
516        let Some(object_ids) = self
517            .gc_manager
518            .try_take_may_delete_object_ids(MIN_MINOR_GC_OBJECT_COUNT)
519        else {
520            return Ok(());
521        };
522        // Objects pinned by either meta backup or time travel should be filtered out.
523        let backup_pinned: HashSet<_> = backup_manager.list_pinned_object_ids().await;
524        // The version_pinned is obtained after the candidate object_ids for deletion, which is new enough for filtering purpose.
525        let version_pinned = {
526            let versioning = self.versioning.read().await;
527            versioning
528                .get_tracked_object_ids(self.context_info.read().await.min_pinned_version_id())
529        };
530        let object_ids = object_ids
531            .into_iter()
532            .filter(|s| !version_pinned.contains(s) && !backup_pinned.contains(&s.as_raw()));
533        let filter_by_time_travel_start_time = Instant::now();
534        let object_ids = self.filter_out_objects_by_time_travel(object_ids).await?;
535        tracing::info!(elapsed = ?filter_by_time_travel_start_time.elapsed(), "filter out objects by time travel in minor GC");
536        // Retry is not necessary. Full GC will handle these objects eventually.
537        self.delete_objects(object_ids.into_iter().collect())
538            .await?;
539        Ok(())
540    }
541}
542
543async fn collect_min_uncommitted_object_id(
544    metadata_manager: &MetadataManager,
545    client_pool: &StreamClientPool,
546) -> Result<HummockRawObjectId> {
547    let futures = metadata_manager
548        .list_active_streaming_compute_nodes()
549        .await
550        .map_err(|err| Error::MetaStore(err.into()))?
551        .into_iter()
552        .map(|worker_node| async move {
553            let client = client_pool.get(&worker_node).await?;
554            let request = GetMinUncommittedObjectIdRequest {};
555            client.get_min_uncommitted_object_id(request).await
556        });
557    let min_watermark = try_join_all(futures)
558        .await
559        .map_err(|err| Error::Internal(err.into()))?
560        .into_iter()
561        .map(|resp| resp.min_uncommitted_object_id)
562        .min()
563        .unwrap_or(u64::MAX.into());
564    Ok(min_watermark)
565}
566
567pub struct FullGcState {
568    is_started: AtomicBool,
569}
570
571impl FullGcState {
572    pub fn new() -> Self {
573        Self {
574            is_started: AtomicBool::new(false),
575        }
576    }
577
578    pub fn try_start(&self) -> bool {
579        self.is_started
580            .compare_exchange(false, true, Ordering::SeqCst, Ordering::SeqCst)
581            .is_ok()
582    }
583
584    pub fn stop(&self) {
585        self.is_started.store(false, Ordering::SeqCst);
586    }
587}
588
589#[cfg(test)]
590mod tests {
591    use std::sync::Arc;
592    use std::time::Duration;
593
594    use itertools::Itertools;
595    use risingwave_hummock_sdk::HummockObjectId;
596    use risingwave_hummock_sdk::compaction_group::StaticCompactionGroupId;
597    use risingwave_rpc_client::HummockMetaClient;
598
599    use crate::hummock::MockHummockMetaClient;
600    use crate::hummock::test_utils::{add_test_tables, setup_compute_env};
601
602    #[tokio::test]
603    async fn test_full_gc() {
604        let (_env, hummock_manager, _cluster_manager, worker_id) = setup_compute_env(80).await;
605        let hummock_meta_client: Arc<dyn HummockMetaClient> = Arc::new(MockHummockMetaClient::new(
606            hummock_manager.clone(),
607            worker_id as _,
608        ));
609        let compaction_group_id = StaticCompactionGroupId::StateDefault;
610        hummock_manager
611            .start_full_gc(
612                Duration::from_secs(hummock_manager.env.opts.min_sst_retention_time_sec + 1),
613                None,
614                None,
615            )
616            .await
617            .unwrap();
618
619        // Empty input results immediate return, without waiting heartbeat.
620        hummock_manager
621            .complete_gc_batch(vec![].into_iter().collect(), None)
622            .await
623            .unwrap();
624
625        // LSMtree is empty. All input object ids should be treated as garbage.
626        // Use fake object ids, because they'll be written to GC history and they shouldn't affect later commit.
627        assert_eq!(
628            3,
629            hummock_manager
630                .complete_gc_batch(
631                    [i64::MAX as u64 - 2, i64::MAX as u64 - 1, i64::MAX as u64]
632                        .into_iter()
633                        .map(|id| HummockObjectId::Sstable(id.into()))
634                        .collect(),
635                    None,
636                )
637                .await
638                .unwrap()
639        );
640
641        // All committed SST ids should be excluded from GC.
642        let sst_infos = add_test_tables(
643            hummock_manager.as_ref(),
644            hummock_meta_client.clone(),
645            compaction_group_id,
646        )
647        .await;
648        let committed_object_ids = sst_infos
649            .into_iter()
650            .flatten()
651            .map(|s| s.object_id)
652            .sorted()
653            .collect_vec();
654        assert!(!committed_object_ids.is_empty());
655        let max_committed_object_id = *committed_object_ids.iter().max().unwrap();
656        assert_eq!(
657            1,
658            hummock_manager
659                .complete_gc_batch(
660                    [committed_object_ids, vec![max_committed_object_id + 1]]
661                        .concat()
662                        .into_iter()
663                        .map(HummockObjectId::Sstable)
664                        .collect(),
665                    None,
666                )
667                .await
668                .unwrap()
669        );
670    }
671}