1use std::cmp::Reverse;
16use std::collections::{HashMap, HashSet};
17use std::future::Future;
18use std::sync::Arc;
19use std::time::Duration;
20
21use futures::stream::{BoxStream, FuturesUnordered};
22use futures::{StreamExt, pin_mut};
23use itertools::Itertools;
24use risingwave_hummock_sdk::CompactionGroupId;
25use risingwave_pb::hummock::compact_task::{self, TaskStatus};
26use risingwave_pb::hummock::level_handler::RunningCompactTask;
27use rw_futures_util::select_all;
28use thiserror_ext::AsReport;
29use tokio::sync::oneshot::{Receiver, Sender};
30use tokio::sync::watch;
31use tokio::task::JoinHandle;
32use tokio_stream::wrappers::IntervalStream;
33use tracing::warn;
34
35use crate::backup_restore::BackupManagerRef;
36use crate::hummock::metrics_utils::{trigger_lsm_stat, trigger_mv_stat};
37use crate::hummock::{HummockManager, TASK_NORMAL};
38
39#[derive(Clone, Copy)]
40enum SchedulingEvent {
41 DynamicCompaction,
42 SpaceReclaimCompaction,
43 TtlCompaction,
44 TombstoneCompaction,
45 GroupSplit,
46 GroupMerge,
47}
48
49#[derive(Clone, Copy)]
50enum MaintenanceEvent {
51 CheckDeadTask,
52 FullGc,
53}
54
55fn interval_stream<E>(period: Duration, event: E) -> BoxStream<'static, E>
56where
57 E: Clone + Send + 'static,
58{
59 let mut interval = tokio::time::interval(period);
60 interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
61 interval.reset();
62 Box::pin(IntervalStream::new(interval).map(move |_| event.clone()))
63}
64
65fn spawn_periodic_loop<F, Fut>(
70 name: &'static str,
71 period: Duration,
72 shutdown_rx: watch::Receiver<bool>,
73 mut handler: F,
74) -> JoinHandle<()>
75where
76 F: FnMut() -> Fut + Send + 'static,
77 Fut: Future<Output = ()> + Send + 'static,
78{
79 spawn_serial_timer_lane(
80 name,
81 vec![interval_stream(period, ())],
82 shutdown_rx,
83 move |()| handler(),
84 )
85}
86
87fn spawn_serial_timer_lane<E, F, Fut>(
89 name: &'static str,
90 triggers: Vec<BoxStream<'static, E>>,
91 mut shutdown_rx: watch::Receiver<bool>,
92 mut handler: F,
93) -> JoinHandle<()>
94where
95 E: Send + 'static,
96 F: FnMut(E) -> Fut + Send + 'static,
97 Fut: Future<Output = ()> + Send + 'static,
98{
99 tokio::spawn(async move {
100 if triggers.is_empty() {
101 tracing::info!(
102 "Hummock timer handler loop [{}] is disabled (no intervals configured)",
103 name
104 );
105 return;
106 }
107
108 let event_stream = select_all(triggers);
109 pin_mut!(event_stream);
110
111 loop {
112 tokio::select! {
113 maybe_event = event_stream.next() => {
114 let Some(event) = maybe_event else {
115 tracing::warn!(
116 "Hummock timer handler loop [{}] event stream ended unexpectedly",
117 name
118 );
119 break;
120 };
121 handler(event).await;
122 }
123 changed = shutdown_rx.changed() => {
124 if changed.is_err() || *shutdown_rx.borrow() {
125 tracing::info!("Hummock timer handler loop [{}] is stopped", name);
126 break;
127 }
128 }
129 }
130 }
131 })
132}
133
134fn spawn_scheduling_loop(
135 hummock_manager: Arc<HummockManager>,
136 split_interval_sec: u64,
137 merge_interval_sec: u64,
138 shutdown_rx: watch::Receiver<bool>,
139) -> JoinHandle<()> {
140 let mut triggers = vec![
141 interval_stream(
142 Duration::from_secs(hummock_manager.env.opts.periodic_compaction_interval_sec),
143 SchedulingEvent::DynamicCompaction,
144 ),
145 interval_stream(
146 Duration::from_secs(
147 hummock_manager
148 .env
149 .opts
150 .periodic_space_reclaim_compaction_interval_sec,
151 ),
152 SchedulingEvent::SpaceReclaimCompaction,
153 ),
154 interval_stream(
155 Duration::from_secs(
156 hummock_manager
157 .env
158 .opts
159 .periodic_ttl_reclaim_compaction_interval_sec,
160 ),
161 SchedulingEvent::TtlCompaction,
162 ),
163 interval_stream(
164 Duration::from_secs(
165 hummock_manager
166 .env
167 .opts
168 .periodic_tombstone_reclaim_compaction_interval_sec,
169 ),
170 SchedulingEvent::TombstoneCompaction,
171 ),
172 ];
173
174 if split_interval_sec > 0 {
175 triggers.push(interval_stream(
176 Duration::from_secs(split_interval_sec),
177 SchedulingEvent::GroupSplit,
178 ));
179 }
180 if merge_interval_sec > 0 {
181 triggers.push(interval_stream(
182 Duration::from_secs(merge_interval_sec),
183 SchedulingEvent::GroupMerge,
184 ));
185 }
186
187 spawn_serial_timer_lane("scheduling", triggers, shutdown_rx, move |event| {
188 let hummock_manager = hummock_manager.clone();
189 async move {
190 if hummock_manager.env.opts.compaction_deterministic_test {
191 return;
192 }
193
194 match event {
195 SchedulingEvent::DynamicCompaction => {
196 hummock_manager
197 .on_handle_trigger_multi_group(compact_task::TaskType::Dynamic)
198 .await;
199 }
200 SchedulingEvent::SpaceReclaimCompaction => {
201 hummock_manager
202 .on_handle_trigger_multi_group(compact_task::TaskType::SpaceReclaim)
203 .await;
204 hummock_manager
206 .on_handle_trigger_multi_group(compact_task::TaskType::VnodeWatermark)
207 .await;
208 }
209 SchedulingEvent::TtlCompaction => {
210 hummock_manager
211 .on_handle_trigger_multi_group(compact_task::TaskType::Ttl)
212 .await;
213 }
214 SchedulingEvent::TombstoneCompaction => {
215 hummock_manager
216 .on_handle_trigger_multi_group(compact_task::TaskType::Tombstone)
217 .await;
218 }
219 SchedulingEvent::GroupSplit => {
220 hummock_manager.on_handle_schedule_group_split().await;
221 }
222 SchedulingEvent::GroupMerge => {
223 hummock_manager.on_handle_schedule_group_merge().await;
224 }
225 }
226 }
227 })
228}
229
230fn spawn_maintenance_loop(
231 hummock_manager: Arc<HummockManager>,
232 backup_manager: Option<BackupManagerRef>,
233 shutdown_rx: watch::Receiver<bool>,
234) -> JoinHandle<()> {
235 const CHECK_PENDING_TASK_PERIOD_SEC: u64 = 300;
236
237 let triggers = vec![
238 interval_stream(
239 Duration::from_secs(CHECK_PENDING_TASK_PERIOD_SEC),
240 MaintenanceEvent::CheckDeadTask,
241 ),
242 interval_stream(
243 Duration::from_secs(hummock_manager.env.opts.full_gc_interval_sec),
244 MaintenanceEvent::FullGc,
245 ),
246 ];
247
248 spawn_serial_timer_lane("maintenance", triggers, shutdown_rx, move |event| {
249 let hummock_manager = hummock_manager.clone();
250 let backup_manager = backup_manager.clone();
251 async move {
252 match event {
253 MaintenanceEvent::CheckDeadTask => {
254 if hummock_manager.env.opts.compaction_deterministic_test {
255 return;
256 }
257 hummock_manager.check_dead_task().await;
258 }
259 MaintenanceEvent::FullGc => {
260 let retention_sec = hummock_manager.env.opts.min_sst_retention_time_sec;
261 tokio::task::spawn(async move {
262 let _ = hummock_manager
263 .start_full_gc(Duration::from_secs(retention_sec), None, backup_manager)
264 .await
265 .inspect_err(|e| warn!(error = %e.as_report(), "Failed to start GC."));
266 });
267 }
268 }
269 }
270 })
271}
272
273async fn supervise_child_loops(
274 mut shutdown_rx: Receiver<()>,
275 child_shutdown_tx: watch::Sender<bool>,
276 child_handles: Vec<(&'static str, JoinHandle<()>)>,
277) {
278 let mut child_handles = child_handles
279 .into_iter()
280 .map(|(name, handle)| async move { (name, handle.await) })
281 .collect::<FuturesUnordered<_>>();
282
283 tokio::select! {
286 _ = &mut shutdown_rx => {
287 let _ = child_shutdown_tx.send(true);
288 while let Some((name, result)) = child_handles.next().await {
289 if let Err(e) = result {
290 warn!(
291 handler = name,
292 error = %e.as_report(),
293 "Hummock timer handler loop failed during shutdown"
294 );
295 }
296 }
297 }
298 Some((name, result)) = child_handles.next() => {
299 match result {
300 Ok(()) => {
301 panic!("Hummock timer handler loop [{name}] exited unexpectedly");
302 }
303 Err(e) if e.is_cancelled() => {
306 let _ = child_shutdown_tx.send(true);
307 while let Some((name, result)) = child_handles.next().await {
308 if let Err(e) = result {
309 warn!(
310 handler = name,
311 error = %e.as_report(),
312 "Hummock timer handler loop failed during shutdown"
313 );
314 }
315 }
316 }
317 Err(e) => {
318 panic!("Hummock timer handler loop [{name}] failed: {}", e.as_report());
319 }
320 }
321 }
322 }
323}
324
325impl HummockManager {
326 pub fn hummock_timer_task(
327 hummock_manager: Arc<Self>,
328 backup_manager: Option<BackupManagerRef>,
329 ) -> (JoinHandle<()>, Sender<()>) {
330 let (shutdown_tx, shutdown_rx) = tokio::sync::oneshot::channel();
331 let join_handle = tokio::spawn(async move {
332 const STAT_REPORT_PERIOD_SEC: u64 = 20;
333 const COMPACTION_HEARTBEAT_PERIOD_SEC: u64 = 1;
334
335 let split_interval_sec = hummock_manager
336 .env
337 .opts
338 .periodic_scheduling_compaction_group_split_interval_sec;
339 let merge_interval_sec = hummock_manager
340 .env
341 .opts
342 .periodic_scheduling_compaction_group_merge_interval_sec;
343
344 tracing::info!(
345 split_interval_sec,
346 merge_interval_sec,
347 report_interval_sec = STAT_REPORT_PERIOD_SEC,
348 heartbeat_interval_sec = COMPACTION_HEARTBEAT_PERIOD_SEC,
349 "Starting Hummock timer task"
350 );
351
352 let (child_shutdown_tx, child_shutdown_rx) = watch::channel(false);
353 let mut child_handles: Vec<(&'static str, JoinHandle<()>)> = Vec::new();
354
355 {
358 let hummock_manager = hummock_manager.clone();
359 child_handles.push((
360 "heartbeat",
361 spawn_periodic_loop(
362 "heartbeat",
363 Duration::from_secs(COMPACTION_HEARTBEAT_PERIOD_SEC),
364 child_shutdown_rx.clone(),
365 move || {
366 let hummock_manager = hummock_manager.clone();
367 async move {
368 hummock_manager
369 .handle_timer_compaction_heartbeat_expired_check()
370 .await;
371 }
372 },
373 ),
374 ));
375 }
376
377 {
379 let hummock_manager = hummock_manager.clone();
380 child_handles.push((
381 "report",
382 spawn_periodic_loop(
383 "report",
384 Duration::from_secs(STAT_REPORT_PERIOD_SEC),
385 child_shutdown_rx.clone(),
386 move || {
387 let hummock_manager = hummock_manager.clone();
388 async move {
389 hummock_manager.handle_timer_report().await;
390 }
391 },
392 ),
393 ));
394 }
395
396 child_handles.push((
400 "scheduling",
401 spawn_scheduling_loop(
402 hummock_manager.clone(),
403 split_interval_sec,
404 merge_interval_sec,
405 child_shutdown_rx.clone(),
406 ),
407 ));
408
409 child_handles.push((
412 "maintenance",
413 spawn_maintenance_loop(
414 hummock_manager.clone(),
415 backup_manager,
416 child_shutdown_rx.clone(),
417 ),
418 ));
419
420 supervise_child_loops(shutdown_rx, child_shutdown_tx, child_handles).await;
421
422 tracing::info!("Hummock timer loop is stopped");
423 });
424 (join_handle, shutdown_tx)
425 }
426
427 async fn handle_timer_report(&self) {
428 let mv_id_to_all_table_ids = self
429 .metadata_manager
430 .get_job_id_to_internal_table_ids_mapping()
431 .await;
432 let table_write_throughput_statistic_manager =
433 self.table_write_throughput_statistic_manager.read().clone();
434 let id_to_config = self.get_compaction_group_map().await;
435
436 let (current_version, version_stats) = {
440 let versioning_guard = self
441 .versioning
442 .read_with_process_name("handle_timer_report")
443 .await;
444 (
445 versioning_guard.current_version.clone(),
446 versioning_guard.version_stats.clone(),
447 )
448 };
449
450 if let Some(mv_id_to_all_table_ids) = mv_id_to_all_table_ids {
451 trigger_mv_stat(&self.metrics, &version_stats, mv_id_to_all_table_ids);
452 }
453
454 let compaction_group_count = current_version.levels.len();
455 self.metrics
456 .compaction_group_count
457 .set(compaction_group_count as i64);
458
459 let max_statistic_expired_time = std::cmp::max(
460 self.env.opts.table_stat_throuput_window_seconds_for_split,
461 self.env.opts.table_stat_throuput_window_seconds_for_merge,
462 );
463
464 for (group_id, group_levels) in ¤t_version.levels {
465 let Some(compaction_group_config) = id_to_config.get(group_id) else {
466 warn!(
467 compaction_group_id = %group_id,
468 "Missing compaction group config when reporting metrics. Skip."
469 );
470 continue;
471 };
472
473 trigger_lsm_stat(
474 &self.metrics,
475 compaction_group_config.compaction_config(),
476 group_levels,
477 *group_id,
478 );
479
480 let member_table_ids = current_version
481 .state_table_info
482 .compaction_group_member_table_ids(*group_id);
483 let group_size = member_table_ids
484 .iter()
485 .map(|table_id| {
486 version_stats
487 .table_stats
488 .get(table_id)
489 .map(|stats| stats.total_key_size + stats.total_value_size)
490 .unwrap_or(0)
491 .max(0) as u64
492 })
493 .sum::<u64>();
494 self.metrics
495 .compaction_group_size
496 .with_label_values(&[&group_id.to_string()])
497 .set(group_size as _);
498
499 let mut avg_throughput = 0;
500 for table_id in member_table_ids {
501 avg_throughput += table_write_throughput_statistic_manager
502 .avg_write_throughput(*table_id, max_statistic_expired_time as i64)
503 as u64;
504 }
505
506 self.metrics
507 .compaction_group_throughput
508 .with_label_values(&[&group_id.to_string()])
509 .set(avg_throughput as _);
510
511 let file_count = group_levels.count_ssts();
512 self.metrics
513 .compaction_group_file_count
514 .with_label_values(&[&group_id.to_string()])
515 .set(file_count as _);
516 }
517 }
518
519 async fn handle_timer_compaction_heartbeat_expired_check(&self) {
520 let expired_tasks: Vec<u64> = self
527 .compactor_manager
528 .get_heartbeat_expired_tasks()
529 .into_iter()
530 .map(|task| task.task_id)
531 .collect();
532 if !expired_tasks.is_empty() {
533 tracing::info!(
534 expired_tasks = ?expired_tasks,
535 "Heartbeat expired compaction tasks detected. Attempting to cancel tasks.",
536 );
537 if let Err(e) = self
538 .cancel_compact_tasks(expired_tasks.clone(), TaskStatus::HeartbeatCanceled)
539 .await
540 {
541 tracing::error!(
542 expired_tasks = ?expired_tasks,
543 error = %e.as_report(),
544 "Attempt to remove compaction task due to elapsed heartbeat failed. We will continue to track its heartbeat
545 until we can successfully report its status",
546 );
547 }
548 }
549 }
550}
551
552impl HummockManager {
553 async fn maybe_normalize_compaction_groups_before_merge(&self) {
554 if !self.env.opts.enable_compaction_group_normalize {
555 return;
556 }
557
558 match self
559 .normalize_overlapping_compaction_groups_with_limit(
560 self.env
561 .opts
562 .max_normalize_splits_per_round
563 .try_into()
564 .unwrap_or(usize::MAX),
565 )
566 .await
567 {
568 Ok(split_count) => {
569 if split_count > 0 {
570 tracing::info!(
571 "normalize compaction groups finished with {} split(s) before merge scheduling",
572 split_count
573 );
574 }
575 }
576 Err(e) => {
577 tracing::warn!(
578 error = %e.as_report(),
579 "failed to normalize compaction groups before merge scheduling"
580 );
581 }
582 }
583 }
584
585 async fn check_dead_task(&self) {
586 const MAX_COMPACTION_L0_MULTIPLIER: u64 = 32;
587 const MAX_COMPACTION_DURATION_SEC: u64 = 20 * 60;
588 let slowdown_groups = {
589 let versioning_guard = self
590 .versioning
591 .read_with_process_name("check_dead_task")
592 .await;
593 let compaction_group_manager = self
594 .compaction_group_manager
595 .read_with_process_name("check_dead_task")
596 .await;
597 let mut slowdown_groups: HashMap<CompactionGroupId, u64> = HashMap::default();
598
599 for (group_id, group_levels) in &versioning_guard.current_version.levels {
600 let compaction_group_config = compaction_group_manager
601 .try_get_compaction_group_config(*group_id)
602 .expect(
603 "compaction group config should exist for every group in current version",
604 );
605 let l0_file_size = group_levels
606 .l0
607 .sub_levels
608 .iter()
609 .map(|level| level.total_file_size)
610 .sum::<u64>();
611 if l0_file_size
612 > MAX_COMPACTION_L0_MULTIPLIER
613 * compaction_group_config
614 .compaction_config
615 .max_bytes_for_level_base
616 {
617 slowdown_groups.insert(*group_id, l0_file_size);
618 }
619 }
620
621 slowdown_groups
622 };
623 if slowdown_groups.is_empty() {
624 return;
625 }
626 let mut pending_tasks: HashMap<u64, (CompactionGroupId, usize, RunningCompactTask)> =
627 HashMap::default();
628 {
629 let compaction_guard = self
630 .compaction
631 .read_with_process_name("check_dead_task")
632 .await;
633 for group_id in slowdown_groups.keys() {
634 if let Some(status) = compaction_guard.compaction_statuses.get(group_id) {
635 for (idx, level_handler) in status.level_handlers.iter().enumerate() {
636 let tasks = level_handler.pending_tasks().to_vec();
637 if tasks.is_empty() {
638 continue;
639 }
640 for task in tasks {
641 pending_tasks.insert(task.task_id, (*group_id, idx, task));
642 }
643 }
644 }
645 }
646 }
647 let task_ids = pending_tasks.keys().cloned().collect_vec();
648 let task_infos = self
649 .compactor_manager
650 .check_tasks_status(&task_ids, Duration::from_secs(MAX_COMPACTION_DURATION_SEC));
651 for (task_id, (compact_time, status)) in task_infos {
652 if status == TASK_NORMAL {
653 continue;
654 }
655 if let Some((group_id, level_id, task)) = pending_tasks.get(&task_id) {
656 let group_size = *slowdown_groups.get(group_id).unwrap();
657 warn!(
658 "COMPACTION SLOW: the task-{} of group-{}(size: {}MB) level-{} has not finished after {:?}, {}, it may cause pending sstable files({:?}) blocking other task.",
659 task_id,
660 group_id,
661 group_size / 1024 / 1024,
662 *level_id,
663 compact_time,
664 status,
665 task.ssts
666 );
667 }
668 }
669 }
670
671 async fn on_handle_schedule_group_split(&self) {
676 let table_write_throughput = self.table_write_throughput_statistic_manager.read().clone();
677
678 let mut group_infos = self.calculate_compaction_group_statistic().await;
679 group_infos.sort_by_key(|group| Reverse(group.group_size));
680
681 for group in group_infos {
682 if group.table_statistic.len() == 1 {
683 continue;
685 }
686
687 self.try_split_compaction_group(&table_write_throughput, group)
688 .await;
689 }
690 }
691
692 #[cfg(test)]
693 pub async fn schedule_group_split_for_test(&self) {
694 self.on_handle_schedule_group_split().await;
695 }
696
697 #[cfg(test)]
698 pub async fn schedule_group_merge_for_test(&self) {
699 self.on_handle_schedule_group_merge().await;
700 }
701
702 async fn on_handle_trigger_multi_group(&self, task_type: compact_task::TaskType) {
703 for cg_id in self.compaction_group_ids().await {
704 self.compaction_state.try_sched_compaction(
705 cg_id,
706 task_type,
707 super::compaction::ScheduleTrigger::Periodic,
708 );
709 }
710 }
711
712 async fn on_handle_schedule_group_merge(&self) {
718 self.maybe_normalize_compaction_groups_before_merge().await;
719
720 let created_tables = match self.metadata_manager.get_created_table_ids().await {
721 Ok(created_tables) => HashSet::from_iter(created_tables),
722 Err(err) => {
723 tracing::warn!(error = %err.as_report(), "failed to fetch created table ids");
724 return;
725 }
726 };
727 let table_write_throughput_statistic_manager =
728 self.table_write_throughput_statistic_manager.read().clone();
729 let mut group_infos = self.calculate_compaction_group_statistic().await;
730 group_infos.sort_by_key(|group| group.table_statistic.keys().next().copied());
732
733 let group_count = group_infos.len();
734 if group_count < 2 {
735 return;
736 }
737
738 let mut base = 0;
739 let mut candidate = 1;
740
741 while candidate < group_count {
742 let group = &group_infos[base];
743 let next_group = &group_infos[candidate];
744 match self
745 .try_merge_compaction_group(
746 &table_write_throughput_statistic_manager,
747 group,
748 next_group,
749 &created_tables,
750 )
751 .await
752 {
753 Ok(_) => candidate += 1,
754 Err(e) => {
755 tracing::debug!(
756 error = %e.as_report(),
757 "Failed to merge compaction group",
758 );
759 base = candidate;
760 candidate = base + 1;
761 }
762 }
763 }
764 }
765}
766
767#[cfg(test)]
768mod tests {
769 mod merge_scheduling {
770 use std::sync::Arc;
771 use std::time::Duration;
772
773 use itertools::Itertools;
774 use risingwave_common::catalog::TableId;
775 use risingwave_hummock_sdk::CompactionGroupId;
776 use risingwave_hummock_sdk::version::HummockVersion;
777 use risingwave_meta_model::WorkerId;
778 use risingwave_pb::common::worker_node::Property;
779 use risingwave_pb::common::{HostAddress, WorkerType};
780
781 use crate::controller::catalog::CatalogController;
782 use crate::controller::cluster::{ClusterController, ClusterControllerRef};
783 use crate::hummock::compaction::compaction_config::CompactionConfigBuilder;
784 use crate::hummock::{CompactorManager, HummockManager, HummockManagerRef};
785 use crate::manager::{MetaOpts, MetaSrvEnv};
786
787 async fn setup_compute_env_with_meta_opts(
788 port: i32,
789 opts: MetaOpts,
790 ) -> (
791 MetaSrvEnv,
792 HummockManagerRef,
793 ClusterControllerRef,
794 WorkerId,
795 ) {
796 let env = MetaSrvEnv::for_test_opts(opts, |_| ()).await;
797 let cluster_ctl = Arc::new(
798 ClusterController::new(env.clone(), Duration::from_secs(1))
799 .await
800 .unwrap(),
801 );
802 let catalog_ctl = Arc::new(CatalogController::new(env.clone()).await.unwrap());
803 let compactor_manager = Arc::new(CompactorManager::for_test());
804 let (compactor_streams_change_tx, _compactor_streams_change_rx) =
805 tokio::sync::mpsc::unbounded_channel();
806 let config = CompactionConfigBuilder::new()
807 .level0_tier_compact_file_number(1)
808 .level0_max_compact_file_number(130)
809 .level0_sub_level_compact_level_count(1)
810 .level0_overlapping_sub_level_compact_level_count(1)
811 .build();
812 let hummock_manager = HummockManager::with_config(
813 env.clone(),
814 cluster_ctl.clone(),
815 catalog_ctl,
816 Arc::new(Default::default()),
817 compactor_manager,
818 config,
819 compactor_streams_change_tx,
820 )
821 .await;
822
823 let worker_id = cluster_ctl
824 .add_worker(
825 WorkerType::ComputeNode,
826 HostAddress {
827 host: "127.0.0.1".to_owned(),
828 port,
829 },
830 Property {
831 is_streaming: true,
832 is_serving: true,
833 parallelism: 4,
834 ..Default::default()
835 },
836 Default::default(),
837 )
838 .await
839 .unwrap();
840 (env, hummock_manager, cluster_ctl, worker_id)
841 }
842
843 async fn get_compaction_group_id_by_table_id(
844 hummock_manager: HummockManagerRef,
845 table_id: u32,
846 ) -> CompactionGroupId {
847 hummock_manager
848 .get_current_version()
849 .await
850 .state_table_info
851 .info()
852 .get(&TableId::new(table_id))
853 .unwrap()
854 .compaction_group_id
855 }
856
857 fn member_table_ids(version: &HummockVersion, group_id: CompactionGroupId) -> Vec<u32> {
858 version
859 .state_table_info
860 .compaction_group_member_table_ids(group_id)
861 .iter()
862 .map(|table_id| table_id.as_raw_id())
863 .collect_vec()
864 }
865
866 fn assert_no_group_overlap(version: &HummockVersion) {
867 let mut ranges = version
868 .levels
869 .keys()
870 .filter_map(|group_id| {
871 let members = member_table_ids(version, *group_id);
872 (!members.is_empty())
873 .then(|| (*members.first().unwrap(), *members.last().unwrap()))
874 })
875 .collect_vec();
876 ranges.sort_by_key(|(min_table_id, _)| *min_table_id);
877 assert!(ranges.windows(2).all(|window| window[0].1 < window[1].0));
878 }
879
880 #[tokio::test]
881 async fn test_merge_scheduling_normalizes_when_split_scheduling_is_disabled() {
882 let mut opts = MetaOpts::test(false);
883 opts.enable_compaction_group_normalize = true;
884 opts.periodic_scheduling_compaction_group_split_interval_sec = 0;
885
886 let (_env, hummock_manager, _, _worker_id) =
887 setup_compute_env_with_meta_opts(80, opts).await;
888 hummock_manager
889 .register_table_ids_for_test(&[(64, 2.into()), (80, 2.into())])
890 .await
891 .unwrap();
892 hummock_manager
893 .register_table_ids_for_test(&[(65, 3.into()), (81, 3.into()), (83, 3.into())])
894 .await
895 .unwrap();
896
897 let cg_64 = get_compaction_group_id_by_table_id(hummock_manager.clone(), 64).await;
898 let cg_65 = get_compaction_group_id_by_table_id(hummock_manager.clone(), 65).await;
899
900 hummock_manager.on_handle_schedule_group_merge().await;
901
902 let version = hummock_manager.get_current_version().await;
903 assert_eq!(member_table_ids(&version, cg_64), vec![64]);
904 assert_eq!(member_table_ids(&version, cg_65), vec![65]);
905 assert_no_group_overlap(&version);
906 }
907 }
908
909 mod periodic_loop {
910 use std::sync::Arc;
911 use std::sync::atomic::{AtomicUsize, Ordering};
912 use std::time::Duration;
913
914 use tokio::sync::watch;
915
916 use super::super::{interval_stream, spawn_serial_timer_lane};
917
918 #[tokio::test(start_paused = true)]
919 async fn test_serial_timer_lane_does_not_run_handlers_concurrently() {
920 #[derive(Clone, Copy)]
921 enum TestEvent {
922 Split,
923 Merge,
924 }
925
926 struct DecrOnDrop(Arc<AtomicUsize>);
927
928 impl Drop for DecrOnDrop {
929 fn drop(&mut self) {
930 self.0.fetch_sub(1, Ordering::SeqCst);
931 }
932 }
933
934 let (shutdown_tx, shutdown_rx) = watch::channel(false);
935 let in_handler = Arc::new(AtomicUsize::new(0));
936 let max_in_handler = Arc::new(AtomicUsize::new(0));
937 let handled = Arc::new(AtomicUsize::new(0));
938
939 let handle = spawn_serial_timer_lane(
940 "test_serial_lane",
941 vec![
942 interval_stream(Duration::from_millis(10), TestEvent::Split),
943 interval_stream(Duration::from_millis(10), TestEvent::Merge),
944 ],
945 shutdown_rx,
946 {
947 let in_handler = in_handler.clone();
948 let max_in_handler = max_in_handler.clone();
949 let handled = handled.clone();
950 move |_event| {
951 let in_handler = in_handler.clone();
952 let max_in_handler = max_in_handler.clone();
953 let handled = handled.clone();
954 async move {
955 handled.fetch_add(1, Ordering::SeqCst);
956 let now = in_handler.fetch_add(1, Ordering::SeqCst) + 1;
957 let _decr = DecrOnDrop(in_handler);
958 max_in_handler.fetch_max(now, Ordering::SeqCst);
959 tokio::time::sleep(Duration::from_millis(15)).await;
960 }
961 }
962 },
963 );
964
965 tokio::time::sleep(Duration::from_millis(120)).await;
966 shutdown_tx.send(true).unwrap();
967 handle.await.unwrap();
968
969 assert!(handled.load(Ordering::SeqCst) > 0);
970 assert_eq!(max_in_handler.load(Ordering::SeqCst), 1);
971 }
972 }
973}