From fa0cf5ec4bc08ee77facb42e70ce07453124664f Mon Sep 17 00:00:00 2001 From: "Lei, HUANG" Date: Thu, 23 Jul 2026 17:39:20 +0800 Subject: [PATCH] refactor(mito2): make compaction scheduling methods synchronous schedule_compaction, handle_pending_compaction_request and schedule_next_compaction no longer await anything after compaction planning became fire-and-forget. Drop the async signature to make the no-suspension-point invariant explicit: these methods always run to completion on the worker loop without reentrancy. Signed-off-by: Lei, HUANG --- src/mito2/src/compaction.rs | 108 ++++++++-------------- src/mito2/src/worker/handle_compaction.rs | 62 +++++-------- 2 files changed, 65 insertions(+), 105 deletions(-) diff --git a/src/mito2/src/compaction.rs b/src/mito2/src/compaction.rs index 50baba56e5..8a52ed2398 100644 --- a/src/mito2/src/compaction.rs +++ b/src/mito2/src/compaction.rs @@ -249,7 +249,7 @@ impl CompactionScheduler { /// Schedules a compaction for the region. /// Returns whether a compaction is scheduled. #[allow(clippy::too_many_arguments)] - pub(crate) async fn schedule_compaction( + pub(crate) fn schedule_compaction( &mut self, region_id: RegionId, compact_options: compact_request::Options, @@ -321,7 +321,7 @@ impl CompactionScheduler { // Handle pending manual compaction request for the region. // // Returns true if should early return, false otherwise. - pub(crate) async fn handle_pending_compaction_request( + pub(crate) fn handle_pending_compaction_request( &mut self, region_id: RegionId, manifest_ctx: &ManifestContextRef, @@ -386,14 +386,11 @@ impl CompactionScheduler { // If there a pending compaction request, handle it first // and defer returning the pending DDL requests to the caller. - if self - .handle_pending_compaction_request( - region_id, - manifest_ctx, - schema_metadata_manager.clone(), - ) - .await - { + if self.handle_pending_compaction_request( + region_id, + manifest_ctx, + schema_metadata_manager.clone(), + ) { return Vec::new(); } @@ -408,8 +405,7 @@ impl CompactionScheduler { } if status.regular_replan_pending { - self.schedule_next_compaction(region_id, manifest_ctx, schema_metadata_manager) - .await; + self.schedule_next_compaction(region_id, manifest_ctx, schema_metadata_manager); return Vec::new(); } @@ -484,7 +480,7 @@ impl CompactionScheduler { /// Schedules next compaction upon a finished compaction. /// Returns whether the compaction is scheduled. - pub(crate) async fn schedule_next_compaction( + pub(crate) fn schedule_next_compaction( &mut self, region_id: RegionId, manifest_ctx: &ManifestContextRef, @@ -919,20 +915,16 @@ impl CompactionScheduler { } } - if self - .handle_pending_compaction_request( - region_id, - manifest_ctx, - schema_metadata_manager.clone(), - ) - .await - { + if self.handle_pending_compaction_request( + region_id, + manifest_ctx, + schema_metadata_manager.clone(), + ) { return Vec::new(); } if self.region_status[®ion_id].regular_replan_pending { - self.schedule_next_compaction(region_id, manifest_ctx, schema_metadata_manager) - .await; + self.schedule_next_compaction(region_id, manifest_ctx, schema_metadata_manager); return Vec::new(); } @@ -2055,7 +2047,6 @@ mod tests { schema_metadata_manager.clone(), 1, ) - .await .unwrap() ); let finished = recv_compaction_pick_finished(rx).await; @@ -2202,7 +2193,6 @@ mod tests { schema_metadata_manager.clone(), 1, ) - .await .unwrap() ); let (manual_tx, _manual_rx) = oneshot::channel(); @@ -2218,7 +2208,6 @@ mod tests { schema_metadata_manager, 1, ) - .await .unwrap() ); @@ -2500,7 +2489,6 @@ mod tests { schema_metadata_manager.clone(), 1, ) - .await .unwrap(); assert!(scheduled); let finished = recv_compaction_pick_finished(&mut rx).await; @@ -2531,7 +2519,6 @@ mod tests { schema_metadata_manager.clone(), 1, ) - .await .unwrap(); assert!(scheduled); let finished = recv_compaction_pick_finished(&mut rx).await; @@ -2589,7 +2576,6 @@ mod tests { schema_metadata_manager.clone(), 1, ) - .await .unwrap(); // The boolean result is what the worker uses to decide whether to update @@ -2658,7 +2644,6 @@ mod tests { schema_metadata_manager, 1, ) - .await .unwrap(); (scheduled, scheduler) } @@ -2694,7 +2679,6 @@ mod tests { schema_metadata_manager.clone(), 1, ) - .await .unwrap() ); let (manual_tx, _manual_rx) = oneshot::channel(); @@ -2710,7 +2694,6 @@ mod tests { schema_metadata_manager, 1, ) - .await .unwrap() ); let status = scheduler.region_status.get(®ion_id).unwrap(); @@ -3050,7 +3033,6 @@ mod tests { schema_metadata_manager.clone(), 1, ) - .await .unwrap() ); assert!( @@ -3065,7 +3047,6 @@ mod tests { schema_metadata_manager.clone(), 1, ) - .await .unwrap() ); let (ddl_tx, mut ddl_rx) = oneshot::channel(); @@ -3198,7 +3179,6 @@ mod tests { schema_metadata_manager.clone(), 1, ) - .await .unwrap() ); let (ddl_tx, mut ddl_rx) = oneshot::channel(); @@ -3251,7 +3231,6 @@ mod tests { schema_metadata_manager.clone(), 1, ) - .await .unwrap() ); postfence_waiters.push(waiter_rx); @@ -3268,7 +3247,6 @@ mod tests { schema_metadata_manager.clone(), 1, ) - .await .unwrap() ); @@ -4124,7 +4102,6 @@ mod tests { schema_metadata_manager.clone(), 1, ) - .await .unwrap(); assert!(scheduled); @@ -4188,7 +4165,6 @@ mod tests { schema_metadata_manager.clone(), 1, ) - .await .unwrap(); // Should schedule 1 compaction. assert!(scheduled); @@ -4230,7 +4206,6 @@ mod tests { schema_metadata_manager.clone(), 1, ) - .await .unwrap(); assert!(!scheduled); assert_eq!(1, scheduler.region_status.len()); @@ -4248,9 +4223,11 @@ mod tests { scheduler .on_compaction_finished(region_id, &manifest_ctx, schema_metadata_manager.clone()) .await; - let scheduled = scheduler - .schedule_next_compaction(region_id, &manifest_ctx, schema_metadata_manager.clone()) - .await; + let scheduled = scheduler.schedule_next_compaction( + region_id, + &manifest_ctx, + schema_metadata_manager.clone(), + ); assert!(scheduled); assert_eq!(1, scheduler.region_status.len()); assert_eq!(1, job_scheduler.num_jobs()); @@ -4284,7 +4261,6 @@ mod tests { schema_metadata_manager, 1, ) - .await .unwrap(); assert!(!scheduled); assert_eq!(2, job_scheduler.num_jobs()); @@ -4372,7 +4348,6 @@ mod tests { schema_metadata_manager.clone(), 1, ) - .await .unwrap(); assert!(scheduled); @@ -4408,7 +4383,6 @@ mod tests { schema_metadata_manager.clone(), 1, ) - .await .unwrap(); assert!(!scheduled); assert_eq!(1, job_scheduler.num_jobs()); @@ -4426,7 +4400,6 @@ mod tests { schema_metadata_manager.clone(), 1, ) - .await .unwrap(); assert!(scheduled); } @@ -4466,18 +4439,16 @@ mod tests { ) .await; - let result = scheduler - .schedule_compaction( - region_id, - compact_request::Options::Regular(Default::default()), - &version_control, - &env.access_layer, - OptionOutputTx::none(), - &manifest_ctx, - schema_metadata_manager.clone(), - 1, - ) - .await; + let result = scheduler.schedule_compaction( + region_id, + compact_request::Options::Regular(Default::default()), + &version_control, + &env.access_layer, + OptionOutputTx::none(), + &manifest_ctx, + schema_metadata_manager.clone(), + 1, + ); assert!(result.unwrap()); assert!(scheduler.region_status.contains_key(®ion_id)); @@ -4551,7 +4522,6 @@ mod tests { schema_metadata_manager.clone(), 1, ) - .await .unwrap(); // Should schedule 1 compaction. assert_eq!(1, scheduler.region_status.len()); @@ -4587,7 +4557,6 @@ mod tests { schema_metadata_manager.clone(), 1, ) - .await .unwrap(); assert_eq!(1, scheduler.region_status.len()); // Current job num should be 1 since compaction is in progress. @@ -4666,7 +4635,6 @@ mod tests { schema_metadata_manager, 1, ) - .await .unwrap(); let result = rx.await.unwrap(); @@ -5196,9 +5164,11 @@ mod tests { let (schema_metadata_manager, _kv_backend) = mock_schema_metadata_manager(); // With no compactable files, next scheduling returns false and removes // the status without creating a background task. - let scheduled = scheduler - .schedule_next_compaction(region_id, &manifest_ctx, schema_metadata_manager.clone()) - .await; + let scheduled = scheduler.schedule_next_compaction( + region_id, + &manifest_ctx, + schema_metadata_manager.clone(), + ); assert!(scheduled); let finished = recv_compaction_pick_finished(&mut rx).await; scheduler @@ -5250,9 +5220,11 @@ mod tests { let (schema_metadata_manager, _kv_backend) = mock_schema_metadata_manager(); // The failing scheduler simulates a submit error; callers must see false. - let scheduled = scheduler - .schedule_next_compaction(region_id, &manifest_ctx, schema_metadata_manager.clone()) - .await; + let scheduled = scheduler.schedule_next_compaction( + region_id, + &manifest_ctx, + schema_metadata_manager.clone(), + ); assert!(scheduled); let finished = recv_compaction_pick_finished(&mut rx).await; scheduler diff --git a/src/mito2/src/worker/handle_compaction.rs b/src/mito2/src/worker/handle_compaction.rs index 1b12729dea..722fa47b01 100644 --- a/src/mito2/src/worker/handle_compaction.rs +++ b/src/mito2/src/worker/handle_compaction.rs @@ -81,20 +81,16 @@ impl RegionWorkerLoop { COMPACTION_REQUEST_COUNT.inc(); let parallelism = req.parallelism.unwrap_or(1) as usize; self.listener.on_compaction_schedule_attempt(region_id); - if let Err(e) = self - .compaction_scheduler - .schedule_compaction( - region.region_id, - req.options, - ®ion.version_control, - ®ion.access_layer, - sender, - ®ion.manifest_ctx, - self.schema_metadata_manager.clone(), - parallelism, - ) - .await - { + if let Err(e) = self.compaction_scheduler.schedule_compaction( + region.region_id, + req.options, + ®ion.version_control, + ®ion.access_layer, + sender, + ®ion.manifest_ctx, + self.schema_metadata_manager.clone(), + parallelism, + ) { error!(e; "Failed to schedule compaction task for region: {}", region_id); } else { info!( @@ -176,15 +172,11 @@ impl RegionWorkerLoop { "minimal compaction interval time {:?} has passed, scheduling next compaction", self.config.min_compaction_interval ); - if self - .compaction_scheduler - .schedule_next_compaction( - region_id, - ®ion.manifest_ctx, - self.schema_metadata_manager.clone(), - ) - .await - { + if self.compaction_scheduler.schedule_next_compaction( + region_id, + ®ion.manifest_ctx, + self.schema_metadata_manager.clone(), + ) { region.update_schedule_compaction_millis(); } } else { @@ -257,20 +249,16 @@ impl RegionWorkerLoop { ); self.listener .on_compaction_schedule_attempt(region.region_id); - match self - .compaction_scheduler - .schedule_compaction( - region.region_id, - compact_request::Options::Regular(Default::default()), - ®ion.version_control, - ®ion.access_layer, - OptionOutputTx::none(), - ®ion.manifest_ctx, - self.schema_metadata_manager.clone(), - 1, // Default for automatic compaction - ) - .await - { + match self.compaction_scheduler.schedule_compaction( + region.region_id, + compact_request::Options::Regular(Default::default()), + ®ion.version_control, + ®ion.access_layer, + OptionOutputTx::none(), + ®ion.manifest_ctx, + self.schema_metadata_manager.clone(), + 1, // Default for automatic compaction + ) { Ok(true) => region.update_schedule_compaction_millis(), Ok(false) => {} Err(e) => {