Skip to main content

risingwave_meta/hummock/manager/compaction/compaction_group_schedule/
normalize.rs

1// Copyright 2026 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
15//! Normalize overlapping table ranges with the existing plan, recheck, and apply loop.
16
17use std::collections::HashMap;
18use std::ops::DerefMut;
19use std::sync::Arc;
20
21use bytes::Bytes;
22use itertools::Itertools;
23use risingwave_common::hash::VirtualNode;
24use risingwave_hummock_sdk::CompactionGroupId;
25use risingwave_hummock_sdk::compact_task::{ReportTask, is_compaction_task_expired};
26use risingwave_hummock_sdk::compaction_group::{StateTableId, group_split};
27use risingwave_hummock_sdk::version::{GroupDelta, GroupDeltas, HummockVersion};
28use risingwave_pb::hummock::compact_task::{TaskStatus, TaskType};
29use risingwave_pb::hummock::{CompatibilityVersion, PbGroupConstruct, PbStateTableInfoDelta};
30
31use super::super::compaction_group_manager::CompactionGroupManager;
32use super::CompactionGroupStatistic;
33use crate::hummock::error::{Error, Result};
34use crate::hummock::manager::transaction::HummockVersionTransaction;
35use crate::hummock::manager::{HummockManager, commit_multi_var};
36use crate::hummock::sequence::{next_compaction_group_id, next_sstable_id};
37
38#[derive(Debug, PartialEq, Eq)]
39struct NormalizePlan {
40    parent_group_id: CompactionGroupId,
41    parent_table_ids: Vec<StateTableId>,
42    boundary_table_id: StateTableId,
43}
44
45impl NormalizePlan {
46    fn split_key(&self) -> Bytes {
47        group_split::build_split_full_key(self.boundary_table_id, VirtualNode::ZERO)
48            .encode()
49            .into()
50    }
51
52    fn split_table_ids(&self) -> (Vec<StateTableId>, Vec<StateTableId>) {
53        let split_full_key =
54            group_split::build_split_full_key(self.boundary_table_id, VirtualNode::ZERO);
55        let (table_ids_left, table_ids_right) =
56            group_split::split_table_ids_with_table_id_and_vnode(
57                &self.parent_table_ids,
58                split_full_key.user_key.table_id,
59                split_full_key.user_key.get_vnode_id(),
60            );
61        assert!(!table_ids_left.is_empty() && !table_ids_right.is_empty());
62        (table_ids_left, table_ids_right)
63    }
64}
65
66fn gen_normalize_plan(
67    left: &CompactionGroupStatistic,
68    right: &CompactionGroupStatistic,
69) -> Option<NormalizePlan> {
70    let left_table_ids = left.table_statistic.keys().copied().collect_vec();
71
72    if left_table_ids.len() <= 1 {
73        return None;
74    }
75
76    let left_max = *left_table_ids.last().unwrap();
77    let right_min = *right.table_statistic.keys().next().unwrap();
78    if left_max < right_min {
79        return None;
80    }
81
82    let boundary_index = left_table_ids.partition_point(|&table_id| table_id < right_min);
83    if boundary_index == 0 || boundary_index >= left_table_ids.len() {
84        return None;
85    }
86    let boundary_table_id = left_table_ids[boundary_index];
87
88    Some(NormalizePlan {
89        parent_group_id: left.group_id,
90        parent_table_ids: left_table_ids,
91        boundary_table_id,
92    })
93}
94
95fn build_normalize_plan_from_group_statistics(
96    groups: &[CompactionGroupStatistic],
97) -> Option<NormalizePlan> {
98    // `calculate_compaction_group_statistic()` iterates all version levels, so newly created or
99    // transiently empty groups can appear here without any member tables.
100    let mut groups = groups
101        .iter()
102        .filter(|group| !group.table_statistic.is_empty())
103        .collect_vec();
104    groups.sort_by_key(|group| *group.table_statistic.keys().next().unwrap());
105
106    groups
107        .split(|group| {
108            group
109                .compaction_group_config
110                .compaction_config
111                .disable_auto_group_scheduling
112                .unwrap_or(false)
113        })
114        .find_map(|segment| {
115            segment
116                .windows(2)
117                .find_map(|pair| gen_normalize_plan(pair[0], pair[1]))
118        })
119}
120
121fn collect_normalize_group_statistics(
122    version: &HummockVersion,
123    compaction_group_manager: &CompactionGroupManager,
124) -> Result<Vec<CompactionGroupStatistic>> {
125    let mut groups = vec![];
126    for group_id in version.levels.keys() {
127        let table_ids = version
128            .state_table_info
129            .compaction_group_member_table_ids(*group_id)
130            .iter()
131            .copied()
132            .collect_vec();
133        if table_ids.is_empty() {
134            continue;
135        }
136
137        let group_config = compaction_group_manager
138            .try_get_compaction_group_config(*group_id)
139            .ok_or_else(|| {
140                Error::CompactionGroup(format!(
141                    "group {} config not found during normalize",
142                    group_id
143                ))
144            })?;
145        groups.push(CompactionGroupStatistic {
146            group_id: *group_id,
147            group_size: 0,
148            table_statistic: table_ids
149                .into_iter()
150                .map(|table_id| (table_id, 0))
151                .collect(),
152            compaction_group_config: group_config,
153        });
154    }
155    Ok(groups)
156}
157
158impl HummockManager {
159    async fn build_normalize_plan(&self) -> Option<NormalizePlan> {
160        let groups = self.calculate_compaction_group_statistic().await;
161        build_normalize_plan_from_group_statistics(&groups)
162    }
163
164    async fn apply_normalize_plan(&self, plan: &NormalizePlan) -> Result<bool> {
165        let (table_ids_right, boundary_table_id, new_compaction_group_id) = {
166            let mut versioning_guard = self
167                .versioning
168                .write_with_process_name("apply_normalize_plan")
169                .await;
170            let versioning = versioning_guard.deref_mut();
171            let mut compaction_group_manager = self
172                .compaction_group_manager
173                .write_with_process_name("apply_normalize_plan")
174                .await;
175
176            let groups = collect_normalize_group_statistics(
177                &versioning.current_version,
178                &compaction_group_manager,
179            )?;
180            let Some(current_plan) = build_normalize_plan_from_group_statistics(&groups) else {
181                return Ok(false);
182            };
183
184            if &current_plan != plan {
185                return Ok(false);
186            }
187
188            let (_table_ids_left, table_ids_right) = plan.split_table_ids();
189
190            let config = compaction_group_manager
191                .try_get_compaction_group_config(plan.parent_group_id)
192                .ok_or_else(|| {
193                    Error::CompactionGroup(format!(
194                        "parent group {} config not found",
195                        plan.parent_group_id
196                    ))
197                })?
198                .compaction_config()
199                .as_ref()
200                .clone();
201
202            let mut compaction_groups_txn = compaction_group_manager.start_compaction_groups_txn();
203            let mut version = HummockVersionTransaction::new(
204                &mut versioning.current_version,
205                &mut versioning.hummock_version_deltas,
206                &mut versioning.table_change_log,
207                self.env.notification_manager(),
208                None,
209                &self.metrics,
210                &self.env.opts,
211                &self.version_stat_tx,
212            );
213            let mut new_version_delta = version.new_delta();
214            let split_key = plan.split_key();
215            let split_sst_count = new_version_delta
216                .latest_version()
217                .count_new_ssts_in_group_split(plan.parent_group_id, split_key.clone());
218            let new_sst_start_id = next_sstable_id(&self.env, split_sst_count).await?;
219            let new_compaction_group_id = next_compaction_group_id(&self.env).await?;
220
221            #[expect(deprecated)]
222            new_version_delta.group_deltas.insert(
223                new_compaction_group_id,
224                GroupDeltas {
225                    group_deltas: vec![GroupDelta::GroupConstruct(Box::new(PbGroupConstruct {
226                        group_config: Some(config.clone()),
227                        group_id: new_compaction_group_id,
228                        parent_group_id: plan.parent_group_id,
229                        new_sst_start_id,
230                        table_ids: vec![],
231                        version: CompatibilityVersion::LATEST as _,
232                        split_key: Some(split_key.into()),
233                    }))],
234                },
235            );
236
237            new_version_delta.with_latest_version(|version, new_version_delta| {
238                for &table_id in &table_ids_right {
239                    let info = version
240                        .state_table_info
241                        .info()
242                        .get(&table_id)
243                        .expect("table should exist before normalize split");
244                    assert!(
245                        new_version_delta
246                            .state_table_info_delta
247                            .insert(
248                                table_id,
249                                PbStateTableInfoDelta {
250                                    committed_epoch: info.committed_epoch,
251                                    compaction_group_id: new_compaction_group_id,
252                                }
253                            )
254                            .is_none()
255                    );
256                }
257            });
258            new_version_delta.pre_apply();
259            compaction_groups_txn
260                .create_compaction_groups(new_compaction_group_id, Arc::new(config));
261
262            commit_multi_var!(self.meta_store_ref(), version, compaction_groups_txn)?;
263            versioning.mark_next_time_travel_version_snapshot();
264            if !self.env.opts.compaction_deterministic_test {
265                for group_id in [plan.parent_group_id, new_compaction_group_id] {
266                    self.try_send_compaction_request(group_id, TaskType::Dynamic);
267                }
268            }
269
270            (
271                table_ids_right,
272                plan.boundary_table_id,
273                new_compaction_group_id,
274            )
275        };
276
277        self.cancel_expired_normalize_split_tasks(plan.parent_group_id)
278            .await?;
279        self.try_update_write_limits(&[plan.parent_group_id, new_compaction_group_id])
280            .await;
281        self.metrics
282            .split_compaction_group_count
283            .with_label_values(&[&plan.parent_group_id.to_string()])
284            .inc();
285        tracing::info!(target: super::TRACE_TARGET,
286            "normalize split success: parent_group={} boundary_table_id={} moved_tables={:?} new_group_id={}",
287            plan.parent_group_id,
288            boundary_table_id,
289            table_ids_right,
290            new_compaction_group_id
291        );
292
293        Ok(true)
294    }
295
296    async fn cancel_expired_normalize_split_tasks(
297        &self,
298        parent_group_id: CompactionGroupId,
299    ) -> Result<()> {
300        let mut canceled_tasks = vec![];
301        let compaction_guard = self
302            .compaction
303            .write_with_process_name("cancel_expired_normalize_split_tasks")
304            .await;
305        let mut versioning_guard = self
306            .versioning
307            .write_with_process_name("cancel_expired_normalize_split_tasks")
308            .await;
309        let versioning = versioning_guard.deref_mut();
310        let compact_task_assignments =
311            compaction_guard.get_compact_task_assignments_by_group_id(parent_group_id);
312        let Some(levels) = versioning.current_version.levels.get(&parent_group_id) else {
313            return Ok(());
314        };
315        compact_task_assignments
316            .into_iter()
317            .for_each(|task_assignment| {
318                let task = &task_assignment.compact_task;
319                if is_compaction_task_expired(
320                    task.compaction_group_version_id,
321                    levels.compaction_group_version_id,
322                ) {
323                    canceled_tasks.push(ReportTask {
324                        task_id: task.task_id,
325                        task_status: TaskStatus::ManualCanceled,
326                        table_stats_change: HashMap::default(),
327                        sorted_output_ssts: vec![],
328                        object_timestamps: HashMap::default(),
329                    });
330                }
331            });
332        canceled_tasks.sort_by_key(|task| task.task_id);
333        canceled_tasks.dedup_by_key(|task| task.task_id);
334
335        if !canceled_tasks.is_empty() {
336            self.report_compact_tasks_impl(canceled_tasks, compaction_guard, versioning_guard)
337                .await?;
338        }
339
340        Ok(())
341    }
342
343    /// Normalize overlapping adjacent compaction groups by split only.
344    ///
345    /// The algorithm repeatedly scans adjacent groups by `min(table_id)` and if
346    /// `max(left) >= min(right)`, it splits `left` at the first table id `>= min(right)`.
347    /// Each step is planned from a read snapshot, then revalidated and applied with a short write
348    /// transaction.
349    pub async fn normalize_overlapping_compaction_groups(&self) -> Result<usize> {
350        self.normalize_overlapping_compaction_groups_with_limit(usize::MAX)
351            .await
352    }
353
354    pub async fn normalize_overlapping_compaction_groups_with_limit(
355        &self,
356        max_splits: usize,
357    ) -> Result<usize> {
358        let mut split_count = 0usize;
359        while split_count < max_splits {
360            let Some(plan) = self.build_normalize_plan().await else {
361                break;
362            };
363
364            if !self.apply_normalize_plan(&plan).await? {
365                tracing::debug!(target: super::TRACE_TARGET,
366                    parent_group_id = %plan.parent_group_id,
367                    boundary_table_id = %plan.boundary_table_id,
368                    "normalize plan became stale before apply"
369                );
370                break;
371            }
372            split_count += 1;
373        }
374
375        Ok(split_count)
376    }
377}
378
379#[cfg(test)]
380mod tests {
381    use super::super::tests::group;
382    use super::{NormalizePlan, build_normalize_plan_from_group_statistics, gen_normalize_plan};
383
384    #[test]
385    fn test_gen_normalize_plan_returns_none_for_single_table_group() {
386        let left = group(1.into(), &[10], false);
387        let right = group(2.into(), &[5, 20], false);
388
389        assert_eq!(None, gen_normalize_plan(&left, &right));
390    }
391
392    #[test]
393    fn test_gen_normalize_plan_returns_none_for_non_overlapping_groups() {
394        let left = group(1.into(), &[1, 2, 3], false);
395        let right = group(2.into(), &[4, 5, 6], false);
396
397        assert_eq!(None, gen_normalize_plan(&left, &right));
398    }
399
400    #[test]
401    fn test_gen_normalize_plan_returns_none_when_boundary_cannot_split_parent() {
402        let left = group(1.into(), &[5, 6, 7], false);
403        let right = group(2.into(), &[4, 8], false);
404
405        assert_eq!(None, gen_normalize_plan(&left, &right));
406    }
407
408    #[test]
409    fn test_gen_normalize_plan_generates_expected_boundary() {
410        let left = group(1.into(), &[1, 4, 7], false);
411        let right = group(2.into(), &[2, 5, 8], false);
412
413        assert_eq!(
414            Some(NormalizePlan {
415                parent_group_id: 1.into(),
416                parent_table_ids: vec![1.into(), 4.into(), 7.into()],
417                boundary_table_id: 4.into(),
418            }),
419            gen_normalize_plan(&left, &right)
420        );
421    }
422
423    #[test]
424    fn test_build_normalize_plan_skips_disabled_boundary_and_continues_later_segment() {
425        let groups = vec![
426            group(1.into(), &[1, 4, 7], false),
427            group(2.into(), &[2, 5, 8], true),
428            group(3.into(), &[10, 13, 16], false),
429            group(4.into(), &[11, 14, 17], false),
430        ];
431
432        assert_eq!(
433            Some(NormalizePlan {
434                parent_group_id: 3.into(),
435                parent_table_ids: vec![10.into(), 13.into(), 16.into()],
436                boundary_table_id: 13.into(),
437            }),
438            build_normalize_plan_from_group_statistics(&groups)
439        );
440    }
441}