Skip to main content

mito2/worker/
handle_compaction.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 api::v1::region::compact_request;
16use common_telemetry::{debug, error, info};
17use store_api::logstore::LogStore;
18use store_api::region_request::RegionCompactRequest;
19use store_api::storage::RegionId;
20
21use crate::compaction::{CompactionPickFinished, CompactionTransition};
22use crate::config::IndexBuildMode;
23use crate::error::{RegionNotFoundSnafu, StaleCompactionExecutionSnafu};
24use crate::metrics::COMPACTION_REQUEST_COUNT;
25use crate::region::MitoRegionRef;
26use crate::request::{
27    BuildIndexRequest, CompactionCancelled, CompactionFailed, CompactionFinished, OnFailure,
28    OptionOutputTx,
29};
30use crate::sst::index::IndexBuildType;
31use crate::worker::RegionWorkerLoop;
32
33fn made_progress(files_to_add: usize, files_to_remove: usize) -> bool {
34    files_to_add > 0 || files_to_remove > files_to_add
35}
36
37impl<S> RegionWorkerLoop<S> {
38    pub(crate) async fn handle_compaction_pick_finished(
39        &mut self,
40        region_id: RegionId,
41        request: CompactionPickFinished,
42    ) where
43        S: LogStore,
44    {
45        let Some(region) = self.regions.get_region(region_id) else {
46            return;
47        };
48        // A terminal pick (canceled, no plan, or failed) may remove the
49        // compaction status and release DDLs fenced behind its picking phase.
50        // Such a pick never produces an execution callback, so the worker must
51        // execute the returned DDLs from this notification.
52        let transition = self
53            .compaction_scheduler
54            .handle_compaction_pick_finished(
55                request,
56                &region.manifest_ctx,
57                self.schema_metadata_manager.clone(),
58            )
59            .await;
60        match transition {
61            CompactionTransition::AutomaticFollowupScheduled => {
62                // There will be a followup compaction so we need to update the
63                // last schedule time to avoid frequent compactions.
64                region.update_schedule_compaction_millis();
65            }
66            CompactionTransition::NoAction => {}
67            CompactionTransition::DdlReady(mut pending_ddls) => {
68                if !pending_ddls.is_empty() {
69                    // Preserve the lifecycle order observed by listeners: compaction
70                    // termination is visible before its dependent DDLs are dispatched.
71                    self.listener.on_compaction_result_notified(region_id).await;
72                    self.handle_ddl_requests(&mut pending_ddls).await;
73                }
74            }
75        }
76    }
77
78    /// Handles compaction request submitted to region worker.
79    pub(crate) async fn handle_compaction_request(
80        &mut self,
81        region_id: RegionId,
82        req: RegionCompactRequest,
83        mut sender: OptionOutputTx,
84    ) {
85        let Some(region) = self.regions.writable_region_or(region_id, &mut sender) else {
86            return;
87        };
88        COMPACTION_REQUEST_COUNT.inc();
89        let parallelism = req.parallelism.unwrap_or(1) as usize;
90        match self.compaction_scheduler.schedule_manual_compaction(
91            req.options,
92            &region.version_control,
93            &region.access_layer,
94            sender,
95            &region.manifest_ctx,
96            self.schema_metadata_manager.clone(),
97            parallelism,
98            req.time_range,
99        ) {
100            // Ok(false) means the request was merged, queued or rejected; the
101            // waiter was already notified or will be completed by the queued
102            // cycle, so only a newly scheduled task is logged as success.
103            Ok(true) => info!(
104                "Successfully scheduled compaction task for region: {}",
105                region_id
106            ),
107            Ok(false) => {}
108            Err(e) => {
109                error!(e; "Failed to schedule compaction task for region: {}", region_id);
110            }
111        }
112    }
113
114    /// Handles compaction finished, update region version and manifest, deleted compacted files.
115    pub(crate) async fn handle_compaction_finished(
116        &mut self,
117        region_id: RegionId,
118        mut request: CompactionFinished,
119    ) where
120        S: LogStore,
121    {
122        let region = match self.regions.get_region(region_id) {
123            Some(region) => region,
124            None => {
125                request.on_failure(RegionNotFoundSnafu { region_id }.build());
126                return;
127            }
128        };
129        // Reject stale terminal results before applying their manifest edit.
130        if !self
131            .compaction_scheduler
132            .is_current_execution(region_id, &request.execution)
133        {
134            request.on_failure(StaleCompactionExecutionSnafu { region_id }.build());
135            return;
136        }
137        let execution = request.execution.clone();
138        let made_progress = made_progress(
139            request.edit.files_to_add.len(),
140            request.edit.files_to_remove.len(),
141        );
142
143        region.version_control.apply_edit(
144            Some(request.edit.clone()),
145            &[],
146            region.file_purger.clone(),
147        );
148
149        let index_build_file_metas = std::mem::take(&mut request.edit.files_to_add);
150
151        // compaction finished.
152        request.on_success();
153        self.listener.on_compaction_result_notified(region_id).await;
154
155        // In async mode, create indexes after compact if new files are created.
156        if self.config.index.build_mode == IndexBuildMode::Async
157            && !index_build_file_metas.is_empty()
158        {
159            self.handle_rebuild_index(
160                BuildIndexRequest {
161                    region_id,
162                    build_type: IndexBuildType::Compact,
163                    file_metas: index_build_file_metas,
164                },
165                OptionOutputTx::new(None),
166            )
167            .await;
168        }
169
170        // Schedule next compaction if necessary.
171        let transition = self
172            .compaction_scheduler
173            .on_execution_finished(
174                region_id,
175                &execution,
176                &region.manifest_ctx,
177                self.schema_metadata_manager.clone(),
178                made_progress,
179            )
180            .await;
181        match transition {
182            CompactionTransition::AutomaticFollowupScheduled => {
183                region.update_schedule_compaction_millis();
184            }
185            CompactionTransition::NoAction => {}
186            CompactionTransition::DdlReady(mut pending_ddls) => {
187                self.handle_ddl_requests(&mut pending_ddls).await;
188            }
189        }
190    }
191
192    pub(crate) async fn handle_compaction_cancelled(
193        &mut self,
194        region_id: RegionId,
195        request: CompactionCancelled,
196    ) where
197        S: LogStore,
198    {
199        let execution = request.execution.clone();
200        let is_current = self.regions.get_region(region_id).is_some_and(|_| {
201            self.compaction_scheduler
202                .is_current_execution(region_id, &execution)
203        });
204        request.on_success();
205
206        if !is_current {
207            return;
208        }
209
210        // Reuse the scheduler's finish path to wake pending DDLs after a cooperative stop.
211        let mut pending_ddls = self
212            .compaction_scheduler
213            .on_execution_cancelled(region_id, &execution)
214            .await;
215        if !pending_ddls.is_empty() {
216            self.listener.on_compaction_result_notified(region_id).await;
217        }
218
219        self.handle_ddl_requests(&mut pending_ddls).await;
220    }
221
222    /// When compaction fails, we simply log the error.
223    pub(crate) async fn handle_compaction_failure(&mut self, req: CompactionFailed) {
224        if self.regions.get_region(req.region_id).is_none() {
225            return;
226        }
227        if !self
228            .compaction_scheduler
229            .is_current_execution(req.region_id, &req.execution)
230        {
231            debug!(
232                "Ignores stale compaction failure for region {}: {:?}",
233                req.region_id, req.err
234            );
235            return;
236        }
237
238        error!(req.err; "Failed to compact region: {}", req.region_id);
239        self.compaction_scheduler
240            .on_execution_failed(req.region_id, &req.execution, req.err);
241    }
242
243    /// Schedule compaction for the region if necessary.
244    pub(crate) async fn schedule_compaction(&mut self, region: &MitoRegionRef) {
245        if region.is_staging() || region.is_enter_staging() {
246            info!(
247                "Region {} is staging or entering staging, skip compaction",
248                region.region_id
249            );
250            return;
251        }
252        let now = self.time_provider.current_time_millis();
253        if now - region.last_schedule_compaction_millis()
254            >= self.config.min_compaction_interval.as_millis() as i64
255        {
256            debug!(
257                "minimal compaction interval time {:?} has passed, scheduling next compaction",
258                self.config.min_compaction_interval
259            );
260            match self.compaction_scheduler.schedule_automatic_compaction(
261                compact_request::Options::Regular(Default::default()),
262                &region.version_control,
263                &region.access_layer,
264                &region.manifest_ctx,
265                self.schema_metadata_manager.clone(),
266            ) {
267                Ok(true) => region.update_schedule_compaction_millis(),
268                Ok(false) => {}
269                Err(e) => {
270                    error!(e; "Failed to schedule compaction for region: {}", region.region_id)
271                }
272            }
273        }
274    }
275}
276
277#[cfg(test)]
278mod tests {
279    use super::made_progress;
280
281    #[test]
282    fn test_nonempty_output_or_file_reduction_is_progress() {
283        for (files_to_add, files_to_remove, expected) in [
284            (3, 3, true),  // Equal-count rewrite.
285            (1, 3, true),  // Ordinary reduction.
286            (0, 3, true),  // Zero-output reduction.
287            (0, 0, false), // Empty edit.
288            (3, 2, true),  // Growth rewrite with output.
289        ] {
290            assert_eq!(
291                expected,
292                made_progress(files_to_add, files_to_remove),
293                "files_to_remove: {files_to_remove}, files_to_add: {files_to_add}"
294            );
295        }
296    }
297}