1use 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 let fetched_created_tables;
76 let created_tables = if let Some(created_tables) = created_tables {
77 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 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 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 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 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 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 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 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 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 {
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 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 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 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 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 {
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 {
362 if let Err(err) = compaction_groups_txn.update_compaction_config(
363 &[left_group_id],
364 &[MutableConfig::SplitWeightByVnode(0)], ) {
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 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 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 versioning.mark_next_time_travel_version_snapshot();
395
396 let mut canceled_tasks = vec![];
398 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 self.metrics
429 .merge_compaction_group_count
430 .with_label_values(&[&left_group_id.to_string()])
431 .inc();
432
433 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 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 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 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 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 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 let new_compaction_group_id = next_compaction_group_id(&self.env).await?;
555 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 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 _, 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 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 for (cg_id, table_ids) in &result {
641 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 versioning.mark_next_time_travel_version_snapshot();
668
669 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 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 {
781 cg_id_to_table_ids
782 .iter()
783 .for_each(|(_cg_id, table_ids)| assert!(table_ids.is_sorted()));
784 }
785
786 {
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 {
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 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 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}