1use 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 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 ¤t_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 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}