1use std::collections::{BTreeMap, 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 {
298 let mut table_write_throughput_statistic_manager =
302 self.table_write_throughput_statistic_manager.write();
303 for &table_id in table_ids.iter().unique() {
304 table_write_throughput_statistic_manager.remove_table(table_id);
305 }
306 }
307
308 let mut versioning_guard = self
309 .versioning
310 .write_with_process_name("unregister_table_ids")
311 .await;
312 let versioning = versioning_guard.deref_mut();
313 let mut version = HummockVersionTransaction::new(
314 &mut versioning.current_version,
315 &mut versioning.hummock_version_deltas,
316 &mut versioning.table_change_log,
317 self.env.notification_manager(),
318 None,
319 &self.metrics,
320 &self.env.opts,
321 &self.version_stat_tx,
322 );
323 let mut new_version_delta = version.new_delta();
324 struct UnregisterGroupChange {
325 remaining_member_count: usize,
326 removed_table_ids: HashSet<TableId>,
327 }
328 let mut group_changes: HashMap<CompactionGroupId, UnregisterGroupChange> = HashMap::new();
329 for table_id in table_ids.into_iter().unique() {
331 let version = new_version_delta.latest_version();
332 let Some(info) = version.state_table_info.info().get(&table_id) else {
333 continue;
334 };
335 let compaction_group_id = info.compaction_group_id;
336
337 let group_change =
338 group_changes
339 .entry(compaction_group_id)
340 .or_insert_with(|| UnregisterGroupChange {
341 remaining_member_count: version
342 .state_table_info
343 .compaction_group_member_tables()
344 .get(&compaction_group_id)
345 .expect("should exist")
346 .len(),
347 removed_table_ids: HashSet::new(),
348 });
349 group_change.remaining_member_count = group_change
350 .remaining_member_count
351 .checked_sub(1)
352 .expect("member table count should be positive");
353 assert!(group_change.removed_table_ids.insert(table_id));
354 new_version_delta.removed_table_ids.insert(table_id);
355 }
356
357 for (group_id, change) in group_changes {
358 if change.remaining_member_count == 0 && group_id > StaticCompactionGroupId::End {
359 let max_level = new_version_delta
360 .latest_version()
361 .get_compaction_group_levels(group_id)
362 .levels
363 .len();
364 new_version_delta
365 .group_deltas
366 .entry(group_id)
367 .or_default()
368 .group_deltas
369 .push(GroupDelta::GroupDestroy(PbGroupDestroy {}));
370 remove_compaction_group_metrics(&self.metrics, group_id, max_level);
371 self.compaction_state.remove_compaction_group(group_id);
373 } else {
374 new_version_delta
375 .group_deltas
376 .entry(group_id)
377 .or_default()
378 .group_deltas
379 .push(GroupDelta::PruneTableIdsFromSsts(change.removed_table_ids));
380 }
381 }
382
383 new_version_delta.pre_apply();
384
385 let mut compaction_group_manager = self
388 .compaction_group_manager
389 .write_with_process_name("unregister_table_ids")
390 .await;
391 let mut compaction_groups_txn = compaction_group_manager.start_compaction_groups_txn();
392
393 compaction_groups_txn.purge(HashSet::from_iter(get_compaction_group_ids(
394 version.latest_version(),
395 )));
396 commit_multi_var!(self.meta_store_ref(), version, compaction_groups_txn)?;
397
398 Ok(())
400 }
401
402 pub async fn update_compaction_config(
403 &self,
404 compaction_group_ids: &[CompactionGroupId],
405 config_to_update: &[MutableConfig],
406 ) -> Result<()> {
407 {
408 let mut compaction_group_manager = self
410 .compaction_group_manager
411 .write_with_process_name("update_compaction_config")
412 .await;
413 let mut compaction_groups_txn = compaction_group_manager.start_compaction_groups_txn();
414 compaction_groups_txn
415 .update_compaction_config(compaction_group_ids, config_to_update)?;
416 commit_multi_var!(self.meta_store_ref(), compaction_groups_txn)?;
417 }
418
419 if config_to_update
420 .iter()
421 .any(|c| matches!(c, MutableConfig::Level0StopWriteThresholdSubLevelNumber(_)))
422 {
423 self.try_update_write_limits(compaction_group_ids).await;
425 }
426
427 Ok(())
428 }
429
430 pub async fn list_compaction_group(&self) -> Vec<CompactionGroupInfo> {
433 let mut versioning_guard = self
434 .versioning
435 .write_with_process_name("list_compaction_group")
436 .await;
437 let versioning = versioning_guard.deref_mut();
438 let current_version = &versioning.current_version;
439 let mut results = vec![];
440 let compaction_group_manager = self
441 .compaction_group_manager
442 .read_with_process_name("list_compaction_group")
443 .await;
444
445 for levels in current_version.levels.values() {
446 let compaction_config = compaction_group_manager
447 .try_get_compaction_group_config(levels.group_id)
448 .unwrap()
449 .compaction_config
450 .as_ref()
451 .clone();
452 let group = CompactionGroupInfo {
453 id: levels.group_id,
454 parent_id: levels.parent_group_id,
455 member_table_ids: current_version
456 .state_table_info
457 .compaction_group_member_table_ids(levels.group_id)
458 .iter()
459 .copied()
460 .collect_vec(),
461 compaction_config: Some(compaction_config),
462 };
463 results.push(group);
464 }
465 results
466 }
467
468 pub(crate) fn calculate_compaction_group_statistic_from_snapshot(
469 current_version: &HummockVersion,
470 version_stats: &HummockVersionStats,
471 id_to_config: &BTreeMap<CompactionGroupId, CompactionGroup>,
472 ) -> Vec<CompactionGroupStatistic> {
473 let mut infos = vec![];
474 for group_id in current_version.levels.keys() {
475 let compaction_group_config = id_to_config
476 .get(group_id)
477 .expect("compaction group config should exist for every group in current version")
478 .clone();
479 let mut group_info = CompactionGroupStatistic {
480 group_id: *group_id,
481 compaction_group_config,
482 ..Default::default()
483 };
484
485 for table_id in current_version
486 .state_table_info
487 .compaction_group_member_table_ids(*group_id)
488 {
489 let stats_size = version_stats
490 .table_stats
491 .get(table_id)
492 .map(|stats| stats.total_key_size + stats.total_value_size)
493 .unwrap_or(0);
494 let table_size = stats_size.max(0) as u64;
495 group_info.group_size += table_size;
496 group_info.table_statistic.insert(*table_id, table_size);
497 }
498
499 infos.push(group_info);
500 }
501 infos
502 }
503
504 pub async fn calculate_compaction_group_statistic(&self) -> Vec<CompactionGroupStatistic> {
505 let versioning_guard = self
506 .versioning
507 .read_with_process_name("calculate_compaction_group_statistic")
508 .await;
509 let manager = self
510 .compaction_group_manager
511 .read_with_process_name("calculate_compaction_group_statistic")
512 .await;
513 Self::calculate_compaction_group_statistic_from_snapshot(
514 &versioning_guard.current_version,
515 &versioning_guard.version_stats,
516 &manager.compaction_groups,
517 )
518 }
519
520 pub(crate) async fn initial_compaction_group_config_after_load(
521 &self,
522 versioning_guard: &Versioning,
523 compaction_group_manager: &mut CompactionGroupManager,
524 ) -> Result<()> {
525 let current_version = &versioning_guard.current_version;
527 let all_group_ids = get_compaction_group_ids(current_version).collect_vec();
528 let default_config = compaction_group_manager.default_compaction_config();
529 let mut compaction_groups_txn = compaction_group_manager.start_compaction_groups_txn();
530 compaction_groups_txn.try_create_compaction_groups(&all_group_ids, default_config);
531 commit_multi_var!(self.meta_store_ref(), compaction_groups_txn)?;
532
533 Ok(())
534 }
535}
536
537pub(crate) struct CompactionGroupManager {
545 compaction_groups: BTreeMap<CompactionGroupId, CompactionGroup>,
546 default_config: Arc<CompactionConfig>,
547 pub write_limit: HashMap<CompactionGroupId, WriteLimit>,
549}
550
551impl CompactionGroupManager {
552 pub fn start_compaction_groups_txn(&mut self) -> CompactionGroupTransaction<'_> {
554 CompactionGroupTransaction::new(&mut self.compaction_groups)
555 }
556
557 #[expect(clippy::type_complexity)]
558 pub fn start_owned_compaction_groups_txn<P: DerefMut<Target = Self>>(
559 inner: P,
560 ) -> BTreeMapTransactionInner<
561 CompactionGroupId,
562 CompactionGroup,
563 DerefMutForward<
564 Self,
565 BTreeMap<CompactionGroupId, CompactionGroup>,
566 P,
567 impl Fn(&Self) -> &BTreeMap<CompactionGroupId, CompactionGroup>,
568 impl Fn(&mut Self) -> &mut BTreeMap<CompactionGroupId, CompactionGroup>,
569 >,
570 > {
571 BTreeMapTransactionInner::new(DerefMutForward::new(
572 inner,
573 |mgr| &mgr.compaction_groups,
574 |mgr| &mut mgr.compaction_groups,
575 ))
576 }
577
578 pub(crate) fn try_get_compaction_group_config(
580 &self,
581 compaction_group_id: impl Into<CompactionGroupId>,
582 ) -> Option<CompactionGroup> {
583 self.compaction_groups
584 .get(&compaction_group_id.into())
585 .cloned()
586 }
587
588 pub(crate) fn default_compaction_config(&self) -> Arc<CompactionConfig> {
590 self.default_config.clone()
591 }
592}
593
594fn update_compaction_config(target: &mut CompactionConfig, items: &[MutableConfig]) -> Result<()> {
595 for item in items {
596 match item {
597 MutableConfig::MaxBytesForLevelBase(c) => {
598 target.max_bytes_for_level_base = *c;
599 }
600 MutableConfig::MaxBytesForLevelMultiplier(c) => {
601 target.max_bytes_for_level_multiplier = *c;
602 }
603 MutableConfig::MaxCompactionBytes(c) => {
604 target.max_compaction_bytes = *c;
605 }
606 MutableConfig::SubLevelMaxCompactionBytes(c) => {
607 target.sub_level_max_compaction_bytes = *c;
608 }
609 MutableConfig::Level0TierCompactFileNumber(c) => {
610 target.level0_tier_compact_file_number = *c;
611 }
612 MutableConfig::TargetFileSizeBase(c) => {
613 target.target_file_size_base = *c;
614 }
615 MutableConfig::CompactionFilterMask(c) => {
616 target.compaction_filter_mask = *c;
617 }
618 MutableConfig::MaxSubCompaction(c) => {
619 target.max_sub_compaction = *c;
620 }
621 MutableConfig::Level0StopWriteThresholdSubLevelNumber(c) => {
622 target.level0_stop_write_threshold_sub_level_number = *c;
623 }
624 MutableConfig::Level0SubLevelCompactLevelCount(c) => {
625 target.level0_sub_level_compact_level_count = *c;
626 }
627 MutableConfig::Level0OverlappingSubLevelCompactLevelCount(c) => {
628 target.level0_overlapping_sub_level_compact_level_count = *c;
629 }
630 MutableConfig::MaxSpaceReclaimBytes(c) => {
631 target.max_space_reclaim_bytes = *c;
632 }
633 MutableConfig::Level0MaxCompactFileNumber(c) => {
634 target.level0_max_compact_file_number = *c;
635 }
636 MutableConfig::EnableEmergencyPicker(c) => {
637 target.enable_emergency_picker = *c;
638 }
639 MutableConfig::TombstoneReclaimRatio(c) => {
640 target.tombstone_reclaim_ratio = *c;
641 }
642 MutableConfig::CompressionAlgorithm(c) => {
643 let level = c.get_level();
644 let max_level = try_u32_max_level(target.max_level)?;
645 if level > max_level {
646 return Err(Error::CompactionGroup(format!(
647 "invalid compression_algorithm level {}, max_level is {}",
648 level, target.max_level
649 )));
650 }
651
652 let Some(algorithm) = target.compression_algorithm.get_mut(level as usize) else {
653 return Err(Error::CompactionGroup(format!(
654 "invalid compression_algorithm level {}, compression_algorithm len is {}",
655 level,
656 target.compression_algorithm.len()
657 )));
658 };
659 algorithm.clone_from(&c.compression_algorithm);
660 }
661 MutableConfig::ResetCompressionAlgorithm(reset) => {
662 if *reset {
663 target.compression_algorithm =
664 default_compaction_config::compression_algorithm_vec(try_u32_max_level(
665 target.max_level,
666 )?);
667 }
668 }
669 MutableConfig::MaxL0CompactLevelCount(c) => {
670 target.max_l0_compact_level_count = Some(*c);
671 }
672 MutableConfig::SstAllowedTrivialMoveMinSize(c) => {
673 target.sst_allowed_trivial_move_min_size = Some(*c);
674 }
675 MutableConfig::SplitWeightByVnode(c) => {
676 target.split_weight_by_vnode = *c;
677 }
678 MutableConfig::DisableAutoGroupScheduling(c) => {
679 target.disable_auto_group_scheduling = Some(*c);
680 }
681 MutableConfig::MaxOverlappingLevelSize(c) => {
682 target.max_overlapping_level_size = Some(*c);
683 }
684 MutableConfig::SstAllowedTrivialMoveMaxCount(c) => {
685 target.sst_allowed_trivial_move_max_count = Some(*c);
686 }
687 MutableConfig::EmergencyLevel0SstFileCount(c) => {
688 target.emergency_level0_sst_file_count = Some(*c);
689 }
690 MutableConfig::EmergencyLevel0SubLevelPartition(c) => {
691 target.emergency_level0_sub_level_partition = Some(*c);
692 }
693 MutableConfig::Level0StopWriteThresholdMaxSstCount(c) => {
694 target.level0_stop_write_threshold_max_sst_count = Some(*c);
695 }
696 MutableConfig::Level0StopWriteThresholdMaxSize(c) => {
697 target.level0_stop_write_threshold_max_size = Some(*c);
698 }
699 MutableConfig::EnableOptimizeL0IntervalSelection(c) => {
700 target.enable_optimize_l0_interval_selection = Some(*c);
701 }
702 #[expect(deprecated)]
703 MutableConfig::VnodeAlignedLevelSizeThreshold(_) => {
704 }
706 MutableConfig::MaxKvCountForXor16(c) => {
707 target.max_kv_count_for_xor16 = optional_non_sentinel_u64_config(*c);
708 }
709 MutableConfig::MaxVnodeKeyRangeBytes(c) => {
710 target.max_vnode_key_range_bytes = optional_positive_u64_config(*c);
711 }
712 MutableConfig::SstableFilterType(c) => {
713 parse_sstable_filter_type(&c.filter_type).map_err(Error::CompactionGroup)?;
714 if target.sstable_filter_type.is_empty() {
715 target.sstable_filter_type = default_compaction_config::sstable_filter_type();
716 target
717 .sstable_filter_type
718 .resize(target.max_level as usize + 1, "xor16".to_owned());
719 }
720 let idx = c.get_level() as usize;
721 let level_entry = target.sstable_filter_type.get_mut(idx).ok_or_else(|| {
722 Error::CompactionGroup(format!(
723 "sstable_filter_type level {} is out of range",
724 idx
725 ))
726 })?;
727 level_entry.clone_from(&c.filter_type);
728 }
729 MutableConfig::SstableFilterLayout(c) => {
730 parse_sstable_filter_layout(&c.layout).map_err(Error::CompactionGroup)?;
731 if target.sstable_filter_layout.is_empty() {
732 target.sstable_filter_layout =
733 default_compaction_config::sstable_filter_layout();
734 target
735 .sstable_filter_layout
736 .resize(target.max_level as usize + 1, "blocked".to_owned());
737 }
738 let idx = c.get_level() as usize;
739 let level_entry = target.sstable_filter_layout.get_mut(idx).ok_or_else(|| {
740 Error::CompactionGroup(format!(
741 "sstable_filter_layout level {} is out of range",
742 idx
743 ))
744 })?;
745 level_entry.clone_from(&c.layout);
746 }
747 }
748 }
749 Ok(())
750}
751
752fn optional_u64_config(value: u64) -> Option<u64> {
753 (value != u64::MIN && value != u64::MAX).then_some(value)
754}
755
756fn optional_non_sentinel_u64_config(value: u64) -> Option<u64> {
757 (value != u64::MAX).then_some(value)
758}
759
760fn optional_positive_u64_config(value: u64) -> Option<u64> {
761 optional_u64_config(value).filter(|value| *value > 0)
762}
763
764fn try_u32_max_level(max_level: u64) -> Result<u32> {
765 u32::try_from(max_level).map_err(|_| {
766 Error::CompactionGroup(format!(
767 "invalid max_level {}, expect <= {}",
768 max_level,
769 u32::MAX
770 ))
771 })
772}
773
774impl CompactionGroupTransaction<'_> {
775 pub fn try_create_compaction_groups(
777 &mut self,
778 compaction_group_ids: &[CompactionGroupId],
779 config: Arc<CompactionConfig>,
780 ) -> bool {
781 let mut trivial = true;
782 for id in compaction_group_ids {
783 if self.contains_key(id) {
784 continue;
785 }
786 let new_entry = CompactionGroup::new(*id, config.as_ref().clone());
787 self.insert(*id, new_entry);
788
789 trivial = false;
790 }
791
792 !trivial
793 }
794
795 pub fn create_compaction_groups(
796 &mut self,
797 compaction_group_id: CompactionGroupId,
798 config: Arc<CompactionConfig>,
799 ) {
800 self.try_create_compaction_groups(&[compaction_group_id], config);
801 }
802
803 pub(crate) fn try_get_compaction_group_config(
805 &self,
806 compaction_group_id: CompactionGroupId,
807 ) -> Option<&CompactionGroup> {
808 self.get(&compaction_group_id)
809 }
810
811 pub fn purge(&mut self, existing_groups: HashSet<CompactionGroupId>) {
813 let stale_group = self
814 .tree_ref()
815 .keys()
816 .cloned()
817 .filter(|k| !existing_groups.contains(k))
818 .collect_vec();
819 if stale_group.is_empty() {
820 return;
821 }
822 for group in stale_group {
823 self.remove(group);
824 }
825 }
826
827 pub(crate) fn update_compaction_config(
828 &mut self,
829 compaction_group_ids: &[CompactionGroupId],
830 config_to_update: &[MutableConfig],
831 ) -> Result<HashMap<CompactionGroupId, CompactionGroup>> {
832 let mut results = HashMap::default();
833 for compaction_group_id in compaction_group_ids.iter().unique() {
834 let group = self.get(compaction_group_id).ok_or_else(|| {
835 Error::CompactionGroup(format!("invalid group {}", *compaction_group_id))
836 })?;
837 let mut config = group.compaction_config.as_ref().clone();
838 update_compaction_config(&mut config, config_to_update)?;
839 if let Err(reason) = validate_compaction_config(&config) {
840 return Err(Error::CompactionGroup(reason));
841 }
842 let mut new_group = group.clone();
843 new_group.compaction_config = Arc::new(config);
844 self.insert(*compaction_group_id, new_group.clone());
845 results.insert(new_group.group_id(), new_group);
846 }
847
848 Ok(results)
849 }
850}
851
852#[cfg(test)]
853mod tests {
854 use std::collections::{BTreeMap, HashSet};
855 use std::sync::Arc;
856
857 use itertools::Itertools;
858 use risingwave_common::id::JobId;
859 use risingwave_hummock_sdk::CompactionGroupId;
860 use risingwave_pb::hummock::rise_ctl_update_compaction_config_request::mutable_config::MutableConfig;
861 use risingwave_pb::hummock::rise_ctl_update_compaction_config_request::{
862 CompressionAlgorithm, SstableFilterLayout, SstableFilterType,
863 };
864
865 use crate::controller::SqlMetaStore;
866 use crate::hummock::commit_multi_var;
867 use crate::hummock::compaction::compaction_config::CompactionConfigBuilder;
868 use crate::hummock::error::Result;
869 use crate::hummock::manager::compaction_group_manager::CompactionGroupManager;
870 use crate::hummock::test_utils::setup_compute_env;
871 use crate::model::{Fragment, StreamJobFragments};
872
873 #[test]
874 fn test_update_compaction_config_filter_type_layout_backward_compat() {
875 let mut config = CompactionConfigBuilder::new().build();
876 config.sstable_filter_type.clear();
877 config.sstable_filter_layout.clear();
878
879 super::update_compaction_config(
880 &mut config,
881 &[MutableConfig::SstableFilterType(SstableFilterType {
882 level: 0,
883 filter_type: "xor8".to_owned(),
884 })],
885 )
886 .unwrap();
887 assert_eq!(
888 config.sstable_filter_type.len(),
889 config.max_level as usize + 1
890 );
891 assert_eq!(config.sstable_filter_type[0], "xor8");
892 assert_eq!(config.sstable_filter_type[5], "xor8");
893
894 super::update_compaction_config(
895 &mut config,
896 &[MutableConfig::SstableFilterType(SstableFilterType {
897 level: 1,
898 filter_type: "none".to_owned(),
899 })],
900 )
901 .unwrap();
902 assert_eq!(config.sstable_filter_type[1], "none");
903
904 super::update_compaction_config(
905 &mut config,
906 &[MutableConfig::SstableFilterLayout(SstableFilterLayout {
907 level: 1,
908 layout: "plain".to_owned(),
909 })],
910 )
911 .unwrap();
912 assert_eq!(
913 config.sstable_filter_layout.len(),
914 config.max_level as usize + 1
915 );
916 assert_eq!(config.sstable_filter_layout[1], "plain");
917 assert_eq!(config.sstable_filter_layout[2], "blocked");
918 }
919
920 #[test]
921 fn test_update_compaction_config_rejects_out_of_range_level() {
922 let mut config = CompactionConfigBuilder::new().build();
923 let oob = config.max_level as u32 + 1;
924
925 assert!(
926 super::update_compaction_config(
927 &mut config,
928 &[MutableConfig::SstableFilterType(SstableFilterType {
929 level: oob,
930 filter_type: "xor8".to_owned(),
931 })],
932 )
933 .is_err()
934 );
935
936 assert!(
937 super::update_compaction_config(
938 &mut config,
939 &[MutableConfig::SstableFilterLayout(SstableFilterLayout {
940 level: oob,
941 layout: "plain".to_owned(),
942 })],
943 )
944 .is_err()
945 );
946
947 assert!(
948 super::update_compaction_config(
949 &mut config,
950 &[MutableConfig::CompressionAlgorithm(CompressionAlgorithm {
951 level: oob,
952 compression_algorithm: "Zstd".to_owned(),
953 })],
954 )
955 .is_err()
956 );
957 }
958
959 #[test]
960 fn test_update_compaction_config_rejects_invalid_filter_metadata() {
961 let mut config = CompactionConfigBuilder::new().build();
962
963 assert!(
964 super::update_compaction_config(
965 &mut config,
966 &[MutableConfig::SstableFilterType(SstableFilterType {
967 level: 0,
968 filter_type: "unknown".to_owned(),
969 })],
970 )
971 .is_err()
972 );
973
974 let mut config = CompactionConfigBuilder::new().build();
975
976 assert!(
977 super::update_compaction_config(
978 &mut config,
979 &[MutableConfig::SstableFilterLayout(SstableFilterLayout {
980 level: 0,
981 layout: "unknown".to_owned(),
982 })],
983 )
984 .is_err()
985 );
986 }
987
988 #[test]
989 fn test_reset_compression_algorithm_false_is_noop() {
990 let mut config = CompactionConfigBuilder::new().build();
991 config.compression_algorithm[3] = "Zstd".to_owned();
992
993 super::update_compaction_config(
994 &mut config,
995 &[MutableConfig::ResetCompressionAlgorithm(false)],
996 )
997 .unwrap();
998
999 assert_eq!(config.compression_algorithm[3], "Zstd");
1000 }
1001
1002 #[tokio::test]
1003 async fn test_inner() {
1004 let (env, ..) = setup_compute_env(8080).await;
1005 let mut inner = CompactionGroupManager::new(&env).await.unwrap();
1006 assert_eq!(inner.compaction_groups.len(), 2);
1007
1008 async fn update_compaction_config(
1009 meta: &SqlMetaStore,
1010 inner: &mut CompactionGroupManager,
1011 cg_ids: &[impl Into<CompactionGroupId> + Copy],
1012 config_to_update: &[MutableConfig],
1013 ) -> Result<()> {
1014 let cg_ids = cg_ids.iter().copied().map_into().collect_vec();
1015 let mut compaction_groups_txn = inner.start_compaction_groups_txn();
1016 compaction_groups_txn.update_compaction_config(&cg_ids, config_to_update)?;
1017 commit_multi_var!(meta, compaction_groups_txn)
1018 }
1019
1020 async fn insert_compaction_group_configs(
1021 meta: &SqlMetaStore,
1022 inner: &mut CompactionGroupManager,
1023 cg_ids: &[u64],
1024 ) {
1025 let default_config = inner.default_compaction_config();
1026 let mut compaction_groups_txn = inner.start_compaction_groups_txn();
1027 if compaction_groups_txn.try_create_compaction_groups(
1028 &cg_ids.iter().copied().map_into().collect_vec(),
1029 default_config,
1030 ) {
1031 commit_multi_var!(meta, compaction_groups_txn).unwrap();
1032 }
1033 }
1034
1035 async fn insert_compaction_group_config_with_max_level(
1036 meta: &SqlMetaStore,
1037 inner: &mut CompactionGroupManager,
1038 cg_id: u64,
1039 max_level: u64,
1040 ) {
1041 let mut config = inner.default_compaction_config().as_ref().clone();
1042 config.max_level = max_level;
1043 config.compression_algorithm =
1044 super::default_compaction_config::compression_algorithm_vec(
1045 super::try_u32_max_level(max_level).expect("max_level should fit u32 in test"),
1046 );
1047 let mut compaction_groups_txn = inner.start_compaction_groups_txn();
1048 compaction_groups_txn.create_compaction_groups(cg_id.into(), Arc::new(config));
1049 commit_multi_var!(meta, compaction_groups_txn).unwrap();
1050 }
1051
1052 update_compaction_config(env.meta_store_ref(), &mut inner, &[100, 200], &[])
1053 .await
1054 .unwrap_err();
1055 insert_compaction_group_configs(env.meta_store_ref(), &mut inner, &[100, 200]).await;
1056 assert_eq!(inner.compaction_groups.len(), 4);
1057 let mut inner = CompactionGroupManager::new(&env).await.unwrap();
1058 assert_eq!(inner.compaction_groups.len(), 4);
1059
1060 update_compaction_config(
1061 env.meta_store_ref(),
1062 &mut inner,
1063 &[100, 200],
1064 &[MutableConfig::MaxSubCompaction(123)],
1065 )
1066 .await
1067 .unwrap();
1068 assert_eq!(inner.compaction_groups.len(), 4);
1069 assert_eq!(
1070 inner
1071 .try_get_compaction_group_config(100)
1072 .unwrap()
1073 .compaction_config
1074 .max_sub_compaction,
1075 123
1076 );
1077 assert_eq!(
1078 inner
1079 .try_get_compaction_group_config(200)
1080 .unwrap()
1081 .compaction_config
1082 .max_sub_compaction,
1083 123
1084 );
1085
1086 insert_compaction_group_config_with_max_level(env.meta_store_ref(), &mut inner, 300, 4)
1087 .await;
1088 update_compaction_config(
1089 env.meta_store_ref(),
1090 &mut inner,
1091 &[300],
1092 &[MutableConfig::ResetCompressionAlgorithm(true)],
1093 )
1094 .await
1095 .unwrap();
1096 assert_eq!(
1097 inner
1098 .try_get_compaction_group_config(300)
1099 .unwrap()
1100 .compaction_config
1101 .compression_algorithm,
1102 super::default_compaction_config::compression_algorithm_vec(4)
1103 );
1104 let err = update_compaction_config(
1105 env.meta_store_ref(),
1106 &mut inner,
1107 &[300],
1108 &[MutableConfig::CompressionAlgorithm(CompressionAlgorithm {
1109 level: 6,
1110 compression_algorithm: "Zstd".to_owned(),
1111 })],
1112 )
1113 .await
1114 .unwrap_err();
1115 assert!(
1116 err.to_string()
1117 .contains("invalid compression_algorithm level 6")
1118 );
1119
1120 update_compaction_config(
1121 env.meta_store_ref(),
1122 &mut inner,
1123 &[100],
1124 &[MutableConfig::MaxKvCountForXor16(0)],
1125 )
1126 .await
1127 .unwrap();
1128 assert_eq!(
1129 inner
1130 .try_get_compaction_group_config(100)
1131 .unwrap()
1132 .compaction_config
1133 .max_kv_count_for_xor16,
1134 Some(0)
1135 );
1136 update_compaction_config(
1137 env.meta_store_ref(),
1138 &mut inner,
1139 &[100],
1140 &[MutableConfig::MaxKvCountForXor16(1024)],
1141 )
1142 .await
1143 .unwrap();
1144 assert_eq!(
1145 inner
1146 .try_get_compaction_group_config(100)
1147 .unwrap()
1148 .compaction_config
1149 .max_kv_count_for_xor16,
1150 Some(1024)
1151 );
1152 update_compaction_config(
1153 env.meta_store_ref(),
1154 &mut inner,
1155 &[100],
1156 &[MutableConfig::MaxKvCountForXor16(u64::MAX)],
1157 )
1158 .await
1159 .unwrap();
1160 assert_eq!(
1161 inner
1162 .try_get_compaction_group_config(100)
1163 .unwrap()
1164 .compaction_config
1165 .max_kv_count_for_xor16,
1166 None
1167 );
1168
1169 update_compaction_config(
1170 env.meta_store_ref(),
1171 &mut inner,
1172 &[100],
1173 &[MutableConfig::MaxVnodeKeyRangeBytes(0)],
1174 )
1175 .await
1176 .unwrap();
1177 assert_eq!(
1178 inner
1179 .try_get_compaction_group_config(100)
1180 .unwrap()
1181 .compaction_config
1182 .max_vnode_key_range_bytes,
1183 None
1184 );
1185 update_compaction_config(
1186 env.meta_store_ref(),
1187 &mut inner,
1188 &[100],
1189 &[MutableConfig::MaxVnodeKeyRangeBytes(1024)],
1190 )
1191 .await
1192 .unwrap();
1193 assert_eq!(
1194 inner
1195 .try_get_compaction_group_config(100)
1196 .unwrap()
1197 .compaction_config
1198 .max_vnode_key_range_bytes,
1199 Some(1024)
1200 );
1201 }
1202
1203 #[tokio::test]
1204 async fn test_manager() {
1205 let (_, compaction_group_manager, ..) = setup_compute_env(8080).await;
1206 let table_fragment_1 = StreamJobFragments::for_test(
1207 JobId::new(10),
1208 BTreeMap::from([(
1209 1.into(),
1210 Fragment {
1211 fragment_id: 1.into(),
1212 state_table_ids: vec![10.into(), 11.into(), 12.into(), 13.into()],
1213 ..Default::default()
1214 },
1215 )]),
1216 );
1217 let table_fragment_2 = StreamJobFragments::for_test(
1218 JobId::new(20),
1219 BTreeMap::from([(
1220 2.into(),
1221 Fragment {
1222 fragment_id: 2.into(),
1223 state_table_ids: vec![20.into(), 21.into(), 22.into(), 23.into()],
1224 ..Default::default()
1225 },
1226 )]),
1227 );
1228
1229 let registered_number = || async {
1231 compaction_group_manager
1232 .list_compaction_group()
1233 .await
1234 .iter()
1235 .map(|cg| cg.member_table_ids.len())
1236 .sum::<usize>()
1237 };
1238 let group_number =
1239 || async { compaction_group_manager.list_compaction_group().await.len() };
1240 assert_eq!(registered_number().await, 0);
1241
1242 compaction_group_manager
1243 .register_table_fragments(
1244 Some(table_fragment_1.stream_job_id().as_mv_table_id()),
1245 table_fragment_1
1246 .internal_table_ids()
1247 .into_iter()
1248 .map_into()
1249 .collect(),
1250 )
1251 .await
1252 .unwrap();
1253 assert_eq!(registered_number().await, 4);
1254 compaction_group_manager
1255 .register_table_fragments(
1256 Some(table_fragment_2.stream_job_id().as_mv_table_id()),
1257 table_fragment_2
1258 .internal_table_ids()
1259 .into_iter()
1260 .map_into()
1261 .collect(),
1262 )
1263 .await
1264 .unwrap();
1265 assert_eq!(registered_number().await, 8);
1266
1267 compaction_group_manager
1269 .unregister_table_fragments_vec(std::slice::from_ref(&table_fragment_1))
1270 .await;
1271 assert_eq!(registered_number().await, 4);
1272
1273 compaction_group_manager
1275 .purge(&table_fragment_2.all_table_ids().collect())
1276 .await
1277 .unwrap();
1278 assert_eq!(registered_number().await, 4);
1279 compaction_group_manager
1280 .purge(&HashSet::new())
1281 .await
1282 .unwrap();
1283 assert_eq!(registered_number().await, 0);
1284
1285 assert_eq!(group_number().await, 2);
1286
1287 compaction_group_manager
1288 .register_table_fragments(
1289 Some(table_fragment_1.stream_job_id().as_mv_table_id()),
1290 table_fragment_1
1291 .internal_table_ids()
1292 .into_iter()
1293 .map_into()
1294 .collect(),
1295 )
1296 .await
1297 .unwrap();
1298 assert_eq!(registered_number().await, 4);
1299 assert_eq!(group_number().await, 2);
1300
1301 compaction_group_manager
1302 .unregister_table_fragments_vec(&[table_fragment_1])
1303 .await;
1304 assert_eq!(registered_number().await, 0);
1305 assert_eq!(group_number().await, 2);
1306 }
1307}