Skip to main content

mito2/worker/
handle_truncate.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
15//! Handling truncate related requests.
16
17use common_telemetry::{debug, info, warn};
18use store_api::logstore::LogStore;
19use store_api::region_request::RegionTruncateRequest;
20use store_api::storage::RegionId;
21
22use crate::error::RegionNotFoundSnafu;
23use crate::manifest::action::{RegionTruncate, TruncateKind};
24use crate::region::RegionLeaderState;
25use crate::request::{DdlRequest, DiscardUnflushedResult, OptionOutputTx, TruncateResult};
26use crate::worker::RegionWorkerLoop;
27
28impl<S: LogStore> RegionWorkerLoop<S> {
29    pub(crate) async fn handle_truncate_request(
30        &mut self,
31        region_id: RegionId,
32        req: RegionTruncateRequest,
33        sender: OptionOutputTx,
34    ) {
35        let region = match self.regions.writable_non_staging_region(region_id) {
36            Ok(region) => region,
37            Err(e) => {
38                sender.send(Err(e));
39                return;
40            }
41        };
42
43        let (sender, req) = match self.flush_scheduler.try_cancel_and_add_ddl(
44            region_id,
45            sender,
46            req,
47            DdlRequest::Truncate,
48        ) {
49            Ok(()) => {
50                self.listener.on_flush_cancel_requested(region_id);
51                return;
52            }
53            Err(request) => request,
54        };
55
56        let (sender, req) = match self.compaction_scheduler.try_cancel_and_add_ddl(
57            region_id,
58            sender,
59            req,
60            DdlRequest::Truncate,
61        ) {
62            Ok(()) => {
63                self.listener.on_compaction_cancel_requested(region_id);
64                return;
65            }
66            Err(request) => request,
67        };
68
69        let version_data = region.version_control.current();
70
71        match req {
72            RegionTruncateRequest::All => {
73                info!("Try to fully truncate region {}", region_id);
74
75                let truncated_entry_id = version_data.last_entry_id;
76                let truncated_sequence = version_data.committed_sequence;
77
78                // Write region truncated to manifest.
79                let truncate = RegionTruncate {
80                    region_id,
81                    kind: TruncateKind::All {
82                        truncated_entry_id,
83                        truncated_sequence,
84                    },
85                    timestamp_ms: None,
86                };
87
88                self.handle_manifest_truncate_action(region, truncate, sender);
89            }
90            RegionTruncateRequest::Unflushed => {
91                let memtables = &version_data.version.memtables;
92                if memtables.is_empty()
93                    && version_data.version.flushed_entry_id == version_data.last_entry_id
94                    && version_data.version.flushed_sequence == version_data.committed_sequence
95                {
96                    self.update_topic_latest_entry_id(&region);
97                    let result = self
98                        .wal
99                        .obsolete(
100                            region_id,
101                            version_data.version.flushed_entry_id,
102                            &region.provider,
103                        )
104                        .await
105                        .map(|_| 0);
106                    sender.send(result);
107                    return;
108                }
109
110                let discarded_rows = memtables.num_rows();
111                let discarded_bytes =
112                    (memtables.mutable_usage() + memtables.immutables_usage()) as u64;
113                warn!(
114                    "Discarding unflushed data from region {}, entry_id: {}, sequence: {}, estimated rows: {}, estimated bytes: {}",
115                    region_id,
116                    version_data.last_entry_id,
117                    version_data.committed_sequence,
118                    discarded_rows,
119                    discarded_bytes,
120                );
121
122                self.handle_manifest_discard_unflushed_action(
123                    region,
124                    version_data.last_entry_id,
125                    version_data.committed_sequence,
126                    discarded_rows,
127                    discarded_bytes,
128                    sender,
129                );
130            }
131            RegionTruncateRequest::ByTimeRanges { time_ranges } => {
132                info!(
133                    "Try to partially truncate region {} by time ranges: {:?}",
134                    region_id, time_ranges
135                );
136                // find all files that are fully contained in the time ranges
137                let mut files_to_truncate = Vec::new();
138                for level in version_data.version.ssts.levels() {
139                    for file in level.files() {
140                        let file_time_range = file.time_range();
141                        // TODO(discord9): This is a naive way to check if is contained, we should
142                        // optimize it later.
143                        let is_subset = time_ranges.iter().any(|(start, end)| {
144                            file_time_range.0 >= *start && file_time_range.1 <= *end
145                        });
146
147                        if is_subset {
148                            files_to_truncate.push(file.meta_ref().clone());
149                        }
150                    }
151                }
152                debug!(
153                    "Found {} files to partially truncate in region {}",
154                    files_to_truncate.len(),
155                    region_id
156                );
157
158                // this could happen if all files are not fully contained in the time ranges
159                if files_to_truncate.is_empty() {
160                    info!("No files to truncate in region {}", region_id);
161                    // directly send success back as no files to truncate
162                    // and no background notify is needed
163                    sender.send(Ok(0));
164                    return;
165                }
166                let truncate = RegionTruncate {
167                    region_id,
168                    kind: TruncateKind::Partial {
169                        files_to_remove: files_to_truncate,
170                    },
171                    timestamp_ms: None,
172                };
173                self.handle_manifest_truncate_action(region, truncate, sender);
174            }
175        };
176    }
177
178    /// Handles truncate result.
179    pub(crate) async fn handle_truncate_result(&mut self, truncate_result: TruncateResult) {
180        let region_id = truncate_result.region_id;
181        let Some(region) = self.regions.get_region(region_id) else {
182            truncate_result.sender.send(
183                RegionNotFoundSnafu {
184                    region_id: truncate_result.region_id,
185                }
186                .fail(),
187            );
188            return;
189        };
190
191        // We are already in the worker loop so we can set the state first.
192        region.switch_state_to_writable(RegionLeaderState::Truncating);
193
194        match truncate_result.result {
195            Ok(()) => {
196                // Applies the truncate action to the region.
197                region
198                    .version_control
199                    .truncate(truncate_result.kind.clone());
200            }
201            Err(e) => {
202                // Unable to truncate the region.
203                truncate_result.sender.send(Err(e));
204                return;
205            }
206        }
207
208        // Notifies flush scheduler.
209        self.flush_scheduler.on_region_truncated(region_id);
210        // Notifies compaction scheduler.
211        self.compaction_scheduler.on_region_truncated(region_id);
212        // Notifies index build scheduler.
213        self.index_build_scheduler
214            .on_region_truncated(region_id)
215            .await;
216
217        if let TruncateKind::All {
218            truncated_entry_id,
219            truncated_sequence: _,
220        } = &truncate_result.kind
221        {
222            // Make all data obsolete.
223            if let Err(e) = self
224                .wal
225                .obsolete(region_id, *truncated_entry_id, &region.provider)
226                .await
227            {
228                truncate_result.sender.send(Err(e));
229                return;
230            }
231        }
232
233        info!(
234            "Complete truncating region: {}, kind: {:?}.",
235            region_id, truncate_result.kind
236        );
237
238        truncate_result.sender.send(Ok(0));
239    }
240
241    /// Handles the result of advancing the replay frontier for unflushed data.
242    pub(crate) async fn handle_discard_unflushed_result(&mut self, result: DiscardUnflushedResult) {
243        let region_id = result.region_id;
244        let Some(region) = self.regions.get_region(region_id) else {
245            result.sender.send(RegionNotFoundSnafu { region_id }.fail());
246            return;
247        };
248
249        region.switch_state_to_writable(RegionLeaderState::Truncating);
250
251        if let Err(e) = result.result {
252            result.sender.send(Err(e));
253            return;
254        }
255
256        region
257            .version_control
258            .discard_unflushed(result.discarded_entry_id, result.discarded_sequence);
259
260        // The flush and compaction schedulers are already quiesced because the request
261        // queued itself behind any running job, so these two are usually no-ops. Index
262        // builds aren't queued that way, so pending builds for surviving files are
263        // retired here. They would abort anyway while the region sits in `Truncating`,
264        // and those files may need a later index rebuild.
265        self.flush_scheduler.on_region_truncated(region_id);
266        self.compaction_scheduler.on_region_truncated(region_id);
267        self.index_build_scheduler
268            .on_region_truncated(region_id)
269            .await;
270        self.update_topic_latest_entry_id(&region);
271
272        // Discarding memtables releases write buffer memory. Notify other workers and
273        // retry requests stalled on this worker before WAL cleanup, which may fail even
274        // though the memory has already been released.
275        self.notify_group();
276        self.handle_stalled_requests().await;
277
278        if let Err(e) = self
279            .wal
280            .obsolete(region_id, result.discarded_entry_id, &region.provider)
281            .await
282        {
283            result.sender.send(Err(e));
284            return;
285        }
286
287        warn!(
288            "Discarded unflushed data from region {}, entry_id: {}, sequence: {}, estimated rows: {}, estimated bytes: {}",
289            region_id,
290            result.discarded_entry_id,
291            result.discarded_sequence,
292            result.discarded_rows,
293            result.discarded_bytes,
294        );
295        result.sender.send(Ok(0));
296    }
297}