1mod planning;
16mod state;
17
18use std::collections::HashMap;
19use std::sync::Arc;
20
21use api::v1::region::compact_request;
22use common_base::Plugins;
23use common_memory_manager::OnExhaustedPolicy;
24use common_meta::key::SchemaMetadataManagerRef;
25use common_telemetry::{debug, error, info};
26use common_time::range::TimestampRange;
27pub(crate) use planning::CompactionPickFinished;
28pub use planning::CompactionRequest;
29pub(crate) use state::CompactionExecution;
30use state::{CompactionStatus, CompactionTrigger, PendingCompaction};
31use store_api::storage::RegionId;
32use tokio::sync::mpsc::Sender;
33
34use crate::access_layer::AccessLayerRef;
35use crate::cache::CacheManagerRef;
36use crate::compaction::memory_manager::CompactionMemoryManager;
37use crate::compaction::task::MAX_PARALLEL_COMPACTION;
38use crate::config::MitoConfig;
39use crate::error::{
40 CompactionCancelledSnafu, Error, ManualCompactionAlreadyRunningSnafu, RegionClosedSnafu,
41 RegionDroppedSnafu, RegionTruncatedSnafu, Result,
42};
43use crate::region::version::VersionControlRef;
44use crate::region::{ManifestContextRef, RegionLeaderState, RegionRoleState};
45use crate::request::{DdlRequest, OptionOutputTx, SenderDdlRequest, WorkerRequestWithTime};
46#[cfg(test)]
47use crate::schedule::RequestCancelResult;
48use crate::schedule::scheduler::SchedulerRef;
49use crate::worker::WorkerListener;
50
51pub(crate) struct CompactionScheduler {
53 scheduler: SchedulerRef,
54 region_status: HashMap<RegionId, CompactionStatus>,
56 request_sender: Sender<WorkerRequestWithTime>,
58 cache_manager: CacheManagerRef,
59 engine_config: Arc<MitoConfig>,
60 memory_manager: Arc<CompactionMemoryManager>,
61 memory_policy: OnExhaustedPolicy,
62 listener: WorkerListener,
63 plugins: Plugins,
65 next_plan_id: u64,
68}
69
70#[derive(Debug)]
72pub(crate) enum CompactionTransition {
73 NoAction,
75 AutomaticFollowupScheduled,
77 DdlReady(Vec<SenderDdlRequest>),
79}
80
81impl CompactionTransition {
82 fn from_pending_ddls(pending_ddls: Vec<SenderDdlRequest>) -> Self {
83 if pending_ddls.is_empty() {
84 Self::NoAction
85 } else {
86 Self::DdlReady(pending_ddls)
87 }
88 }
89
90 #[cfg(test)]
91 fn is_empty(&self) -> bool {
92 !matches!(self, Self::DdlReady(pending_ddls) if !pending_ddls.is_empty())
93 }
94
95 #[cfg(test)]
96 fn len(&self) -> usize {
97 match self {
98 Self::DdlReady(pending_ddls) => pending_ddls.len(),
99 Self::NoAction | Self::AutomaticFollowupScheduled => 0,
100 }
101 }
102}
103
104impl CompactionScheduler {
106 #[allow(clippy::too_many_arguments)]
118 pub(crate) fn new(
119 scheduler: SchedulerRef,
120 request_sender: Sender<WorkerRequestWithTime>,
121 cache_manager: CacheManagerRef,
122 engine_config: Arc<MitoConfig>,
123 listener: WorkerListener,
124 plugins: Plugins,
125 memory_manager: Arc<CompactionMemoryManager>,
126 memory_policy: OnExhaustedPolicy,
127 ) -> Self {
128 Self {
129 scheduler,
130 region_status: HashMap::new(),
131 request_sender,
132 cache_manager,
133 engine_config,
134 memory_manager,
135 memory_policy,
136 listener,
137 plugins,
138 next_plan_id: 0,
139 }
140 }
141
142 pub(crate) fn schedule_automatic_compaction(
157 &mut self,
158 compact_options: compact_request::Options,
159 version_control: &VersionControlRef,
160 access_layer: &AccessLayerRef,
161 manifest_ctx: &ManifestContextRef,
162 schema_metadata_manager: SchemaMetadataManagerRef,
163 ) -> Result<bool> {
164 self.schedule_compaction(
165 CompactionTrigger::Automatic,
166 compact_options,
167 version_control,
168 access_layer,
169 OptionOutputTx::none(),
170 manifest_ctx,
171 schema_metadata_manager,
172 1, None,
174 )
175 }
176
177 #[allow(clippy::too_many_arguments)]
195 pub(crate) fn schedule_manual_compaction(
196 &mut self,
197 compact_options: compact_request::Options,
198 version_control: &VersionControlRef,
199 access_layer: &AccessLayerRef,
200 waiter: OptionOutputTx,
201 manifest_ctx: &ManifestContextRef,
202 schema_metadata_manager: SchemaMetadataManagerRef,
203 max_parallelism: usize,
204 time_range: Option<TimestampRange>,
205 ) -> Result<bool> {
206 self.schedule_compaction(
207 CompactionTrigger::Manual,
208 compact_options,
209 version_control,
210 access_layer,
211 waiter,
212 manifest_ctx,
213 schema_metadata_manager,
214 max_parallelism,
215 time_range,
216 )
217 }
218
219 pub(crate) async fn handle_compaction_pick_finished(
234 &mut self,
235 finished: CompactionPickFinished,
236 manifest_ctx: &ManifestContextRef,
237 schema_metadata_manager: SchemaMetadataManagerRef,
238 ) -> CompactionTransition {
239 self.handle_compaction_pick_finished_inner(finished, manifest_ctx, schema_metadata_manager)
240 .await
241 }
242
243 pub(crate) fn is_current_execution(
254 &self,
255 region_id: RegionId,
256 execution: &CompactionExecution,
257 ) -> bool {
258 self.region_status
259 .get(®ion_id)
260 .is_some_and(|status| status.matches_execution(execution))
261 }
262
263 pub(crate) async fn on_execution_finished(
277 &mut self,
278 region_id: RegionId,
279 execution: &CompactionExecution,
280 manifest_ctx: &ManifestContextRef,
281 schema_metadata_manager: SchemaMetadataManagerRef,
282 ) -> CompactionTransition {
283 if !self.is_current_execution(region_id, execution) {
284 return CompactionTransition::NoAction;
285 }
286 self.on_compaction_finished(region_id, manifest_ctx, schema_metadata_manager)
287 .await
288 }
289
290 pub(crate) async fn on_execution_cancelled(
302 &mut self,
303 region_id: RegionId,
304 execution: &CompactionExecution,
305 ) -> Vec<SenderDdlRequest> {
306 if !self.is_current_execution(region_id, execution) {
307 return Vec::new();
308 }
309 self.on_compaction_cancelled(region_id).await
310 }
311
312 pub(crate) fn on_execution_failed(
324 &mut self,
325 region_id: RegionId,
326 execution: &CompactionExecution,
327 err: Arc<Error>,
328 ) {
329 if !self.is_current_execution(region_id, execution) {
330 return;
331 }
332 self.on_compaction_failed(region_id, err);
333 }
334
335 pub(crate) fn on_region_dropped(&mut self, region_id: RegionId) {
345 self.remove_region_on_failure(
346 region_id,
347 Arc::new(RegionDroppedSnafu { region_id }.build()),
348 );
349 }
350
351 pub(crate) fn on_region_closed(&mut self, region_id: RegionId) {
361 self.remove_region_on_failure(region_id, Arc::new(RegionClosedSnafu { region_id }.build()));
362 }
363
364 pub(crate) fn on_region_truncated(&mut self, region_id: RegionId) {
374 self.remove_region_on_failure(
375 region_id,
376 Arc::new(RegionTruncatedSnafu { region_id }.build()),
377 );
378 }
379
380 pub(crate) fn try_cancel_and_add_ddl<T>(
396 &mut self,
397 region_id: RegionId,
398 sender: OptionOutputTx,
399 request: T,
400 into_ddl_request: impl FnOnce(T) -> DdlRequest,
401 ) -> std::result::Result<(), (OptionOutputTx, T)> {
402 let Some(status) = self.region_status.get_mut(®ion_id) else {
403 return Err((sender, request));
404 };
405 status.request_cancel();
406
407 let request = SenderDdlRequest {
408 region_id,
409 sender,
410 request: into_ddl_request(request),
411 };
412 debug!(
413 "Added pending DDL request for region: {}, ddl: {:?}",
414 request.region_id, request.request
415 );
416 status.pending_ddl_requests.push(request);
417 Ok(())
418 }
419}
420
421impl CompactionScheduler {
423 fn next_plan_id(counter: &mut u64) -> u64 {
428 let plan_id = *counter;
429 *counter = counter.wrapping_add(1);
430 plan_id
431 }
432
433 #[allow(clippy::too_many_arguments)]
434 fn schedule_compaction(
435 &mut self,
436 trigger: CompactionTrigger,
437 compact_options: compact_request::Options,
438 version_control: &VersionControlRef,
439 access_layer: &AccessLayerRef,
440 waiter: OptionOutputTx,
441 manifest_ctx: &ManifestContextRef,
442 schema_metadata_manager: SchemaMetadataManagerRef,
443 max_parallelism: usize,
444 time_range: Option<TimestampRange>,
445 ) -> Result<bool> {
446 let region_id = version_control.region_id();
447 let current_state = manifest_ctx.current_state();
448 if current_state == RegionRoleState::Leader(RegionLeaderState::Staging) {
449 info!(
450 "Skipping compaction for region {} in staging mode, options: {:?}",
451 region_id, compact_options
452 );
453 waiter.send(Ok(0));
454 return Ok(false);
455 }
456
457 if let Some(status) = self.region_status.get_mut(®ion_id) {
458 if !status.pending_ddl_requests.is_empty() {
461 waiter.send(CompactionCancelledSnafu.fail());
462 info!(
463 "Region {} has pending DDL requests, ignoring compaction: {:?}",
464 region_id, compact_options
465 );
466 return Ok(false);
467 }
468
469 match trigger {
470 CompactionTrigger::Automatic => status.mark_automatic_trigger(),
471 CompactionTrigger::Manual if status.is_manual_compaction() => {
472 waiter.send(ManualCompactionAlreadyRunningSnafu { region_id }.fail());
473 info!(
474 "Region {} already has a manually triggered compaction running",
475 region_id
476 );
477 }
478 CompactionTrigger::Manual => {
479 status.set_pending_request(PendingCompaction {
480 options: compact_options,
481 waiter,
482 max_parallelism,
483 time_range,
484 });
485 info!(
486 "Region {} is running an automatic compaction; manual compaction will be re-scheduled",
487 region_id
488 );
489 }
490 }
491 return Ok(false);
492 }
493
494 let plan_id = Self::next_plan_id(&mut self.next_plan_id);
496 let mut status = CompactionStatus::new(
497 region_id,
498 version_control.clone(),
499 access_layer.clone(),
500 plan_id,
501 trigger,
502 );
503 let request = status.new_compaction_request(
504 self.request_sender.clone(),
505 self.engine_config.clone(),
506 self.cache_manager.clone(),
507 manifest_ctx,
508 self.listener.clone(),
509 schema_metadata_manager,
510 max_parallelism,
511 );
512 status.merge_waiter(waiter);
513 self.region_status.insert(region_id, status);
514 self.dispatch_compaction_planning(plan_id, request, compact_options, time_range);
515 self.listener.on_compaction_scheduled(region_id);
516 Ok(true)
517 }
518
519 fn handle_pending_compaction_request(
523 &mut self,
524 region_id: RegionId,
525 manifest_ctx: &ManifestContextRef,
526 schema_metadata_manager: SchemaMetadataManagerRef,
527 ) -> bool {
528 let Some(status) = self.region_status.get_mut(®ion_id) else {
529 return true;
530 };
531
532 let Some(pending_request) = std::mem::take(&mut status.pending_request) else {
535 return false;
536 };
537
538 let PendingCompaction {
539 options,
540 waiter,
541 max_parallelism,
542 time_range,
543 } = pending_request;
544
545 let request = status.new_compaction_request(
546 self.request_sender.clone(),
547 self.engine_config.clone(),
548 self.cache_manager.clone(),
549 manifest_ctx,
550 self.listener.clone(),
551 schema_metadata_manager,
552 max_parallelism,
553 );
554 status.merge_waiter(waiter);
555 let plan_id = Self::next_plan_id(&mut self.next_plan_id);
559 status.start_picking_with_trigger(plan_id, CompactionTrigger::Manual);
560 self.dispatch_compaction_planning(plan_id, request, options, time_range);
561 debug!(
562 "Successfully scheduled manual compaction planning for region id: {}",
563 region_id
564 );
565 true
566 }
567
568 async fn on_compaction_finished(
570 &mut self,
571 region_id: RegionId,
572 manifest_ctx: &ManifestContextRef,
573 schema_metadata_manager: SchemaMetadataManagerRef,
574 ) -> CompactionTransition {
575 if !self.region_status.contains_key(®ion_id) {
576 return CompactionTransition::NoAction;
577 }
578
579 if self.handle_pending_compaction_request(
580 region_id,
581 manifest_ctx,
582 schema_metadata_manager.clone(),
583 ) {
584 return CompactionTransition::NoAction;
585 }
586
587 let Some(status) = self.region_status.get_mut(®ion_id) else {
590 return CompactionTransition::NoAction;
591 };
592 for waiter in status.take_waiters() {
593 waiter.send(Ok(0));
594 }
595
596 let pending_ddl_requests = std::mem::take(&mut status.pending_ddl_requests);
600 if !pending_ddl_requests.is_empty() {
601 self.region_status.remove(®ion_id);
603 return CompactionTransition::DdlReady(pending_ddl_requests);
606 }
607
608 if status.active.reset_automatic_followup()
609 && self.schedule_automatic_followup(region_id, manifest_ctx, schema_metadata_manager)
610 {
611 return CompactionTransition::AutomaticFollowupScheduled;
612 }
613 self.region_status.remove(®ion_id);
614 CompactionTransition::NoAction
615 }
616
617 fn schedule_automatic_followup(
618 &mut self,
619 region_id: RegionId,
620 manifest_ctx: &ManifestContextRef,
621 schema_metadata_manager: SchemaMetadataManagerRef,
622 ) -> bool {
623 let Some(status) = self.region_status.get_mut(®ion_id) else {
624 return false;
625 };
626 let request = status.new_compaction_request(
628 self.request_sender.clone(),
629 self.engine_config.clone(),
630 self.cache_manager.clone(),
631 manifest_ctx,
632 self.listener.clone(),
633 schema_metadata_manager,
634 MAX_PARALLEL_COMPACTION,
635 );
636 let plan_id = Self::next_plan_id(&mut self.next_plan_id);
640 status.start_regular_picking(plan_id);
641 self.dispatch_compaction_planning(
642 plan_id,
643 request,
644 compact_request::Options::Regular(Default::default()),
645 None,
646 );
647 debug!(
648 "Successfully scheduled next compaction planning for region id: {}",
649 region_id
650 );
651 true
652 }
653
654 async fn on_compaction_cancelled(&mut self, region_id: RegionId) -> Vec<SenderDdlRequest> {
656 self.remove_region_on_cancel(region_id)
657 }
658
659 fn on_compaction_failed(&mut self, region_id: RegionId, err: Arc<Error>) {
661 error!(err; "Region {} failed to compact, cancel all pending tasks", region_id);
662 self.remove_region_on_failure(region_id, err);
663 }
664
665 #[cfg(test)]
666 fn add_ddl_request_to_pending(&mut self, request: SenderDdlRequest) {
667 self.region_status
668 .get_mut(&request.region_id)
669 .unwrap()
670 .pending_ddl_requests
671 .push(request);
672 }
673
674 #[cfg(test)]
675 fn has_pending_ddls(&self, region_id: RegionId) -> bool {
676 let has_pending = self
677 .region_status
678 .get(®ion_id)
679 .map(|status| !status.pending_ddl_requests.is_empty())
680 .unwrap_or(false);
681 debug!(
682 "Checked pending DDL requests for region: {}, has_pending: {}",
683 region_id, has_pending
684 );
685 has_pending
686 }
687
688 #[cfg(test)]
689 fn request_cancel(&mut self, region_id: RegionId) -> RequestCancelResult {
690 self.region_status
691 .get_mut(®ion_id)
692 .unwrap()
693 .request_cancel()
694 }
695
696 fn remove_region_on_failure(&mut self, region_id: RegionId, err: Arc<Error>) {
697 let Some(status) = self.region_status.remove(®ion_id) else {
699 return;
700 };
701
702 status.on_failure(err);
704 }
705
706 fn remove_region_on_cancel(&mut self, region_id: RegionId) -> Vec<SenderDdlRequest> {
707 let Some(status) = self.region_status.remove(®ion_id) else {
708 return Vec::new();
709 };
710
711 status.on_cancel()
712 }
713}
714
715impl Drop for CompactionScheduler {
716 fn drop(&mut self) {
717 for (region_id, status) in self.region_status.drain() {
718 status.on_failure(Arc::new(RegionClosedSnafu { region_id }.build()));
720 }
721 }
722}
723
724#[cfg(test)]
725#[path = "scheduler_test.rs"]
726mod tests;