Skip to main content

risingwave_meta/hummock/manager/compaction/
compaction_group_manager.rs

1// Copyright 2024 RisingWave Labs
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use 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    /// Should not be called inside [`HummockManager`], because it requests locks internally.
99    /// The implementation acquires `versioning` lock.
100    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    /// The implementation acquires `compaction_group_manager` lock.
112    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    /// Registers `table_fragments` to compaction groups.
122    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            // materialized_view
133            pairs.push((mv_table, StaticCompactionGroupId::MaterializedView));
134        }
135        // internal states
136        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    /// Unregisters `table_fragments` from compaction groups
145    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    /// Unregisters stale members and groups
155    /// The caller should ensure `table_fragments_list` remain unchanged during `purge`.
156    /// Currently `purge` is only called during meta service start ups.
157    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        // As we have released versioning lock, the version that `to_unregister` is calculated from
171        // may not be the same as the one used in unregister_table_ids. It is OK.
172        self.unregister_table_ids(to_unregister).await
173    }
174
175    /// The implementation acquires `versioning` lock.
176    ///
177    /// The method name is temporarily added with a `_for_test` prefix to mark
178    /// that it's currently only used in test.
179    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        // All NewCompactionGroup pairs are mapped to one new compaction group.
209        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            // Remove table write throughput statistics
299            // The Caller acquires `Send`, so we should safely use `write` lock before the await point.
300            // The table write throughput statistic accepts data inconsistencies (unregister table ids fail), so we can clean it up in advance.
301            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        // Remove member tables
330        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                // clean up compaction schedule state for the removed group
372                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        // Purge may cause write to meta store. If it hurts performance while holding versioning
386        // lock, consider to make it in batch.
387        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        // No need to handle DeltaType::GroupDestroy during time travel.
399        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            // Avoid lock conflicts with `try_update_write_limits``
409            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            // Update write limits with lock
424            self.try_update_write_limits(compaction_group_ids).await;
425        }
426
427        Ok(())
428    }
429
430    /// Gets complete compaction group info.
431    /// It is the aggregate of `HummockVersion` and `CompactionGroupConfig`
432    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        // 1. Due to version compatibility, we fix some of the configuration of older versions after hummock starts.
526        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
537/// We muse ensure there is an entry exists in [`CompactionGroupManager`] for any
538/// compaction group found in current hummock version. That's done by invoking
539/// `get_or_insert_compaction_group_config` or `get_or_insert_compaction_group_configs` before
540/// adding any group in current hummock version:
541/// 1. initialize default static compaction group.
542/// 2. register new table to new compaction group.
543/// 3. move existent table to new compaction group.
544pub(crate) struct CompactionGroupManager {
545    compaction_groups: BTreeMap<CompactionGroupId, CompactionGroup>,
546    default_config: Arc<CompactionConfig>,
547    /// Tables that write limit is trigger for.
548    pub write_limit: HashMap<CompactionGroupId, WriteLimit>,
549}
550
551impl CompactionGroupManager {
552    /// Starts a transaction to update compaction group configs.
553    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    /// Tries to get compaction group config for `compaction_group_id`.
579    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    /// Tries to get compaction group config for `compaction_group_id`.
589    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                // Deprecated. Keep accepting the field for old clients but do not apply it.
705            }
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    /// Inserts compaction group configs if they do not exist.
776    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    /// Tries to get compaction group config for `compaction_group_id`.
804    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    /// Removes stale group configs.
812    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        // Test register_table_fragments
1230        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        // Test unregister_table_fragments
1268        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        // Test purge_stale_members: table fragments
1274        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}