1use 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 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 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 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 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 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 pub async fn delete_version_deltas(&self) -> Result<usize> {
188 let mut versioning_guard = self
189 .versioning
190 .write_with_process_name("delete_version_deltas")
191 .await;
192 let versioning = versioning_guard.deref_mut();
193 let context_info = self
194 .context_info
195 .read_with_process_name("delete_version_deltas")
196 .await;
197 if !context_info.version_safe_points.is_empty() {
200 return Ok(0);
201 }
202 let version_id = versioning.checkpoint.version.id;
205 let res = hummock_version_delta::Entity::delete_many()
206 .filter(hummock_version_delta::Column::Id.lte(version_id))
207 .exec(&self.env.meta_store_ref().conn)
208 .await?;
209 tracing::debug!(rows_affected = res.rows_affected, "Deleted version deltas");
210 versioning
211 .hummock_version_deltas
212 .retain(|id, _| *id > version_id);
213 #[cfg(test)]
214 {
215 drop(context_info);
216 drop(versioning_guard);
217 self.check_state_consistency().await;
218 }
219 Ok(res.rows_affected as usize)
220 }
221
222 pub async fn finalize_objects_to_delete(
224 &self,
225 object_ids: impl Iterator<Item = HummockObjectId>,
226 ) -> Result<Vec<HummockObjectId>> {
227 let versioning = self
229 .versioning
230 .read_with_process_name("finalize_objects_to_delete")
231 .await;
232 let min_pinned_version_id = self
233 .context_info
234 .read_with_process_name("finalize_objects_to_delete")
235 .await
236 .min_pinned_version_id();
237 let tracked_object_ids: HashSet<HummockObjectId> =
238 versioning.get_tracked_object_ids(min_pinned_version_id);
239 let to_delete = object_ids
240 .filter(|object_id| !tracked_object_ids.contains(object_id))
241 .collect_vec();
242 self.write_gc_history(to_delete.iter().copied()).await?;
245 Ok(to_delete)
246 }
247
248 pub async fn start_full_gc(
251 &self,
252 object_retention_time: Duration,
253 prefix: Option<String>,
254 backup_manager: Option<BackupManagerRef>,
255 ) -> Result<()> {
256 if !self.full_gc_state.try_start() {
257 return Err(anyhow::anyhow!("failed to start GC due to an ongoing process").into());
258 }
259 let _guard = scopeguard::guard(self.full_gc_state.clone(), |full_gc_state| {
260 full_gc_state.stop()
261 });
262 self.metrics.full_gc_trigger_count.inc();
263 let object_retention_time = cmp::max(
264 object_retention_time,
265 Duration::from_secs(self.env.opts.min_sst_retention_time_sec),
266 );
267 let limit = self.env.opts.full_gc_object_limit;
268 let mut start_after = None;
269 let object_retention_watermark = self
270 .now()
271 .await?
272 .saturating_sub(object_retention_time.as_secs());
273 let mut total_object_count = 0;
274 let mut total_object_size = 0;
275 tracing::info!(
276 retention_sec = object_retention_time.as_secs(),
277 prefix,
278 limit,
279 "Start GC."
280 );
281 loop {
282 tracing::debug!(
283 retention_sec = object_retention_time.as_secs(),
284 prefix,
285 start_after,
286 limit,
287 "Start a GC batch."
288 );
289 let (object_ids, batch_object_count, batch_object_size, next_start_after) = self
290 .gc_manager
291 .list_objects(
292 object_retention_watermark,
293 prefix.clone(),
294 start_after.clone(),
295 Some(limit),
296 )
297 .await?;
298 total_object_count += batch_object_count;
299 total_object_size += batch_object_size;
300 tracing::debug!(
301 ?object_ids,
302 batch_object_count,
303 batch_object_size,
304 "Finish listing a GC batch."
305 );
306 self.complete_gc_batch(object_ids, backup_manager.clone())
307 .await?;
308 if next_start_after.is_none() {
309 break;
310 }
311 start_after = next_start_after;
312 }
313 tracing::info!(total_object_count, total_object_size, "Finish GC");
314 self.metrics.total_object_size.set(total_object_size as _);
315 self.metrics.total_object_count.set(total_object_count as _);
316 match self.time_travel_pinned_object_count().await {
317 Ok(count) => {
318 self.metrics.time_travel_object_count.set(count as _);
319 }
320 Err(err) => {
321 use thiserror_ext::AsReport;
322 tracing::warn!(error = %err.as_report(), "Failed to count time travel objects.");
323 }
324 }
325 Ok(())
326 }
327
328 pub(crate) async fn complete_gc_batch(
331 &self,
332 object_ids: HashSet<HummockObjectId>,
333 backup_manager: Option<BackupManagerRef>,
334 ) -> Result<usize> {
335 if object_ids.is_empty() {
336 return Ok(0);
337 }
338 let pinned_by_metadata_backup = match backup_manager.as_ref() {
339 Some(b) => b.list_pinned_object_ids().await,
340 None => HashSet::default(),
341 };
342 let min_object_id = collect_min_uncommitted_object_id(
346 &self.metadata_manager,
347 self.env.stream_client_pool(),
348 )
349 .await?;
350 let metrics = &self.metrics;
351 let candidate_object_number = object_ids.len();
352 metrics
353 .full_gc_candidate_object_count
354 .observe(candidate_object_number as _);
355 let object_ids = object_ids
357 .into_iter()
358 .filter(|s| !pinned_by_metadata_backup.contains(&s.as_raw()))
359 .collect_vec();
360 let after_metadata_backup = object_ids.len();
361 let filter_by_time_travel_start_time = Instant::now();
363 let object_ids = self
364 .filter_out_objects_by_time_travel(object_ids.into_iter())
365 .await?;
366 tracing::info!(elapsed = ?filter_by_time_travel_start_time.elapsed(), "filter out objects by time travel in full GC");
367 let after_time_travel = object_ids.len();
368 let object_ids = object_ids
370 .into_iter()
371 .filter(|id| id.as_raw() < min_object_id)
372 .collect_vec();
373 let after_min_object_id = object_ids.len();
374 let after_version = self
376 .finalize_objects_to_delete(object_ids.into_iter())
377 .await?;
378 let after_version_count = after_version.len();
379 metrics
380 .full_gc_selected_object_count
381 .observe(after_version_count as _);
382 tracing::info!(
383 candidate_object_number,
384 after_metadata_backup,
385 after_time_travel,
386 after_min_object_id,
387 after_version_count,
388 "complete gc batch"
389 );
390 self.delete_objects(after_version).await?;
391 Ok(after_version_count)
392 }
393
394 pub async fn now(&self) -> Result<u64> {
395 let mut guard = self.now.lock().await;
396 let new_now = SystemTime::now()
397 .duration_since(SystemTime::UNIX_EPOCH)
398 .expect("Clock may have gone backwards")
399 .as_secs();
400 if new_now < *guard {
401 return Err(anyhow::anyhow!(format!(
402 "unexpected decreasing now, old={}, new={}",
403 *guard, new_now
404 ))
405 .into());
406 }
407 *guard = new_now;
408 drop(guard);
409 let m = hummock_sequence::ActiveModel {
411 name: ActiveValue::Set(HUMMOCK_NOW.into()),
412 seq: ActiveValue::Set(new_now.try_into().unwrap()),
413 };
414 hummock_sequence::Entity::insert(m)
415 .on_conflict(
416 OnConflict::column(hummock_sequence::Column::Name)
417 .update_column(hummock_sequence::Column::Seq)
418 .to_owned(),
419 )
420 .exec(&self.env.meta_store_ref().conn)
421 .await?;
422 Ok(new_now)
423 }
424
425 pub(crate) async fn load_now(&self) -> Result<Option<u64>> {
426 let now = hummock_sequence::Entity::find_by_id(HUMMOCK_NOW.to_owned())
427 .one(&self.env.meta_store_ref().conn)
428 .await?
429 .map(|m| m.seq.try_into().unwrap());
430 Ok(now)
431 }
432
433 async fn write_gc_history(
434 &self,
435 object_ids: impl Iterator<Item = HummockObjectId>,
436 ) -> Result<()> {
437 if self.env.opts.gc_history_retention_time_sec == 0 {
438 return Ok(());
439 }
440 let now = self.now().await?;
441 let dt = DateTime::from_timestamp(now.try_into().unwrap(), 0).unwrap();
442 let mut models = object_ids.map(|o| hummock_gc_history::ActiveModel {
443 object_id: Set(o.as_raw().into()),
444 mark_delete_at: Set(dt.naive_utc()),
445 });
446 let db = &self.meta_store_ref().conn;
447 let gc_history_low_watermark = DateTime::from_timestamp(
448 now.saturating_sub(self.env.opts.gc_history_retention_time_sec)
449 .try_into()
450 .unwrap(),
451 0,
452 )
453 .unwrap();
454 hummock_gc_history::Entity::delete_many()
455 .filter(hummock_gc_history::Column::MarkDeleteAt.lt(gc_history_low_watermark))
456 .exec(db)
457 .await?;
458 let mut is_finished = false;
459 while !is_finished {
460 let mut batch = vec![];
461 let mut count: usize = self.env.opts.hummock_gc_history_insert_batch_size;
462 while count > 0 {
463 let Some(m) = models.next() else {
464 is_finished = true;
465 break;
466 };
467 count -= 1;
468 batch.push(m);
469 }
470 if batch.is_empty() {
471 break;
472 }
473 hummock_gc_history::Entity::insert_many(batch)
474 .on_conflict_do_nothing()
475 .exec(db)
476 .await?;
477 }
478 Ok(())
479 }
480
481 pub async fn delete_time_travel_metadata(
482 &self,
483 pinned_snapshot_epochs: HashMap<TableId, HashSet<HummockEpoch>>,
484 ) -> MetaResult<()> {
485 let current_epoch_time = Epoch::now().physical_time();
486 let epoch_watermark = Epoch::from_physical_time(
487 current_epoch_time.saturating_sub(
488 self.env
489 .system_params_reader()
490 .await
491 .time_travel_retention_ms(),
492 ),
493 )
494 .0;
495 self.truncate_time_travel_metadata(epoch_watermark, pinned_snapshot_epochs)
496 .await?;
497 Ok(())
498 }
499
500 pub async fn delete_objects(&self, objects_to_delete: Vec<HummockObjectId>) -> Result<usize> {
505 let total = objects_to_delete.len();
506 for objects in objects_to_delete.chunks(1000) {
507 if self.env.opts.vacuum_spin_interval_ms != 0 {
508 tokio::time::sleep(Duration::from_millis(self.env.opts.vacuum_spin_interval_ms))
509 .await;
510 }
511 let delete_batch: HashSet<_> = objects.iter().copied().collect();
512 tracing::info!(?delete_batch, "Attempt to delete objects.");
513 self.gc_manager
514 .delete_objects(delete_batch.iter().copied())
515 .await?;
516 tracing::debug!(deleted_object_ids = ?delete_batch, "Finish deleting objects.");
517 }
518 Ok(total)
519 }
520
521 pub async fn try_start_minor_gc(&self, backup_manager: BackupManagerRef) -> Result<()> {
523 const MIN_MINOR_GC_OBJECT_COUNT: usize = 1000;
524 let Some(object_ids) = self
525 .gc_manager
526 .try_take_may_delete_object_ids(MIN_MINOR_GC_OBJECT_COUNT)
527 else {
528 return Ok(());
529 };
530 let backup_pinned: HashSet<_> = backup_manager.list_pinned_object_ids().await;
532 let version_pinned = {
534 let versioning = self
535 .versioning
536 .read_with_process_name("try_start_minor_gc")
537 .await;
538 let min_pinned_version_id = self
539 .context_info
540 .read_with_process_name("try_start_minor_gc")
541 .await
542 .min_pinned_version_id();
543 versioning.get_tracked_object_ids(min_pinned_version_id)
544 };
545 let object_ids = object_ids
546 .into_iter()
547 .filter(|s| !version_pinned.contains(s) && !backup_pinned.contains(&s.as_raw()));
548 let filter_by_time_travel_start_time = Instant::now();
549 let object_ids = self.filter_out_objects_by_time_travel(object_ids).await?;
550 tracing::info!(elapsed = ?filter_by_time_travel_start_time.elapsed(), "filter out objects by time travel in minor GC");
551 self.delete_objects(object_ids.into_iter().collect())
553 .await?;
554 Ok(())
555 }
556}
557
558async fn collect_min_uncommitted_object_id(
559 metadata_manager: &MetadataManager,
560 client_pool: &StreamClientPool,
561) -> Result<HummockRawObjectId> {
562 let futures = metadata_manager
563 .list_active_streaming_compute_nodes()
564 .await
565 .map_err(|err| Error::MetaStore(err.into()))?
566 .into_iter()
567 .map(|worker_node| async move {
568 let client = client_pool.get(&worker_node).await?;
569 let request = GetMinUncommittedObjectIdRequest {};
570 client.get_min_uncommitted_object_id(request).await
571 });
572 let min_watermark = try_join_all(futures)
573 .await
574 .map_err(|err| Error::Internal(err.into()))?
575 .into_iter()
576 .map(|resp| resp.min_uncommitted_object_id)
577 .min()
578 .unwrap_or(u64::MAX.into());
579 Ok(min_watermark)
580}
581
582pub struct FullGcState {
583 is_started: AtomicBool,
584}
585
586impl FullGcState {
587 pub fn new() -> Self {
588 Self {
589 is_started: AtomicBool::new(false),
590 }
591 }
592
593 pub fn try_start(&self) -> bool {
594 self.is_started
595 .compare_exchange(false, true, Ordering::SeqCst, Ordering::SeqCst)
596 .is_ok()
597 }
598
599 pub fn stop(&self) {
600 self.is_started.store(false, Ordering::SeqCst);
601 }
602}
603
604#[cfg(test)]
605mod tests {
606 use std::sync::Arc;
607 use std::time::Duration;
608
609 use itertools::Itertools;
610 use risingwave_hummock_sdk::HummockObjectId;
611 use risingwave_hummock_sdk::compaction_group::StaticCompactionGroupId;
612 use risingwave_rpc_client::HummockMetaClient;
613
614 use crate::hummock::MockHummockMetaClient;
615 use crate::hummock::test_utils::{add_test_tables, setup_compute_env};
616
617 #[tokio::test]
618 async fn test_full_gc() {
619 let (_env, hummock_manager, _cluster_manager, worker_id) = setup_compute_env(80).await;
620 let hummock_meta_client: Arc<dyn HummockMetaClient> = Arc::new(MockHummockMetaClient::new(
621 hummock_manager.clone(),
622 worker_id as _,
623 ));
624 let compaction_group_id = StaticCompactionGroupId::StateDefault;
625 hummock_manager
626 .start_full_gc(
627 Duration::from_secs(hummock_manager.env.opts.min_sst_retention_time_sec + 1),
628 None,
629 None,
630 )
631 .await
632 .unwrap();
633
634 hummock_manager
636 .complete_gc_batch(vec![].into_iter().collect(), None)
637 .await
638 .unwrap();
639
640 assert_eq!(
643 3,
644 hummock_manager
645 .complete_gc_batch(
646 [i64::MAX as u64 - 2, i64::MAX as u64 - 1, i64::MAX as u64]
647 .into_iter()
648 .map(|id| HummockObjectId::Sstable(id.into()))
649 .collect(),
650 None,
651 )
652 .await
653 .unwrap()
654 );
655
656 let sst_infos = add_test_tables(
658 hummock_manager.as_ref(),
659 hummock_meta_client.clone(),
660 compaction_group_id,
661 )
662 .await;
663 let committed_object_ids = sst_infos
664 .into_iter()
665 .flatten()
666 .map(|s| s.object_id)
667 .sorted()
668 .collect_vec();
669 assert!(!committed_object_ids.is_empty());
670 let max_committed_object_id = *committed_object_ids.iter().max().unwrap();
671 assert_eq!(
672 1,
673 hummock_manager
674 .complete_gc_batch(
675 [committed_object_ids, vec![max_committed_object_id + 1]]
676 .concat()
677 .into_iter()
678 .map(HummockObjectId::Sstable)
679 .collect(),
680 None,
681 )
682 .await
683 .unwrap()
684 );
685 }
686}