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;
422 }
423
424 let group_config = {
428 let config_manager = self
429 .compaction_group_manager
430 .read_with_process_name("get_compact_tasks_impl")
431 .await;
432
433 match config_manager.try_get_compaction_group_config(compaction_group_id) {
434 Some(config) => config,
435 None => continue,
436 }
437 };
438
439 let task_id = self
441 .next_compaction_task_id_with_prefetch(
442 self.env.opts.compaction_task_id_refill_capacity,
443 )
444 .await?;
445
446 if !compaction_statuses.contains_key(&compaction_group_id) {
447 compaction_statuses.insert(
449 compaction_group_id,
450 CompactStatus::new(
451 compaction_group_id,
452 group_config.compaction_config.max_level,
453 ),
454 );
455 }
456 let mut compact_status = compaction_statuses.get_mut(compaction_group_id).unwrap();
457
458 let mut stats = LocalSelectorStatistic::default();
459 let member_table_ids: Vec<_> = version
460 .latest_version()
461 .state_table_info
462 .compaction_group_member_table_ids(compaction_group_id)
463 .iter()
464 .copied()
465 .collect();
466
467 let mut table_id_to_option: HashMap<TableId, _> = HashMap::default();
468
469 {
470 let guard = self.table_id_to_table_option.read();
471 for table_id in &member_table_ids {
472 if let Some(opts) = guard.get(table_id) {
473 table_id_to_option.insert(*table_id, *opts);
474 }
475 }
476 }
477
478 let in_progress_compactions = InProgressCompactionView::for_group(
479 compact_task_assignment.tree_ref().values(),
480 compaction_group_id,
481 );
482
483 while let Some(picked_task) = compact_status.get_compact_task(
484 version
485 .latest_version()
486 .get_compaction_group_levels(compaction_group_id),
487 version
488 .latest_version()
489 .state_table_info
490 .compaction_group_member_table_ids(compaction_group_id),
491 task_id as HummockCompactionTaskId,
492 &group_config,
493 &mut stats,
494 selector,
495 &table_id_to_option,
496 developer_config.clone(),
497 &version.latest_version().table_watermarks,
498 &version.latest_version().state_table_info,
499 &in_progress_compactions,
500 ) {
501 let compaction_group_levels = version
502 .latest_version()
503 .get_compaction_group_levels(compaction_group_id);
504 let target_level_id = picked_task.input.target_level as u32;
505 let is_target_level_last = compaction_group_levels.is_last_level(target_level_id);
506 let table_options = table_id_to_option
507 .iter()
508 .map(|(table_id, table_option)| (*table_id, TableOption::from(table_option)))
509 .collect();
510 let built_compact_task = self.build_ready_compact_task(
511 picked_task,
512 CompactTaskBuildContext {
513 task_id,
514 compaction_group_id: group_config.group_id,
515 compaction_group_version_id: compaction_group_levels
516 .compaction_group_version_id,
517 existing_table_ids: member_table_ids.clone(),
518 table_options,
519 is_target_level_last,
520 compaction_config: group_config.compaction_config.clone(),
521 current_epoch_time: Epoch::now().0,
522 },
523 &version.latest_version().table_watermarks,
524 &all_versioned_table_schemas,
525 );
526
527 match built_compact_task {
528 BuiltCompactTask::MetaFinished(compact_task) => {
529 let label = compact_task.task_label();
530 tracing::debug!(
531 "{} for compaction group {}: input: {:?}, cost time: {:?}",
532 label,
533 compact_task.compaction_group_id,
534 compact_task.input_ssts,
535 start_time.elapsed()
536 );
537 compact_status.report_compact_task(&compact_task);
538 update_table_stats_for_vnode_watermark_trivial_reclaim(
539 &mut version_stats.table_stats,
540 &compact_task,
541 );
542 self.metrics
543 .compact_frequency
544 .with_label_values(&[
545 label,
546 &compact_task.compaction_group_id.to_string(),
547 selector.task_type().as_str_name(),
548 "SUCCESS",
549 ])
550 .inc();
551
552 version.apply_compact_task(&compact_task);
553 trivial_tasks.push(compact_task);
554 if trivial_tasks.len() >= self.env.opts.max_trivial_move_task_count_per_loop
555 {
556 break 'outside;
557 }
558 }
559 BuiltCompactTask::PendingAssignment(compact_task) => {
560 compact_task_assignment.insert(
561 compact_task.task_id,
562 CompactTaskAssignment {
563 compact_task: compact_task.clone(),
564 context_id: META_NODE_ID, },
566 );
567
568 pick_tasks.push(compact_task);
569 break;
570 }
571 }
572
573 stats.report_to_metrics(compaction_group_id, self.metrics.as_ref());
574 stats = LocalSelectorStatistic::default();
575 }
576 if pick_tasks
577 .last()
578 .map(|task| task.compaction_group_id != compaction_group_id)
579 .unwrap_or(true)
580 {
581 unschedule_groups.push(compaction_group_id);
582 }
583 stats.report_to_metrics(compaction_group_id, self.metrics.as_ref());
584 }
585
586 if !trivial_tasks.is_empty() {
587 commit_multi_var!(
588 self.meta_store_ref(),
589 compaction_statuses,
590 compact_task_assignment,
591 version,
592 version_stats
593 )?;
594 self.metrics
595 .compact_task_batch_count
596 .with_label_values(&["batch_trivial_move"])
597 .observe(trivial_tasks.len() as f64);
598
599 for trivial_task in &trivial_tasks {
600 self.metrics
601 .compact_task_trivial_move_sst_count
602 .with_label_values(&[&trivial_task.compaction_group_id.to_string()])
603 .observe(trivial_task.input_ssts[0].table_infos.len() as _);
604 }
605
606 drop(versioning_guard);
607 } else {
608 drop(versioning_guard);
611 commit_multi_var!(
612 self.meta_store_ref(),
613 compaction_statuses,
614 compact_task_assignment
615 )?;
616 }
617 drop(compaction_guard);
618 if !pick_tasks.is_empty() {
619 self.metrics
620 .compact_task_batch_count
621 .with_label_values(&["batch_get_compact_task"])
622 .observe(pick_tasks.len() as f64);
623 }
624
625 for compact_task in &mut pick_tasks {
626 let compaction_group_id = compact_task.compaction_group_id;
627
628 self.compactor_manager
630 .initiate_task_heartbeat(compact_task.clone());
631
632 compact_task.task_status = TaskStatus::Pending;
634 let compact_task_statistics = statistics_compact_task(compact_task);
635
636 let level_type_label = build_compact_task_level_type_metrics_label(
637 compact_task.input_ssts[0].level_idx as usize,
638 compact_task.input_ssts.last().unwrap().level_idx as usize,
639 );
640
641 let level_count = compact_task.input_ssts.len();
642 if compact_task.input_ssts[0].level_idx == 0 {
643 self.metrics
644 .l0_compact_level_count
645 .with_label_values(&[&compaction_group_id.to_string(), &level_type_label])
646 .observe(level_count as _);
647 }
648
649 self.metrics
650 .compact_task_size
651 .with_label_values(&[&compaction_group_id.to_string(), &level_type_label])
652 .observe(compact_task_statistics.total_file_size as _);
653
654 self.metrics
655 .compact_task_size
656 .with_label_values(&[
657 &compaction_group_id.to_string(),
658 &format!("{} uncompressed", level_type_label),
659 ])
660 .observe(compact_task_statistics.total_uncompressed_file_size as _);
661
662 self.metrics
663 .compact_task_file_count
664 .with_label_values(&[&compaction_group_id.to_string(), &level_type_label])
665 .observe(compact_task_statistics.total_file_count as _);
666
667 tracing::trace!(
668 "For compaction group {}: pick up {} {} sub_level in level {} to compact to target {}. cost time: {:?} compact_task_statistics {:?}",
669 compaction_group_id,
670 level_count,
671 compact_task.input_ssts[0].level_type.as_str_name(),
672 compact_task.input_ssts[0].level_idx,
673 compact_task.target_level,
674 start_time.elapsed(),
675 compact_task_statistics
676 );
677 }
678
679 #[cfg(test)]
680 {
681 self.check_state_consistency().await;
682 }
683 pick_tasks.extend(trivial_tasks);
684 Ok((pick_tasks, unschedule_groups))
685 }
686
687 pub async fn cancel_compact_task(&self, task_id: u64, task_status: TaskStatus) -> Result<bool> {
689 fail_point!("fp_cancel_compact_task", |_| Err(Error::MetaStore(
690 anyhow::anyhow!("failpoint metastore err")
691 )));
692 let ret = self
693 .cancel_compact_task_impl(vec![task_id], task_status)
694 .await?;
695 Ok(ret[0])
696 }
697
698 pub async fn cancel_compact_tasks(
699 &self,
700 tasks: Vec<u64>,
701 task_status: TaskStatus,
702 ) -> Result<Vec<bool>> {
703 self.cancel_compact_task_impl(tasks, task_status).await
704 }
705
706 async fn cancel_compact_task_impl(
707 &self,
708 task_ids: Vec<u64>,
709 task_status: TaskStatus,
710 ) -> Result<Vec<bool>> {
711 assert!(CANCEL_STATUS_SET.contains(&task_status));
712 let tasks = task_ids
713 .into_iter()
714 .map(|task_id| ReportTask {
715 task_id,
716 task_status,
717 sorted_output_ssts: vec![],
718 table_stats_change: HashMap::default(),
719 object_timestamps: HashMap::default(),
720 })
721 .collect_vec();
722 let rets = self.report_compact_tasks(tasks).await?;
723 #[cfg(test)]
724 {
725 self.check_state_consistency().await;
726 }
727 Ok(rets)
728 }
729
730 async fn get_compact_tasks(
731 &self,
732 mut compaction_groups: Vec<CompactionGroupId>,
733 max_select_count: usize,
734 selector: &mut dyn CompactionSelector,
735 ) -> Result<(Vec<CompactTask>, Vec<CompactionGroupId>)> {
736 fail_point!("fp_get_compact_task", |_| Err(Error::MetaStore(
737 anyhow::anyhow!("failpoint metastore error")
738 )));
739 compaction_groups.shuffle(&mut thread_rng());
740 let (mut tasks, groups) = self
741 .get_compact_tasks_impl(compaction_groups, max_select_count, selector)
742 .await?;
743 tasks.retain(|task| {
744 if task.task_status == TaskStatus::Success {
745 debug_assert!(task.is_trivial_reclaim() || task.is_trivial_move_task());
746 false
747 } else {
748 true
749 }
750 });
751 Ok((tasks, groups))
752 }
753
754 pub async fn get_compact_task(
755 &self,
756 compaction_group_id: CompactionGroupId,
757 selector: &mut dyn CompactionSelector,
758 ) -> Result<Option<CompactTask>> {
759 fail_point!("fp_get_compact_task", |_| Err(Error::MetaStore(
760 anyhow::anyhow!("failpoint metastore error")
761 )));
762
763 let (normal_tasks, _) = self
764 .get_compact_tasks_impl(vec![compaction_group_id], 1, selector)
765 .await?;
766 for task in normal_tasks {
767 if task.task_status != TaskStatus::Success {
768 return Ok(Some(task));
769 }
770 debug_assert!(task.is_trivial_reclaim() || task.is_trivial_move_task());
771 }
772 Ok(None)
773 }
774
775 pub async fn manual_get_compact_task(
776 &self,
777 compaction_group_id: CompactionGroupId,
778 manual_compaction_option: ManualCompactionOption,
779 ) -> Result<Option<CompactTask>> {
780 let (task, _) = self
781 .manual_get_compact_task_with_info(compaction_group_id, manual_compaction_option)
782 .await?;
783 Ok(task)
784 }
785
786 pub async fn manual_get_compact_task_with_info(
787 &self,
788 compaction_group_id: CompactionGroupId,
789 manual_compaction_option: ManualCompactionOption,
790 ) -> Result<(Option<CompactTask>, bool)> {
791 let mut selector = ManualCompactionSelector::new(manual_compaction_option);
792 let task = self
793 .get_compact_task(compaction_group_id, &mut selector)
794 .await?;
795 if let Some(err) = selector.validation_error() {
796 return Err(Error::InvalidManualCompactionOption(err.to_owned()));
797 }
798 Ok((task, selector.blocked_by_pending()))
799 }
800
801 pub async fn report_compact_task(
802 &self,
803 task_id: u64,
804 task_status: TaskStatus,
805 sorted_output_ssts: Vec<SstableInfo>,
806 table_stats_change: Option<PbTableStatsMap>,
807 object_timestamps: HashMap<HummockSstableObjectId, u64>,
808 ) -> Result<bool> {
809 let rets = self
810 .report_compact_tasks(vec![ReportTask {
811 task_id,
812 task_status,
813 sorted_output_ssts,
814 table_stats_change: table_stats_change.unwrap_or_default(),
815 object_timestamps,
816 }])
817 .await?;
818 Ok(rets[0])
819 }
820
821 pub async fn report_compact_tasks(&self, report_tasks: Vec<ReportTask>) -> Result<Vec<bool>> {
822 let compaction_guard = self
823 .compaction
824 .write_with_process_name("report_compact_tasks")
825 .await;
826 let versioning_guard = self
827 .versioning
828 .write_with_process_name("report_compact_tasks")
829 .await;
830
831 self.report_compact_tasks_impl(report_tasks, compaction_guard, versioning_guard)
832 .await
833 }
834
835 pub async fn report_compact_tasks_impl(
843 &self,
844 report_tasks: Vec<ReportTask>,
845 mut compaction_guard: impl DerefMut<Target = Compaction>,
846 mut versioning_guard: impl DerefMut<Target = Versioning>,
847 ) -> Result<Vec<bool>> {
848 let deterministic_mode = self.env.opts.compaction_deterministic_test;
849 let compaction: &mut Compaction = &mut compaction_guard;
850 let start_time = Instant::now();
851 let original_keys = compaction.compaction_statuses.keys().cloned().collect_vec();
852 let mut compact_statuses = BTreeMapTransaction::new(&mut compaction.compaction_statuses);
853 let mut rets = vec![false; report_tasks.len()];
854 let mut compact_task_assignment =
855 BTreeMapTransaction::new(&mut compaction.compact_task_assignment);
856 let versioning: &mut Versioning = &mut versioning_guard;
858
859 for group_id in original_keys {
861 if !versioning.current_version.levels.contains_key(&group_id) {
862 compact_statuses.remove(group_id);
863 }
864 }
865 let mut tasks = vec![];
866
867 let mut version = HummockVersionTransaction::new(
868 &mut versioning.current_version,
869 &mut versioning.hummock_version_deltas,
870 &mut versioning.table_change_log,
871 self.env.notification_manager(),
872 None,
873 &self.metrics,
874 &self.env.opts,
875 &self.version_stat_tx,
876 );
877
878 if deterministic_mode {
879 version.disable_apply_to_txn();
880 }
881
882 let mut version_stats = HummockVersionStatsTransaction::new(
883 &mut versioning.version_stats,
884 self.env.notification_manager(),
885 );
886 let mut success_count = 0;
887 let mut report_results = Vec::with_capacity(rets.len());
888 for (idx, task) in report_tasks.into_iter().enumerate() {
889 rets[idx] = true;
890 let task_id = task.task_id;
891 let mut task_status = task.task_status;
892 let mut compact_task = match compact_task_assignment.remove(task.task_id) {
893 Some(compact_task_assignment) => compact_task_assignment.compact_task,
894 None => {
895 tracing::warn!("{}", format!("compact task {} not found", task.task_id));
896 rets[idx] = false;
897 report_results.push(CompactionTaskReportResult {
898 task_id,
899 task_status,
900 reported: false,
901 });
902 continue;
903 }
904 };
905
906 {
907 compact_task.task_status = task.task_status;
909 compact_task.sorted_output_ssts = task.sorted_output_ssts;
910 }
911
912 match compact_statuses.get_mut(compact_task.compaction_group_id) {
913 Some(mut compact_status) => {
914 compact_status.report_compact_task(&compact_task);
915 }
916 None => {
917 compact_task.task_status = TaskStatus::InvalidGroupCanceled;
923 }
924 }
925
926 let is_success = if let TaskStatus::Success = compact_task.task_status {
927 match self
928 .report_compaction_sanity_check(&task.object_timestamps)
929 .await
930 {
931 Err(e) => {
932 warn!(
933 "failed to commit compaction task {} {}",
934 compact_task.task_id,
935 e.as_report()
936 );
937 compact_task.task_status = TaskStatus::RetentionTimeRejected;
938 false
939 }
940 _ => {
941 let group = version
942 .latest_version()
943 .levels
944 .get(&compact_task.compaction_group_id)
945 .unwrap();
946 let is_expired = compact_task.is_expired(group.compaction_group_version_id);
947 if is_expired {
948 compact_task.task_status = TaskStatus::InputOutdatedCanceled;
949 warn!(
950 "The task may be expired because of group split, task:\n {:?}",
951 compact_task_to_string(&compact_task)
952 );
953 }
954 !is_expired
955 }
956 }
957 } else {
958 false
959 };
960 if is_success {
961 success_count += 1;
962 version.apply_compact_task(&compact_task);
963 if purge_prost_table_stats(
964 &mut version_stats.table_stats,
965 version.latest_version(),
966 &HashSet::default(),
967 ) {
968 self.metrics.version_stats.reset();
969 versioning.local_metrics.clear();
970 }
971 add_prost_table_stats_map(&mut version_stats.table_stats, &task.table_stats_change);
972 trigger_local_table_stat(
973 &self.metrics,
974 &mut versioning.local_metrics,
975 &version_stats,
976 &task.table_stats_change,
977 );
978 }
979 task_status = compact_task.task_status;
980 report_results.push(CompactionTaskReportResult {
981 task_id,
982 task_status,
983 reported: rets[idx],
984 });
985 tasks.push(compact_task);
986 }
987 if success_count > 0 {
988 commit_multi_var!(
989 self.meta_store_ref(),
990 compact_statuses,
991 compact_task_assignment,
992 version,
993 version_stats
994 )?;
995
996 self.metrics
997 .compact_task_batch_count
998 .with_label_values(&["batch_report_task"])
999 .observe(success_count as f64);
1000 } else {
1001 commit_multi_var!(
1003 self.meta_store_ref(),
1004 compact_statuses,
1005 compact_task_assignment
1006 )?;
1007 }
1008
1009 self.notify_compaction_task_report_waiters(report_results);
1010
1011 let mut success_groups = vec![];
1012 for compact_task in &tasks {
1013 self.compactor_manager
1014 .remove_task_heartbeat(compact_task.task_id);
1015 tracing::trace!(
1016 "Reported compaction task. {}. cost time: {:?}",
1017 compact_task_to_string(compact_task),
1018 start_time.elapsed(),
1019 );
1020
1021 if !deterministic_mode
1022 && (matches!(compact_task.task_type, compact_task::TaskType::Dynamic)
1023 || matches!(compact_task.task_type, compact_task::TaskType::Emergency))
1024 {
1025 self.try_send_compaction_request(
1027 compact_task.compaction_group_id,
1028 compact_task::TaskType::Dynamic,
1029 );
1030 }
1031
1032 if compact_task.task_status == TaskStatus::Success {
1033 success_groups.push(compact_task.compaction_group_id);
1034 }
1035 }
1036
1037 trigger_compact_tasks_stat(
1038 &self.metrics,
1039 &tasks,
1040 &compaction.compaction_statuses,
1041 &versioning_guard.current_version,
1042 );
1043 drop(versioning_guard);
1044 if !success_groups.is_empty() {
1045 self.try_update_write_limits(&success_groups).await;
1046 }
1047 Ok(rets)
1048 }
1049
1050 pub async fn trigger_compaction_deterministic(
1053 &self,
1054 _base_version_id: HummockVersionId,
1055 compaction_groups: Vec<CompactionGroupId>,
1056 ) -> Result<()> {
1057 self.on_current_version(|old_version| {
1058 tracing::info!(
1059 "Trigger compaction for version {}, groups {:?}",
1060 old_version.id,
1061 compaction_groups
1062 );
1063 })
1064 .await;
1065
1066 if compaction_groups.is_empty() {
1067 return Ok(());
1068 }
1069 for compaction_group in compaction_groups {
1070 self.try_send_compaction_request(compaction_group, compact_task::TaskType::Dynamic);
1071 }
1072 Ok(())
1073 }
1074
1075 pub async fn trigger_manual_compaction(
1076 &self,
1077 compaction_group: CompactionGroupId,
1078 manual_compaction_option: ManualCompactionOption,
1079 ) -> Result<ManualCompactionTriggerResult> {
1080 let start_time = Instant::now();
1081 let exclusive = manual_compaction_option.exclusive;
1082
1083 let compactor = match self.compactor_manager.next_compactor() {
1085 Some(compactor) => compactor,
1086 None => {
1087 tracing::warn!("trigger_manual_compaction No compactor is available.");
1088 return Err(anyhow::anyhow!(
1089 "trigger_manual_compaction No compactor is available. compaction_group {}",
1090 compaction_group
1091 )
1092 .into());
1093 }
1094 };
1095
1096 let compact_task = self
1098 .manual_get_compact_task_with_info(compaction_group, manual_compaction_option)
1099 .await;
1100 let (compact_task, blocked_by_pending) = match compact_task {
1101 Ok((compact_task, blocked_by_pending)) => (compact_task, blocked_by_pending),
1102 Err(err) => {
1103 tracing::warn!(error = %err.as_report(), "Failed to get compaction task");
1104 if matches!(err, Error::InvalidManualCompactionOption(_)) {
1105 return Err(err);
1106 }
1107
1108 return Err(anyhow::anyhow!(err)
1109 .context(format!(
1110 "Failed to get compaction task for compaction_group {}",
1111 compaction_group,
1112 ))
1113 .into());
1114 }
1115 };
1116 let compact_task = match compact_task {
1117 Some(compact_task) => compact_task,
1118 None => {
1119 if exclusive && blocked_by_pending {
1120 return Ok(ManualCompactionTriggerResult::Retry);
1121 }
1122 return Err(anyhow::anyhow!(
1124 "trigger_manual_compaction No compaction_task is available. compaction_group {}",
1125 compaction_group
1126 )
1127 .into());
1128 }
1129 };
1130
1131 let task_id = compact_task.task_id;
1133 let compact_task_string = compact_task_to_string(&compact_task);
1134 tracing::info!(
1135 compact_task_string,
1136 duration = ?start_time.elapsed(),
1137 "Triggered manual compaction task."
1138 );
1139
1140 let report_rx = self.register_compaction_task_report_waiter(task_id);
1141 if let Err(err) = compactor
1142 .send_event(ResponseEvent::CompactTask(compact_task.into()))
1143 .with_context(|| {
1144 format!(
1145 "Failed to trigger compaction task for compaction_group {}",
1146 compaction_group,
1147 )
1148 })
1149 {
1150 self.remove_compaction_task_report_waiter(task_id);
1151 return Err(err.into());
1152 }
1153
1154 let report_result = match report_rx.await {
1155 Ok(result) => result,
1156 Err(_) => {
1157 self.remove_compaction_task_report_waiter(task_id);
1158 return Err(anyhow::anyhow!(
1159 "trigger_manual_compaction wait report failed. compaction_group {}",
1160 compaction_group
1161 )
1162 .into());
1163 }
1164 };
1165 if !report_result.reported {
1166 return Err(anyhow::anyhow!(
1167 "trigger_manual_compaction report not accepted. task_id {}",
1168 report_result.task_id
1169 )
1170 .into());
1171 }
1172
1173 if report_result.task_status == TaskStatus::NoAvailCpuResourceCanceled
1174 || report_result.task_status == TaskStatus::NoAvailMemoryResourceCanceled
1175 {
1176 return Ok(ManualCompactionTriggerResult::Retry);
1177 }
1178
1179 tracing::info!(
1180 ?report_result,
1181 duration = ?start_time.elapsed(),
1182 "Completed manual compaction task."
1183 );
1184
1185 Ok(ManualCompactionTriggerResult::Submitted)
1186 }
1187
1188 pub fn try_send_compaction_request(
1190 &self,
1191 compaction_group: CompactionGroupId,
1192 task_type: compact_task::TaskType,
1193 ) -> bool {
1194 self.compaction_state.try_sched_compaction(
1195 compaction_group,
1196 task_type,
1197 ScheduleTrigger::NewData,
1198 )
1199 }
1200
1201 fn apply_split_weight_by_vnode_partition(
1204 &self,
1205 compact_task: &mut CompactTask,
1206 compaction_config: &CompactionConfig,
1207 compact_table_ids: &[TableId],
1208 ) {
1209 if compaction_config.split_weight_by_vnode > 0 {
1210 for table_id in compact_table_ids {
1211 compact_task
1212 .table_vnode_partition
1213 .insert(*table_id, compact_task.split_weight_by_vnode);
1214 }
1215
1216 return;
1217 }
1218
1219 let mut table_size_info: HashMap<TableId, u64> = HashMap::default();
1221 for input_ssts in &compact_task.input_ssts {
1222 for sst in &input_ssts.table_infos {
1223 for table_id in &sst.table_ids {
1224 *table_size_info.entry(*table_id).or_default() +=
1225 sst.sst_size / (sst.table_ids.len() as u64);
1226 }
1227 }
1228 }
1229
1230 let hybrid_vnode_count = self.env.opts.hybrid_partition_node_count;
1231 let default_partition_count = self.env.opts.partition_vnode_count;
1232 let compact_task_table_size_partition_threshold_low = self
1233 .env
1234 .opts
1235 .compact_task_table_size_partition_threshold_low;
1236 let compact_task_table_size_partition_threshold_high = self
1237 .env
1238 .opts
1239 .compact_task_table_size_partition_threshold_high;
1240
1241 let table_write_throughput_statistic_manager =
1243 self.table_write_throughput_statistic_manager.read();
1244 let timestamp = chrono::Utc::now().timestamp();
1245
1246 for (table_id, compact_table_size) in table_size_info {
1247 let write_throughput = table_write_throughput_statistic_manager
1248 .get_table_throughput_descending(table_id, timestamp)
1249 .peekable()
1250 .peek()
1251 .map(|item| item.throughput)
1252 .unwrap_or(0);
1253
1254 if compact_table_size > compact_task_table_size_partition_threshold_high
1255 && default_partition_count > 0
1256 {
1257 compact_task
1258 .table_vnode_partition
1259 .insert(table_id, default_partition_count);
1260 } else if (compact_table_size > compact_task_table_size_partition_threshold_low
1261 || (write_throughput > self.env.opts.table_high_write_throughput_threshold
1262 && compact_table_size > compaction_config.target_file_size_base))
1263 && hybrid_vnode_count > 0
1264 {
1265 compact_task
1266 .table_vnode_partition
1267 .insert(table_id, hybrid_vnode_count);
1268 } else if compact_table_size > compaction_config.target_file_size_base {
1269 compact_task.table_vnode_partition.insert(table_id, 1);
1270 }
1271 }
1272
1273 compact_task
1274 .table_vnode_partition
1275 .retain(|table_id, _| compact_table_ids.contains(table_id));
1276 }
1277
1278 pub(crate) fn calculate_vnode_partition(
1279 &self,
1280 compact_task: &mut CompactTask,
1281 compaction_config: &CompactionConfig,
1282 compact_table_ids: &[TableId],
1283 ) {
1284 if compact_task.target_level > compact_task.base_level {
1289 return;
1290 }
1291
1292 self.apply_split_weight_by_vnode_partition(
1294 compact_task,
1295 compaction_config,
1296 compact_table_ids,
1297 );
1298 }
1299
1300 fn build_ready_compact_task(
1301 &self,
1302 picked_task: PickedCompactionTask,
1303 context: CompactTaskBuildContext,
1304 table_watermarks: &HashMap<TableId, Arc<TableWatermarks>>,
1305 all_versioned_table_schemas: &HashMap<TableId, Vec<i32>>,
1306 ) -> BuiltCompactTask {
1307 let compaction_config = context.compaction_config.clone();
1308 let (mut compact_task, compact_table_ids) = build_base_compact_task(picked_task, context);
1309
1310 if compact_task.is_trivial_reclaim() {
1311 compact_task.task_status = TaskStatus::Success;
1312 compact_task.sorted_output_ssts.clear();
1313 return BuiltCompactTask::MetaFinished(compact_task);
1314 }
1315
1316 if compact_task.is_trivial_move_task() {
1317 compact_task.task_status = TaskStatus::Success;
1318 compact_task.sorted_output_ssts = compact_task.input_ssts[0]
1319 .read_sstable_infos()
1320 .cloned()
1321 .collect();
1322 return BuiltCompactTask::MetaFinished(compact_task);
1323 }
1324
1325 self.prepare_compact_task_for_assignment(
1326 &mut compact_task,
1327 compaction_config.as_ref(),
1328 &compact_table_ids,
1329 safe_epoch_table_watermarks_impl(table_watermarks, &compact_table_ids),
1330 all_versioned_table_schemas,
1331 );
1332
1333 BuiltCompactTask::PendingAssignment(compact_task)
1334 }
1335
1336 fn prepare_compact_task_for_assignment(
1337 &self,
1338 compact_task: &mut CompactTask,
1339 compaction_config: &CompactionConfig,
1340 compact_table_ids: &[TableId],
1341 table_watermarks: BTreeMap<TableId, TableWatermarks>,
1342 all_versioned_table_schemas: &HashMap<TableId, Vec<i32>>,
1343 ) {
1344 self.calculate_vnode_partition(compact_task, compaction_config, compact_table_ids);
1345 attach_compact_task_table_metadata(
1346 compact_task,
1347 compact_table_ids,
1348 table_watermarks,
1349 all_versioned_table_schemas,
1350 );
1351 }
1352
1353 pub fn compactor_manager_ref(&self) -> crate::hummock::CompactorManagerRef {
1354 self.compactor_manager.clone()
1355 }
1356
1357 fn register_compaction_task_report_waiter(
1358 &self,
1359 task_id: HummockCompactionTaskId,
1360 ) -> Receiver<CompactionTaskReportResult> {
1361 let (tx, rx) = tokio::sync::oneshot::channel();
1362 self.compaction_task_report_notifiers
1363 .lock()
1364 .register(task_id, tx);
1365 rx
1366 }
1367
1368 fn remove_compaction_task_report_waiter(&self, task_id: HummockCompactionTaskId) {
1369 self.compaction_task_report_notifiers.lock().remove(task_id);
1370 }
1371
1372 fn notify_compaction_task_report_waiters(&self, results: Vec<CompactionTaskReportResult>) {
1373 let mut guard = self.compaction_task_report_notifiers.lock();
1374 for result in results {
1375 guard.notify(result);
1376 }
1377 }
1378}
1379
1380#[cfg(any(test, feature = "test"))]
1381impl HummockManager {
1382 pub async fn compaction_task_from_assignment_for_test(
1383 &self,
1384 task_id: u64,
1385 ) -> Option<CompactTaskAssignment> {
1386 let compaction_guard = self
1387 .compaction
1388 .read_with_process_name("compaction_task_from_assignment_for_test")
1389 .await;
1390 let assignment_ref = &compaction_guard.compact_task_assignment;
1391 assignment_ref.get(&task_id).cloned()
1392 }
1393
1394 pub async fn report_compact_task_for_test(
1395 &self,
1396 task_id: u64,
1397 compact_task: Option<CompactTask>,
1398 task_status: TaskStatus,
1399 sorted_output_ssts: Vec<SstableInfo>,
1400 table_stats_change: Option<PbTableStatsMap>,
1401 ) -> Result<()> {
1402 if let Some(task) = compact_task {
1403 let mut guard = self
1404 .compaction
1405 .write_with_process_name("report_compact_task_for_test")
1406 .await;
1407 guard.compact_task_assignment.insert(
1408 task_id,
1409 CompactTaskAssignment {
1410 compact_task: task,
1411 context_id: 0.into(),
1412 },
1413 );
1414 }
1415
1416 self.report_compact_tasks(vec![ReportTask {
1419 task_id,
1420 task_status,
1421 sorted_output_ssts,
1422 table_stats_change: table_stats_change.unwrap_or_default(),
1423 object_timestamps: HashMap::default(),
1424 }])
1425 .await?;
1426 Ok(())
1427 }
1428}
1429
1430#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1432pub enum ScheduleTrigger {
1433 NewData,
1435 Periodic,
1437}
1438
1439pub struct CompactionScheduleSnapshot {
1444 scheduled: HashSet<(CompactionGroupId, compact_task::TaskType)>,
1445 snapshot_time: Instant,
1446}
1447
1448impl CompactionScheduleSnapshot {
1449 const TASK_TYPE_PRIORITY: &[TaskType] = &[
1451 TaskType::Dynamic,
1452 TaskType::SpaceReclaim,
1453 TaskType::Ttl,
1454 TaskType::Tombstone,
1455 TaskType::VnodeWatermark,
1456 ];
1457
1458 pub fn snapshot_time(&self) -> Instant {
1459 self.snapshot_time
1460 }
1461
1462 pub fn pick_compaction_groups_and_type(&self) -> Option<(Vec<CompactionGroupId>, TaskType)> {
1467 let group_ids = self.group_ids_shuffled();
1468 let mut normal_groups = vec![];
1469 for cg_id in group_ids {
1470 if let Some(pick_type) = self.pick_type(cg_id) {
1471 if pick_type == TaskType::Dynamic {
1472 normal_groups.push(cg_id);
1473 } else if normal_groups.is_empty() {
1474 return Some((vec![cg_id], pick_type));
1475 }
1476 }
1477 }
1478 if normal_groups.is_empty() {
1479 None
1480 } else {
1481 Some((normal_groups, TaskType::Dynamic))
1482 }
1483 }
1484
1485 fn group_ids_shuffled(&self) -> Vec<CompactionGroupId> {
1486 let mut group_ids: Vec<_> = self.scheduled.iter().map(|(g, _)| *g).unique().collect();
1487 group_ids.shuffle(&mut thread_rng());
1488 group_ids
1489 }
1490
1491 fn pick_type(&self, group: CompactionGroupId) -> Option<TaskType> {
1492 Self::TASK_TYPE_PRIORITY
1493 .iter()
1494 .find(|t| self.scheduled.contains(&(group, **t)))
1495 .copied()
1496 }
1497}
1498
1499#[derive(Debug, Default)]
1504pub struct CompactionState {
1505 inner: Mutex<CompactionStateInner>,
1506}
1507
1508#[derive(Debug, Default)]
1509struct CompactionStateInner {
1510 scheduled: HashSet<(CompactionGroupId, compact_task::TaskType)>,
1511 dynamic_cooldown: HashSet<CompactionGroupId>,
1513 last_new_data_time: HashMap<CompactionGroupId, Instant>,
1515}
1516
1517impl CompactionState {
1518 pub fn new() -> Self {
1519 Self {
1520 inner: Default::default(),
1521 }
1522 }
1523
1524 pub fn try_sched_compaction(
1528 &self,
1529 compaction_group: CompactionGroupId,
1530 task_type: TaskType,
1531 trigger: ScheduleTrigger,
1532 ) -> bool {
1533 let mut guard = self.inner.lock();
1534 if task_type == TaskType::Dynamic {
1535 match trigger {
1536 ScheduleTrigger::NewData => {
1537 guard.dynamic_cooldown.remove(&compaction_group);
1538 guard
1539 .last_new_data_time
1540 .insert(compaction_group, Instant::now());
1541 }
1542 ScheduleTrigger::Periodic => {
1543 if guard.dynamic_cooldown.contains(&compaction_group) {
1544 return false;
1545 }
1546 }
1547 }
1548 }
1549 guard.scheduled.insert((compaction_group, task_type))
1550 }
1551
1552 pub fn unschedule(
1555 &self,
1556 compaction_group: CompactionGroupId,
1557 task_type: compact_task::TaskType,
1558 snapshot_time: Instant,
1559 ) {
1560 let mut guard = self.inner.lock();
1561 guard.scheduled.remove(&(compaction_group, task_type));
1562 if task_type == TaskType::Dynamic {
1563 let has_new_data = guard
1564 .last_new_data_time
1565 .get(&compaction_group)
1566 .is_some_and(|t| *t > snapshot_time);
1567 if !has_new_data {
1568 guard.dynamic_cooldown.insert(compaction_group);
1569 }
1570 }
1571 }
1572
1573 pub fn snapshot(&self) -> CompactionScheduleSnapshot {
1575 let guard = self.inner.lock();
1576 let snapshot_time = Instant::now();
1578 CompactionScheduleSnapshot {
1579 scheduled: guard.scheduled.clone(),
1580 snapshot_time,
1581 }
1582 }
1583
1584 pub fn remove_compaction_group(&self, compaction_group: CompactionGroupId) {
1586 let mut guard = self.inner.lock();
1587 guard
1588 .scheduled
1589 .retain(|(group, _)| *group != compaction_group);
1590 guard.dynamic_cooldown.remove(&compaction_group);
1591 guard.last_new_data_time.remove(&compaction_group);
1592 }
1593}
1594
1595impl Compaction {
1596 pub fn get_compact_task_assignments_by_group_id(
1597 &self,
1598 compaction_group_id: CompactionGroupId,
1599 ) -> Vec<CompactTaskAssignment> {
1600 self.compact_task_assignment
1601 .values()
1602 .filter_map(|assignment| {
1603 if assignment.compact_task.compaction_group_id == compaction_group_id {
1604 Some(assignment.clone())
1605 } else {
1606 None
1607 }
1608 })
1609 .collect()
1610 }
1611}
1612
1613#[derive(Clone, Default)]
1614pub struct CompactionGroupStatistic {
1615 pub group_id: CompactionGroupId,
1616 pub group_size: u64,
1617 pub table_statistic: BTreeMap<StateTableId, u64>,
1618 pub compaction_group_config: CompactionGroup,
1619}
1620
1621fn update_table_stats_for_vnode_watermark_trivial_reclaim(
1623 table_stats: &mut PbTableStatsMap,
1624 task: &CompactTask,
1625) {
1626 if task.task_type != TaskType::VnodeWatermark {
1627 return;
1628 }
1629 let mut deleted_table_keys: HashMap<TableId, u64> = HashMap::default();
1630 for s in task.input_ssts.iter().flat_map(|l| l.table_infos.iter()) {
1631 assert_eq!(s.table_ids.len(), 1);
1632 let e = deleted_table_keys.entry(s.table_ids[0]).or_insert(0);
1633 *e += s.total_key_count;
1634 }
1635 for (table_id, delete_count) in deleted_table_keys {
1636 let Some(stats) = table_stats.get_mut(&table_id) else {
1637 continue;
1638 };
1639 if stats.total_key_count == 0 {
1640 continue;
1641 }
1642 let new_total_key_count = stats.total_key_count.saturating_sub(delete_count as i64);
1643 let ratio = new_total_key_count as f64 / stats.total_key_count as f64;
1644 stats.total_key_count = new_total_key_count;
1646 stats.total_key_size = (stats.total_key_size as f64 * ratio).ceil() as i64;
1648 stats.total_value_size = (stats.total_value_size as f64 * ratio).ceil() as i64;
1649 }
1650}
1651
1652#[derive(Debug, Clone)]
1653pub enum GroupState {
1654 Normal,
1656
1657 Emergency(String), WriteStop(String), }
1663
1664impl GroupState {
1665 pub fn is_write_stop(&self) -> bool {
1666 matches!(self, Self::WriteStop(_))
1667 }
1668
1669 pub fn is_emergency(&self) -> bool {
1670 matches!(self, Self::Emergency(_))
1671 }
1672
1673 pub fn reason(&self) -> Option<&str> {
1674 match self {
1675 Self::Emergency(reason) | Self::WriteStop(reason) => Some(reason),
1676 _ => None,
1677 }
1678 }
1679}
1680
1681#[derive(Clone, Default)]
1682pub struct GroupStateValidator;
1683
1684impl GroupStateValidator {
1685 pub fn write_stop_sub_level_count(
1686 level_count: usize,
1687 compaction_config: &CompactionConfig,
1688 ) -> bool {
1689 let threshold = compaction_config.level0_stop_write_threshold_sub_level_number as usize;
1690 level_count > threshold
1691 }
1692
1693 pub fn write_stop_l0_size(l0_size: u64, compaction_config: &CompactionConfig) -> bool {
1694 l0_size
1695 > compaction_config
1696 .level0_stop_write_threshold_max_size
1697 .unwrap_or(compaction_config::level0_stop_write_threshold_max_size())
1698 }
1699
1700 pub fn write_stop_l0_file_count(
1701 l0_file_count: usize,
1702 compaction_config: &CompactionConfig,
1703 ) -> bool {
1704 l0_file_count
1705 > compaction_config
1706 .level0_stop_write_threshold_max_sst_count
1707 .unwrap_or(compaction_config::level0_stop_write_threshold_max_sst_count())
1708 as usize
1709 }
1710
1711 pub fn emergency_l0_file_count(
1712 l0_file_count: usize,
1713 compaction_config: &CompactionConfig,
1714 ) -> bool {
1715 l0_file_count
1716 > compaction_config
1717 .emergency_level0_sst_file_count
1718 .unwrap_or(compaction_config::emergency_level0_sst_file_count())
1719 as usize
1720 }
1721
1722 pub fn emergency_l0_partition_count(
1723 last_l0_sub_level_partition_count: usize,
1724 compaction_config: &CompactionConfig,
1725 ) -> bool {
1726 last_l0_sub_level_partition_count
1727 > compaction_config
1728 .emergency_level0_sub_level_partition
1729 .unwrap_or(compaction_config::emergency_level0_sub_level_partition())
1730 as usize
1731 }
1732
1733 pub fn check_single_group_write_stop(
1734 levels: &Levels,
1735 compaction_config: &CompactionConfig,
1736 ) -> GroupState {
1737 if Self::write_stop_sub_level_count(levels.l0.sub_levels.len(), compaction_config) {
1738 return GroupState::WriteStop(format!(
1739 "WriteStop(l0_level_count: {}, threshold: {}) too many L0 sub levels",
1740 levels.l0.sub_levels.len(),
1741 compaction_config.level0_stop_write_threshold_sub_level_number
1742 ));
1743 }
1744
1745 if Self::write_stop_l0_file_count(
1746 levels
1747 .l0
1748 .sub_levels
1749 .iter()
1750 .map(|l| l.table_infos.len())
1751 .sum(),
1752 compaction_config,
1753 ) {
1754 return GroupState::WriteStop(format!(
1755 "WriteStop(l0_sst_count: {}, threshold: {}) too many L0 sst files",
1756 levels
1757 .l0
1758 .sub_levels
1759 .iter()
1760 .map(|l| l.table_infos.len())
1761 .sum::<usize>(),
1762 compaction_config
1763 .level0_stop_write_threshold_max_sst_count
1764 .unwrap_or(compaction_config::level0_stop_write_threshold_max_sst_count())
1765 ));
1766 }
1767
1768 if Self::write_stop_l0_size(levels.l0.total_file_size, compaction_config) {
1769 return GroupState::WriteStop(format!(
1770 "WriteStop(l0_size: {}, threshold: {}) too large L0 size",
1771 levels.l0.total_file_size,
1772 compaction_config
1773 .level0_stop_write_threshold_max_size
1774 .unwrap_or(compaction_config::level0_stop_write_threshold_max_size())
1775 ));
1776 }
1777
1778 GroupState::Normal
1779 }
1780
1781 pub fn check_single_group_emergency(
1782 levels: &Levels,
1783 compaction_config: &CompactionConfig,
1784 ) -> GroupState {
1785 if Self::emergency_l0_file_count(
1786 levels
1787 .l0
1788 .sub_levels
1789 .iter()
1790 .map(|l| l.table_infos.len())
1791 .sum(),
1792 compaction_config,
1793 ) {
1794 return GroupState::Emergency(format!(
1795 "Emergency(l0_sst_count: {}, threshold: {}) too many L0 sst files",
1796 levels
1797 .l0
1798 .sub_levels
1799 .iter()
1800 .map(|l| l.table_infos.len())
1801 .sum::<usize>(),
1802 compaction_config
1803 .emergency_level0_sst_file_count
1804 .unwrap_or(compaction_config::emergency_level0_sst_file_count())
1805 ));
1806 }
1807
1808 if Self::emergency_l0_partition_count(
1809 levels
1810 .l0
1811 .sub_levels
1812 .first()
1813 .map(|l| l.table_infos.len())
1814 .unwrap_or(0),
1815 compaction_config,
1816 ) {
1817 return GroupState::Emergency(format!(
1818 "Emergency(l0_partition_count: {}, threshold: {}) too many L0 partitions",
1819 levels
1820 .l0
1821 .sub_levels
1822 .first()
1823 .map(|l| l.table_infos.len())
1824 .unwrap_or(0),
1825 compaction_config
1826 .emergency_level0_sub_level_partition
1827 .unwrap_or(compaction_config::emergency_level0_sub_level_partition())
1828 ));
1829 }
1830
1831 GroupState::Normal
1832 }
1833
1834 pub fn group_state(levels: &Levels, compaction_config: &CompactionConfig) -> GroupState {
1835 let state = Self::check_single_group_write_stop(levels, compaction_config);
1836 if state.is_write_stop() {
1837 return state;
1838 }
1839
1840 Self::check_single_group_emergency(levels, compaction_config)
1841 }
1842}
1843
1844#[cfg(test)]
1845mod compaction_state_tests {
1846 use risingwave_pb::hummock::compact_task::TaskType;
1847
1848 use super::*;
1849
1850 #[test]
1851 fn test_basic_schedule_and_unschedule() {
1852 let state = CompactionState::new();
1853 let group_id: CompactionGroupId = 1.into();
1854
1855 assert!(state.try_sched_compaction(group_id, TaskType::Dynamic, ScheduleTrigger::NewData));
1857 assert!(!state.try_sched_compaction(group_id, TaskType::Dynamic, ScheduleTrigger::NewData));
1859 assert!(state.try_sched_compaction(group_id, TaskType::Ttl, ScheduleTrigger::Periodic));
1861
1862 let snapshot = state.snapshot();
1864 assert!(snapshot.scheduled.contains(&(group_id, TaskType::Dynamic)));
1865 assert!(snapshot.scheduled.contains(&(group_id, TaskType::Ttl)));
1866
1867 state.unschedule(group_id, TaskType::Dynamic, snapshot.snapshot_time());
1869 let snapshot2 = state.snapshot();
1870 assert!(!snapshot2.scheduled.contains(&(group_id, TaskType::Dynamic)));
1871 assert!(snapshot2.scheduled.contains(&(group_id, TaskType::Ttl)));
1872 }
1873
1874 #[test]
1875 fn test_cooldown_blocks_periodic_trigger() {
1876 let state = CompactionState::new();
1877 let group_id: CompactionGroupId = 1.into();
1878
1879 assert!(state.try_sched_compaction(group_id, TaskType::Dynamic, ScheduleTrigger::NewData));
1881 let snapshot = state.snapshot();
1882 state.unschedule(group_id, TaskType::Dynamic, snapshot.snapshot_time());
1883
1884 assert!(state.inner.lock().dynamic_cooldown.contains(&group_id));
1886
1887 assert!(!state.try_sched_compaction(
1889 group_id,
1890 TaskType::Dynamic,
1891 ScheduleTrigger::Periodic
1892 ));
1893 }
1894
1895 #[test]
1896 fn test_new_data_clears_cooldown() {
1897 let state = CompactionState::new();
1898 let group_id: CompactionGroupId = 1.into();
1899
1900 assert!(state.try_sched_compaction(group_id, TaskType::Dynamic, ScheduleTrigger::NewData));
1902 let snapshot = state.snapshot();
1903 state.unschedule(group_id, TaskType::Dynamic, snapshot.snapshot_time());
1904 assert!(state.inner.lock().dynamic_cooldown.contains(&group_id));
1905
1906 assert!(state.try_sched_compaction(group_id, TaskType::Dynamic, ScheduleTrigger::NewData));
1908 assert!(!state.inner.lock().dynamic_cooldown.contains(&group_id));
1909 }
1910
1911 #[test]
1912 fn test_cooldown_only_affects_dynamic_type() {
1913 let state = CompactionState::new();
1914 let group_id: CompactionGroupId = 1.into();
1915
1916 assert!(state.try_sched_compaction(group_id, TaskType::Dynamic, ScheduleTrigger::NewData));
1918 let snapshot = state.snapshot();
1919 state.unschedule(group_id, TaskType::Dynamic, snapshot.snapshot_time());
1920
1921 let group_id_2: CompactionGroupId = 2.into();
1923 assert!(state.try_sched_compaction(group_id_2, TaskType::Ttl, ScheduleTrigger::Periodic));
1924 let snapshot2 = state.snapshot();
1925 state.unschedule(group_id_2, TaskType::Ttl, snapshot2.snapshot_time());
1926 assert!(!state.inner.lock().dynamic_cooldown.contains(&group_id_2));
1927
1928 assert!(state.try_sched_compaction(group_id, TaskType::Ttl, ScheduleTrigger::Periodic));
1930 assert!(state.try_sched_compaction(
1931 group_id,
1932 TaskType::SpaceReclaim,
1933 ScheduleTrigger::Periodic
1934 ));
1935 }
1936
1937 #[test]
1938 fn test_race_condition_new_data_after_snapshot() {
1939 let state = CompactionState::new();
1940 let group_id: CompactionGroupId = 1.into();
1941
1942 assert!(state.try_sched_compaction(group_id, TaskType::Dynamic, ScheduleTrigger::NewData));
1943 let snapshot = state.snapshot();
1944
1945 {
1947 let mut guard = state.inner.lock();
1948 guard.last_new_data_time.insert(group_id, Instant::now());
1949 }
1950
1951 state.unschedule(group_id, TaskType::Dynamic, snapshot.snapshot_time());
1953 assert!(
1954 !state.inner.lock().dynamic_cooldown.contains(&group_id),
1955 "Should skip cooldown when new data arrived after snapshot"
1956 );
1957 }
1958
1959 #[test]
1960 fn test_remove_compaction_group_cleans_all_state() {
1961 let state = CompactionState::new();
1962 let group_id: CompactionGroupId = 1.into();
1963
1964 assert!(state.try_sched_compaction(group_id, TaskType::Dynamic, ScheduleTrigger::NewData));
1966 assert!(state.try_sched_compaction(group_id, TaskType::Ttl, ScheduleTrigger::Periodic));
1967 state.inner.lock().dynamic_cooldown.insert(group_id);
1968
1969 state.remove_compaction_group(group_id);
1971
1972 let guard = state.inner.lock();
1974 assert!(!guard.scheduled.contains(&(group_id, TaskType::Dynamic)));
1975 assert!(!guard.scheduled.contains(&(group_id, TaskType::Ttl)));
1976 assert!(!guard.dynamic_cooldown.contains(&group_id));
1977 assert!(!guard.last_new_data_time.contains_key(&group_id));
1978 }
1979
1980 #[test]
1981 fn test_snapshot_pick_type_priority() {
1982 let state = CompactionState::new();
1983 let group_id: CompactionGroupId = 1.into();
1984
1985 assert_eq!(state.snapshot().pick_type(group_id), None);
1987
1988 state.try_sched_compaction(
1990 group_id,
1991 TaskType::VnodeWatermark,
1992 ScheduleTrigger::Periodic,
1993 );
1994 assert_eq!(
1995 state.snapshot().pick_type(group_id),
1996 Some(TaskType::VnodeWatermark)
1997 );
1998
1999 state.try_sched_compaction(group_id, TaskType::Tombstone, ScheduleTrigger::Periodic);
2000 assert_eq!(
2001 state.snapshot().pick_type(group_id),
2002 Some(TaskType::Tombstone)
2003 );
2004
2005 state.try_sched_compaction(group_id, TaskType::Ttl, ScheduleTrigger::Periodic);
2006 assert_eq!(state.snapshot().pick_type(group_id), Some(TaskType::Ttl));
2007
2008 state.try_sched_compaction(group_id, TaskType::SpaceReclaim, ScheduleTrigger::Periodic);
2009 assert_eq!(
2010 state.snapshot().pick_type(group_id),
2011 Some(TaskType::SpaceReclaim)
2012 );
2013
2014 state.try_sched_compaction(group_id, TaskType::Dynamic, ScheduleTrigger::NewData);
2015 assert_eq!(
2016 state.snapshot().pick_type(group_id),
2017 Some(TaskType::Dynamic)
2018 );
2019 }
2020
2021 #[test]
2022 fn test_multiple_groups_independent_cooldown() {
2023 let state = CompactionState::new();
2024 let g1: CompactionGroupId = 1.into();
2025 let g2: CompactionGroupId = 2.into();
2026
2027 state.try_sched_compaction(g1, TaskType::Dynamic, ScheduleTrigger::NewData);
2028 state.try_sched_compaction(g2, TaskType::Dynamic, ScheduleTrigger::NewData);
2029 let snapshot = state.snapshot();
2030
2031 state.unschedule(g1, TaskType::Dynamic, snapshot.snapshot_time());
2033
2034 let guard = state.inner.lock();
2035 assert!(guard.dynamic_cooldown.contains(&g1));
2036 assert!(!guard.dynamic_cooldown.contains(&g2));
2037 }
2038
2039 #[test]
2040 fn test_pick_compaction_groups_empty() {
2041 let state = CompactionState::new();
2042 let snapshot = state.snapshot();
2043 assert!(snapshot.pick_compaction_groups_and_type().is_none());
2045 }
2046
2047 #[test]
2048 fn test_pick_compaction_groups_mixed_types() {
2049 let state = CompactionState::new();
2050 let g1: CompactionGroupId = 1.into();
2051 let g2: CompactionGroupId = 2.into();
2052 let g3: CompactionGroupId = 3.into();
2053
2054 state.try_sched_compaction(g1, TaskType::Dynamic, ScheduleTrigger::NewData);
2056 state.try_sched_compaction(g2, TaskType::Ttl, ScheduleTrigger::Periodic);
2057 state.try_sched_compaction(g3, TaskType::Dynamic, ScheduleTrigger::NewData);
2058
2059 let snapshot = state.snapshot();
2060 let (groups, task_type) = snapshot.pick_compaction_groups_and_type().unwrap();
2061
2062 if task_type == TaskType::Dynamic {
2067 assert!(groups.contains(&g1));
2068 assert!(groups.contains(&g3));
2069 assert!(!groups.contains(&g2)); } else {
2071 assert_eq!(task_type, TaskType::Ttl);
2072 assert_eq!(groups, vec![g2]);
2073 }
2074 }
2075
2076 #[test]
2077 fn test_pick_compaction_groups_all_dynamic() {
2078 let state = CompactionState::new();
2079 let g1: CompactionGroupId = 1.into();
2080 let g2: CompactionGroupId = 2.into();
2081
2082 state.try_sched_compaction(g1, TaskType::Dynamic, ScheduleTrigger::NewData);
2083 state.try_sched_compaction(g2, TaskType::Dynamic, ScheduleTrigger::NewData);
2084
2085 let snapshot = state.snapshot();
2086 let (groups, task_type) = snapshot.pick_compaction_groups_and_type().unwrap();
2087 assert_eq!(task_type, TaskType::Dynamic);
2088 assert!(groups.contains(&g1));
2089 assert!(groups.contains(&g2));
2090 }
2091
2092 #[test]
2093 fn test_pick_compaction_groups_single_non_dynamic() {
2094 let state = CompactionState::new();
2095 let g1: CompactionGroupId = 1.into();
2096
2097 state.try_sched_compaction(g1, TaskType::SpaceReclaim, ScheduleTrigger::Periodic);
2098
2099 let snapshot = state.snapshot();
2100 let (groups, task_type) = snapshot.pick_compaction_groups_and_type().unwrap();
2101 assert_eq!(task_type, TaskType::SpaceReclaim);
2102 assert_eq!(groups, vec![g1]);
2103 }
2104}