1use std::collections::{BTreeMap, HashMap, HashSet};
16use std::ops::DerefMut;
17use std::sync::{Arc, LazyLock};
18use std::time::Instant;
19
20use anyhow::Context;
21use compaction_event_loop::{
22 HummockCompactionEventDispatcher, HummockCompactionEventHandler, HummockCompactionEventLoop,
23 HummockCompactorDedicatedEventLoop,
24};
25use fail::fail_point;
26use itertools::Itertools;
27use parking_lot::Mutex;
28use rand::rng as thread_rng;
29use rand::seq::SliceRandom;
30use risingwave_common::catalog::TableId;
31use risingwave_common::config::meta::default::compaction_config;
32use risingwave_common::util::epoch::Epoch;
33use risingwave_hummock_sdk::compact_task::{CompactTask, CompactTaskAssignment, ReportTask};
34use risingwave_hummock_sdk::compaction_group::StateTableId;
35use risingwave_hummock_sdk::compaction_group::hummock_version_ext::safe_epoch_table_watermarks_impl;
36use risingwave_hummock_sdk::level::Levels;
37use risingwave_hummock_sdk::sstable_info::SstableInfo;
38use risingwave_hummock_sdk::table_stats::{
39 PbTableStatsMap, add_prost_table_stats_map, purge_prost_table_stats,
40};
41use risingwave_hummock_sdk::table_watermark::TableWatermarks;
42use risingwave_hummock_sdk::version::{GroupDelta, IntraLevelDelta};
43use risingwave_hummock_sdk::{
44 CompactionGroupId, HummockCompactionTaskId, HummockContextId, HummockSstableId,
45 HummockSstableObjectId, HummockVersionId, compact_task_to_string, statistics_compact_task,
46};
47use risingwave_meta_model::hummock_sequence::COMPACTION_TASK_ID;
48use risingwave_pb::hummock::compact_task::{TaskStatus, TaskType};
49use risingwave_pb::hummock::subscribe_compaction_event_response::Event as ResponseEvent;
50use risingwave_pb::hummock::{
51 CompactTaskAssignment as PbCompactTaskAssignment, CompactionConfig, PbCompactStatus,
52 SubscribeCompactionEventRequest, TableOption, compact_task,
53};
54use thiserror_ext::AsReport;
55use tokio::sync::mpsc::UnboundedReceiver;
56use tokio::sync::oneshot::{Receiver, Sender};
57use tokio::task::JoinHandle;
58use tonic::Streaming;
59use tracing::warn;
60
61use crate::hummock::compaction::in_progress_compaction::InProgressCompactionView;
62use crate::hummock::compaction::selector::level_selector::PickerInfo;
63use crate::hummock::compaction::selector::{
64 DynamicLevelSelector, DynamicLevelSelectorCore, LocalSelectorStatistic, ManualCompactionOption,
65 ManualCompactionSelector, SpaceReclaimCompactionSelector, TombstoneCompactionSelector,
66 TtlCompactionSelector, VnodeWatermarkCompactionSelector,
67};
68use crate::hummock::compaction::{
69 CompactStatus, CompactionDeveloperConfig, CompactionSelector,
70 CompactionTask as PickedCompactionTask,
71};
72use crate::hummock::error::{Error, Result};
73use crate::hummock::manager::CompactionTaskReportResult;
74use crate::hummock::manager::compaction::compact_task_builder::{
75 CompactTaskBuildContext, attach_compact_task_table_metadata, build_base_compact_task,
76};
77use crate::hummock::manager::transaction::{
78 HummockVersionStatsTransaction, HummockVersionTransaction,
79};
80use crate::hummock::manager::versioning::Versioning;
81use crate::hummock::metrics_utils::{
82 build_compact_task_level_type_metrics_label, trigger_compact_tasks_stat,
83 trigger_local_table_stat,
84};
85use crate::hummock::model::CompactionGroup;
86use crate::hummock::{HummockManager, commit_multi_var};
87use crate::manager::META_NODE_ID;
88use crate::model::BTreeMapTransaction;
89
90#[derive(Debug, Clone, Copy, PartialEq, Eq)]
91pub enum ManualCompactionTriggerResult {
92 Submitted,
93 Retry,
94}
95
96mod compact_task_builder;
97pub mod compaction_event_loop;
98pub mod compaction_group_manager;
99pub mod compaction_group_schedule;
100
101static CANCEL_STATUS_SET: LazyLock<HashSet<TaskStatus>> = LazyLock::new(|| {
102 [
103 TaskStatus::ManualCanceled,
104 TaskStatus::SendFailCanceled,
105 TaskStatus::AssignFailCanceled,
106 TaskStatus::HeartbeatCanceled,
107 TaskStatus::InvalidGroupCanceled,
108 TaskStatus::NoAvailMemoryResourceCanceled,
109 TaskStatus::NoAvailCpuResourceCanceled,
110 TaskStatus::HeartbeatProgressCanceled,
111 ]
112 .into_iter()
113 .collect()
114});
115
116fn init_selectors() -> HashMap<compact_task::TaskType, Box<dyn CompactionSelector>> {
117 let mut compaction_selectors: HashMap<compact_task::TaskType, Box<dyn CompactionSelector>> =
118 HashMap::default();
119 compaction_selectors.insert(
120 compact_task::TaskType::Dynamic,
121 Box::<DynamicLevelSelector>::default(),
122 );
123 compaction_selectors.insert(
124 compact_task::TaskType::SpaceReclaim,
125 Box::<SpaceReclaimCompactionSelector>::default(),
126 );
127 compaction_selectors.insert(
128 compact_task::TaskType::Ttl,
129 Box::<TtlCompactionSelector>::default(),
130 );
131 compaction_selectors.insert(
132 compact_task::TaskType::Tombstone,
133 Box::<TombstoneCompactionSelector>::default(),
134 );
135 compaction_selectors.insert(
136 compact_task::TaskType::VnodeWatermark,
137 Box::<VnodeWatermarkCompactionSelector>::default(),
138 );
139 compaction_selectors
140}
141
142enum BuiltCompactTask {
143 MetaFinished(CompactTask),
144 PendingAssignment(CompactTask),
145}
146
147impl HummockVersionTransaction<'_> {
148 fn apply_compact_task(&mut self, compact_task: &CompactTask) {
149 let mut version_delta = self.new_delta();
150 let trivial_move = compact_task.is_trivial_move_task();
151 version_delta.trivial_move = trivial_move;
152
153 let group_deltas = &mut version_delta
154 .group_deltas
155 .entry(compact_task.compaction_group_id)
156 .or_default()
157 .group_deltas;
158 let mut removed_table_ids_map: BTreeMap<u32, HashSet<HummockSstableId>> =
159 BTreeMap::default();
160
161 for level in &compact_task.input_ssts {
162 let level_idx = level.level_idx;
163
164 removed_table_ids_map
165 .entry(level_idx)
166 .or_default()
167 .extend(level.table_infos.iter().map(|sst| sst.sst_id));
168 }
169
170 for (level_idx, removed_table_ids) in removed_table_ids_map {
171 let group_delta = GroupDelta::IntraLevel(IntraLevelDelta::new(
172 level_idx,
173 0, removed_table_ids,
175 vec![], 0, compact_task.compaction_group_version_id,
178 ));
179
180 group_deltas.push(group_delta);
181 }
182
183 let group_delta = GroupDelta::IntraLevel(IntraLevelDelta::new(
184 compact_task.target_level,
185 compact_task.target_sub_level_id,
186 HashSet::new(), compact_task.sorted_output_ssts.clone(),
188 compact_task.split_weight_by_vnode,
189 compact_task.compaction_group_version_id,
190 ));
191
192 group_deltas.push(group_delta);
193 version_delta.pre_apply();
194 }
195}
196
197#[derive(Default)]
198pub struct Compaction {
199 pub compact_task_assignment: BTreeMap<HummockCompactionTaskId, CompactTaskAssignment>,
201 pub compaction_statuses: BTreeMap<CompactionGroupId, CompactStatus>,
203
204 pub _deterministic_mode: bool,
205}
206
207impl HummockManager {
208 pub async fn get_assigned_compact_task_num(&self) -> u64 {
209 self.compaction
210 .read_with_process_name("get_assigned_compact_task_num")
211 .await
212 .compact_task_assignment
213 .len() as u64
214 }
215
216 pub async fn list_compaction_status(
217 &self,
218 ) -> (Vec<PbCompactStatus>, Vec<PbCompactTaskAssignment>) {
219 let (compaction_statuses, compact_task_assignments) = {
220 let compaction = self
221 .compaction
222 .read_with_process_name("list_compaction_status")
223 .await;
224 (
225 compaction
226 .compaction_statuses
227 .values()
228 .map_into()
229 .collect_vec(),
230 compaction
231 .compact_task_assignment
232 .values()
233 .cloned()
234 .collect_vec(),
235 )
236 };
237
238 (
239 compaction_statuses,
240 compact_task_assignments
241 .into_iter()
242 .map(PbCompactTaskAssignment::from)
243 .collect(),
244 )
245 }
246
247 pub async fn get_compaction_scores(
248 &self,
249 compaction_group_id: CompactionGroupId,
250 ) -> Vec<PickerInfo> {
251 let (status, levels, group) = {
252 let compaction = self
253 .compaction
254 .read_with_process_name("get_compaction_scores")
255 .await;
256 let versioning = self
257 .versioning
258 .read_with_process_name("get_compaction_scores")
259 .await;
260 let config_manager = self
261 .compaction_group_manager
262 .read_with_process_name("get_compaction_scores")
263 .await;
264 match (
265 compaction.compaction_statuses.get(&compaction_group_id),
266 versioning.current_version.levels.get(&compaction_group_id),
267 config_manager.try_get_compaction_group_config(compaction_group_id),
268 ) {
269 (Some(cs), Some(v), Some(cf)) => (cs.to_owned(), v.to_owned(), cf),
270 _ => {
271 return vec![];
272 }
273 }
274 };
275 let dynamic_level_core = DynamicLevelSelectorCore::new(
276 group.compaction_config,
277 Arc::new(CompactionDeveloperConfig::default()),
278 );
279 let ctx = dynamic_level_core.get_priority_levels(&levels, &status.level_handlers);
280 ctx.score_levels
281 }
282}
283
284impl HummockManager {
285 pub fn compaction_event_loop(
286 hummock_manager: Arc<Self>,
287 compactor_streams_change_rx: UnboundedReceiver<(
288 HummockContextId,
289 Streaming<SubscribeCompactionEventRequest>,
290 )>,
291 ) -> Vec<(JoinHandle<()>, Sender<()>)> {
292 let mut join_handle_vec = Vec::default();
293
294 let hummock_compaction_event_handler =
295 HummockCompactionEventHandler::new(hummock_manager.clone());
296
297 let dedicated_event_loop = HummockCompactorDedicatedEventLoop::new(
298 hummock_manager.clone(),
299 hummock_compaction_event_handler.clone(),
300 );
301
302 let (dedicated_event_loop_join_handle, event_tx, shutdown_tx) = dedicated_event_loop.run();
303 join_handle_vec.push((dedicated_event_loop_join_handle, shutdown_tx));
304
305 let hummock_compaction_event_dispatcher = HummockCompactionEventDispatcher::new(
306 hummock_manager.env.opts.clone(),
307 hummock_compaction_event_handler,
308 Some(event_tx),
309 );
310
311 let event_loop = HummockCompactionEventLoop::new(
312 hummock_compaction_event_dispatcher,
313 hummock_manager.metrics.clone(),
314 compactor_streams_change_rx,
315 );
316
317 let (event_loop_join_handle, event_loop_shutdown_tx) = event_loop.run();
318 join_handle_vec.push((event_loop_join_handle, event_loop_shutdown_tx));
319
320 join_handle_vec
321 }
322
323 pub fn add_compactor_stream(
324 &self,
325 context_id: HummockContextId,
326 req_stream: Streaming<SubscribeCompactionEventRequest>,
327 ) {
328 self.compactor_streams_change_tx
329 .send((context_id, req_stream))
330 .unwrap();
331 }
332}
333
334impl HummockManager {
335 async fn next_compaction_task_id_with_prefetch(&self, refill_capacity: u32) -> Result<u64> {
338 self.prefetched_compaction_task_ids
339 .next(refill_capacity, |count| async move {
340 self.env
341 .hummock_seq
342 .next_interval(COMPACTION_TASK_ID, count)
343 .await
344 })
345 .await
346 }
347
348 pub async fn get_compact_tasks_impl(
349 &self,
350 compaction_groups: Vec<CompactionGroupId>,
351 max_select_count: usize,
352 selector: &mut dyn CompactionSelector,
353 ) -> Result<(Vec<CompactTask>, Vec<CompactionGroupId>)> {
354 let deterministic_mode = self.env.opts.compaction_deterministic_test;
355
356 let mut compaction_guard = self
357 .compaction
358 .write_with_process_name("get_compact_tasks_impl")
359 .await;
360 let mut versioning_guard = self
361 .versioning
362 .write_with_process_name("get_compact_tasks_impl")
363 .await;
364 let compaction: &mut Compaction = &mut compaction_guard;
365 let versioning: &mut Versioning = &mut versioning_guard;
366
367 let start_time = Instant::now();
368 let mut compaction_statuses = BTreeMapTransaction::new(&mut compaction.compaction_statuses);
369
370 let mut compact_task_assignment =
371 BTreeMapTransaction::new(&mut compaction.compact_task_assignment);
372
373 let mut version = HummockVersionTransaction::new(
374 &mut versioning.current_version,
375 &mut versioning.hummock_version_deltas,
376 &mut versioning.table_change_log,
377 self.env.notification_manager(),
378 None,
379 &self.metrics,
380 &self.env.opts,
381 &self.version_stat_tx,
382 );
383 let mut version_stats = HummockVersionStatsTransaction::new(
385 &mut versioning.version_stats,
386 self.env.notification_manager(),
387 );
388
389 if deterministic_mode {
390 version.disable_apply_to_txn();
391 }
392 let all_versioned_table_schemas = if self.env.opts.enable_dropped_column_reclaim {
393 self.metadata_manager
394 .catalog_controller
395 .get_versioned_table_schemas()
396 .await
397 .map_err(|e| Error::Internal(e.into()))?
398 } else {
399 HashMap::default()
400 };
401 let mut unschedule_groups = vec![];
402 let mut trivial_tasks = vec![];
403 let mut pick_tasks = vec![];
404 let developer_config = Arc::new(CompactionDeveloperConfig::new_from_meta_opts(
405 &self.env.opts,
406 ));
407 'outside: for compaction_group_id in compaction_groups {
412 if pick_tasks.len() >= max_select_count {
413 break;
414 }
415
416 if !version
417 .latest_version()
418 .levels
419 .contains_key(&compaction_group_id)
420 {
421 continue;
423 }
424
425 let group_config = {
429 let config_manager = self
430 .compaction_group_manager
431 .read_with_process_name("get_compact_tasks_impl")
432 .await;
433
434 match config_manager.try_get_compaction_group_config(compaction_group_id) {
435 Some(config) => config,
436 None => continue,
437 }
438 };
439
440 let task_id = self
442 .next_compaction_task_id_with_prefetch(
443 self.env.opts.compaction_task_id_refill_capacity,
444 )
445 .await?;
446
447 if !compaction_statuses.contains_key(&compaction_group_id) {
448 compaction_statuses.insert(
450 compaction_group_id,
451 CompactStatus::new(
452 compaction_group_id,
453 group_config.compaction_config.max_level,
454 ),
455 );
456 }
457 let mut compact_status = compaction_statuses.get_mut(compaction_group_id).unwrap();
458
459 let mut stats = LocalSelectorStatistic::default();
460 let member_table_ids: Vec<_> = version
461 .latest_version()
462 .state_table_info
463 .compaction_group_member_table_ids(compaction_group_id)
464 .iter()
465 .copied()
466 .collect();
467
468 let mut table_id_to_option: HashMap<TableId, _> = HashMap::default();
469
470 {
471 let guard = self.table_id_to_table_option.read();
472 for table_id in &member_table_ids {
473 if let Some(opts) = guard.get(table_id) {
474 table_id_to_option.insert(*table_id, *opts);
475 }
476 }
477 }
478
479 let in_progress_compactions = InProgressCompactionView::for_group(
480 compact_task_assignment.tree_ref().values(),
481 compaction_group_id,
482 );
483
484 while let Some(picked_task) = compact_status.get_compact_task(
485 version
486 .latest_version()
487 .get_compaction_group_levels(compaction_group_id),
488 version
489 .latest_version()
490 .state_table_info
491 .compaction_group_member_table_ids(compaction_group_id),
492 task_id as HummockCompactionTaskId,
493 &group_config,
494 &mut stats,
495 selector,
496 &table_id_to_option,
497 developer_config.clone(),
498 &version.latest_version().table_watermarks,
499 &version.latest_version().state_table_info,
500 &in_progress_compactions,
501 ) {
502 let compaction_group_levels = version
503 .latest_version()
504 .get_compaction_group_levels(compaction_group_id);
505 let target_level_id = picked_task.input.target_level as u32;
506 let is_target_level_last = compaction_group_levels.is_last_level(target_level_id);
507 let table_options = table_id_to_option
508 .iter()
509 .map(|(table_id, table_option)| (*table_id, TableOption::from(table_option)))
510 .collect();
511 let built_compact_task = self.build_ready_compact_task(
512 picked_task,
513 CompactTaskBuildContext {
514 task_id,
515 compaction_group_id: group_config.group_id,
516 compaction_group_version_id: compaction_group_levels
517 .compaction_group_version_id,
518 existing_table_ids: member_table_ids.clone(),
519 table_options,
520 is_target_level_last,
521 compaction_config: group_config.compaction_config.clone(),
522 current_epoch_time: Epoch::now().0,
523 },
524 &version.latest_version().table_watermarks,
525 &all_versioned_table_schemas,
526 );
527
528 match built_compact_task {
529 BuiltCompactTask::MetaFinished(compact_task) => {
530 let label = compact_task.task_label();
531 tracing::debug!(
532 "{} for compaction group {}: input: {:?}, cost time: {:?}",
533 label,
534 compact_task.compaction_group_id,
535 compact_task.input_ssts,
536 start_time.elapsed()
537 );
538 compact_status.report_compact_task(&compact_task);
539 update_table_stats_for_vnode_watermark_trivial_reclaim(
540 &mut version_stats.table_stats,
541 &compact_task,
542 );
543 self.metrics
544 .compact_frequency
545 .with_label_values(&[
546 label,
547 &compact_task.compaction_group_id.to_string(),
548 selector.task_type().as_str_name(),
549 "SUCCESS",
550 ])
551 .inc();
552
553 version.apply_compact_task(&compact_task);
554 trivial_tasks.push(compact_task);
555 if trivial_tasks.len() >= self.env.opts.max_trivial_move_task_count_per_loop
556 {
557 break 'outside;
558 }
559 }
560 BuiltCompactTask::PendingAssignment(compact_task) => {
561 compact_task_assignment.insert(
562 compact_task.task_id,
563 CompactTaskAssignment {
564 compact_task: compact_task.clone(),
565 context_id: META_NODE_ID, },
567 );
568
569 pick_tasks.push(compact_task);
570 break;
571 }
572 }
573
574 stats.report_to_metrics(compaction_group_id, self.metrics.as_ref());
575 stats = LocalSelectorStatistic::default();
576 }
577 if pick_tasks
578 .last()
579 .map(|task| task.compaction_group_id != compaction_group_id)
580 .unwrap_or(true)
581 {
582 unschedule_groups.push(compaction_group_id);
583 }
584 stats.report_to_metrics(compaction_group_id, self.metrics.as_ref());
585 }
586
587 if !trivial_tasks.is_empty() {
588 commit_multi_var!(
589 self.meta_store_ref(),
590 compaction_statuses,
591 compact_task_assignment,
592 version,
593 version_stats
594 )?;
595 self.metrics
596 .compact_task_batch_count
597 .with_label_values(&["batch_trivial_move"])
598 .observe(trivial_tasks.len() as f64);
599
600 for trivial_task in &trivial_tasks {
601 self.metrics
602 .compact_task_trivial_move_sst_count
603 .with_label_values(&[&trivial_task.compaction_group_id.to_string()])
604 .observe(trivial_task.input_ssts[0].table_infos.len() as _);
605 }
606
607 drop(versioning_guard);
608 } else {
609 drop(versioning_guard);
612 commit_multi_var!(
613 self.meta_store_ref(),
614 compaction_statuses,
615 compact_task_assignment
616 )?;
617 }
618 drop(compaction_guard);
619 if !pick_tasks.is_empty() {
620 self.metrics
621 .compact_task_batch_count
622 .with_label_values(&["batch_get_compact_task"])
623 .observe(pick_tasks.len() as f64);
624 }
625
626 for compact_task in &mut pick_tasks {
627 let compaction_group_id = compact_task.compaction_group_id;
628
629 self.compactor_manager
631 .initiate_task_heartbeat(compact_task.clone());
632
633 compact_task.task_status = TaskStatus::Pending;
635 let compact_task_statistics = statistics_compact_task(compact_task);
636
637 let level_type_label = build_compact_task_level_type_metrics_label(
638 compact_task.input_ssts[0].level_idx as usize,
639 compact_task.input_ssts.last().unwrap().level_idx as usize,
640 );
641
642 let level_count = compact_task.input_ssts.len();
643 if compact_task.input_ssts[0].level_idx == 0 {
644 self.metrics
645 .l0_compact_level_count
646 .with_label_values(&[&compaction_group_id.to_string(), &level_type_label])
647 .observe(level_count as _);
648 }
649
650 self.metrics
651 .compact_task_size
652 .with_label_values(&[&compaction_group_id.to_string(), &level_type_label])
653 .observe(compact_task_statistics.total_file_size as _);
654
655 self.metrics
656 .compact_task_size
657 .with_label_values(&[
658 &compaction_group_id.to_string(),
659 &format!("{} uncompressed", level_type_label),
660 ])
661 .observe(compact_task_statistics.total_uncompressed_file_size as _);
662
663 self.metrics
664 .compact_task_file_count
665 .with_label_values(&[&compaction_group_id.to_string(), &level_type_label])
666 .observe(compact_task_statistics.total_file_count as _);
667
668 tracing::trace!(
669 "For compaction group {}: pick up {} {} sub_level in level {} to compact to target {}. cost time: {:?} compact_task_statistics {:?}",
670 compaction_group_id,
671 level_count,
672 compact_task.input_ssts[0].level_type.as_str_name(),
673 compact_task.input_ssts[0].level_idx,
674 compact_task.target_level,
675 start_time.elapsed(),
676 compact_task_statistics
677 );
678 }
679
680 #[cfg(test)]
681 {
682 self.check_state_consistency().await;
683 }
684 pick_tasks.extend(trivial_tasks);
685 Ok((pick_tasks, unschedule_groups))
686 }
687
688 pub async fn cancel_compact_task(&self, task_id: u64, task_status: TaskStatus) -> Result<bool> {
690 fail_point!("fp_cancel_compact_task", |_| Err(Error::MetaStore(
691 anyhow::anyhow!("failpoint metastore err")
692 )));
693 let ret = self
694 .cancel_compact_task_impl(vec![task_id], task_status)
695 .await?;
696 Ok(ret[0])
697 }
698
699 pub async fn cancel_compact_tasks(
700 &self,
701 tasks: Vec<u64>,
702 task_status: TaskStatus,
703 ) -> Result<Vec<bool>> {
704 self.cancel_compact_task_impl(tasks, task_status).await
705 }
706
707 async fn cancel_compact_task_impl(
708 &self,
709 task_ids: Vec<u64>,
710 task_status: TaskStatus,
711 ) -> Result<Vec<bool>> {
712 assert!(CANCEL_STATUS_SET.contains(&task_status));
713 let tasks = task_ids
714 .into_iter()
715 .map(|task_id| ReportTask {
716 task_id,
717 task_status,
718 sorted_output_ssts: vec![],
719 table_stats_change: HashMap::default(),
720 object_timestamps: HashMap::default(),
721 })
722 .collect_vec();
723 let rets = self.report_compact_tasks(tasks).await?;
724 #[cfg(test)]
725 {
726 self.check_state_consistency().await;
727 }
728 Ok(rets)
729 }
730
731 async fn get_compact_tasks(
732 &self,
733 mut compaction_groups: Vec<CompactionGroupId>,
734 max_select_count: usize,
735 selector: &mut dyn CompactionSelector,
736 ) -> Result<(Vec<CompactTask>, Vec<CompactionGroupId>)> {
737 fail_point!("fp_get_compact_task", |_| Err(Error::MetaStore(
738 anyhow::anyhow!("failpoint metastore error")
739 )));
740 compaction_groups.shuffle(&mut thread_rng());
741 let (mut tasks, groups) = self
742 .get_compact_tasks_impl(compaction_groups, max_select_count, selector)
743 .await?;
744 tasks.retain(|task| {
745 if task.task_status == TaskStatus::Success {
746 debug_assert!(task.is_trivial_reclaim() || task.is_trivial_move_task());
747 false
748 } else {
749 true
750 }
751 });
752 Ok((tasks, groups))
753 }
754
755 pub async fn get_compact_task(
756 &self,
757 compaction_group_id: CompactionGroupId,
758 selector: &mut dyn CompactionSelector,
759 ) -> Result<Option<CompactTask>> {
760 fail_point!("fp_get_compact_task", |_| Err(Error::MetaStore(
761 anyhow::anyhow!("failpoint metastore error")
762 )));
763
764 let (normal_tasks, _) = self
765 .get_compact_tasks_impl(vec![compaction_group_id], 1, selector)
766 .await?;
767 for task in normal_tasks {
768 if task.task_status != TaskStatus::Success {
769 return Ok(Some(task));
770 }
771 debug_assert!(task.is_trivial_reclaim() || task.is_trivial_move_task());
772 }
773 Ok(None)
774 }
775
776 pub async fn manual_get_compact_task(
777 &self,
778 compaction_group_id: CompactionGroupId,
779 manual_compaction_option: ManualCompactionOption,
780 ) -> Result<Option<CompactTask>> {
781 let (task, _) = self
782 .manual_get_compact_task_with_info(compaction_group_id, manual_compaction_option)
783 .await?;
784 Ok(task)
785 }
786
787 pub async fn manual_get_compact_task_with_info(
788 &self,
789 compaction_group_id: CompactionGroupId,
790 manual_compaction_option: ManualCompactionOption,
791 ) -> Result<(Option<CompactTask>, bool)> {
792 let mut selector = ManualCompactionSelector::new(manual_compaction_option);
793 let task = self
794 .get_compact_task(compaction_group_id, &mut selector)
795 .await?;
796 if let Some(err) = selector.validation_error() {
797 return Err(Error::InvalidManualCompactionOption(err.to_owned()));
798 }
799 Ok((task, selector.blocked_by_pending()))
800 }
801
802 pub async fn report_compact_task(
803 &self,
804 task_id: u64,
805 task_status: TaskStatus,
806 sorted_output_ssts: Vec<SstableInfo>,
807 table_stats_change: Option<PbTableStatsMap>,
808 object_timestamps: HashMap<HummockSstableObjectId, u64>,
809 ) -> Result<bool> {
810 let rets = self
811 .report_compact_tasks(vec![ReportTask {
812 task_id,
813 task_status,
814 sorted_output_ssts,
815 table_stats_change: table_stats_change.unwrap_or_default(),
816 object_timestamps,
817 }])
818 .await?;
819 Ok(rets[0])
820 }
821
822 pub async fn report_compact_tasks(&self, report_tasks: Vec<ReportTask>) -> Result<Vec<bool>> {
823 let compaction_guard = self
824 .compaction
825 .write_with_process_name("report_compact_tasks")
826 .await;
827 let versioning_guard = self
828 .versioning
829 .write_with_process_name("report_compact_tasks")
830 .await;
831
832 self.report_compact_tasks_impl(report_tasks, compaction_guard, versioning_guard)
833 .await
834 }
835
836 pub async fn report_compact_tasks_impl(
844 &self,
845 report_tasks: Vec<ReportTask>,
846 mut compaction_guard: impl DerefMut<Target = Compaction>,
847 mut versioning_guard: impl DerefMut<Target = Versioning>,
848 ) -> Result<Vec<bool>> {
849 let deterministic_mode = self.env.opts.compaction_deterministic_test;
850 let compaction: &mut Compaction = &mut compaction_guard;
851 let start_time = Instant::now();
852 let original_keys = compaction.compaction_statuses.keys().cloned().collect_vec();
853 let mut compact_statuses = BTreeMapTransaction::new(&mut compaction.compaction_statuses);
854 let mut rets = vec![false; report_tasks.len()];
855 let mut compact_task_assignment =
856 BTreeMapTransaction::new(&mut compaction.compact_task_assignment);
857 let versioning: &mut Versioning = &mut versioning_guard;
859
860 for group_id in original_keys {
862 if !versioning.current_version.levels.contains_key(&group_id) {
863 compact_statuses.remove(group_id);
864 }
865 }
866 let mut tasks = vec![];
867
868 let mut version = HummockVersionTransaction::new(
869 &mut versioning.current_version,
870 &mut versioning.hummock_version_deltas,
871 &mut versioning.table_change_log,
872 self.env.notification_manager(),
873 None,
874 &self.metrics,
875 &self.env.opts,
876 &self.version_stat_tx,
877 );
878
879 if deterministic_mode {
880 version.disable_apply_to_txn();
881 }
882
883 let mut version_stats = HummockVersionStatsTransaction::new(
884 &mut versioning.version_stats,
885 self.env.notification_manager(),
886 );
887 let mut success_count = 0;
888 let mut report_results = Vec::with_capacity(rets.len());
889 for (idx, task) in report_tasks.into_iter().enumerate() {
890 rets[idx] = true;
891 let task_id = task.task_id;
892 let mut task_status = task.task_status;
893 let mut compact_task = match compact_task_assignment.remove(task.task_id) {
894 Some(compact_task_assignment) => compact_task_assignment.compact_task,
895 None => {
896 tracing::warn!("{}", format!("compact task {} not found", task.task_id));
897 rets[idx] = false;
898 report_results.push(CompactionTaskReportResult {
899 task_id,
900 task_status,
901 reported: false,
902 });
903 continue;
904 }
905 };
906
907 {
908 compact_task.task_status = task.task_status;
910 compact_task.sorted_output_ssts = task.sorted_output_ssts;
911 }
912
913 match compact_statuses.get_mut(compact_task.compaction_group_id) {
914 Some(mut compact_status) => {
915 compact_status.report_compact_task(&compact_task);
916 }
917 None => {
918 compact_task.task_status = TaskStatus::InvalidGroupCanceled;
924 }
925 }
926
927 let is_success = if let TaskStatus::Success = compact_task.task_status {
928 match self
929 .report_compaction_sanity_check(&task.object_timestamps)
930 .await
931 {
932 Err(e) => {
933 warn!(
934 "failed to commit compaction task {} {}",
935 compact_task.task_id,
936 e.as_report()
937 );
938 compact_task.task_status = TaskStatus::RetentionTimeRejected;
939 false
940 }
941 _ => {
942 let group = version
943 .latest_version()
944 .levels
945 .get(&compact_task.compaction_group_id)
946 .unwrap();
947 let is_expired = compact_task.is_expired(group.compaction_group_version_id);
948 if is_expired {
949 compact_task.task_status = TaskStatus::InputOutdatedCanceled;
950 warn!(
951 "The task may be expired because of group split, task:\n {:?}",
952 compact_task_to_string(&compact_task)
953 );
954 }
955 !is_expired
956 }
957 }
958 } else {
959 false
960 };
961 if is_success {
962 success_count += 1;
963 version.apply_compact_task(&compact_task);
964 if purge_prost_table_stats(
965 &mut version_stats.table_stats,
966 version.latest_version(),
967 &HashSet::default(),
968 ) {
969 self.metrics.version_stats.reset();
970 versioning.local_metrics.clear();
971 }
972 add_prost_table_stats_map(&mut version_stats.table_stats, &task.table_stats_change);
973 trigger_local_table_stat(
974 &self.metrics,
975 &mut versioning.local_metrics,
976 &version_stats,
977 &task.table_stats_change,
978 );
979 }
980 task_status = compact_task.task_status;
981 report_results.push(CompactionTaskReportResult {
982 task_id,
983 task_status,
984 reported: rets[idx],
985 });
986 tasks.push(compact_task);
987 }
988 if success_count > 0 {
989 commit_multi_var!(
990 self.meta_store_ref(),
991 compact_statuses,
992 compact_task_assignment,
993 version,
994 version_stats
995 )?;
996
997 self.metrics
998 .compact_task_batch_count
999 .with_label_values(&["batch_report_task"])
1000 .observe(success_count as f64);
1001 } else {
1002 commit_multi_var!(
1004 self.meta_store_ref(),
1005 compact_statuses,
1006 compact_task_assignment
1007 )?;
1008 }
1009
1010 self.notify_compaction_task_report_waiters(report_results);
1011
1012 let mut success_groups = vec![];
1013 for compact_task in &tasks {
1014 self.compactor_manager
1015 .remove_task_heartbeat(compact_task.task_id);
1016 tracing::trace!(
1017 "Reported compaction task. {}. cost time: {:?}",
1018 compact_task_to_string(compact_task),
1019 start_time.elapsed(),
1020 );
1021
1022 if !deterministic_mode
1023 && versioning_guard
1024 .current_version
1025 .levels
1026 .contains_key(&compact_task.compaction_group_id)
1027 && (matches!(compact_task.task_type, compact_task::TaskType::Dynamic)
1028 || matches!(compact_task.task_type, compact_task::TaskType::Emergency))
1029 {
1030 self.try_send_compaction_request(
1032 compact_task.compaction_group_id,
1033 compact_task::TaskType::Dynamic,
1034 );
1035 }
1036
1037 if compact_task.task_status == TaskStatus::Success {
1038 success_groups.push(compact_task.compaction_group_id);
1039 }
1040 }
1041
1042 trigger_compact_tasks_stat(
1043 &self.metrics,
1044 &tasks,
1045 &compaction.compaction_statuses,
1046 &versioning_guard.current_version,
1047 );
1048 drop(versioning_guard);
1049 if !success_groups.is_empty() {
1050 self.try_update_write_limits(&success_groups).await;
1051 }
1052 Ok(rets)
1053 }
1054
1055 pub async fn trigger_compaction_deterministic(
1058 &self,
1059 _base_version_id: HummockVersionId,
1060 compaction_groups: Vec<CompactionGroupId>,
1061 ) -> Result<()> {
1062 self.on_current_version(|old_version| {
1063 tracing::info!(
1064 "Trigger compaction for version {}, groups {:?}",
1065 old_version.id,
1066 compaction_groups
1067 );
1068 for &compaction_group in &compaction_groups {
1069 if old_version.levels.contains_key(&compaction_group) {
1070 self.try_send_compaction_request(
1071 compaction_group,
1072 compact_task::TaskType::Dynamic,
1073 );
1074 }
1075 }
1076 })
1077 .await;
1078
1079 Ok(())
1080 }
1081
1082 pub async fn trigger_manual_compaction(
1083 &self,
1084 compaction_group: CompactionGroupId,
1085 manual_compaction_option: ManualCompactionOption,
1086 ) -> Result<ManualCompactionTriggerResult> {
1087 let start_time = Instant::now();
1088 let exclusive = manual_compaction_option.exclusive;
1089
1090 let compactor = match self.compactor_manager.next_compactor() {
1092 Some(compactor) => compactor,
1093 None => {
1094 tracing::warn!("trigger_manual_compaction No compactor is available.");
1095 return Err(anyhow::anyhow!(
1096 "trigger_manual_compaction No compactor is available. compaction_group {}",
1097 compaction_group
1098 )
1099 .into());
1100 }
1101 };
1102
1103 let compact_task = self
1105 .manual_get_compact_task_with_info(compaction_group, manual_compaction_option)
1106 .await;
1107 let (compact_task, blocked_by_pending) = match compact_task {
1108 Ok((compact_task, blocked_by_pending)) => (compact_task, blocked_by_pending),
1109 Err(err) => {
1110 tracing::warn!(error = %err.as_report(), "Failed to get compaction task");
1111 if matches!(err, Error::InvalidManualCompactionOption(_)) {
1112 return Err(err);
1113 }
1114
1115 return Err(anyhow::anyhow!(err)
1116 .context(format!(
1117 "Failed to get compaction task for compaction_group {}",
1118 compaction_group,
1119 ))
1120 .into());
1121 }
1122 };
1123 let compact_task = match compact_task {
1124 Some(compact_task) => compact_task,
1125 None => {
1126 if exclusive && blocked_by_pending {
1127 return Ok(ManualCompactionTriggerResult::Retry);
1128 }
1129 return Err(anyhow::anyhow!(
1131 "trigger_manual_compaction No compaction_task is available. compaction_group {}",
1132 compaction_group
1133 )
1134 .into());
1135 }
1136 };
1137
1138 let task_id = compact_task.task_id;
1140 let compact_task_string = compact_task_to_string(&compact_task);
1141 tracing::info!(
1142 compact_task_string,
1143 duration = ?start_time.elapsed(),
1144 "Triggered manual compaction task."
1145 );
1146
1147 let report_rx = self.register_compaction_task_report_waiter(task_id);
1148 if let Err(err) = compactor
1149 .send_event(ResponseEvent::CompactTask(compact_task.into()))
1150 .with_context(|| {
1151 format!(
1152 "Failed to trigger compaction task for compaction_group {}",
1153 compaction_group,
1154 )
1155 })
1156 {
1157 self.remove_compaction_task_report_waiter(task_id);
1158 return Err(err.into());
1159 }
1160
1161 let report_result = match report_rx.await {
1162 Ok(result) => result,
1163 Err(_) => {
1164 self.remove_compaction_task_report_waiter(task_id);
1165 return Err(anyhow::anyhow!(
1166 "trigger_manual_compaction wait report failed. compaction_group {}",
1167 compaction_group
1168 )
1169 .into());
1170 }
1171 };
1172 if !report_result.reported {
1173 return Err(anyhow::anyhow!(
1174 "trigger_manual_compaction report not accepted. task_id {}",
1175 report_result.task_id
1176 )
1177 .into());
1178 }
1179
1180 if report_result.task_status == TaskStatus::NoAvailCpuResourceCanceled
1181 || report_result.task_status == TaskStatus::NoAvailMemoryResourceCanceled
1182 {
1183 return Ok(ManualCompactionTriggerResult::Retry);
1184 }
1185
1186 tracing::info!(
1187 ?report_result,
1188 duration = ?start_time.elapsed(),
1189 "Completed manual compaction task."
1190 );
1191
1192 Ok(ManualCompactionTriggerResult::Submitted)
1193 }
1194
1195 pub(crate) fn try_send_compaction_request(
1199 &self,
1200 compaction_group: CompactionGroupId,
1201 task_type: compact_task::TaskType,
1202 ) -> bool {
1203 self.compaction_state.try_sched_compaction(
1204 compaction_group,
1205 task_type,
1206 ScheduleTrigger::NewData,
1207 )
1208 }
1209
1210 pub async fn trigger_compaction_for_all_groups(
1212 &self,
1213 task_type: TaskType,
1214 trigger: ScheduleTrigger,
1215 ) {
1216 let versioning = self
1218 .versioning
1219 .read_with_process_name("on_handle_trigger_multi_group")
1220 .await;
1221 #[cfg(test)]
1222 compaction_state_tests::before_trigger().await;
1223 for &group_id in versioning.current_version.levels.keys() {
1224 self.compaction_state
1225 .try_sched_compaction(group_id, task_type, trigger);
1226 }
1227 }
1228
1229 fn apply_split_weight_by_vnode_partition(
1232 &self,
1233 compact_task: &mut CompactTask,
1234 compaction_config: &CompactionConfig,
1235 compact_table_ids: &[TableId],
1236 ) {
1237 if compaction_config.split_weight_by_vnode > 0 {
1238 for table_id in compact_table_ids {
1239 compact_task
1240 .table_vnode_partition
1241 .insert(*table_id, compact_task.split_weight_by_vnode);
1242 }
1243
1244 return;
1245 }
1246
1247 let mut table_size_info: HashMap<TableId, u64> = HashMap::default();
1249 for input_ssts in &compact_task.input_ssts {
1250 for sst in &input_ssts.table_infos {
1251 for table_id in &sst.table_ids {
1252 *table_size_info.entry(*table_id).or_default() +=
1253 sst.sst_size / (sst.table_ids.len() as u64);
1254 }
1255 }
1256 }
1257
1258 let hybrid_vnode_count = self.env.opts.hybrid_partition_node_count;
1259 let default_partition_count = self.env.opts.partition_vnode_count;
1260 let compact_task_table_size_partition_threshold_low = self
1261 .env
1262 .opts
1263 .compact_task_table_size_partition_threshold_low;
1264 let compact_task_table_size_partition_threshold_high = self
1265 .env
1266 .opts
1267 .compact_task_table_size_partition_threshold_high;
1268
1269 let table_write_throughput_statistic_manager =
1271 self.table_write_throughput_statistic_manager.read();
1272
1273 for (table_id, compact_table_size) in table_size_info {
1274 let write_throughput = table_write_throughput_statistic_manager
1275 .latest_table_throughput(table_id)
1276 .unwrap_or(0);
1277
1278 if compact_table_size > compact_task_table_size_partition_threshold_high
1279 && default_partition_count > 0
1280 {
1281 compact_task
1282 .table_vnode_partition
1283 .insert(table_id, default_partition_count);
1284 } else if (compact_table_size > compact_task_table_size_partition_threshold_low
1285 || (write_throughput > self.env.opts.table_high_write_throughput_threshold
1286 && compact_table_size > compaction_config.target_file_size_base))
1287 && hybrid_vnode_count > 0
1288 {
1289 compact_task
1290 .table_vnode_partition
1291 .insert(table_id, hybrid_vnode_count);
1292 } else if compact_table_size > compaction_config.target_file_size_base {
1293 compact_task.table_vnode_partition.insert(table_id, 1);
1294 }
1295 }
1296
1297 compact_task
1298 .table_vnode_partition
1299 .retain(|table_id, _| compact_table_ids.contains(table_id));
1300 }
1301
1302 pub(crate) fn calculate_vnode_partition(
1303 &self,
1304 compact_task: &mut CompactTask,
1305 compaction_config: &CompactionConfig,
1306 compact_table_ids: &[TableId],
1307 ) {
1308 if compact_task.target_level > compact_task.base_level {
1313 return;
1314 }
1315
1316 self.apply_split_weight_by_vnode_partition(
1318 compact_task,
1319 compaction_config,
1320 compact_table_ids,
1321 );
1322 }
1323
1324 fn build_ready_compact_task(
1325 &self,
1326 picked_task: PickedCompactionTask,
1327 context: CompactTaskBuildContext,
1328 table_watermarks: &HashMap<TableId, Arc<TableWatermarks>>,
1329 all_versioned_table_schemas: &HashMap<TableId, Vec<i32>>,
1330 ) -> BuiltCompactTask {
1331 let compaction_config = context.compaction_config.clone();
1332 let (mut compact_task, compact_table_ids) = build_base_compact_task(picked_task, context);
1333
1334 if compact_task.is_trivial_reclaim() {
1335 compact_task.task_status = TaskStatus::Success;
1336 compact_task.sorted_output_ssts.clear();
1337 return BuiltCompactTask::MetaFinished(compact_task);
1338 }
1339
1340 if compact_task.is_trivial_move_task() {
1341 compact_task.task_status = TaskStatus::Success;
1342 compact_task.sorted_output_ssts = compact_task.input_ssts[0]
1343 .read_sstable_infos()
1344 .cloned()
1345 .collect();
1346 return BuiltCompactTask::MetaFinished(compact_task);
1347 }
1348
1349 self.prepare_compact_task_for_assignment(
1350 &mut compact_task,
1351 compaction_config.as_ref(),
1352 &compact_table_ids,
1353 safe_epoch_table_watermarks_impl(table_watermarks, &compact_table_ids),
1354 all_versioned_table_schemas,
1355 );
1356
1357 BuiltCompactTask::PendingAssignment(compact_task)
1358 }
1359
1360 fn prepare_compact_task_for_assignment(
1361 &self,
1362 compact_task: &mut CompactTask,
1363 compaction_config: &CompactionConfig,
1364 compact_table_ids: &[TableId],
1365 table_watermarks: BTreeMap<TableId, TableWatermarks>,
1366 all_versioned_table_schemas: &HashMap<TableId, Vec<i32>>,
1367 ) {
1368 self.calculate_vnode_partition(compact_task, compaction_config, compact_table_ids);
1369 attach_compact_task_table_metadata(
1370 compact_task,
1371 compact_table_ids,
1372 table_watermarks,
1373 all_versioned_table_schemas,
1374 );
1375 }
1376
1377 pub fn compactor_manager_ref(&self) -> crate::hummock::CompactorManagerRef {
1378 self.compactor_manager.clone()
1379 }
1380
1381 fn register_compaction_task_report_waiter(
1382 &self,
1383 task_id: HummockCompactionTaskId,
1384 ) -> Receiver<CompactionTaskReportResult> {
1385 let (tx, rx) = tokio::sync::oneshot::channel();
1386 self.compaction_task_report_notifiers
1387 .lock()
1388 .register(task_id, tx);
1389 rx
1390 }
1391
1392 fn remove_compaction_task_report_waiter(&self, task_id: HummockCompactionTaskId) {
1393 self.compaction_task_report_notifiers.lock().remove(task_id);
1394 }
1395
1396 fn notify_compaction_task_report_waiters(&self, results: Vec<CompactionTaskReportResult>) {
1397 let mut guard = self.compaction_task_report_notifiers.lock();
1398 for result in results {
1399 guard.notify(result);
1400 }
1401 }
1402}
1403
1404#[cfg(any(test, feature = "test"))]
1405impl HummockManager {
1406 pub async fn compaction_task_from_assignment_for_test(
1407 &self,
1408 task_id: u64,
1409 ) -> Option<CompactTaskAssignment> {
1410 let compaction_guard = self
1411 .compaction
1412 .read_with_process_name("compaction_task_from_assignment_for_test")
1413 .await;
1414 let assignment_ref = &compaction_guard.compact_task_assignment;
1415 assignment_ref.get(&task_id).cloned()
1416 }
1417
1418 pub async fn report_compact_task_for_test(
1419 &self,
1420 task_id: u64,
1421 compact_task: Option<CompactTask>,
1422 task_status: TaskStatus,
1423 sorted_output_ssts: Vec<SstableInfo>,
1424 table_stats_change: Option<PbTableStatsMap>,
1425 ) -> Result<()> {
1426 if let Some(task) = compact_task {
1427 let mut guard = self
1428 .compaction
1429 .write_with_process_name("report_compact_task_for_test")
1430 .await;
1431 guard.compact_task_assignment.insert(
1432 task_id,
1433 CompactTaskAssignment {
1434 compact_task: task,
1435 context_id: 0.into(),
1436 },
1437 );
1438 }
1439
1440 self.report_compact_tasks(vec![ReportTask {
1443 task_id,
1444 task_status,
1445 sorted_output_ssts,
1446 table_stats_change: table_stats_change.unwrap_or_default(),
1447 object_timestamps: HashMap::default(),
1448 }])
1449 .await?;
1450 Ok(())
1451 }
1452}
1453
1454#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1456pub enum ScheduleTrigger {
1457 NewData,
1459 Periodic,
1461}
1462
1463pub struct CompactionScheduleSnapshot {
1468 scheduled: HashSet<(CompactionGroupId, compact_task::TaskType)>,
1469 generation: u64,
1470}
1471
1472impl CompactionScheduleSnapshot {
1473 const TASK_TYPE_PRIORITY: &[TaskType] = &[
1475 TaskType::Dynamic,
1476 TaskType::SpaceReclaim,
1477 TaskType::Ttl,
1478 TaskType::Tombstone,
1479 TaskType::VnodeWatermark,
1480 ];
1481
1482 pub fn generation(&self) -> u64 {
1483 self.generation
1484 }
1485
1486 pub fn pick_compaction_groups_and_type(&self) -> Option<(Vec<CompactionGroupId>, TaskType)> {
1491 let group_ids = self.group_ids_shuffled();
1492 let mut normal_groups = vec![];
1493 for cg_id in group_ids {
1494 if let Some(pick_type) = self.pick_type(cg_id) {
1495 if pick_type == TaskType::Dynamic {
1496 normal_groups.push(cg_id);
1497 } else if normal_groups.is_empty() {
1498 return Some((vec![cg_id], pick_type));
1499 }
1500 }
1501 }
1502 if normal_groups.is_empty() {
1503 None
1504 } else {
1505 Some((normal_groups, TaskType::Dynamic))
1506 }
1507 }
1508
1509 fn group_ids_shuffled(&self) -> Vec<CompactionGroupId> {
1510 let mut group_ids: Vec<_> = self.scheduled.iter().map(|(g, _)| *g).unique().collect();
1511 group_ids.shuffle(&mut thread_rng());
1512 group_ids
1513 }
1514
1515 fn pick_type(&self, group: CompactionGroupId) -> Option<TaskType> {
1516 Self::TASK_TYPE_PRIORITY
1517 .iter()
1518 .find(|t| self.scheduled.contains(&(group, **t)))
1519 .copied()
1520 }
1521}
1522
1523#[derive(Debug, Default)]
1528pub struct CompactionState {
1529 inner: Mutex<CompactionStateInner>,
1530}
1531
1532#[derive(Debug, Default)]
1533struct CompactionStateInner {
1534 generation: u64,
1535 scheduled: HashSet<(CompactionGroupId, compact_task::TaskType)>,
1536 dynamic_cooldown: HashSet<CompactionGroupId>,
1538 last_new_data_generation: HashMap<CompactionGroupId, u64>,
1540}
1541
1542impl CompactionState {
1543 pub fn new() -> Self {
1544 Self {
1545 inner: Default::default(),
1546 }
1547 }
1548
1549 pub fn try_sched_compaction(
1553 &self,
1554 compaction_group: CompactionGroupId,
1555 task_type: TaskType,
1556 trigger: ScheduleTrigger,
1557 ) -> bool {
1558 let mut guard = self.inner.lock();
1559 if task_type == TaskType::Dynamic {
1560 match trigger {
1561 ScheduleTrigger::NewData => {
1562 guard.dynamic_cooldown.remove(&compaction_group);
1563 guard.generation += 1;
1565 let generation = guard.generation;
1566 guard
1567 .last_new_data_generation
1568 .insert(compaction_group, generation);
1569 }
1570 ScheduleTrigger::Periodic => {
1571 if guard.dynamic_cooldown.contains(&compaction_group) {
1572 return false;
1573 }
1574 }
1575 }
1576 }
1577 guard.scheduled.insert((compaction_group, task_type))
1578 }
1579
1580 pub fn unschedule(
1583 &self,
1584 compaction_group: CompactionGroupId,
1585 task_type: compact_task::TaskType,
1586 generation: u64,
1587 ) {
1588 let mut guard = self.inner.lock();
1589 if task_type == TaskType::Dynamic
1590 && guard
1591 .last_new_data_generation
1592 .get(&compaction_group)
1593 .is_some_and(|g| *g > generation)
1594 {
1595 return;
1596 }
1597 if guard.scheduled.remove(&(compaction_group, task_type)) && task_type == TaskType::Dynamic
1599 {
1600 guard.dynamic_cooldown.insert(compaction_group);
1601 }
1602 }
1603
1604 pub fn snapshot(&self) -> CompactionScheduleSnapshot {
1606 let guard = self.inner.lock();
1607 let generation = guard.generation;
1608 CompactionScheduleSnapshot {
1609 scheduled: guard.scheduled.clone(),
1610 generation,
1611 }
1612 }
1613
1614 pub fn remove_compaction_group(&self, compaction_group: CompactionGroupId) {
1616 let mut guard = self.inner.lock();
1617 guard
1618 .scheduled
1619 .retain(|(group, _)| *group != compaction_group);
1620 guard.dynamic_cooldown.remove(&compaction_group);
1621 guard.last_new_data_generation.remove(&compaction_group);
1622 }
1623}
1624
1625impl Compaction {
1626 pub fn get_compact_task_assignments_by_group_id(
1627 &self,
1628 compaction_group_id: CompactionGroupId,
1629 ) -> Vec<CompactTaskAssignment> {
1630 self.compact_task_assignment
1631 .values()
1632 .filter_map(|assignment| {
1633 if assignment.compact_task.compaction_group_id == compaction_group_id {
1634 Some(assignment.clone())
1635 } else {
1636 None
1637 }
1638 })
1639 .collect()
1640 }
1641}
1642
1643#[derive(Clone, Default)]
1644pub struct CompactionGroupStatistic {
1645 pub group_id: CompactionGroupId,
1646 pub group_size: u64,
1647 pub table_statistic: BTreeMap<StateTableId, u64>,
1648 pub compaction_group_config: CompactionGroup,
1649}
1650
1651fn update_table_stats_for_vnode_watermark_trivial_reclaim(
1653 table_stats: &mut PbTableStatsMap,
1654 task: &CompactTask,
1655) {
1656 if task.task_type != TaskType::VnodeWatermark {
1657 return;
1658 }
1659 let mut deleted_table_keys: HashMap<TableId, u64> = HashMap::default();
1660 for s in task.input_ssts.iter().flat_map(|l| l.table_infos.iter()) {
1661 assert_eq!(s.table_ids.len(), 1);
1662 let e = deleted_table_keys.entry(s.table_ids[0]).or_insert(0);
1663 *e += s.total_key_count;
1664 }
1665 for (table_id, delete_count) in deleted_table_keys {
1666 let Some(stats) = table_stats.get_mut(&table_id) else {
1667 continue;
1668 };
1669 if stats.total_key_count == 0 {
1670 continue;
1671 }
1672 let new_total_key_count = stats.total_key_count.saturating_sub(delete_count as i64);
1673 let ratio = new_total_key_count as f64 / stats.total_key_count as f64;
1674 stats.total_key_count = new_total_key_count;
1676 stats.total_key_size = (stats.total_key_size as f64 * ratio).ceil() as i64;
1678 stats.total_value_size = (stats.total_value_size as f64 * ratio).ceil() as i64;
1679 }
1680}
1681
1682#[derive(Debug, Clone)]
1683pub enum GroupState {
1684 Normal,
1686
1687 Emergency(String), WriteStop(String), }
1693
1694impl GroupState {
1695 pub fn is_write_stop(&self) -> bool {
1696 matches!(self, Self::WriteStop(_))
1697 }
1698
1699 pub fn is_emergency(&self) -> bool {
1700 matches!(self, Self::Emergency(_))
1701 }
1702
1703 pub fn reason(&self) -> Option<&str> {
1704 match self {
1705 Self::Emergency(reason) | Self::WriteStop(reason) => Some(reason),
1706 _ => None,
1707 }
1708 }
1709}
1710
1711#[derive(Clone, Default)]
1712pub struct GroupStateValidator;
1713
1714impl GroupStateValidator {
1715 pub fn write_stop_sub_level_count(
1716 level_count: usize,
1717 compaction_config: &CompactionConfig,
1718 ) -> bool {
1719 let threshold = compaction_config.level0_stop_write_threshold_sub_level_number as usize;
1720 level_count > threshold
1721 }
1722
1723 pub fn write_stop_l0_size(l0_size: u64, compaction_config: &CompactionConfig) -> bool {
1724 l0_size
1725 > compaction_config
1726 .level0_stop_write_threshold_max_size
1727 .unwrap_or(compaction_config::level0_stop_write_threshold_max_size())
1728 }
1729
1730 pub fn write_stop_l0_file_count(
1731 l0_file_count: usize,
1732 compaction_config: &CompactionConfig,
1733 ) -> bool {
1734 l0_file_count
1735 > compaction_config
1736 .level0_stop_write_threshold_max_sst_count
1737 .unwrap_or(compaction_config::level0_stop_write_threshold_max_sst_count())
1738 as usize
1739 }
1740
1741 pub fn emergency_l0_file_count(
1742 l0_file_count: usize,
1743 compaction_config: &CompactionConfig,
1744 ) -> bool {
1745 l0_file_count
1746 > compaction_config
1747 .emergency_level0_sst_file_count
1748 .unwrap_or(compaction_config::emergency_level0_sst_file_count())
1749 as usize
1750 }
1751
1752 pub fn emergency_l0_partition_count(
1753 last_l0_sub_level_partition_count: usize,
1754 compaction_config: &CompactionConfig,
1755 ) -> bool {
1756 last_l0_sub_level_partition_count
1757 > compaction_config
1758 .emergency_level0_sub_level_partition
1759 .unwrap_or(compaction_config::emergency_level0_sub_level_partition())
1760 as usize
1761 }
1762
1763 pub fn check_single_group_write_stop(
1764 levels: &Levels,
1765 compaction_config: &CompactionConfig,
1766 ) -> GroupState {
1767 if Self::write_stop_sub_level_count(levels.l0.sub_levels.len(), compaction_config) {
1768 return GroupState::WriteStop(format!(
1769 "WriteStop(l0_level_count: {}, threshold: {}) too many L0 sub levels",
1770 levels.l0.sub_levels.len(),
1771 compaction_config.level0_stop_write_threshold_sub_level_number
1772 ));
1773 }
1774
1775 if Self::write_stop_l0_file_count(
1776 levels
1777 .l0
1778 .sub_levels
1779 .iter()
1780 .map(|l| l.table_infos.len())
1781 .sum(),
1782 compaction_config,
1783 ) {
1784 return GroupState::WriteStop(format!(
1785 "WriteStop(l0_sst_count: {}, threshold: {}) too many L0 sst files",
1786 levels
1787 .l0
1788 .sub_levels
1789 .iter()
1790 .map(|l| l.table_infos.len())
1791 .sum::<usize>(),
1792 compaction_config
1793 .level0_stop_write_threshold_max_sst_count
1794 .unwrap_or(compaction_config::level0_stop_write_threshold_max_sst_count())
1795 ));
1796 }
1797
1798 if Self::write_stop_l0_size(levels.l0.total_file_size, compaction_config) {
1799 return GroupState::WriteStop(format!(
1800 "WriteStop(l0_size: {}, threshold: {}) too large L0 size",
1801 levels.l0.total_file_size,
1802 compaction_config
1803 .level0_stop_write_threshold_max_size
1804 .unwrap_or(compaction_config::level0_stop_write_threshold_max_size())
1805 ));
1806 }
1807
1808 GroupState::Normal
1809 }
1810
1811 pub fn check_single_group_emergency(
1812 levels: &Levels,
1813 compaction_config: &CompactionConfig,
1814 ) -> GroupState {
1815 if Self::emergency_l0_file_count(
1816 levels
1817 .l0
1818 .sub_levels
1819 .iter()
1820 .map(|l| l.table_infos.len())
1821 .sum(),
1822 compaction_config,
1823 ) {
1824 return GroupState::Emergency(format!(
1825 "Emergency(l0_sst_count: {}, threshold: {}) too many L0 sst files",
1826 levels
1827 .l0
1828 .sub_levels
1829 .iter()
1830 .map(|l| l.table_infos.len())
1831 .sum::<usize>(),
1832 compaction_config
1833 .emergency_level0_sst_file_count
1834 .unwrap_or(compaction_config::emergency_level0_sst_file_count())
1835 ));
1836 }
1837
1838 if Self::emergency_l0_partition_count(
1839 levels
1840 .l0
1841 .sub_levels
1842 .first()
1843 .map(|l| l.table_infos.len())
1844 .unwrap_or(0),
1845 compaction_config,
1846 ) {
1847 return GroupState::Emergency(format!(
1848 "Emergency(l0_partition_count: {}, threshold: {}) too many L0 partitions",
1849 levels
1850 .l0
1851 .sub_levels
1852 .first()
1853 .map(|l| l.table_infos.len())
1854 .unwrap_or(0),
1855 compaction_config
1856 .emergency_level0_sub_level_partition
1857 .unwrap_or(compaction_config::emergency_level0_sub_level_partition())
1858 ));
1859 }
1860
1861 GroupState::Normal
1862 }
1863
1864 pub fn group_state(levels: &Levels, compaction_config: &CompactionConfig) -> GroupState {
1865 let state = Self::check_single_group_write_stop(levels, compaction_config);
1866 if state.is_write_stop() {
1867 return state;
1868 }
1869
1870 Self::check_single_group_emergency(levels, compaction_config)
1871 }
1872}
1873
1874#[cfg(test)]
1875mod compaction_state_tests {
1876 use risingwave_pb::hummock::compact_task::TaskType;
1877
1878 use super::*;
1879
1880 tokio::task_local! {
1881 static BEFORE_TRIGGER: Arc<tokio::sync::Barrier>;
1882 }
1883
1884 pub(super) async fn before_trigger() {
1885 if let Ok(barrier) = BEFORE_TRIGGER.try_with(Arc::clone) {
1886 barrier.wait().await;
1887 barrier.wait().await;
1888 }
1889 }
1890
1891 #[tokio::test]
1892 #[cfg(not(madsim))]
1893 async fn test_trigger_all_groups_serializes_with_deletion() {
1894 use std::time::Duration;
1895
1896 use crate::hummock::test_utils::setup_compute_env;
1897
1898 for trigger in [ScheduleTrigger::NewData, ScheduleTrigger::Periodic] {
1899 let (_, manager, _, _) = setup_compute_env(80).await;
1900 manager
1901 .register_table_ids_for_test(&[(100, 2.into()), (101, 2.into())])
1902 .await
1903 .unwrap();
1904 let (group, _) = manager
1905 .move_state_tables_to_dedicated_compaction_group(2.into(), &[101.into()], None)
1906 .await
1907 .unwrap();
1908 let snapshot = manager.compaction_state.snapshot();
1909 manager
1910 .compaction_state
1911 .unschedule(2.into(), TaskType::Dynamic, snapshot.generation());
1912 let barrier = Arc::new(tokio::sync::Barrier::new(2));
1913 let delete = async {
1914 barrier.wait().await;
1915 {
1916 let write = tokio::task::unconstrained(manager.versioning.write());
1917 tokio::pin!(write);
1918 assert!(
1919 futures::poll!(write.as_mut()).is_pending(),
1920 "group deletion must wait until candidate publication finishes"
1921 );
1922 }
1923 barrier.wait().await;
1924 manager.unregister_table_ids([101.into()]).await.unwrap();
1925 };
1926 tokio::time::timeout(Duration::from_secs(10), async {
1927 tokio::join!(
1928 BEFORE_TRIGGER.scope(
1929 barrier.clone(),
1930 manager.trigger_compaction_for_all_groups(TaskType::Dynamic, trigger),
1931 ),
1932 delete,
1933 )
1934 })
1935 .await
1936 .unwrap();
1937 assert!(!manager.compaction_group_ids().await.contains(&group));
1938 assert_eq!(
1939 manager
1940 .compaction_state
1941 .snapshot()
1942 .scheduled
1943 .contains(&(2.into(), TaskType::Dynamic)),
1944 trigger == ScheduleTrigger::NewData,
1945 "registration wakes cooled groups; periodic scans must respect cooldown"
1946 );
1947 assert!(
1948 !manager
1949 .compaction_state
1950 .snapshot()
1951 .scheduled
1952 .iter()
1953 .any(|(id, _)| *id == group)
1954 );
1955 }
1956 }
1957
1958 #[test]
1959 fn test_basic_schedule_and_unschedule() {
1960 let state = CompactionState::new();
1961 let group_id: CompactionGroupId = 1.into();
1962
1963 assert!(state.try_sched_compaction(group_id, TaskType::Dynamic, ScheduleTrigger::NewData));
1965 assert!(!state.try_sched_compaction(group_id, TaskType::Dynamic, ScheduleTrigger::NewData));
1967 assert!(state.try_sched_compaction(group_id, TaskType::Ttl, ScheduleTrigger::Periodic));
1969
1970 let snapshot = state.snapshot();
1972 assert!(snapshot.scheduled.contains(&(group_id, TaskType::Dynamic)));
1973 assert!(snapshot.scheduled.contains(&(group_id, TaskType::Ttl)));
1974
1975 state.unschedule(group_id, TaskType::Dynamic, snapshot.generation());
1977 let snapshot2 = state.snapshot();
1978 assert!(!snapshot2.scheduled.contains(&(group_id, TaskType::Dynamic)));
1979 assert!(snapshot2.scheduled.contains(&(group_id, TaskType::Ttl)));
1980 }
1981
1982 #[test]
1983 fn test_cooldown_blocks_periodic_trigger() {
1984 let state = CompactionState::new();
1985 let group_id: CompactionGroupId = 1.into();
1986
1987 assert!(state.try_sched_compaction(group_id, TaskType::Dynamic, ScheduleTrigger::NewData));
1989 let snapshot = state.snapshot();
1990 state.unschedule(group_id, TaskType::Dynamic, snapshot.generation());
1991
1992 assert!(state.inner.lock().dynamic_cooldown.contains(&group_id));
1994
1995 assert!(!state.try_sched_compaction(
1997 group_id,
1998 TaskType::Dynamic,
1999 ScheduleTrigger::Periodic
2000 ));
2001 }
2002
2003 #[test]
2004 fn test_new_data_clears_cooldown() {
2005 let state = CompactionState::new();
2006 let group_id: CompactionGroupId = 1.into();
2007
2008 assert!(state.try_sched_compaction(group_id, TaskType::Dynamic, ScheduleTrigger::NewData));
2010 let snapshot = state.snapshot();
2011 state.unschedule(group_id, TaskType::Dynamic, snapshot.generation());
2012 assert!(state.inner.lock().dynamic_cooldown.contains(&group_id));
2013
2014 assert!(state.try_sched_compaction(group_id, TaskType::Dynamic, ScheduleTrigger::NewData));
2016 assert!(!state.inner.lock().dynamic_cooldown.contains(&group_id));
2017 }
2018
2019 #[test]
2020 fn test_cooldown_only_affects_dynamic_type() {
2021 let state = CompactionState::new();
2022 let group_id: CompactionGroupId = 1.into();
2023
2024 assert!(state.try_sched_compaction(group_id, TaskType::Dynamic, ScheduleTrigger::NewData));
2026 let snapshot = state.snapshot();
2027 state.unschedule(group_id, TaskType::Dynamic, snapshot.generation());
2028
2029 let group_id_2: CompactionGroupId = 2.into();
2031 assert!(state.try_sched_compaction(group_id_2, TaskType::Ttl, ScheduleTrigger::Periodic));
2032 let snapshot2 = state.snapshot();
2033 state.unschedule(group_id_2, TaskType::Ttl, snapshot2.generation());
2034 assert!(!state.inner.lock().dynamic_cooldown.contains(&group_id_2));
2035
2036 assert!(state.try_sched_compaction(group_id, TaskType::Ttl, ScheduleTrigger::Periodic));
2038 assert!(state.try_sched_compaction(
2039 group_id,
2040 TaskType::SpaceReclaim,
2041 ScheduleTrigger::Periodic
2042 ));
2043 }
2044
2045 #[test]
2046 fn test_race_condition_new_data_after_snapshot() {
2047 let state = CompactionState::new();
2048 let group_id: CompactionGroupId = 1.into();
2049
2050 assert!(state.try_sched_compaction(group_id, TaskType::Dynamic, ScheduleTrigger::NewData));
2051 let snapshot = state.snapshot();
2052
2053 assert!(!state.try_sched_compaction(group_id, TaskType::Dynamic, ScheduleTrigger::NewData));
2055 let latest = state.snapshot();
2056 state.unschedule(group_id, TaskType::Dynamic, snapshot.generation());
2057 assert!(
2058 state
2059 .snapshot()
2060 .scheduled
2061 .contains(&(group_id, TaskType::Dynamic))
2062 );
2063 assert!(!state.inner.lock().dynamic_cooldown.contains(&group_id));
2064
2065 state.try_sched_compaction(2.into(), TaskType::Dynamic, ScheduleTrigger::NewData);
2067 state.unschedule(group_id, TaskType::Dynamic, latest.generation());
2068 assert!(
2069 !state
2070 .snapshot()
2071 .scheduled
2072 .contains(&(group_id, TaskType::Dynamic))
2073 );
2074 assert!(!state.try_sched_compaction(
2075 group_id,
2076 TaskType::Dynamic,
2077 ScheduleTrigger::Periodic
2078 ));
2079
2080 assert!(state.try_sched_compaction(group_id, TaskType::Dynamic, ScheduleTrigger::NewData));
2083 state.unschedule(group_id, TaskType::Dynamic, latest.generation());
2084 assert!(
2085 state
2086 .snapshot()
2087 .scheduled
2088 .contains(&(group_id, TaskType::Dynamic))
2089 );
2090 assert!(!state.inner.lock().dynamic_cooldown.contains(&group_id));
2091 }
2092
2093 #[test]
2094 fn test_remove_compaction_group_cleans_all_state() {
2095 let state = CompactionState::new();
2096 let group_id: CompactionGroupId = 1.into();
2097
2098 assert!(state.try_sched_compaction(group_id, TaskType::Dynamic, ScheduleTrigger::NewData));
2100 assert!(state.try_sched_compaction(group_id, TaskType::Ttl, ScheduleTrigger::Periodic));
2101 state.inner.lock().dynamic_cooldown.insert(group_id);
2102
2103 let snapshot = state.snapshot();
2105 state.remove_compaction_group(group_id);
2106
2107 let guard = state.inner.lock();
2109 assert!(!guard.scheduled.contains(&(group_id, TaskType::Dynamic)));
2110 assert!(!guard.scheduled.contains(&(group_id, TaskType::Ttl)));
2111 assert!(!guard.dynamic_cooldown.contains(&group_id));
2112 assert!(!guard.last_new_data_generation.contains_key(&group_id));
2113 drop(guard);
2114
2115 state.unschedule(group_id, TaskType::Dynamic, snapshot.generation());
2117 let guard = state.inner.lock();
2118 assert!(guard.scheduled.is_empty());
2119 assert!(guard.dynamic_cooldown.is_empty());
2120 assert!(guard.last_new_data_generation.is_empty());
2121 }
2122
2123 #[test]
2124 fn test_snapshot_pick_type_priority() {
2125 let state = CompactionState::new();
2126 let group_id: CompactionGroupId = 1.into();
2127
2128 assert_eq!(state.snapshot().pick_type(group_id), None);
2130
2131 state.try_sched_compaction(
2133 group_id,
2134 TaskType::VnodeWatermark,
2135 ScheduleTrigger::Periodic,
2136 );
2137 assert_eq!(
2138 state.snapshot().pick_type(group_id),
2139 Some(TaskType::VnodeWatermark)
2140 );
2141
2142 state.try_sched_compaction(group_id, TaskType::Tombstone, ScheduleTrigger::Periodic);
2143 assert_eq!(
2144 state.snapshot().pick_type(group_id),
2145 Some(TaskType::Tombstone)
2146 );
2147
2148 state.try_sched_compaction(group_id, TaskType::Ttl, ScheduleTrigger::Periodic);
2149 assert_eq!(state.snapshot().pick_type(group_id), Some(TaskType::Ttl));
2150
2151 state.try_sched_compaction(group_id, TaskType::SpaceReclaim, ScheduleTrigger::Periodic);
2152 assert_eq!(
2153 state.snapshot().pick_type(group_id),
2154 Some(TaskType::SpaceReclaim)
2155 );
2156
2157 state.try_sched_compaction(group_id, TaskType::Dynamic, ScheduleTrigger::NewData);
2158 assert_eq!(
2159 state.snapshot().pick_type(group_id),
2160 Some(TaskType::Dynamic)
2161 );
2162 }
2163
2164 #[test]
2165 fn test_multiple_groups_independent_cooldown() {
2166 let state = CompactionState::new();
2167 let g1: CompactionGroupId = 1.into();
2168 let g2: CompactionGroupId = 2.into();
2169
2170 state.try_sched_compaction(g1, TaskType::Dynamic, ScheduleTrigger::NewData);
2171 state.try_sched_compaction(g2, TaskType::Dynamic, ScheduleTrigger::NewData);
2172 let snapshot = state.snapshot();
2173
2174 state.unschedule(g1, TaskType::Dynamic, snapshot.generation());
2176
2177 let guard = state.inner.lock();
2178 assert!(guard.dynamic_cooldown.contains(&g1));
2179 assert!(!guard.dynamic_cooldown.contains(&g2));
2180 }
2181
2182 #[test]
2183 fn test_pick_compaction_groups_empty() {
2184 let state = CompactionState::new();
2185 let snapshot = state.snapshot();
2186 assert!(snapshot.pick_compaction_groups_and_type().is_none());
2188 }
2189
2190 #[test]
2191 fn test_pick_compaction_groups_mixed_types() {
2192 let state = CompactionState::new();
2193 let g1: CompactionGroupId = 1.into();
2194 let g2: CompactionGroupId = 2.into();
2195 let g3: CompactionGroupId = 3.into();
2196
2197 state.try_sched_compaction(g1, TaskType::Dynamic, ScheduleTrigger::NewData);
2199 state.try_sched_compaction(g2, TaskType::Ttl, ScheduleTrigger::Periodic);
2200 state.try_sched_compaction(g3, TaskType::Dynamic, ScheduleTrigger::NewData);
2201
2202 let snapshot = state.snapshot();
2203 let (groups, task_type) = snapshot.pick_compaction_groups_and_type().unwrap();
2204
2205 if task_type == TaskType::Dynamic {
2210 assert!(groups.contains(&g1));
2211 assert!(groups.contains(&g3));
2212 assert!(!groups.contains(&g2)); } else {
2214 assert_eq!(task_type, TaskType::Ttl);
2215 assert_eq!(groups, vec![g2]);
2216 }
2217 }
2218
2219 #[test]
2220 fn test_pick_compaction_groups_all_dynamic() {
2221 let state = CompactionState::new();
2222 let g1: CompactionGroupId = 1.into();
2223 let g2: CompactionGroupId = 2.into();
2224
2225 state.try_sched_compaction(g1, TaskType::Dynamic, ScheduleTrigger::NewData);
2226 state.try_sched_compaction(g2, TaskType::Dynamic, ScheduleTrigger::NewData);
2227
2228 let snapshot = state.snapshot();
2229 let (groups, task_type) = snapshot.pick_compaction_groups_and_type().unwrap();
2230 assert_eq!(task_type, TaskType::Dynamic);
2231 assert!(groups.contains(&g1));
2232 assert!(groups.contains(&g2));
2233 }
2234
2235 #[test]
2236 fn test_pick_compaction_groups_single_non_dynamic() {
2237 let state = CompactionState::new();
2238 let g1: CompactionGroupId = 1.into();
2239
2240 state.try_sched_compaction(g1, TaskType::SpaceReclaim, ScheduleTrigger::Periodic);
2241
2242 let snapshot = state.snapshot();
2243 let (groups, task_type) = snapshot.pick_compaction_groups_and_type().unwrap();
2244 assert_eq!(task_type, TaskType::SpaceReclaim);
2245 assert_eq!(groups, vec![g1]);
2246 }
2247}