Skip to main content

risingwave_meta/hummock/manager/compaction/compaction_group_schedule/
topology.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//! Split and merge transactions. Lock acquisition, commit, candidate publication, and
16//! post-commit cleanup stay together so their ordering is visible at each operation.
17
18use std::collections::{BTreeMap, BTreeSet, HashMap, HashSet};
19use std::ops::DerefMut;
20use std::sync::Arc;
21
22use bytes::Bytes;
23use itertools::Itertools;
24use risingwave_common::catalog::TableId;
25use risingwave_common::hash::VirtualNode;
26use risingwave_hummock_sdk::compact_task::{ReportTask, is_compaction_task_expired};
27use risingwave_hummock_sdk::compaction_group::{StateTableId, group_split};
28use risingwave_hummock_sdk::version::{GroupDelta, GroupDeltas};
29use risingwave_hummock_sdk::{CompactionGroupId, can_concat};
30use risingwave_pb::hummock::compact_task::{TaskStatus, TaskType};
31use risingwave_pb::hummock::rise_ctl_update_compaction_config_request::mutable_config::MutableConfig;
32use risingwave_pb::hummock::{
33    CompatibilityVersion, PbGroupConstruct, PbGroupMerge, PbStateTableInfoDelta,
34};
35use thiserror_ext::AsReport;
36
37use super::{CompactionGroupStatistic, merge_policy};
38use crate::hummock::error::{Error, Result};
39use crate::hummock::manager::transaction::HummockVersionTransaction;
40use crate::hummock::manager::{HummockManager, commit_multi_var};
41use crate::hummock::metrics_utils::remove_compaction_group_metrics;
42use crate::hummock::sequence::{next_compaction_group_id, next_sstable_id};
43
44impl HummockManager {
45    pub async fn merge_compaction_group(
46        &self,
47        group_1: CompactionGroupId,
48        group_2: CompactionGroupId,
49    ) -> Result<()> {
50        self.merge_compaction_group_impl(group_1, group_2, None, false)
51            .await
52            .map(|_| ())
53    }
54
55    pub async fn merge_compaction_group_for_test(
56        &self,
57        group_1: CompactionGroupId,
58        group_2: CompactionGroupId,
59        created_tables: HashSet<TableId>,
60    ) -> Result<()> {
61        self.merge_compaction_group_impl(group_1, group_2, Some(&created_tables), false)
62            .await
63            .map(|_| ())
64    }
65
66    pub async fn merge_compaction_group_impl(
67        &self,
68        group_1: CompactionGroupId,
69        group_2: CompactionGroupId,
70        created_tables: Option<&HashSet<TableId>>,
71        validate_policy: bool,
72    ) -> Result<CompactionGroupStatistic> {
73        // Catalog access can wait for DDL or the database. Do it before taking Hummock write
74        // locks, then conservatively check this snapshot against the current group members.
75        let fetched_created_tables;
76        let created_tables = if let Some(created_tables) = created_tables {
77            // Reuse the batch snapshot; a newly added table is conservatively rejected by
78            // the current-membership check below, without another catalog query per merge.
79            created_tables
80        } else {
81            fetched_created_tables = match self.metadata_manager.get_created_table_ids().await {
82                Ok(created_tables) => HashSet::from_iter(created_tables),
83                Err(err) => {
84                    tracing::warn!(target: super::TRACE_TARGET, error = %err.as_report(), "failed to fetch created table ids");
85                    return Err(Error::CompactionGroup(format!(
86                        "merge group_1 {} group_2 {} failed to fetch created table ids",
87                        group_1, group_2
88                    )));
89                }
90            };
91            &fetched_created_tables
92        };
93        let compaction_guard = self
94            .compaction
95            .write_with_process_name("merge_compaction_group_impl")
96            .await;
97        let mut versioning_guard = self
98            .versioning
99            .write_with_process_name("merge_compaction_group_impl")
100            .await;
101        let versioning = versioning_guard.deref_mut();
102        // Validate parameters.
103        if !versioning.current_version.levels.contains_key(&group_1) {
104            return Err(Error::CompactionGroup(format!("invalid group {}", group_1)));
105        }
106
107        if !versioning.current_version.levels.contains_key(&group_2) {
108            return Err(Error::CompactionGroup(format!("invalid group {}", group_2)));
109        }
110
111        let state_table_info = &versioning.current_version.state_table_info;
112        let mut member_table_ids_1 = state_table_info
113            .compaction_group_member_table_ids(group_1)
114            .iter()
115            .cloned()
116            .collect_vec();
117
118        if member_table_ids_1.is_empty() {
119            return Err(Error::CompactionGroup(format!(
120                "group_1 {} is empty",
121                group_1
122            )));
123        }
124
125        let mut member_table_ids_2 = state_table_info
126            .compaction_group_member_table_ids(group_2)
127            .iter()
128            .cloned()
129            .collect_vec();
130
131        if member_table_ids_2.is_empty() {
132            return Err(Error::CompactionGroup(format!(
133                "group_2 {} is empty",
134                group_2
135            )));
136        }
137
138        debug_assert!(!member_table_ids_1.is_empty());
139        debug_assert!(!member_table_ids_2.is_empty());
140        assert!(member_table_ids_1.is_sorted());
141        assert!(member_table_ids_2.is_sorted());
142
143        // Do not merge a group containing a table that is still being created.
144        if member_table_ids_1
145            .iter()
146            .chain(&member_table_ids_2)
147            .any(|table_id| !created_tables.contains(table_id))
148        {
149            return Err(Error::CompactionGroup(format!(
150                "Cannot merge creating group {} next_group {} member_table_ids_1 {:?} member_table_ids_2 {:?}",
151                group_1, group_2, member_table_ids_1, member_table_ids_2
152            )));
153        }
154
155        // Make sure `member_table_ids_1` is smaller than `member_table_ids_2`
156        let (left_group_id, right_group_id) =
157            if member_table_ids_1.first().unwrap() < member_table_ids_2.first().unwrap() {
158                (group_1, group_2)
159            } else {
160                std::mem::swap(&mut member_table_ids_1, &mut member_table_ids_2);
161                (group_2, group_1)
162            };
163
164        // We can only merge two groups with non-overlapping member table ids.
165        // After the swap above, member_table_ids_1 has the smaller first element.
166        // If the last element of member_table_ids_1 >= the first element of member_table_ids_2,
167        // the two groups' table id ranges overlap and cannot be merged.
168        if member_table_ids_1.last().unwrap() >= member_table_ids_2.first().unwrap() {
169            return Err(Error::CompactionGroup(format!(
170                "invalid merge group_1 {} group_2 {}: table id ranges overlap",
171                left_group_id, right_group_id
172            )));
173        }
174
175        let combined_member_table_ids = member_table_ids_1
176            .iter()
177            .chain(member_table_ids_2.iter())
178            .collect_vec();
179        assert!(combined_member_table_ids.is_sorted());
180
181        // These members and their sizes are stable under versioning. Reuse the same
182        // per-table sizes for the policy check and the returned survivor statistics.
183        let survivor_tables = combined_member_table_ids
184            .iter()
185            .map(|&&table_id| {
186                let size = versioning
187                    .version_stats
188                    .table_stats
189                    .get(&table_id)
190                    .map(|stats| (stats.total_key_size + stats.total_value_size).max(0) as u64)
191                    .unwrap_or(0);
192                (table_id, size)
193            })
194            .collect_vec();
195
196        let mut compaction_group_manager = self
197            .compaction_group_manager
198            .write_with_process_name("merge_compaction_group_impl")
199            .await;
200        if validate_policy {
201            let group = compaction_group_manager
202                .try_get_compaction_group_config(group_1)
203                .ok_or_else(|| Error::CompactionGroup(format!("invalid group {group_1}")))?;
204            let next_group = compaction_group_manager
205                .try_get_compaction_group_config(group_2)
206                .ok_or_else(|| Error::CompactionGroup(format!("invalid group {group_2}")))?;
207            let current_size = survivor_tables.iter().map(|(_, size)| size).sum();
208            // The scheduling snapshot is only a prefilter. Check all mutable policy inputs
209            // under the locks already needed to apply the merge, before scanning SSTs.
210            merge_policy::validate_group_config(&group, &next_group, current_size, &self.env.opts)?;
211            merge_policy::validate_group_levels(
212                &group,
213                &next_group,
214                &self.env.opts,
215                &versioning.current_version,
216            )?;
217            // Commit hands off to statistics before releasing versioning. This read cannot
218            // pair the new version with throughput from before its commit.
219            if !merge_policy::check_is_low_write_throughput(
220                &self.table_write_throughput_statistic_manager.read(),
221                combined_member_table_ids.iter().map(|&&table_id| table_id),
222                &self.env.opts,
223            ) {
224                return Err(Error::CompactionGroup(format!(
225                    "Cannot merge groups {group_1} and {group_2}: current throughput is not observed cold"
226                )));
227            }
228        }
229
230        // check duplicated sst_id
231        let mut sst_id_set = HashSet::new();
232        for sst_id in versioning
233            .current_version
234            .get_sst_ids_by_group_id(left_group_id)
235            .chain(
236                versioning
237                    .current_version
238                    .get_sst_ids_by_group_id(right_group_id),
239            )
240        {
241            if !sst_id_set.insert(sst_id) {
242                return Err(Error::CompactionGroup(format!(
243                    "invalid merge group_1 {} group_2 {} duplicated sst_id {}",
244                    left_group_id, right_group_id, sst_id
245                )));
246            }
247        }
248
249        // check branched sst on non-overlap level
250        {
251            let left_levels = versioning
252                .current_version
253                .get_compaction_group_levels(left_group_id);
254
255            let right_levels = versioning
256                .current_version
257                .get_compaction_group_levels(right_group_id);
258
259            // we can not check the l0 sub level, because the sub level id will be rewritten when merge
260            // This check will ensure that other non-overlapping level ssts can be concat and that the key_range is correct.
261            let max_level = std::cmp::max(left_levels.levels.len(), right_levels.levels.len());
262            for level_idx in 1..=max_level {
263                let left_level = left_levels.get_level(level_idx);
264                let right_level = right_levels.get_level(level_idx);
265                if left_level.table_infos.is_empty() || right_level.table_infos.is_empty() {
266                    continue;
267                }
268
269                let left_last_sst = left_level.table_infos.last().unwrap().clone();
270                let right_first_sst = right_level.table_infos.first().unwrap().clone();
271                let left_sst_id = left_last_sst.sst_id;
272                let right_sst_id = right_first_sst.sst_id;
273                let left_obj_id = left_last_sst.object_id;
274                let right_obj_id = right_first_sst.object_id;
275
276                // Since the sst key_range within a group is legal, we only need to check the ssts adjacent to the two groups.
277                if !can_concat(&[left_last_sst, right_first_sst]) {
278                    return Err(Error::CompactionGroup(format!(
279                        "invalid merge group_1 {} group_2 {} level_idx {} left_last_sst_id {} right_first_sst_id {} left_obj_id {} right_obj_id {}",
280                        left_group_id,
281                        right_group_id,
282                        level_idx,
283                        left_sst_id,
284                        right_sst_id,
285                        left_obj_id,
286                        right_obj_id
287                    )));
288                }
289            }
290        }
291
292        let mut version = HummockVersionTransaction::new(
293            &mut versioning.current_version,
294            &mut versioning.hummock_version_deltas,
295            &mut versioning.table_change_log,
296            self.env.notification_manager(),
297            None,
298            &self.metrics,
299            &self.env.opts,
300            &self.version_stat_tx,
301        );
302        let mut new_version_delta = version.new_delta();
303
304        let target_compaction_group_id = {
305            // merge right_group_id to left_group_id and remove right_group_id
306            new_version_delta.group_deltas.insert(
307                left_group_id,
308                GroupDeltas {
309                    group_deltas: vec![GroupDelta::GroupMerge(PbGroupMerge {
310                        left_group_id,
311                        right_group_id,
312                    })],
313                },
314            );
315            left_group_id
316        };
317
318        // TODO: remove compaciton group_id from state_table_info
319        // rewrite compaction_group_id for all tables
320        new_version_delta.with_latest_version(|version, new_version_delta| {
321            for &table_id in combined_member_table_ids {
322                let info = version
323                    .state_table_info
324                    .info()
325                    .get(&table_id)
326                    .expect("have check exist previously");
327                assert!(
328                    new_version_delta
329                        .state_table_info_delta
330                        .insert(
331                            table_id,
332                            PbStateTableInfoDelta {
333                                committed_epoch: info.committed_epoch,
334                                compaction_group_id: target_compaction_group_id,
335                            }
336                        )
337                        .is_none()
338                );
339            }
340        });
341
342        let survivor_config = {
343            let mut compaction_groups_txn = compaction_group_manager.start_compaction_groups_txn();
344
345            // for metrics reclaim
346            {
347                let right_group_max_level = new_version_delta
348                    .latest_version()
349                    .get_compaction_group_levels(right_group_id)
350                    .levels
351                    .len();
352
353                remove_compaction_group_metrics(
354                    &self.metrics,
355                    right_group_id,
356                    right_group_max_level,
357                );
358            }
359
360            // clear `partition_vnode_count` for the hybrid group
361            {
362                if let Err(err) = compaction_groups_txn.update_compaction_config(
363                    &[left_group_id],
364                    &[MutableConfig::SplitWeightByVnode(0)], // default
365                ) {
366                    tracing::error!(target: super::TRACE_TARGET,
367                        error = %err.as_report(),
368                        "failed to update compaction config for group-{}",
369                        left_group_id
370                    );
371                }
372            }
373
374            new_version_delta.pre_apply();
375
376            // remove right_group_id
377            compaction_groups_txn.remove(right_group_id);
378            commit_multi_var!(self.meta_store_ref(), version, compaction_groups_txn)?;
379            compaction_group_manager
380                .try_get_compaction_group_config(left_group_id)
381                .expect("merged group config should exist")
382        };
383
384        drop(compaction_group_manager);
385
386        // Update candidates only after commit, while versioning still protects group membership.
387        self.compaction_state
388            .remove_compaction_group(right_group_id);
389        if !self.env.opts.compaction_deterministic_test {
390            self.try_send_compaction_request(left_group_id, TaskType::Dynamic);
391        }
392
393        // Instead of handling DeltaType::GroupConstruct for time travel, simply enforce a version snapshot.
394        versioning.mark_next_time_travel_version_snapshot();
395
396        // cancel tasks
397        let mut canceled_tasks = vec![];
398        // after merge, all tasks in right_group_id should be canceled
399        // Failure of cancel does not cause correctness problems, the report task will have better interception, and the operation here is designed to free up compactor compute resources more quickly.
400        let compact_task_assignments =
401            compaction_guard.get_compact_task_assignments_by_group_id(right_group_id);
402        compact_task_assignments
403            .into_iter()
404            .for_each(|task_assignment| {
405                let task = &task_assignment.compact_task;
406                assert_eq!(task.compaction_group_id, right_group_id);
407                canceled_tasks.push(ReportTask {
408                    task_id: task.task_id,
409                    task_status: TaskStatus::ManualCanceled,
410                    table_stats_change: HashMap::default(),
411                    sorted_output_ssts: vec![],
412                    object_timestamps: HashMap::default(),
413                });
414            });
415
416        if !canceled_tasks.is_empty() {
417            self.report_compact_tasks_impl(canceled_tasks, compaction_guard, versioning_guard)
418                .await?;
419        } else {
420            drop(versioning_guard);
421            drop(compaction_guard);
422        }
423
424        self.try_update_write_limits(&[left_group_id, right_group_id])
425            .await;
426
427        // Count successful merges here for both manual and automatic callers.
428        self.metrics
429            .merge_compaction_group_count
430            .with_label_values(&[&left_group_id.to_string()])
431            .inc();
432
433        // This is the actual survivor at the topology commit, including its current members and
434        // config. Build its map after releasing the write locks, without rescanning other groups.
435        Ok(CompactionGroupStatistic {
436            group_id: left_group_id,
437            group_size: survivor_tables.iter().map(|(_, size)| size).sum(),
438            table_statistic: survivor_tables.into_iter().collect(),
439            compaction_group_config: survivor_config,
440        })
441    }
442}
443
444impl HummockManager {
445    /// Split `table_ids` to a dedicated compaction group.(will be split by the `table_id` and `vnode`.)
446    /// Returns the compaction group id containing the `table_ids` and the mapping of compaction group id to table ids.
447    /// The split will follow the following rules
448    /// 1. ssts with `key_range.left` greater than `split_key` will be split to the right group
449    /// 2. the sst containing `split_key` will be split into two separate ssts and their `key_range` will be changed `sst_1`: [`sst.key_range.left`, `split_key`) `sst_2`: [`split_key`, `sst.key_range.right`]
450    /// 3. currently only `vnode` 0 and `vnode` max is supported. (Due to the above rule, vnode max will be rewritten as `table_id` + 1, `vnode` 0)
451    ///   - `parent_group_id`: the `group_id` to split
452    ///   - `split_table_ids`: the `table_ids` to split, now we still support to split multiple tables to one group at once, pass `split_table_ids` for per `split` operation for checking
453    ///   - `table_id_to_split`: the `table_id` to split
454    ///   - `vnode_to_split`: the `vnode` to split
455    ///   - `partition_vnode_count`: the partition count for the single table group if need
456    async fn split_compaction_group_impl(
457        &self,
458        parent_group_id: CompactionGroupId,
459        split_table_ids: &[StateTableId],
460        table_id_to_split: StateTableId,
461        vnode_to_split: VirtualNode,
462        partition_vnode_count: Option<u32>,
463        validate_policy: bool,
464    ) -> Result<Vec<(CompactionGroupId, Vec<StateTableId>)>> {
465        let mut result = vec![];
466        let compaction_guard = self
467            .compaction
468            .write_with_process_name("split_compaction_group_impl")
469            .await;
470        let mut versioning_guard = self
471            .versioning
472            .write_with_process_name("split_compaction_group_impl")
473            .await;
474        let versioning = versioning_guard.deref_mut();
475        // Validate parameters.
476        if !versioning
477            .current_version
478            .levels
479            .contains_key(&parent_group_id)
480        {
481            return Err(Error::CompactionGroup(format!(
482                "invalid group {}",
483                parent_group_id
484            )));
485        }
486
487        let member_table_ids = versioning
488            .current_version
489            .state_table_info
490            .compaction_group_member_table_ids(parent_group_id)
491            .iter()
492            .copied()
493            .collect::<BTreeSet<_>>();
494
495        if !member_table_ids.contains(&table_id_to_split) {
496            return Err(Error::CompactionGroup(format!(
497                "table {} doesn't in group {}",
498                table_id_to_split, parent_group_id
499            )));
500        }
501
502        let split_full_key = group_split::build_split_full_key(table_id_to_split, vnode_to_split);
503
504        // change to vec for partition
505        let table_ids = member_table_ids.into_iter().collect_vec();
506        if table_ids == split_table_ids {
507            return Err(Error::CompactionGroup(format!(
508                "invalid split attempt for group {}: all member tables are moved",
509                parent_group_id
510            )));
511        }
512        // avoid decode split_key when caller is aware of the table_id and vnode
513        let (table_ids_left, table_ids_right) =
514            group_split::split_table_ids_with_table_id_and_vnode(
515                &table_ids,
516                split_full_key.user_key.table_id,
517                split_full_key.user_key.get_vnode_id(),
518            );
519        if table_ids_left.is_empty() || table_ids_right.is_empty() {
520            // not need to split group if all tables are in the same side
521            if !table_ids_left.is_empty() {
522                result.push((parent_group_id, table_ids_left));
523            }
524
525            if !table_ids_right.is_empty() {
526                result.push((parent_group_id, table_ids_right));
527            }
528            return Ok(result);
529        }
530
531        result.push((parent_group_id, table_ids_left));
532
533        let split_key: Bytes = split_full_key.encode().into();
534
535        let mut version = HummockVersionTransaction::new(
536            &mut versioning.current_version,
537            &mut versioning.hummock_version_deltas,
538            &mut versioning.table_change_log,
539            self.env.notification_manager(),
540            None,
541            &self.metrics,
542            &self.env.opts,
543            &self.version_stat_tx,
544        );
545        let mut new_version_delta = version.new_delta();
546
547        let split_sst_count = new_version_delta
548            .latest_version()
549            .count_new_ssts_in_group_split(parent_group_id, split_key.clone());
550
551        let new_sst_start_id = next_sstable_id(&self.env, split_sst_count).await?;
552        let (new_compaction_group_id, config) = {
553            // All NewCompactionGroup pairs are mapped to one new compaction group.
554            let new_compaction_group_id = next_compaction_group_id(&self.env).await?;
555            // Inherit config from parent group
556            let config = self
557                .compaction_group_manager
558                .read_with_process_name("split_compaction_group_impl")
559                .await
560                .try_get_compaction_group_config(parent_group_id)
561                .ok_or_else(|| {
562                    Error::CompactionGroup(format!(
563                        "parent group {} config not found",
564                        parent_group_id
565                    ))
566                })?
567                .compaction_config()
568                .as_ref()
569                .clone();
570
571            #[expect(deprecated)]
572            // fill the deprecated field with default value
573            new_version_delta.group_deltas.insert(
574                new_compaction_group_id,
575                GroupDeltas {
576                    group_deltas: vec![GroupDelta::GroupConstruct(Box::new(PbGroupConstruct {
577                        group_config: Some(config.clone()),
578                        group_id: new_compaction_group_id,
579                        parent_group_id,
580                        new_sst_start_id,
581                        table_ids: vec![],
582                        version: CompatibilityVersion::LATEST as _, // for compatibility
583                        split_key: Some(split_key.into()),
584                    }))],
585                },
586            );
587            (new_compaction_group_id, config)
588        };
589
590        new_version_delta.with_latest_version(|version, new_version_delta| {
591            for &table_id in &table_ids_right {
592                let info = version
593                    .state_table_info
594                    .info()
595                    .get(&table_id)
596                    .expect("have check exist previously");
597                assert!(
598                    new_version_delta
599                        .state_table_info_delta
600                        .insert(
601                            table_id,
602                            PbStateTableInfoDelta {
603                                committed_epoch: info.committed_epoch,
604                                compaction_group_id: new_compaction_group_id,
605                            }
606                        )
607                        .is_none()
608                );
609            }
610        });
611
612        result.push((new_compaction_group_id, table_ids_right));
613
614        {
615            let mut compaction_group_manager = self
616                .compaction_group_manager
617                .write_with_process_name("split_compaction_group_impl")
618                .await;
619            // Check automatic policy at each committed step, while config changes are
620            // excluded by the same write guard used to create the child group.
621            if validate_policy
622                && compaction_group_manager
623                    .try_get_compaction_group_config(parent_group_id)
624                    .expect("current parent config should exist")
625                    .compaction_config
626                    .disable_auto_group_scheduling
627                    .unwrap_or(false)
628            {
629                return Err(Error::CompactionGroup(format!(
630                    "group {parent_group_id} disables automatic splitting"
631                )));
632            }
633            let mut compaction_groups_txn = compaction_group_manager.start_compaction_groups_txn();
634            compaction_groups_txn
635                .create_compaction_groups(new_compaction_group_id, Arc::new(config));
636
637            // check if need to update the compaction config for the single table group and guarantee the operation atomicity
638            // `partition_vnode_count` only works inside a table, to avoid a lot of slicing sst, we only enable it in groups with high throughput and only one table.
639            // The target `table_ids` might be split to an existing group, so we need to try to update its config
640            for (cg_id, table_ids) in &result {
641                // check the split_tables had been place to the dedicated compaction group
642                if let Some(partition_vnode_count) = partition_vnode_count
643                    && table_ids.len() == 1
644                    && table_ids == split_table_ids
645                    && let Err(err) = compaction_groups_txn.update_compaction_config(
646                        &[*cg_id],
647                        &[MutableConfig::SplitWeightByVnode(partition_vnode_count)],
648                    )
649                {
650                    tracing::error!(target: super::TRACE_TARGET,
651                        error = %err.as_report(),
652                        "failed to update compaction config for group-{}",
653                        cg_id
654                    );
655                }
656            }
657
658            new_version_delta.pre_apply();
659            commit_multi_var!(self.meta_store_ref(), version, compaction_groups_txn)?;
660        }
661        if !self.env.opts.compaction_deterministic_test {
662            for (group_id, _) in &result {
663                self.try_send_compaction_request(*group_id, TaskType::Dynamic);
664            }
665        }
666        // Instead of handling DeltaType::GroupConstruct for time travel, simply enforce a version snapshot.
667        versioning.mark_next_time_travel_version_snapshot();
668
669        // The expired compact tasks will be canceled.
670        // Failure of cancel does not cause correctness problems, the report task will have better interception, and the operation here is designed to free up compactor compute resources more quickly.
671        let mut canceled_tasks = vec![];
672        let compact_task_assignments =
673            compaction_guard.get_compact_task_assignments_by_group_id(parent_group_id);
674        let levels = versioning
675            .current_version
676            .get_compaction_group_levels(parent_group_id);
677        compact_task_assignments
678            .into_iter()
679            .for_each(|task_assignment| {
680                let task = &task_assignment.compact_task;
681                let is_expired = is_compaction_task_expired(
682                    task.compaction_group_version_id,
683                    levels.compaction_group_version_id,
684                );
685                if is_expired {
686                    canceled_tasks.push(ReportTask {
687                        task_id: task.task_id,
688                        task_status: TaskStatus::ManualCanceled,
689                        table_stats_change: HashMap::default(),
690                        sorted_output_ssts: vec![],
691                        object_timestamps: HashMap::default(),
692                    });
693                }
694            });
695
696        if !canceled_tasks.is_empty() {
697            self.report_compact_tasks_impl(canceled_tasks, compaction_guard, versioning_guard)
698                .await?;
699        } else {
700            drop(versioning_guard);
701            drop(compaction_guard);
702        }
703
704        let affected_group_ids = result.iter().map(|(cg_id, _)| *cg_id).collect_vec();
705        self.try_update_write_limits(&affected_group_ids).await;
706
707        self.metrics
708            .split_compaction_group_count
709            .with_label_values(&[&parent_group_id.to_string()])
710            .inc();
711
712        Ok(result)
713    }
714
715    /// Split `table_ids` to a dedicated compaction group.
716    /// Returns the compaction group id containing the `table_ids` and the mapping of compaction group id to table ids.
717    pub async fn move_state_tables_to_dedicated_compaction_group(
718        &self,
719        parent_group_id: CompactionGroupId,
720        table_ids: &[StateTableId],
721        partition_vnode_count: Option<u32>,
722    ) -> Result<(
723        CompactionGroupId,
724        BTreeMap<CompactionGroupId, Vec<StateTableId>>,
725    )> {
726        self.move_state_tables_to_dedicated_compaction_group_impl(
727            parent_group_id,
728            table_ids,
729            partition_vnode_count,
730            false,
731        )
732        .await
733    }
734
735    pub(super) async fn move_state_tables_to_dedicated_compaction_group_impl(
736        &self,
737        parent_group_id: CompactionGroupId,
738        table_ids: &[StateTableId],
739        partition_vnode_count: Option<u32>,
740        validate_policy: bool,
741    ) -> Result<(
742        CompactionGroupId,
743        BTreeMap<CompactionGroupId, Vec<StateTableId>>,
744    )> {
745        if table_ids.is_empty() {
746            return Err(Error::CompactionGroup(
747                "table_ids must not be empty".to_owned(),
748            ));
749        }
750
751        if !table_ids.is_sorted() {
752            return Err(Error::CompactionGroup(
753                "table_ids must be sorted".to_owned(),
754            ));
755        }
756
757        let parent_table_ids = {
758            let versioning_guard = self
759                .versioning
760                .read_with_process_name("move_state_tables_to_dedicated_compaction_group")
761                .await;
762            versioning_guard
763                .current_version
764                .state_table_info
765                .compaction_group_member_table_ids(parent_group_id)
766                .iter()
767                .copied()
768                .collect_vec()
769        };
770
771        if parent_table_ids == table_ids {
772            return Err(Error::CompactionGroup(format!(
773                "invalid split attempt for group {}: all member tables are moved",
774                parent_group_id
775            )));
776        }
777
778        fn check_table_ids_valid(cg_id_to_table_ids: &BTreeMap<CompactionGroupId, Vec<TableId>>) {
779            // 1. table_ids in different cg are sorted.
780            {
781                cg_id_to_table_ids
782                    .iter()
783                    .for_each(|(_cg_id, table_ids)| assert!(table_ids.is_sorted()));
784            }
785
786            // 2.table_ids in different cg are non-overlapping
787            {
788                let mut table_table_ids_vec = cg_id_to_table_ids.values().cloned().collect_vec();
789                table_table_ids_vec.sort_by(|a, b| a[0].cmp(&b[0]));
790                assert!(table_table_ids_vec.concat().is_sorted());
791            }
792
793            // 3.table_ids belong to one and only one cg.
794            {
795                let mut all_table_ids = HashSet::new();
796                for table_ids in cg_id_to_table_ids.values() {
797                    for table_id in table_ids {
798                        assert!(all_table_ids.insert(*table_id));
799                    }
800                }
801            }
802        }
803
804        // move [3,4,5,6]
805        // [1,2,3,4,5,6,7,8,9,10] -> [1,2] [3,4,5,6] [7,8,9,10]
806        // split key
807        // 1. table_id = 3, vnode = 0, epoch = MAX
808        // 2. table_id = 7, vnode = 0, epoch = MAX
809
810        // The new compaction group id is always generate on the right side
811        // Hence, we return the first compaction group id as the result
812        // split 1
813        let mut cg_id_to_table_ids: BTreeMap<CompactionGroupId, Vec<TableId>> = BTreeMap::new();
814        let table_id_to_split = *table_ids.first().unwrap();
815        let mut target_compaction_group_id: CompactionGroupId = 0.into();
816        let result_vec = self
817            .split_compaction_group_impl(
818                parent_group_id,
819                table_ids,
820                table_id_to_split,
821                VirtualNode::ZERO,
822                partition_vnode_count,
823                validate_policy,
824            )
825            .await?;
826        assert!(result_vec.len() <= 2);
827
828        let mut finish_move = false;
829        for (cg_id, table_ids_after_split) in result_vec {
830            if table_ids_after_split.contains(&table_id_to_split) {
831                target_compaction_group_id = cg_id;
832            }
833
834            if table_ids_after_split == table_ids {
835                finish_move = true;
836            }
837
838            cg_id_to_table_ids.insert(cg_id, table_ids_after_split);
839        }
840        check_table_ids_valid(&cg_id_to_table_ids);
841
842        if finish_move {
843            return Ok((target_compaction_group_id, cg_id_to_table_ids));
844        }
845
846        // split 2
847        // See the example above and the split rule in `split_compaction_group_impl`.
848        let table_id_to_split = *table_ids.last().unwrap();
849        let result_vec = self
850            .split_compaction_group_impl(
851                target_compaction_group_id,
852                table_ids,
853                table_id_to_split,
854                VirtualNode::MAX_REPRESENTABLE,
855                partition_vnode_count,
856                validate_policy,
857            )
858            .await?;
859        assert!(result_vec.len() <= 2);
860        for (cg_id, table_ids_after_split) in result_vec {
861            if table_ids_after_split.contains(&table_id_to_split) {
862                target_compaction_group_id = cg_id;
863            }
864            cg_id_to_table_ids.insert(cg_id, table_ids_after_split);
865        }
866        check_table_ids_valid(&cg_id_to_table_ids);
867
868        Ok((target_compaction_group_id, cg_id_to_table_ids))
869    }
870}