Skip to main content

mito2/compaction/scheduler/
planning.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use 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
55/// Region compaction request.
56pub struct CompactionRequest {
57    pub(crate) engine_config: Arc<MitoConfig>,
58    pub(crate) current_version: CompactionVersion,
59    pub(crate) access_layer: AccessLayerRef,
60    /// Sender to send notification to the region worker.
61    pub(crate) request_sender: mpsc::Sender<WorkerRequestWithTime>,
62    /// Start time of compaction task.
63    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
77/// Result returned to the worker after background compaction planning.
78pub(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/// Pure planning completion sent back to the owning region worker.
98#[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    /// Runs the planning future and always sends the planning result back to
137    /// the worker, even if the planning panics.
138    ///
139    /// The worker only leaves the picking phase after it receives the
140    /// `CompactionPickFinished` notification. If a panicked planning task
141    /// swallowed the notification, the region would be stuck in the picking
142    /// phase forever, blocking all future compactions and pending DDLs (e.g.
143    /// entering staging) of the region.
144    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        // The idiomatic way to handle a panic result.
151        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        // Avoid queuing serial outputs from a stale snapshot. Repicking after each
207        // batch lets newly flushed L0 files outrank L1 work that has not started.
208        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            &region_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    /// Applies a background planning result to the current compaction lifecycle.
275    ///
276    /// # Returns
277    ///
278    /// Reports an automatic follow-up dispatched after Picking, or DDL requests
279    /// released when Picking terminates without creating an execution. The
280    /// owning worker must dispatch returned DDLs immediately because no
281    /// execution callback will follow.
282    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(&region_id) else {
291            return CompactionTransition::NoAction;
292        };
293
294        if !status.is_picking(finished.plan_id) {
295            return CompactionTransition::NoAction;
296        }
297        // Cancellation during Picking is completed by this notification. No
298        // execution callback will follow, so return DDLs released by the fence.
299        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, &current.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(&region_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(&region_id) {
341                            status.set_phase(phase);
342                        }
343                        // The execution now owns the terminal transition. Any
344                        // fenced DDLs remain queued until its callback arrives.
345                        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            // These paths never create an execution. Finish the Picking
363            // lifecycle here and release DDLs if no follow-up was scheduled.
364            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(&region_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(&region_id) else {
412            return CompactionTransition::NoAction;
413        };
414
415        // A queued DDL supersedes a retained automatic follow-up, matching the
416        // execution terminal path in `on_compaction_finished`.
417        let pending_ddls = std::mem::take(&mut status.pending_ddl_requests);
418        if !pending_ddls.is_empty() {
419            self.region_status.remove(&region_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(&region_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        // If specified to run compaction remotely, we schedule the compaction job remotely.
450        // It will fall back to local compaction if there is no remote job scheduler.
451        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(&region_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                        // An error may be ambiguous after the remote scheduler consumed
497                        // the notifier. Fence a delayed remote callback from the local fallback.
498                        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        // Check whether this local compaction can ever fit before submitting it.
514        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(&region_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
611/// Estimates compaction memory as the sum of all input files' maximum row-group
612/// uncompressed sizes.
613fn 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
625/// Rebuilds picker output with current SST handles while preserving the picker's grouping.
626///
627/// Picking runs in background on a version snapshot that may be stale by the
628/// time the plan is accepted: a concurrent flush, compaction or index rebuild
629/// can replace a selected file with a new handle carrying updated metadata
630/// (e.g. `index_version`), or remove the file entirely. The handles in the
631/// picker output therefore cannot be used as-is; re-resolving them against the
632/// current version both detects gone files (aborting the plan) and ensures the
633/// execution reads and reserves the up-to-date handle.
634fn 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}