1use std::collections::{BTreeMap, BTreeSet, HashMap, HashSet};
16use std::ops::DerefMut;
17use std::sync::Arc;
18
19use itertools::Itertools;
20use risingwave_common::catalog::TableId;
21use risingwave_common::config::meta::default::compaction_config as default_compaction_config;
22use risingwave_common::util::epoch::INVALID_EPOCH;
23use risingwave_hummock_sdk::CompactionGroupId;
24use risingwave_hummock_sdk::compaction_group::StaticCompactionGroupId;
25use risingwave_hummock_sdk::compaction_group::hummock_version_ext::get_compaction_group_ids;
26use risingwave_hummock_sdk::filter_utils::{
27 parse_sstable_filter_layout, parse_sstable_filter_type,
28};
29use risingwave_hummock_sdk::version::{GroupDelta, HummockVersion};
30use risingwave_meta_model::compaction_config;
31use risingwave_pb::hummock::rise_ctl_update_compaction_config_request::mutable_config::MutableConfig;
32use risingwave_pb::hummock::write_limits::WriteLimit;
33use risingwave_pb::hummock::{
34 CompactionConfig, CompactionGroupInfo, HummockVersionStats, PbGroupConstruct, PbGroupDestroy,
35 PbStateTableInfoDelta,
36};
37use sea_orm::EntityTrait;
38use tokio::sync::OnceCell;
39
40use super::CompactionGroupStatistic;
41use crate::hummock::compaction::compaction_config::{
42 CompactionConfigBuilder, validate_compaction_config,
43};
44use crate::hummock::error::{Error, Result};
45use crate::hummock::manager::transaction::HummockVersionTransaction;
46use crate::hummock::manager::versioning::Versioning;
47use crate::hummock::manager::{HummockManager, commit_multi_var};
48use crate::hummock::metrics_utils::remove_compaction_group_metrics;
49use crate::hummock::model::CompactionGroup;
50use crate::hummock::sequence::next_compaction_group_id;
51use crate::manager::MetaSrvEnv;
52use crate::model::{
53 BTreeMapTransaction, BTreeMapTransactionInner, DerefMutForward, MetadataModelError,
54};
55
56type CompactionGroupTransaction<'a> = BTreeMapTransaction<'a, CompactionGroupId, CompactionGroup>;
57
58impl CompactionGroupManager {
59 pub(crate) async fn new(env: &MetaSrvEnv) -> Result<CompactionGroupManager> {
60 let default_config = match env.opts.compaction_config.as_ref() {
61 None => CompactionConfigBuilder::new().build(),
62 Some(opt) => CompactionConfigBuilder::with_opt(opt).build(),
63 };
64 Self::new_with_config(env, default_config).await
65 }
66
67 pub(crate) async fn new_with_config(
68 env: &MetaSrvEnv,
69 default_config: CompactionConfig,
70 ) -> Result<CompactionGroupManager> {
71 let mut compaction_group_manager = CompactionGroupManager {
72 compaction_groups: BTreeMap::new(),
73 default_config: Arc::new(default_config),
74 write_limit: Default::default(),
75 };
76
77 let loaded_compaction_groups: BTreeMap<CompactionGroupId, CompactionGroup> =
78 compaction_config::Entity::find()
79 .all(&env.meta_store_ref().conn)
80 .await
81 .map_err(MetadataModelError::from)?
82 .into_iter()
83 .map(|m| (m.compaction_group_id as CompactionGroupId, m.into()))
84 .collect();
85
86 compaction_group_manager.init(loaded_compaction_groups);
87 Ok(compaction_group_manager)
88 }
89
90 fn init(&mut self, loaded_compaction_groups: BTreeMap<CompactionGroupId, CompactionGroup>) {
91 if !loaded_compaction_groups.is_empty() {
92 self.compaction_groups = loaded_compaction_groups;
93 }
94 }
95}
96
97impl HummockManager {
98 pub async fn compaction_group_ids(&self) -> Vec<CompactionGroupId> {
101 get_compaction_group_ids(
102 &self
103 .versioning
104 .read_with_process_name("compaction_group_ids")
105 .await
106 .current_version,
107 )
108 .collect_vec()
109 }
110
111 pub async fn get_compaction_group_map(&self) -> BTreeMap<CompactionGroupId, CompactionGroup> {
113 self.compaction_group_manager
114 .read_with_process_name("get_compaction_group_map")
115 .await
116 .compaction_groups
117 .clone()
118 }
119
120 #[cfg(test)]
121 pub async fn register_table_fragments(
123 &self,
124 mv_table: Option<TableId>,
125 mut internal_tables: Vec<TableId>,
126 ) -> Result<()> {
127 let mut pairs = vec![];
128 if let Some(mv_table) = mv_table {
129 if internal_tables.extract_if(.., |t| *t == mv_table).count() > 0 {
130 tracing::warn!("`mv_table` {} found in `internal_tables`", mv_table);
131 }
132 pairs.push((mv_table, StaticCompactionGroupId::MaterializedView));
134 }
135 for table_id in internal_tables {
137 pairs.push((table_id, StaticCompactionGroupId::StateDefault));
138 }
139 self.register_table_ids_for_test(&pairs).await?;
140 Ok(())
141 }
142
143 #[cfg(test)]
144 pub async fn unregister_table_fragments_vec(
146 &self,
147 table_fragments: &[crate::model::StreamJobFragments],
148 ) {
149 self.unregister_table_ids(table_fragments.iter().flat_map(|t| t.all_table_ids()))
150 .await
151 .unwrap();
152 }
153
154 pub async fn purge(&self, valid_ids: &HashSet<TableId>) -> Result<()> {
158 let to_unregister = self
159 .versioning
160 .read_with_process_name("purge")
161 .await
162 .current_version
163 .state_table_info
164 .info()
165 .keys()
166 .cloned()
167 .filter(|table_id| !valid_ids.contains(table_id))
168 .collect_vec();
169
170 self.unregister_table_ids(to_unregister).await
173 }
174
175 pub async fn register_table_ids_for_test(
180 &self,
181 pairs: &[(impl Into<TableId> + Copy, CompactionGroupId)],
182 ) -> Result<()> {
183 if pairs.is_empty() {
184 return Ok(());
185 }
186 let mut versioning_guard = self
187 .versioning
188 .write_with_process_name("register_table_ids_for_test")
189 .await;
190 let versioning = versioning_guard.deref_mut();
191 let mut compaction_group_manager = self
192 .compaction_group_manager
193 .write_with_process_name("register_table_ids_for_test")
194 .await;
195 let current_version = &versioning.current_version;
196 let default_config = compaction_group_manager.default_compaction_config();
197 let mut compaction_groups_txn = compaction_group_manager.start_compaction_groups_txn();
198
199 for (table_id, _) in pairs {
200 let table_id = (*table_id).into();
201 if let Some(info) = current_version.state_table_info.info().get(&table_id) {
202 return Err(Error::CompactionGroup(format!(
203 "table {} already {:?}",
204 table_id, info
205 )));
206 }
207 }
208 let new_compaction_group_id: OnceCell<CompactionGroupId> = OnceCell::new();
210 let mut version = HummockVersionTransaction::new(
211 &mut versioning.current_version,
212 &mut versioning.hummock_version_deltas,
213 &mut versioning.table_change_log,
214 self.env.notification_manager(),
215 None,
216 &self.metrics,
217 &self.env.opts,
218 &self.version_stat_tx,
219 );
220 let mut new_version_delta = version.new_delta();
221
222 let committed_epoch = new_version_delta
223 .latest_version()
224 .state_table_info
225 .info()
226 .values()
227 .map(|info| info.committed_epoch)
228 .max()
229 .unwrap_or(INVALID_EPOCH);
230
231 for (table_id, raw_group_id) in pairs {
232 let table_id = (*table_id).into();
233 let mut group_id = *raw_group_id;
234 if group_id == StaticCompactionGroupId::NewCompactionGroup {
235 let mut is_group_init = false;
236 group_id = *new_compaction_group_id
237 .get_or_try_init(|| async {
238 next_compaction_group_id(&self.env).await.inspect(|_| {
239 is_group_init = true;
240 })
241 })
242 .await?;
243 if is_group_init {
244 let group_deltas = &mut new_version_delta
245 .group_deltas
246 .entry(group_id)
247 .or_default()
248 .group_deltas;
249
250 let config =
251 match compaction_groups_txn.try_get_compaction_group_config(group_id) {
252 Some(config) => config.compaction_config.as_ref().clone(),
253 None => {
254 compaction_groups_txn
255 .create_compaction_groups(group_id, default_config.clone());
256 default_config.as_ref().clone()
257 }
258 };
259
260 let group_delta = GroupDelta::GroupConstruct(Box::new(PbGroupConstruct {
261 group_config: Some(config),
262 group_id,
263 ..Default::default()
264 }));
265
266 group_deltas.push(group_delta);
267 }
268 }
269 assert!(
270 new_version_delta
271 .state_table_info_delta
272 .insert(
273 table_id,
274 PbStateTableInfoDelta {
275 committed_epoch,
276 compaction_group_id: *raw_group_id,
277 }
278 )
279 .is_none()
280 );
281 }
282 new_version_delta.pre_apply();
283 commit_multi_var!(self.meta_store_ref(), version, compaction_groups_txn)?;
284
285 Ok(())
286 }
287
288 pub async fn unregister_table_ids(
289 &self,
290 table_ids: impl IntoIterator<Item = TableId>,
291 ) -> Result<()> {
292 let table_ids = table_ids.into_iter().collect_vec();
293 if table_ids.is_empty() {
294 return Ok(());
295 }
296
297 let mut versioning_guard = self
298 .versioning
299 .write_with_process_name("unregister_table_ids")
300 .await;
301 let versioning = versioning_guard.deref_mut();
302 let mut version = HummockVersionTransaction::new(
303 &mut versioning.current_version,
304 &mut versioning.hummock_version_deltas,
305 &mut versioning.table_change_log,
306 self.env.notification_manager(),
307 None,
308 &self.metrics,
309 &self.env.opts,
310 &self.version_stat_tx,
311 );
312 let mut new_version_delta = version.new_delta();
313 struct UnregisterGroupChange {
314 remaining_member_count: usize,
315 removed_table_ids: HashSet<TableId>,
316 }
317 let mut group_changes: HashMap<CompactionGroupId, UnregisterGroupChange> = HashMap::new();
318 for table_id in table_ids.iter().copied().unique() {
320 let version = new_version_delta.latest_version();
321 let Some(info) = version.state_table_info.info().get(&table_id) else {
322 continue;
323 };
324 let compaction_group_id = info.compaction_group_id;
325
326 let group_change =
327 group_changes
328 .entry(compaction_group_id)
329 .or_insert_with(|| UnregisterGroupChange {
330 remaining_member_count: version
331 .state_table_info
332 .compaction_group_member_tables()
333 .get(&compaction_group_id)
334 .expect("should exist")
335 .len(),
336 removed_table_ids: HashSet::new(),
337 });
338 group_change.remaining_member_count = group_change
339 .remaining_member_count
340 .checked_sub(1)
341 .expect("member table count should be positive");
342 assert!(group_change.removed_table_ids.insert(table_id));
343 new_version_delta.removed_table_ids.insert(table_id);
344 }
345
346 let mut removed_groups = vec![];
349 for (group_id, change) in group_changes {
350 if change.remaining_member_count == 0 && group_id > StaticCompactionGroupId::End {
351 let max_level = new_version_delta
352 .latest_version()
353 .get_compaction_group_levels(group_id)
354 .levels
355 .len();
356 new_version_delta
357 .group_deltas
358 .entry(group_id)
359 .or_default()
360 .group_deltas
361 .push(GroupDelta::GroupDestroy(PbGroupDestroy {}));
362 remove_compaction_group_metrics(&self.metrics, group_id, max_level);
363 removed_groups.push(group_id);
364 } else {
365 new_version_delta
366 .group_deltas
367 .entry(group_id)
368 .or_default()
369 .group_deltas
370 .push(GroupDelta::PruneTableIdsFromSsts(change.removed_table_ids));
371 }
372 }
373
374 new_version_delta.pre_apply();
375
376 let mut compaction_group_manager = self
379 .compaction_group_manager
380 .write_with_process_name("unregister_table_ids")
381 .await;
382 let mut compaction_groups_txn = compaction_group_manager.start_compaction_groups_txn();
383
384 compaction_groups_txn.purge(HashSet::from_iter(get_compaction_group_ids(
385 version.latest_version(),
386 )));
387 commit_multi_var!(self.meta_store_ref(), version, compaction_groups_txn)?;
388
389 for group_id in removed_groups {
390 self.compaction_state.remove_compaction_group(group_id);
391 }
392
393 let mut stats = self.table_write_throughput_statistic_manager.write();
396 drop(compaction_group_manager);
397 drop(versioning_guard);
398 for table_id in table_ids {
399 stats.remove_table(table_id);
400 }
401
402 Ok(())
404 }
405
406 pub async fn update_compaction_config(
407 &self,
408 compaction_group_ids: &[CompactionGroupId],
409 config_to_update: &[MutableConfig],
410 ) -> Result<()> {
411 {
412 let mut compaction_group_manager = self
414 .compaction_group_manager
415 .write_with_process_name("update_compaction_config")
416 .await;
417 let mut compaction_groups_txn = compaction_group_manager.start_compaction_groups_txn();
418 compaction_groups_txn
419 .update_compaction_config(compaction_group_ids, config_to_update)?;
420 commit_multi_var!(self.meta_store_ref(), compaction_groups_txn)?;
421 }
422
423 if config_to_update
424 .iter()
425 .any(|c| matches!(c, MutableConfig::Level0StopWriteThresholdSubLevelNumber(_)))
426 {
427 self.try_update_write_limits(compaction_group_ids).await;
429 }
430
431 Ok(())
432 }
433
434 pub async fn list_compaction_group(&self) -> Vec<CompactionGroupInfo> {
437 let mut versioning_guard = self
438 .versioning
439 .write_with_process_name("list_compaction_group")
440 .await;
441 let versioning = versioning_guard.deref_mut();
442 let current_version = &versioning.current_version;
443 let mut results = vec![];
444 let compaction_group_manager = self
445 .compaction_group_manager
446 .read_with_process_name("list_compaction_group")
447 .await;
448
449 for levels in current_version.levels.values() {
450 let compaction_config = compaction_group_manager
451 .try_get_compaction_group_config(levels.group_id)
452 .unwrap()
453 .compaction_config
454 .as_ref()
455 .clone();
456 let group = CompactionGroupInfo {
457 id: levels.group_id,
458 parent_id: levels.parent_group_id,
459 member_table_ids: current_version
460 .state_table_info
461 .compaction_group_member_table_ids(levels.group_id)
462 .iter()
463 .copied()
464 .collect_vec(),
465 compaction_config: Some(compaction_config),
466 };
467 results.push(group);
468 }
469 results
470 }
471
472 pub(crate) fn calculate_compaction_group_statistic_from_snapshot(
473 current_version: &HummockVersion,
474 version_stats: &HummockVersionStats,
475 id_to_config: &BTreeMap<CompactionGroupId, CompactionGroup>,
476 ) -> Vec<CompactionGroupStatistic> {
477 let mut infos = vec![];
478 for group_id in current_version.levels.keys() {
479 let compaction_group_config = id_to_config
480 .get(group_id)
481 .expect("compaction group config should exist for every group in current version")
482 .clone();
483 let mut group_info = CompactionGroupStatistic {
484 group_id: *group_id,
485 compaction_group_config,
486 ..Default::default()
487 };
488
489 for table_id in current_version
490 .state_table_info
491 .compaction_group_member_table_ids(*group_id)
492 {
493 let stats_size = version_stats
494 .table_stats
495 .get(table_id)
496 .map(|stats| stats.total_key_size + stats.total_value_size)
497 .unwrap_or(0);
498 let table_size = stats_size.max(0) as u64;
499 group_info.group_size += table_size;
500 group_info.table_statistic.insert(*table_id, table_size);
501 }
502
503 infos.push(group_info);
504 }
505 infos
506 }
507
508 pub async fn calculate_compaction_group_statistic(&self) -> Vec<CompactionGroupStatistic> {
509 let versioning_guard = self
510 .versioning
511 .read_with_process_name("calculate_compaction_group_statistic")
512 .await;
513 let manager = self
514 .compaction_group_manager
515 .read_with_process_name("calculate_compaction_group_statistic")
516 .await;
517 Self::calculate_compaction_group_statistic_from_snapshot(
518 &versioning_guard.current_version,
519 &versioning_guard.version_stats,
520 &manager.compaction_groups,
521 )
522 }
523
524 pub(crate) async fn calculate_compaction_group_statistic_for_tables(
525 &self,
526 table_ids: &[TableId],
527 ) -> Vec<CompactionGroupStatistic> {
528 let groups = {
529 let versioning = self
530 .versioning
531 .read_with_process_name("calculate_compaction_group_statistic_for_tables")
532 .await;
533 let manager = self
534 .compaction_group_manager
535 .read_with_process_name("calculate_compaction_group_statistic_for_tables")
536 .await;
537 let version = &versioning.current_version;
538 let group_ids: BTreeSet<_> = table_ids
539 .iter()
540 .filter_map(|table_id| {
541 version
542 .state_table_info
543 .info()
544 .get(table_id)
545 .map(|info| info.compaction_group_id)
546 })
547 .collect();
548 group_ids
549 .into_iter()
550 .map(|group_id| {
551 let config = manager
552 .try_get_compaction_group_config(group_id)
553 .expect("current group config should exist");
554 let tables = version
555 .state_table_info
556 .compaction_group_member_table_ids(group_id)
557 .iter()
558 .map(|table_id| {
559 let size = versioning
560 .version_stats
561 .table_stats
562 .get(table_id)
563 .map(|stats| stats.total_key_size + stats.total_value_size)
564 .unwrap_or(0)
565 .max(0) as u64;
566 (*table_id, size)
567 })
568 .collect_vec();
569 (group_id, config, tables)
570 })
571 .collect_vec()
572 };
573 groups
574 .into_iter()
575 .map(
576 |(group_id, compaction_group_config, tables)| CompactionGroupStatistic {
577 group_id,
578 group_size: tables.iter().map(|(_, size)| size).sum(),
579 table_statistic: tables.into_iter().collect(),
580 compaction_group_config,
581 },
582 )
583 .collect()
584 }
585
586 pub(crate) async fn initial_compaction_group_config_after_load(
587 &self,
588 versioning_guard: &Versioning,
589 compaction_group_manager: &mut CompactionGroupManager,
590 ) -> Result<()> {
591 let current_version = &versioning_guard.current_version;
593 let all_group_ids = get_compaction_group_ids(current_version).collect_vec();
594 let default_config = compaction_group_manager.default_compaction_config();
595 let mut compaction_groups_txn = compaction_group_manager.start_compaction_groups_txn();
596 compaction_groups_txn.try_create_compaction_groups(&all_group_ids, default_config);
597 commit_multi_var!(self.meta_store_ref(), compaction_groups_txn)?;
598
599 Ok(())
600 }
601}
602
603pub(crate) struct CompactionGroupManager {
611 compaction_groups: BTreeMap<CompactionGroupId, CompactionGroup>,
612 default_config: Arc<CompactionConfig>,
613 pub write_limit: HashMap<CompactionGroupId, WriteLimit>,
615}
616
617impl CompactionGroupManager {
618 pub fn start_compaction_groups_txn(&mut self) -> CompactionGroupTransaction<'_> {
620 CompactionGroupTransaction::new(&mut self.compaction_groups)
621 }
622
623 #[expect(clippy::type_complexity)]
624 pub fn start_owned_compaction_groups_txn<P: DerefMut<Target = Self>>(
625 inner: P,
626 ) -> BTreeMapTransactionInner<
627 CompactionGroupId,
628 CompactionGroup,
629 DerefMutForward<
630 Self,
631 BTreeMap<CompactionGroupId, CompactionGroup>,
632 P,
633 impl Fn(&Self) -> &BTreeMap<CompactionGroupId, CompactionGroup>,
634 impl Fn(&mut Self) -> &mut BTreeMap<CompactionGroupId, CompactionGroup>,
635 >,
636 > {
637 BTreeMapTransactionInner::new(DerefMutForward::new(
638 inner,
639 |mgr| &mgr.compaction_groups,
640 |mgr| &mut mgr.compaction_groups,
641 ))
642 }
643
644 pub(crate) fn try_get_compaction_group_config(
646 &self,
647 compaction_group_id: impl Into<CompactionGroupId>,
648 ) -> Option<CompactionGroup> {
649 self.compaction_groups
650 .get(&compaction_group_id.into())
651 .cloned()
652 }
653
654 pub(crate) fn default_compaction_config(&self) -> Arc<CompactionConfig> {
656 self.default_config.clone()
657 }
658}
659
660fn update_compaction_config(target: &mut CompactionConfig, items: &[MutableConfig]) -> Result<()> {
661 for item in items {
662 match item {
663 MutableConfig::MaxBytesForLevelBase(c) => {
664 target.max_bytes_for_level_base = *c;
665 }
666 MutableConfig::MaxBytesForLevelMultiplier(c) => {
667 target.max_bytes_for_level_multiplier = *c;
668 }
669 MutableConfig::MaxCompactionBytes(c) => {
670 target.max_compaction_bytes = *c;
671 }
672 MutableConfig::SubLevelMaxCompactionBytes(c) => {
673 target.sub_level_max_compaction_bytes = *c;
674 }
675 MutableConfig::Level0TierCompactFileNumber(c) => {
676 target.level0_tier_compact_file_number = *c;
677 }
678 MutableConfig::TargetFileSizeBase(c) => {
679 target.target_file_size_base = *c;
680 }
681 MutableConfig::CompactionFilterMask(c) => {
682 target.compaction_filter_mask = *c;
683 }
684 MutableConfig::MaxSubCompaction(c) => {
685 target.max_sub_compaction = *c;
686 }
687 MutableConfig::Level0StopWriteThresholdSubLevelNumber(c) => {
688 target.level0_stop_write_threshold_sub_level_number = *c;
689 }
690 MutableConfig::Level0SubLevelCompactLevelCount(c) => {
691 target.level0_sub_level_compact_level_count = *c;
692 }
693 MutableConfig::Level0OverlappingSubLevelCompactLevelCount(c) => {
694 target.level0_overlapping_sub_level_compact_level_count = *c;
695 }
696 MutableConfig::MaxSpaceReclaimBytes(c) => {
697 target.max_space_reclaim_bytes = *c;
698 }
699 MutableConfig::Level0MaxCompactFileNumber(c) => {
700 target.level0_max_compact_file_number = *c;
701 }
702 MutableConfig::EnableEmergencyPicker(c) => {
703 target.enable_emergency_picker = *c;
704 }
705 MutableConfig::TombstoneReclaimRatio(c) => {
706 target.tombstone_reclaim_ratio = *c;
707 }
708 MutableConfig::CompressionAlgorithm(c) => {
709 let level = c.get_level();
710 let max_level = try_u32_max_level(target.max_level)?;
711 if level > max_level {
712 return Err(Error::CompactionGroup(format!(
713 "invalid compression_algorithm level {}, max_level is {}",
714 level, target.max_level
715 )));
716 }
717
718 let Some(algorithm) = target.compression_algorithm.get_mut(level as usize) else {
719 return Err(Error::CompactionGroup(format!(
720 "invalid compression_algorithm level {}, compression_algorithm len is {}",
721 level,
722 target.compression_algorithm.len()
723 )));
724 };
725 algorithm.clone_from(&c.compression_algorithm);
726 }
727 MutableConfig::ResetCompressionAlgorithm(reset) => {
728 if *reset {
729 target.compression_algorithm =
730 default_compaction_config::compression_algorithm_vec(try_u32_max_level(
731 target.max_level,
732 )?);
733 }
734 }
735 MutableConfig::MaxL0CompactLevelCount(c) => {
736 target.max_l0_compact_level_count = Some(*c);
737 }
738 MutableConfig::SstAllowedTrivialMoveMinSize(c) => {
739 target.sst_allowed_trivial_move_min_size = Some(*c);
740 }
741 MutableConfig::SplitWeightByVnode(c) => {
742 target.split_weight_by_vnode = *c;
743 }
744 MutableConfig::DisableAutoGroupScheduling(c) => {
745 target.disable_auto_group_scheduling = Some(*c);
746 }
747 MutableConfig::MaxOverlappingLevelSize(c) => {
748 target.max_overlapping_level_size = Some(*c);
749 }
750 MutableConfig::SstAllowedTrivialMoveMaxCount(c) => {
751 target.sst_allowed_trivial_move_max_count = Some(*c);
752 }
753 MutableConfig::EmergencyLevel0SstFileCount(c) => {
754 target.emergency_level0_sst_file_count = Some(*c);
755 }
756 MutableConfig::EmergencyLevel0SubLevelPartition(c) => {
757 target.emergency_level0_sub_level_partition = Some(*c);
758 }
759 MutableConfig::Level0StopWriteThresholdMaxSstCount(c) => {
760 target.level0_stop_write_threshold_max_sst_count = Some(*c);
761 }
762 MutableConfig::Level0StopWriteThresholdMaxSize(c) => {
763 target.level0_stop_write_threshold_max_size = Some(*c);
764 }
765 MutableConfig::EnableOptimizeL0IntervalSelection(c) => {
766 target.enable_optimize_l0_interval_selection = Some(*c);
767 }
768 #[expect(deprecated)]
769 MutableConfig::VnodeAlignedLevelSizeThreshold(_) => {
770 }
772 MutableConfig::MaxKvCountForXor16(c) => {
773 target.max_kv_count_for_xor16 = optional_non_sentinel_u64_config(*c);
774 }
775 MutableConfig::MaxVnodeKeyRangeBytes(c) => {
776 target.max_vnode_key_range_bytes = optional_positive_u64_config(*c);
777 }
778 MutableConfig::SstableFilterType(c) => {
779 parse_sstable_filter_type(&c.filter_type).map_err(Error::CompactionGroup)?;
780 if target.sstable_filter_type.is_empty() {
781 target.sstable_filter_type = default_compaction_config::sstable_filter_type();
782 target
783 .sstable_filter_type
784 .resize(target.max_level as usize + 1, "xor16".to_owned());
785 }
786 let idx = c.get_level() as usize;
787 let level_entry = target.sstable_filter_type.get_mut(idx).ok_or_else(|| {
788 Error::CompactionGroup(format!(
789 "sstable_filter_type level {} is out of range",
790 idx
791 ))
792 })?;
793 level_entry.clone_from(&c.filter_type);
794 }
795 MutableConfig::SstableFilterLayout(c) => {
796 parse_sstable_filter_layout(&c.layout).map_err(Error::CompactionGroup)?;
797 if target.sstable_filter_layout.is_empty() {
798 target.sstable_filter_layout =
799 default_compaction_config::sstable_filter_layout();
800 target
801 .sstable_filter_layout
802 .resize(target.max_level as usize + 1, "blocked".to_owned());
803 }
804 let idx = c.get_level() as usize;
805 let level_entry = target.sstable_filter_layout.get_mut(idx).ok_or_else(|| {
806 Error::CompactionGroup(format!(
807 "sstable_filter_layout level {} is out of range",
808 idx
809 ))
810 })?;
811 level_entry.clone_from(&c.layout);
812 }
813 }
814 }
815 Ok(())
816}
817
818fn optional_u64_config(value: u64) -> Option<u64> {
819 (value != u64::MIN && value != u64::MAX).then_some(value)
820}
821
822fn optional_non_sentinel_u64_config(value: u64) -> Option<u64> {
823 (value != u64::MAX).then_some(value)
824}
825
826fn optional_positive_u64_config(value: u64) -> Option<u64> {
827 optional_u64_config(value).filter(|value| *value > 0)
828}
829
830fn try_u32_max_level(max_level: u64) -> Result<u32> {
831 u32::try_from(max_level).map_err(|_| {
832 Error::CompactionGroup(format!(
833 "invalid max_level {}, expect <= {}",
834 max_level,
835 u32::MAX
836 ))
837 })
838}
839
840impl CompactionGroupTransaction<'_> {
841 pub fn try_create_compaction_groups(
843 &mut self,
844 compaction_group_ids: &[CompactionGroupId],
845 config: Arc<CompactionConfig>,
846 ) -> bool {
847 let mut trivial = true;
848 for id in compaction_group_ids {
849 if self.contains_key(id) {
850 continue;
851 }
852 let new_entry = CompactionGroup::new(*id, config.as_ref().clone());
853 self.insert(*id, new_entry);
854
855 trivial = false;
856 }
857
858 !trivial
859 }
860
861 pub fn create_compaction_groups(
862 &mut self,
863 compaction_group_id: CompactionGroupId,
864 config: Arc<CompactionConfig>,
865 ) {
866 self.try_create_compaction_groups(&[compaction_group_id], config);
867 }
868
869 pub(crate) fn try_get_compaction_group_config(
871 &self,
872 compaction_group_id: CompactionGroupId,
873 ) -> Option<&CompactionGroup> {
874 self.get(&compaction_group_id)
875 }
876
877 pub fn purge(&mut self, existing_groups: HashSet<CompactionGroupId>) {
879 let stale_group = self
880 .tree_ref()
881 .keys()
882 .cloned()
883 .filter(|k| !existing_groups.contains(k))
884 .collect_vec();
885 if stale_group.is_empty() {
886 return;
887 }
888 for group in stale_group {
889 self.remove(group);
890 }
891 }
892
893 pub(crate) fn update_compaction_config(
894 &mut self,
895 compaction_group_ids: &[CompactionGroupId],
896 config_to_update: &[MutableConfig],
897 ) -> Result<HashMap<CompactionGroupId, CompactionGroup>> {
898 let mut results = HashMap::default();
899 for compaction_group_id in compaction_group_ids.iter().unique() {
900 let group = self.get(compaction_group_id).ok_or_else(|| {
901 Error::CompactionGroup(format!("invalid group {}", *compaction_group_id))
902 })?;
903 let mut config = group.compaction_config.as_ref().clone();
904 update_compaction_config(&mut config, config_to_update)?;
905 if let Err(reason) = validate_compaction_config(&config) {
906 return Err(Error::CompactionGroup(reason));
907 }
908 let mut new_group = group.clone();
909 new_group.compaction_config = Arc::new(config);
910 self.insert(*compaction_group_id, new_group.clone());
911 results.insert(new_group.group_id(), new_group);
912 }
913
914 Ok(results)
915 }
916}
917
918#[cfg(test)]
919mod tests {
920 use std::collections::{BTreeMap, HashSet};
921 use std::sync::Arc;
922
923 use itertools::Itertools;
924 use risingwave_common::id::JobId;
925 use risingwave_hummock_sdk::CompactionGroupId;
926 use risingwave_pb::hummock::rise_ctl_update_compaction_config_request::mutable_config::MutableConfig;
927 use risingwave_pb::hummock::rise_ctl_update_compaction_config_request::{
928 CompressionAlgorithm, SstableFilterLayout, SstableFilterType,
929 };
930
931 use crate::controller::SqlMetaStore;
932 use crate::hummock::commit_multi_var;
933 use crate::hummock::compaction::compaction_config::CompactionConfigBuilder;
934 use crate::hummock::error::Result;
935 use crate::hummock::manager::compaction_group_manager::CompactionGroupManager;
936 use crate::hummock::test_utils::setup_compute_env;
937 use crate::model::{Fragment, StreamJobFragments};
938
939 #[test]
940 fn test_update_compaction_config_filter_type_layout_backward_compat() {
941 let mut config = CompactionConfigBuilder::new().build();
942 config.sstable_filter_type.clear();
943 config.sstable_filter_layout.clear();
944
945 super::update_compaction_config(
946 &mut config,
947 &[MutableConfig::SstableFilterType(SstableFilterType {
948 level: 0,
949 filter_type: "xor8".to_owned(),
950 })],
951 )
952 .unwrap();
953 assert_eq!(
954 config.sstable_filter_type.len(),
955 config.max_level as usize + 1
956 );
957 assert_eq!(config.sstable_filter_type[0], "xor8");
958 assert_eq!(config.sstable_filter_type[5], "xor8");
959
960 super::update_compaction_config(
961 &mut config,
962 &[MutableConfig::SstableFilterType(SstableFilterType {
963 level: 1,
964 filter_type: "none".to_owned(),
965 })],
966 )
967 .unwrap();
968 assert_eq!(config.sstable_filter_type[1], "none");
969
970 super::update_compaction_config(
971 &mut config,
972 &[MutableConfig::SstableFilterLayout(SstableFilterLayout {
973 level: 1,
974 layout: "plain".to_owned(),
975 })],
976 )
977 .unwrap();
978 assert_eq!(
979 config.sstable_filter_layout.len(),
980 config.max_level as usize + 1
981 );
982 assert_eq!(config.sstable_filter_layout[1], "plain");
983 assert_eq!(config.sstable_filter_layout[2], "blocked");
984 }
985
986 #[test]
987 fn test_update_compaction_config_rejects_out_of_range_level() {
988 let mut config = CompactionConfigBuilder::new().build();
989 let oob = config.max_level as u32 + 1;
990
991 assert!(
992 super::update_compaction_config(
993 &mut config,
994 &[MutableConfig::SstableFilterType(SstableFilterType {
995 level: oob,
996 filter_type: "xor8".to_owned(),
997 })],
998 )
999 .is_err()
1000 );
1001
1002 assert!(
1003 super::update_compaction_config(
1004 &mut config,
1005 &[MutableConfig::SstableFilterLayout(SstableFilterLayout {
1006 level: oob,
1007 layout: "plain".to_owned(),
1008 })],
1009 )
1010 .is_err()
1011 );
1012
1013 assert!(
1014 super::update_compaction_config(
1015 &mut config,
1016 &[MutableConfig::CompressionAlgorithm(CompressionAlgorithm {
1017 level: oob,
1018 compression_algorithm: "Zstd".to_owned(),
1019 })],
1020 )
1021 .is_err()
1022 );
1023 }
1024
1025 #[test]
1026 fn test_update_compaction_config_rejects_invalid_filter_metadata() {
1027 let mut config = CompactionConfigBuilder::new().build();
1028
1029 assert!(
1030 super::update_compaction_config(
1031 &mut config,
1032 &[MutableConfig::SstableFilterType(SstableFilterType {
1033 level: 0,
1034 filter_type: "unknown".to_owned(),
1035 })],
1036 )
1037 .is_err()
1038 );
1039
1040 let mut config = CompactionConfigBuilder::new().build();
1041
1042 assert!(
1043 super::update_compaction_config(
1044 &mut config,
1045 &[MutableConfig::SstableFilterLayout(SstableFilterLayout {
1046 level: 0,
1047 layout: "unknown".to_owned(),
1048 })],
1049 )
1050 .is_err()
1051 );
1052 }
1053
1054 #[test]
1055 fn test_reset_compression_algorithm_false_is_noop() {
1056 let mut config = CompactionConfigBuilder::new().build();
1057 config.compression_algorithm[3] = "Zstd".to_owned();
1058
1059 super::update_compaction_config(
1060 &mut config,
1061 &[MutableConfig::ResetCompressionAlgorithm(false)],
1062 )
1063 .unwrap();
1064
1065 assert_eq!(config.compression_algorithm[3], "Zstd");
1066 }
1067
1068 #[tokio::test]
1069 async fn test_inner() {
1070 let (env, ..) = setup_compute_env(8080).await;
1071 let mut inner = CompactionGroupManager::new(&env).await.unwrap();
1072 assert_eq!(inner.compaction_groups.len(), 2);
1073
1074 async fn update_compaction_config(
1075 meta: &SqlMetaStore,
1076 inner: &mut CompactionGroupManager,
1077 cg_ids: &[impl Into<CompactionGroupId> + Copy],
1078 config_to_update: &[MutableConfig],
1079 ) -> Result<()> {
1080 let cg_ids = cg_ids.iter().copied().map_into().collect_vec();
1081 let mut compaction_groups_txn = inner.start_compaction_groups_txn();
1082 compaction_groups_txn.update_compaction_config(&cg_ids, config_to_update)?;
1083 commit_multi_var!(meta, compaction_groups_txn)
1084 }
1085
1086 async fn insert_compaction_group_configs(
1087 meta: &SqlMetaStore,
1088 inner: &mut CompactionGroupManager,
1089 cg_ids: &[u64],
1090 ) {
1091 let default_config = inner.default_compaction_config();
1092 let mut compaction_groups_txn = inner.start_compaction_groups_txn();
1093 if compaction_groups_txn.try_create_compaction_groups(
1094 &cg_ids.iter().copied().map_into().collect_vec(),
1095 default_config,
1096 ) {
1097 commit_multi_var!(meta, compaction_groups_txn).unwrap();
1098 }
1099 }
1100
1101 async fn insert_compaction_group_config_with_max_level(
1102 meta: &SqlMetaStore,
1103 inner: &mut CompactionGroupManager,
1104 cg_id: u64,
1105 max_level: u64,
1106 ) {
1107 let mut config = inner.default_compaction_config().as_ref().clone();
1108 config.max_level = max_level;
1109 config.compression_algorithm =
1110 super::default_compaction_config::compression_algorithm_vec(
1111 super::try_u32_max_level(max_level).expect("max_level should fit u32 in test"),
1112 );
1113 let mut compaction_groups_txn = inner.start_compaction_groups_txn();
1114 compaction_groups_txn.create_compaction_groups(cg_id.into(), Arc::new(config));
1115 commit_multi_var!(meta, compaction_groups_txn).unwrap();
1116 }
1117
1118 update_compaction_config(env.meta_store_ref(), &mut inner, &[100, 200], &[])
1119 .await
1120 .unwrap_err();
1121 insert_compaction_group_configs(env.meta_store_ref(), &mut inner, &[100, 200]).await;
1122 assert_eq!(inner.compaction_groups.len(), 4);
1123 let mut inner = CompactionGroupManager::new(&env).await.unwrap();
1124 assert_eq!(inner.compaction_groups.len(), 4);
1125
1126 update_compaction_config(
1127 env.meta_store_ref(),
1128 &mut inner,
1129 &[100, 200],
1130 &[MutableConfig::MaxSubCompaction(123)],
1131 )
1132 .await
1133 .unwrap();
1134 assert_eq!(inner.compaction_groups.len(), 4);
1135 assert_eq!(
1136 inner
1137 .try_get_compaction_group_config(100)
1138 .unwrap()
1139 .compaction_config
1140 .max_sub_compaction,
1141 123
1142 );
1143 assert_eq!(
1144 inner
1145 .try_get_compaction_group_config(200)
1146 .unwrap()
1147 .compaction_config
1148 .max_sub_compaction,
1149 123
1150 );
1151
1152 insert_compaction_group_config_with_max_level(env.meta_store_ref(), &mut inner, 300, 4)
1153 .await;
1154 update_compaction_config(
1155 env.meta_store_ref(),
1156 &mut inner,
1157 &[300],
1158 &[MutableConfig::ResetCompressionAlgorithm(true)],
1159 )
1160 .await
1161 .unwrap();
1162 assert_eq!(
1163 inner
1164 .try_get_compaction_group_config(300)
1165 .unwrap()
1166 .compaction_config
1167 .compression_algorithm,
1168 super::default_compaction_config::compression_algorithm_vec(4)
1169 );
1170 let err = update_compaction_config(
1171 env.meta_store_ref(),
1172 &mut inner,
1173 &[300],
1174 &[MutableConfig::CompressionAlgorithm(CompressionAlgorithm {
1175 level: 6,
1176 compression_algorithm: "Zstd".to_owned(),
1177 })],
1178 )
1179 .await
1180 .unwrap_err();
1181 assert!(
1182 err.to_string()
1183 .contains("invalid compression_algorithm level 6")
1184 );
1185
1186 update_compaction_config(
1187 env.meta_store_ref(),
1188 &mut inner,
1189 &[100],
1190 &[MutableConfig::MaxKvCountForXor16(0)],
1191 )
1192 .await
1193 .unwrap();
1194 assert_eq!(
1195 inner
1196 .try_get_compaction_group_config(100)
1197 .unwrap()
1198 .compaction_config
1199 .max_kv_count_for_xor16,
1200 Some(0)
1201 );
1202 update_compaction_config(
1203 env.meta_store_ref(),
1204 &mut inner,
1205 &[100],
1206 &[MutableConfig::MaxKvCountForXor16(1024)],
1207 )
1208 .await
1209 .unwrap();
1210 assert_eq!(
1211 inner
1212 .try_get_compaction_group_config(100)
1213 .unwrap()
1214 .compaction_config
1215 .max_kv_count_for_xor16,
1216 Some(1024)
1217 );
1218 update_compaction_config(
1219 env.meta_store_ref(),
1220 &mut inner,
1221 &[100],
1222 &[MutableConfig::MaxKvCountForXor16(u64::MAX)],
1223 )
1224 .await
1225 .unwrap();
1226 assert_eq!(
1227 inner
1228 .try_get_compaction_group_config(100)
1229 .unwrap()
1230 .compaction_config
1231 .max_kv_count_for_xor16,
1232 None
1233 );
1234
1235 update_compaction_config(
1236 env.meta_store_ref(),
1237 &mut inner,
1238 &[100],
1239 &[MutableConfig::MaxVnodeKeyRangeBytes(0)],
1240 )
1241 .await
1242 .unwrap();
1243 assert_eq!(
1244 inner
1245 .try_get_compaction_group_config(100)
1246 .unwrap()
1247 .compaction_config
1248 .max_vnode_key_range_bytes,
1249 None
1250 );
1251 update_compaction_config(
1252 env.meta_store_ref(),
1253 &mut inner,
1254 &[100],
1255 &[MutableConfig::MaxVnodeKeyRangeBytes(1024)],
1256 )
1257 .await
1258 .unwrap();
1259 assert_eq!(
1260 inner
1261 .try_get_compaction_group_config(100)
1262 .unwrap()
1263 .compaction_config
1264 .max_vnode_key_range_bytes,
1265 Some(1024)
1266 );
1267 }
1268
1269 #[tokio::test]
1270 async fn test_manager() {
1271 let (_, compaction_group_manager, ..) = setup_compute_env(8080).await;
1272 let table_fragment_1 = StreamJobFragments::for_test(
1273 JobId::new(10),
1274 BTreeMap::from([(
1275 1.into(),
1276 Fragment {
1277 fragment_id: 1.into(),
1278 state_table_ids: vec![10.into(), 11.into(), 12.into(), 13.into()],
1279 ..Default::default()
1280 },
1281 )]),
1282 );
1283 let table_fragment_2 = StreamJobFragments::for_test(
1284 JobId::new(20),
1285 BTreeMap::from([(
1286 2.into(),
1287 Fragment {
1288 fragment_id: 2.into(),
1289 state_table_ids: vec![20.into(), 21.into(), 22.into(), 23.into()],
1290 ..Default::default()
1291 },
1292 )]),
1293 );
1294
1295 let registered_number = || async {
1297 compaction_group_manager
1298 .list_compaction_group()
1299 .await
1300 .iter()
1301 .map(|cg| cg.member_table_ids.len())
1302 .sum::<usize>()
1303 };
1304 let group_number =
1305 || async { compaction_group_manager.list_compaction_group().await.len() };
1306 assert_eq!(registered_number().await, 0);
1307
1308 compaction_group_manager
1309 .register_table_fragments(
1310 Some(table_fragment_1.stream_job_id().as_mv_table_id()),
1311 table_fragment_1
1312 .internal_table_ids()
1313 .into_iter()
1314 .map_into()
1315 .collect(),
1316 )
1317 .await
1318 .unwrap();
1319 assert_eq!(registered_number().await, 4);
1320 compaction_group_manager
1321 .register_table_fragments(
1322 Some(table_fragment_2.stream_job_id().as_mv_table_id()),
1323 table_fragment_2
1324 .internal_table_ids()
1325 .into_iter()
1326 .map_into()
1327 .collect(),
1328 )
1329 .await
1330 .unwrap();
1331 assert_eq!(registered_number().await, 8);
1332
1333 compaction_group_manager
1335 .unregister_table_fragments_vec(std::slice::from_ref(&table_fragment_1))
1336 .await;
1337 assert_eq!(registered_number().await, 4);
1338
1339 compaction_group_manager
1341 .purge(&table_fragment_2.all_table_ids().collect())
1342 .await
1343 .unwrap();
1344 assert_eq!(registered_number().await, 4);
1345 compaction_group_manager
1346 .purge(&HashSet::new())
1347 .await
1348 .unwrap();
1349 assert_eq!(registered_number().await, 0);
1350
1351 assert_eq!(group_number().await, 2);
1352
1353 compaction_group_manager
1354 .register_table_fragments(
1355 Some(table_fragment_1.stream_job_id().as_mv_table_id()),
1356 table_fragment_1
1357 .internal_table_ids()
1358 .into_iter()
1359 .map_into()
1360 .collect(),
1361 )
1362 .await
1363 .unwrap();
1364 assert_eq!(registered_number().await, 4);
1365 assert_eq!(group_number().await, 2);
1366
1367 compaction_group_manager
1368 .unregister_table_fragments_vec(&[table_fragment_1])
1369 .await;
1370 assert_eq!(registered_number().await, 0);
1371 assert_eq!(group_number().await, 2);
1372 }
1373}