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.versioning.write().await;
189 let versioning = versioning_guard.deref_mut();
190 let context_info = self.context_info.read().await;
191 if !context_info.version_safe_points.is_empty() {
194 return Ok(0);
195 }
196 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 pub async fn finalize_objects_to_delete(
218 &self,
219 object_ids: impl Iterator<Item = HummockObjectId>,
220 ) -> Result<Vec<HummockObjectId>> {
221 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 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 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 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 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 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 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 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 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 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 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 let backup_pinned: HashSet<_> = backup_manager.list_pinned_object_ids().await;
524 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 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 hummock_manager
621 .complete_gc_batch(vec![].into_iter().collect(), None)
622 .await
623 .unwrap();
624
625 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 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}