1use std::fmt;
16use std::future::Future;
17use std::sync::{Arc, Mutex};
18use std::time::Instant;
19
20use api::v1::region::compact_request;
21use common_base::Plugins;
22use common_meta::key::SchemaMetadataManagerRef;
23use common_telemetry::{debug, error, info, warn};
24use common_time::TimeToLive;
25use common_time::range::TimestampRange;
26use futures::FutureExt;
27use snafu::ResultExt;
28use store_api::storage::RegionId;
29use tokio::sync::mpsc::{self, Sender};
30
31use crate::access_layer::AccessLayerRef;
32use crate::cache::CacheManagerRef;
33use crate::compaction::compactor::{CompactionRegion, CompactionVersion, DefaultCompactor};
34use crate::compaction::picker::{CompactionTask, PickerOutput, new_picker};
35use crate::compaction::scheduler::state::{CompactingFiles, CompactionExecution, CompactionPhase};
36use crate::compaction::scheduler::{CompactionScheduler, CompactionTransition};
37use crate::compaction::task::CompactionTaskImpl;
38use crate::compaction::{CompactionOutput, find_dynamic_options};
39use crate::config::MitoConfig;
40use crate::error::{CompactRegionSnafu, Error, RemoteCompactionSnafu, Result, UnexpectedSnafu};
41use crate::metrics::{
42 COMPACTION_MEMORY_REJECTED, COMPACTION_STAGE_ELAPSED, INFLIGHT_COMPACTION_COUNT,
43};
44use crate::region::ManifestContextRef;
45use crate::region::options::RegionOptions;
46use crate::request::{BackgroundNotify, OutputTx, WorkerRequest, WorkerRequestWithTime};
47use crate::schedule::CancellableTaskState;
48use crate::schedule::remote_job_scheduler::{
49 CompactionJob, DefaultNotifier, RemoteJob, RemoteJobSchedulerRef,
50};
51use crate::sst::file::{FileHandle, UncommittedSsts};
52use crate::sst::version::SstVersion;
53use crate::worker::WorkerListener;
54
55pub struct CompactionRequest {
57 pub(crate) engine_config: Arc<MitoConfig>,
58 pub(crate) current_version: CompactionVersion,
59 pub(crate) access_layer: AccessLayerRef,
60 pub(crate) request_sender: mpsc::Sender<WorkerRequestWithTime>,
62 pub(crate) start_time: Instant,
64 pub(crate) cache_manager: CacheManagerRef,
65 pub(crate) manifest_ctx: ManifestContextRef,
66 pub(crate) listener: WorkerListener,
67 pub(crate) schema_metadata_manager: SchemaMetadataManagerRef,
68 pub(crate) max_parallelism: usize,
69}
70
71impl CompactionRequest {
72 pub(crate) fn region_id(&self) -> RegionId {
73 self.current_version.metadata.region_id
74 }
75}
76
77pub(crate) enum CompactionPlanningResult {
79 Prepared(PreparedCompaction),
80 NoPlan,
81 Error(Arc<Error>),
82}
83
84impl fmt::Debug for CompactionPlanningResult {
85 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
86 match self {
87 Self::Prepared(prepared) => f
88 .debug_tuple("Prepared")
89 .field(&prepared.compaction_region.region_id)
90 .finish(),
91 Self::NoPlan => f.write_str("NoPlan"),
92 Self::Error(err) => f.debug_tuple("Error").field(err).finish(),
93 }
94 }
95}
96
97#[derive(Debug)]
99pub(crate) struct CompactionPickFinished {
100 pub(crate) region_id: RegionId,
101 pub(crate) plan_id: u64,
102 pub(crate) result: CompactionPlanningResult,
103}
104
105pub(crate) struct PreparedCompaction {
106 pub(super) compaction_region: CompactionRegion,
107 pub(super) picker_output: PickerOutput,
108 start_time: Instant,
109 ttl: TimeToLive,
110}
111
112impl CompactionScheduler {
113 pub(super) fn dispatch_compaction_planning(
114 &self,
115 plan_id: u64,
116 request: CompactionRequest,
117 options: compact_request::Options,
118 time_range: Option<TimestampRange>,
119 ) {
120 let plugins = self.plugins.clone();
121 let max_background_compactions = self.engine_config.max_background_compactions;
122 common_runtime::spawn_compact(async move {
123 let region_id = request.region_id();
124 let request_sender = request.request_sender.clone();
125 let planning = Self::prepare_compaction(
126 request,
127 options,
128 plugins,
129 max_background_compactions,
130 time_range,
131 );
132 Self::notify_planning_result(region_id, plan_id, request_sender, planning).await;
133 });
134 }
135
136 pub(super) async fn notify_planning_result(
145 region_id: RegionId,
146 plan_id: u64,
147 request_sender: Sender<WorkerRequestWithTime>,
148 planning: impl Future<Output = CompactionPlanningResult> + Send,
149 ) {
150 let result = std::panic::AssertUnwindSafe(planning).catch_unwind().await.unwrap_or_else(|payload| {
152 let reason = if let Some(message) = payload.as_ref().downcast_ref::<&str>() {
153 message.to_string()
154 } else if let Some(message) = payload.as_ref().downcast_ref::<String>() {
155 message.clone()
156 } else {
157 "unknown panic".to_string()
158 };
159 CompactionPlanningResult::Error(Arc::new(
160 UnexpectedSnafu {
161 reason: format!(
162 "Compaction planning panicked for region {region_id}, plan_id {plan_id}: {reason}"
163 ),
164 }
165 .build(),
166 ))
167 });
168 if let CompactionPlanningResult::Error(err) = &result {
169 error!(err; "Compaction planning failed for region {}, plan_id: {}", region_id, plan_id);
170 }
171 let request = WorkerRequestWithTime::new(WorkerRequest::Background {
172 region_id,
173 notify: BackgroundNotify::CompactionPickFinished(CompactionPickFinished {
174 region_id,
175 plan_id,
176 result,
177 }),
178 });
179 if request_sender.send(request).await.is_err() {
180 warn!("Failed to send compaction planning result for region {region_id}");
181 }
182 }
183
184 async fn prepare_compaction(
185 request: CompactionRequest,
186 options: compact_request::Options,
187 plugins: Plugins,
188 max_background_compactions: usize,
189 time_range: Option<TimestampRange>,
190 ) -> CompactionPlanningResult {
191 let region_id = request.region_id();
192 let (dynamic_compaction_opts, ttl) = find_dynamic_options(
193 region_id,
194 &request.current_version.options,
195 &request.schema_metadata_manager,
196 )
197 .await
198 .unwrap_or_else(|e| {
199 warn!(e; "Failed to find dynamic options for region: {}", region_id);
200 (
201 request.current_version.options.compaction.clone(),
202 request.current_version.options.ttl.unwrap_or_default(),
203 )
204 });
205
206 let max_picker_outputs = max_background_compactions.min(request.max_parallelism.max(1));
209 let region_options = RegionOptions {
210 compaction: dynamic_compaction_opts,
211 ..request.current_version.options.clone()
212 };
213 let picker = new_picker(
214 &options,
215 ®ion_options,
216 Some(max_picker_outputs),
217 time_range,
218 );
219 let region_id = request.region_id();
220 let CompactionRequest {
221 engine_config,
222 current_version,
223 access_layer,
224 request_sender: _,
225 start_time,
226 cache_manager,
227 manifest_ctx,
228 listener,
229 schema_metadata_manager: _,
230 max_parallelism,
231 } = request;
232
233 debug!(
234 "Pick compaction strategy {:?} for region: {}, ttl: {:?}",
235 picker, region_id, ttl
236 );
237
238 let compaction_region = CompactionRegion {
239 region_id,
240 current_version: current_version.clone(),
241 region_options,
242 engine_config: engine_config.clone(),
243 region_metadata: current_version.metadata.clone(),
244 cache_manager: cache_manager.clone(),
245 access_layer: access_layer.clone(),
246 manifest_ctx: manifest_ctx.clone(),
247 file_purger: None,
248 ttl: Some(ttl),
249 max_parallelism,
250 plugins,
251 };
252
253 listener.on_compaction_pick_begin(region_id).await;
254 let _pick_timer = COMPACTION_STAGE_ELAPSED
255 .with_label_values(&["pick"])
256 .start_timer();
257 let picker_output = match picker.pick(&compaction_region).await {
258 Ok(output) => output,
259 Err(err) => return CompactionPlanningResult::Error(Arc::new(err)),
260 };
261
262 let Some(picker_output) = picker_output else {
263 return CompactionPlanningResult::NoPlan;
264 };
265
266 CompactionPlanningResult::Prepared(PreparedCompaction {
267 compaction_region,
268 picker_output,
269 start_time,
270 ttl,
271 })
272 }
273
274 pub(super) async fn handle_compaction_pick_finished_inner(
283 &mut self,
284 finished: CompactionPickFinished,
285 manifest_ctx: &ManifestContextRef,
286 schema_metadata_manager: SchemaMetadataManagerRef,
287 ) -> CompactionTransition {
288 let region_id = finished.region_id;
289 let plan_id = finished.plan_id;
290 let Some(status) = self.region_status.get(®ion_id) else {
291 return CompactionTransition::NoAction;
292 };
293
294 if !status.is_picking(finished.plan_id) {
295 return CompactionTransition::NoAction;
296 }
297 if !status.accept_plan(finished.plan_id) {
300 return CompactionTransition::from_pending_ddls(
301 self.remove_region_on_cancel(region_id),
302 );
303 }
304
305 match finished.result {
306 CompactionPlanningResult::Prepared(mut prepared) => {
307 let current = status.version_control.current().version;
308 let Some(picker_output) =
309 refresh_picker_output(prepared.picker_output, ¤t.ssts)
310 else {
311 return self
312 .finish_compaction_planning(
313 region_id,
314 None,
315 manifest_ctx,
316 schema_metadata_manager,
317 )
318 .await;
319 };
320 let Some(files) = CompactingFiles::try_new(&picker_output) else {
321 return self
322 .finish_compaction_planning(
323 region_id,
324 None,
325 manifest_ctx,
326 schema_metadata_manager,
327 )
328 .await;
329 };
330 prepared.picker_output = picker_output;
331 let Some(status) = self.region_status.get_mut(®ion_id) else {
332 return CompactionTransition::NoAction;
333 };
334 let waiters = status.take_waiters();
335 match self
336 .submit_prepared_compaction(prepared, files, waiters, plan_id)
337 .await
338 {
339 Ok(Some(phase)) => {
340 if let Some(status) = self.region_status.get_mut(®ion_id) {
341 status.set_phase(phase);
342 }
343 CompactionTransition::NoAction
346 }
347 Ok(None) => {
348 self.finish_compaction_planning(
349 region_id,
350 None,
351 manifest_ctx,
352 schema_metadata_manager,
353 )
354 .await
355 }
356 Err(err) => {
357 self.remove_region_on_failure(region_id, Arc::new(err));
358 CompactionTransition::NoAction
359 }
360 }
361 }
362 CompactionPlanningResult::NoPlan => {
365 self.finish_compaction_planning(
366 region_id,
367 None,
368 manifest_ctx,
369 schema_metadata_manager,
370 )
371 .await
372 }
373 CompactionPlanningResult::Error(err) => {
374 self.finish_compaction_planning(
375 region_id,
376 Some(err),
377 manifest_ctx,
378 schema_metadata_manager,
379 )
380 .await
381 }
382 }
383 }
384
385 async fn finish_compaction_planning(
386 &mut self,
387 region_id: RegionId,
388 err: Option<Arc<Error>>,
389 manifest_ctx: &ManifestContextRef,
390 schema_metadata_manager: SchemaMetadataManagerRef,
391 ) -> CompactionTransition {
392 let Some(status) = self.region_status.get_mut(®ion_id) else {
393 return CompactionTransition::NoAction;
394 };
395 for waiter in status.take_waiters() {
396 if let Some(err) = &err {
397 waiter.send(Err(err.clone()).context(CompactRegionSnafu { region_id }));
398 } else {
399 waiter.send(Ok(0));
400 }
401 }
402
403 if self.handle_pending_compaction_request(
404 region_id,
405 manifest_ctx,
406 schema_metadata_manager.clone(),
407 ) {
408 return CompactionTransition::NoAction;
409 }
410
411 let Some(status) = self.region_status.get_mut(®ion_id) else {
412 return CompactionTransition::NoAction;
413 };
414
415 let pending_ddls = std::mem::take(&mut status.pending_ddl_requests);
418 if !pending_ddls.is_empty() {
419 self.region_status.remove(®ion_id);
420 return CompactionTransition::DdlReady(pending_ddls);
421 }
422
423 if status.active.reset_automatic_followup()
424 && self.schedule_automatic_followup(region_id, manifest_ctx, schema_metadata_manager)
425 {
426 return CompactionTransition::AutomaticFollowupScheduled;
427 }
428
429 self.region_status.remove(®ion_id);
430 CompactionTransition::NoAction
431 }
432
433 async fn submit_prepared_compaction(
434 &mut self,
435 prepared: PreparedCompaction,
436 files: CompactingFiles,
437 waiters: Vec<OutputTx>,
438 mut plan_id: u64,
439 ) -> Result<Option<CompactionPhase>> {
440 let PreparedCompaction {
441 compaction_region,
442 picker_output,
443 start_time,
444 ttl,
445 } = prepared;
446 let region_id = compaction_region.region_id;
447 let dynamic_compaction_opts = &compaction_region.region_options.compaction;
448
449 let waiters = if dynamic_compaction_opts.remote_compaction() {
452 if let Some(remote_job_scheduler) = &self.plugins.get::<RemoteJobSchedulerRef>() {
453 let execution = CompactionExecution::new(plan_id, files.clone());
454 let remote_compaction_job = CompactionJob {
455 compaction_region: compaction_region.clone(),
456 picker_output: picker_output.clone(),
457 start_time,
458 waiters,
459 ttl,
460 };
461
462 let result = remote_job_scheduler
463 .schedule(
464 RemoteJob::CompactionJob(remote_compaction_job),
465 Box::new(DefaultNotifier::new(
466 self.request_sender.clone(),
467 execution.clone(),
468 )),
469 )
470 .await;
471
472 match result {
473 Ok(job_id) => {
474 info!(
475 "Scheduled remote compaction job {} for region {}",
476 job_id, region_id
477 );
478 INFLIGHT_COMPACTION_COUNT.inc();
479 return Ok(Some(CompactionPhase::Remote { execution }));
480 }
481 Err(e) => {
482 if !dynamic_compaction_opts.fallback_to_local() {
483 error!(e; "Failed to schedule remote compaction job for region {}", region_id);
484 if let Some(status) = self.region_status.get_mut(®ion_id) {
485 status.extend_waiters(e.waiters);
486 }
487 return RemoteCompactionSnafu {
488 region_id,
489 job_id: None,
490 reason: e.reason,
491 }
492 .fail();
493 }
494
495 error!(e; "Failed to schedule remote compaction job for region {}, fallback to local compaction", region_id);
496 plan_id = Self::next_plan_id(&mut self.next_plan_id);
499 e.waiters
500 }
501 }
502 } else {
503 debug!(
504 "Remote compaction is not enabled, fallback to local compaction for region {}",
505 region_id
506 );
507 waiters
508 }
509 } else {
510 waiters
511 };
512
513 let estimated_bytes = estimate_compaction_bytes(&picker_output);
515 if let Some(limit_bytes) = self.exceeds_compaction_memory_limit(estimated_bytes) {
516 COMPACTION_MEMORY_REJECTED
517 .with_label_values(&["oversized"])
518 .inc();
519 warn!(
520 "Skip compaction for region {} because estimated memory {} bytes exceeds compaction memory limit {} bytes",
521 region_id, estimated_bytes, limit_bytes,
522 );
523 for waiter in waiters {
524 waiter.send(Ok(0));
525 }
526 return Ok(None);
527 }
528
529 let state = CancellableTaskState::new();
530 let cancel_handle = state.cancel_handle();
531 let execution = CompactionExecution::new(plan_id, files);
532 let uncommitted = UncommittedSsts::new(
533 region_id,
534 compaction_region.access_layer.clone(),
535 Some(compaction_region.cache_manager.clone()),
536 );
537 let local_compaction_task = Box::new(CompactionTaskImpl {
538 state: state.clone(),
539 execution: execution.clone(),
540 request_sender: self.request_sender.clone(),
541 waiters,
542 start_time,
543 listener: self.listener.clone(),
544 picker_output,
545 compaction_region,
546 compactor: Arc::new(DefaultCompactor::with_cancel_handle(
547 cancel_handle.clone(),
548 uncommitted.clone(),
549 )),
550 memory_manager: self.memory_manager.clone(),
551 memory_policy: self.memory_policy,
552 estimated_memory_bytes: estimated_bytes,
553 uncommitted,
554 });
555
556 match self.submit_compaction_task(local_compaction_task, region_id) {
557 Ok(()) => Ok(Some(CompactionPhase::Local { state, execution })),
558 Err((err, task)) => {
559 if let (Some(status), Some(mut task)) =
560 (self.region_status.get_mut(®ion_id), task)
561 {
562 status.append_waiters(&mut task.waiters);
563 }
564 Err(err)
565 }
566 }
567 }
568
569 fn submit_compaction_task(
570 &mut self,
571 task: Box<CompactionTaskImpl>,
572 region_id: RegionId,
573 ) -> std::result::Result<(), (Error, Option<Box<CompactionTaskImpl>>)> {
574 let task = Arc::new(Mutex::new(Some(task)));
575 let task_to_run = task.clone();
576 match self.scheduler.schedule(Box::pin(async move {
577 let task = task_to_run
578 .lock()
579 .unwrap_or_else(|poisoned| poisoned.into_inner())
580 .take();
581 if let Some(mut task) = task {
582 INFLIGHT_COMPACTION_COUNT.inc();
583 task.run().await;
584 INFLIGHT_COMPACTION_COUNT.dec();
585 } else {
586 error!("Compaction task was missing when the scheduled job started");
587 }
588 })) {
589 Ok(()) => Ok(()),
590 Err(err) => {
591 error!(err; "Failed to submit compaction request for region {}", region_id);
592 let task = task
593 .lock()
594 .unwrap_or_else(|poisoned| poisoned.into_inner())
595 .take();
596 Err((err, task))
597 }
598 }
599 }
600
601 fn exceeds_compaction_memory_limit(&self, estimated_bytes: u64) -> Option<u64> {
602 let limit_bytes = self.memory_manager.limit_bytes();
603 if limit_bytes > 0 && estimated_bytes > limit_bytes {
604 Some(limit_bytes)
605 } else {
606 None
607 }
608 }
609}
610
611fn estimate_compaction_bytes(picker_output: &PickerOutput) -> u64 {
614 picker_output
615 .outputs
616 .iter()
617 .flat_map(|output| output.inputs.iter())
618 .map(|file: &FileHandle| {
619 let meta = file.meta_ref();
620 meta.max_row_group_uncompressed_size
621 })
622 .sum()
623}
624
625fn refresh_picker_output(output: PickerOutput, current: &SstVersion) -> Option<PickerOutput> {
635 let refresh = |file: FileHandle| {
636 current
637 .file_for_compaction(&file)
638 .filter(|current| !current.is_deleted() && !current.compacting())
639 .cloned()
640 };
641 let outputs = output
642 .outputs
643 .into_iter()
644 .map(|output| {
645 let inputs = output
646 .inputs
647 .into_iter()
648 .map(&refresh)
649 .collect::<Option<Vec<_>>>()?;
650 Some(CompactionOutput { inputs, ..output })
651 })
652 .collect::<Option<Vec<_>>>()?;
653 let expired_ssts = output
654 .expired_ssts
655 .into_iter()
656 .map(refresh)
657 .collect::<Option<Vec<_>>>()?;
658
659 Some(PickerOutput {
660 outputs,
661 expired_ssts,
662 time_window_size: output.time_window_size,
663 max_file_size: output.max_file_size,
664 })
665}