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