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, 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    /// 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        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        // Remove member tables
319        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        // Defer schedule-state cleanup until commit succeeds: it is not transactional and
347        // cannot be restored automatically if the group deletion fails.
348        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        // Purge may cause write to meta store. If it hurts performance while holding versioning
377        // lock, consider to make it in batch.
378        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        // Serialize removal with commit's statistics publication. Cleaning up before taking
394        // versioning could let an in-flight commit recreate a deleted table's history.
395        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        // No need to handle DeltaType::GroupDestroy during time travel.
403        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            // Avoid lock conflicts with `try_update_write_limits``
413            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            // Update write limits with lock
428            self.try_update_write_limits(compaction_group_ids).await;
429        }
430
431        Ok(())
432    }
433
434    /// Gets complete compaction group info.
435    /// It is the aggregate of `HummockVersion` and `CompactionGroupConfig`
436    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        // 1. Due to version compatibility, we fix some of the configuration of older versions after hummock starts.
592        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
603/// We muse ensure there is an entry exists in [`CompactionGroupManager`] for any
604/// compaction group found in current hummock version. That's done by invoking
605/// `get_or_insert_compaction_group_config` or `get_or_insert_compaction_group_configs` before
606/// adding any group in current hummock version:
607/// 1. initialize default static compaction group.
608/// 2. register new table to new compaction group.
609/// 3. move existent table to new compaction group.
610pub(crate) struct CompactionGroupManager {
611    compaction_groups: BTreeMap<CompactionGroupId, CompactionGroup>,
612    default_config: Arc<CompactionConfig>,
613    /// Tables that write limit is trigger for.
614    pub write_limit: HashMap<CompactionGroupId, WriteLimit>,
615}
616
617impl CompactionGroupManager {
618    /// Starts a transaction to update compaction group configs.
619    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    /// Tries to get compaction group config for `compaction_group_id`.
645    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    /// Tries to get compaction group config for `compaction_group_id`.
655    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                // Deprecated. Keep accepting the field for old clients but do not apply it.
771            }
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    /// Inserts compaction group configs if they do not exist.
842    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    /// Tries to get compaction group config for `compaction_group_id`.
870    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    /// Removes stale group configs.
878    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        // Test register_table_fragments
1296        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        // Test unregister_table_fragments
1334        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        // Test purge_stale_members: table fragments
1340        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}