Skip to main content

mito2/worker/
handle_rebuild_index.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//! Handles index build requests.
16
17use std::collections::HashMap;
18use std::sync::Arc;
19
20use api::v1::region::build_index_request;
21use common_telemetry::{debug, error, warn};
22use store_api::region_request::RegionBuildIndexRequest;
23use store_api::storage::{FileId, RegionId};
24use tokio::sync::mpsc;
25
26use crate::cache::CacheStrategy;
27use crate::error::{InvalidRequestSnafu, Result, WorkerStoppedSnafu};
28use crate::manifest::action::RegionEdit;
29use crate::metrics::INDEX_PUBLICATION_STALE_TOTAL;
30use crate::region::version::VersionRef;
31use crate::region::{IndexBuildSource, MitoRegionRef};
32use crate::request::{
33    BuildIndexRequest, IndexBuildFailed, IndexBuildFinished, IndexBuildStopped, OptionOutputTx,
34};
35use crate::sst::file::{FileHandle, FileMeta, RegionFileId, RegionIndexId};
36use crate::sst::index::{
37    IndexBuildOutcome, IndexBuildTask, IndexBuildType, IndexerBuilderImpl, ResultMpscSender,
38    cleanup_stale_index_caches,
39};
40use crate::worker::RegionWorkerLoop;
41
42impl<S> RegionWorkerLoop<S> {
43    pub(crate) fn new_index_build_task(
44        &self,
45        region: &MitoRegionRef,
46        version: &VersionRef,
47        file: FileHandle,
48        file_meta: FileMeta,
49        build_type: IndexBuildType,
50        result_sender: ResultMpscSender,
51    ) -> IndexBuildTask {
52        let access_layer = region.access_layer.clone();
53
54        let puffin_manager = if let Some(write_cache) = self.cache_manager.write_cache() {
55            write_cache.build_puffin_manager()
56        } else {
57            access_layer.build_puffin_manager()
58        };
59
60        let intermediate_manager = if let Some(write_cache) = self.cache_manager.write_cache() {
61            write_cache.intermediate_manager().clone()
62        } else {
63            access_layer.intermediate_manager().clone()
64        };
65
66        let indexer_builder_ref = Arc::new(IndexerBuilderImpl {
67            build_type: build_type.clone(),
68            metadata: version.metadata.clone(),
69            inverted_index_config: self.config.inverted_index.clone(),
70            fulltext_index_config: self.config.fulltext_index.clone(),
71            bloom_filter_index_config: self.config.bloom_filter_index.clone(),
72            index_options: version.options.index_options.clone(),
73            intermediate_manager,
74            puffin_manager,
75            write_cache_enabled: self.cache_manager.write_cache().is_some(),
76        });
77
78        IndexBuildTask {
79            region_id: region.region_id,
80            file: file.clone(),
81            target_region_metadata: version.metadata.clone(),
82            source: IndexBuildSource::new(file_meta, version.metadata.schema_version),
83            reason: build_type,
84            access_layer: access_layer.clone(),
85            listener: self.listener.clone(),
86            manifest_ctx: region.manifest_ctx.clone(),
87            write_cache: self.cache_manager.write_cache().cloned(),
88            cache_manager: Some(self.cache_manager.clone()),
89            file_purger: file.file_purger(),
90            request_sender: self.sender.clone(),
91            indexer_builder: indexer_builder_ref.clone(),
92            result_sender,
93        }
94    }
95
96    /// Handles manual build index requests.
97    pub(crate) async fn handle_build_index_request(
98        &mut self,
99        region_id: RegionId,
100        req: RegionBuildIndexRequest,
101        sender: OptionOutputTx,
102    ) {
103        match req.options {
104            Some(build_index_request::Options::SeriesIndex(_)) => {
105                self.handle_build_series_index_request(region_id, sender);
106            }
107            None | Some(build_index_request::Options::SstIndex(_)) => {
108                self.handle_rebuild_index(
109                    BuildIndexRequest {
110                        region_id,
111                        build_type: IndexBuildType::Manual,
112                        file_metas: Vec::new(),
113                    },
114                    sender,
115                )
116                .await;
117            }
118        }
119    }
120
121    /// Submits manual reconciliation without blocking the region worker on index I/O.
122    fn handle_build_series_index_request(&self, region_id: RegionId, mut sender: OptionOutputTx) {
123        let Some(region) = self.regions.writable_region_or(region_id, &mut sender) else {
124            return;
125        };
126        let Some(state) = &self.series_index_task_state else {
127            sender.send(
128                InvalidRequestSnafu {
129                    region_id,
130                    reason: "series index is disabled; enable experimental_enable_series_index",
131                }
132                .fail(),
133            );
134            return;
135        };
136        let receiver = match state.try_reconcile(region) {
137            Ok(receiver) => receiver,
138            Err(error) => {
139                sender.send(Err(error));
140                return;
141            }
142        };
143        let worker_id = self.id;
144        common_runtime::spawn_global(async move {
145            let result = receiver
146                .await
147                .unwrap_or_else(|_| WorkerStoppedSnafu { id: worker_id }.fail());
148            sender.send(result.map(|_| 0));
149        });
150    }
151
152    pub(crate) async fn handle_rebuild_index(
153        &mut self,
154        request: BuildIndexRequest,
155        mut sender: OptionOutputTx,
156    ) {
157        let region_id = request.region_id;
158        let Some(region) = self.regions.writable_region_or(region_id, &mut sender) else {
159            return;
160        };
161
162        let version_control = region.version_control.clone();
163        let version = version_control.current().version;
164        // A committed index publication may not have reached version control
165        // yet. Use the manifest's FileMeta as the conditional publication
166        // source while keeping the builder's schema generation from `version`.
167        let manifest = region.manifest_ctx.manifest().await;
168        let current_file_meta = |file_id| {
169            let file_meta = manifest.files.get(&file_id);
170            if file_meta.is_none() {
171                debug!(
172                    "Skipping index build because file is absent from manifest, region: {}, file_id: {}",
173                    region_id, file_id
174                );
175            }
176            file_meta.cloned()
177        };
178
179        let all_files: HashMap<FileId, FileHandle> = version
180            .ssts
181            .levels()
182            .iter()
183            .flat_map(|level| level.files.iter())
184            .filter(|(_, handle)| !handle.is_deleted() && !handle.compacting())
185            .map(|(id, handle)| (*id, handle.clone()))
186            .collect();
187
188        let build_tasks = if request.file_metas.is_empty() {
189            // If no specific files are provided, find files whose index is inconsistent with the region metadata.
190            all_files
191                .values()
192                .filter_map(|file| {
193                    let file_meta = current_file_meta(file.meta_ref().file_id)?;
194                    (!file_meta.is_index_consistent_with_region(&version.metadata.column_metadatas))
195                        .then(|| (file.clone(), file_meta))
196                })
197                .collect::<Vec<_>>()
198        } else {
199            request
200                .file_metas
201                .iter()
202                .filter_map(|meta| {
203                    let file = all_files.get(&meta.file_id)?;
204                    let file_meta = current_file_meta(meta.file_id)?;
205                    Some((file.clone(), file_meta))
206                })
207                .collect::<Vec<_>>()
208        };
209
210        if build_tasks.is_empty() {
211            debug!(
212                "No files need to build index for region {}, request: {:?}",
213                region_id, request
214            );
215            sender.send(Ok(0));
216            return;
217        }
218
219        let num_tasks = build_tasks.len();
220        let (tx, mut rx) = mpsc::channel::<Result<IndexBuildOutcome>>(num_tasks);
221
222        for (file_handle, file_meta) in build_tasks {
223            debug!(
224                "Scheduling index build for region {}, file_id {}",
225                region_id,
226                file_handle.meta_ref().file_id
227            );
228
229            if region.should_abort_index() {
230                warn!(
231                    "Region {} is in state {:?}, abort index rebuild process for file_id {}",
232                    region_id,
233                    region.state(),
234                    file_handle.meta_ref().file_id
235                );
236                break;
237            }
238
239            let task = self.new_index_build_task(
240                &region,
241                &version,
242                file_handle.clone(),
243                file_meta,
244                request.build_type.clone(),
245                tx.clone(),
246            );
247            let _ = self
248                .index_build_scheduler
249                .schedule_build(&region.version_control, task)
250                .await;
251        }
252        // Wait for all index build tasks to finish and notify the caller.
253        common_runtime::spawn_global(async move {
254            for _ in 0..num_tasks {
255                if let Some(Err(e)) = rx.recv().await {
256                    warn!(e; "Index build task failed for region: {}", region_id);
257                    sender.send(Err(e));
258                    return;
259                }
260            }
261            sender.send(Ok(0));
262        });
263    }
264
265    pub(crate) async fn handle_index_build_finished(
266        &mut self,
267        region_id: RegionId,
268        request: IndexBuildFinished,
269    ) {
270        let region = match self.regions.get_region(region_id) {
271            Some(region) => region,
272            None => {
273                warn!(
274                    "Region not found for index build finished, region_id: {}",
275                    region_id
276                );
277                return;
278            }
279        };
280
281        let file_meta = request.file_meta;
282        let region_file_id = RegionFileId::new(file_meta.region_id, file_meta.file_id);
283        let index_id = RegionIndexId::new(region_file_id, file_meta.index_version);
284
285        // Clean old puffin-related cache before making the new metadata visible.
286        let cache_strategy = CacheStrategy::EnableAll(self.cache_manager.clone());
287        cache_strategy.evict_puffin_cache(index_id).await;
288
289        let manager = region.manifest_ctx.manifest_manager.read().await;
290        let manifest = manager.manifest();
291        let is_current = manifest.files.get(&file_meta.file_id) == Some(&file_meta);
292        if is_current {
293            region.version_control.apply_edit(
294                Some(RegionEdit {
295                    files_to_add: vec![file_meta.clone()],
296                    files_to_remove: Vec::new(),
297                    timestamp_ms: None,
298                    flushed_sequence: None,
299                    flushed_entry_id: None,
300                    committed_sequence: None,
301                    compaction_time_window: None,
302                }),
303                &[],
304                region.file_purger.clone(),
305            );
306        }
307        drop(manager);
308
309        if !is_current {
310            INDEX_PUBLICATION_STALE_TOTAL
311                .with_label_values(&["worker_apply"])
312                .inc();
313            warn!(
314                "Ignores stale index build result, region: {}, file_id: {}, index_version: {}, committed manifest version: {}",
315                region_id, file_meta.file_id, file_meta.index_version, request.manifest_version
316            );
317            cleanup_stale_index_caches(
318                index_id,
319                &region.access_layer,
320                Some(&self.cache_manager),
321                self.cache_manager.write_cache(),
322            )
323            .await;
324            self.listener.on_index_build_abort(region_file_id).await;
325            return;
326        }
327
328        self.listener.on_index_build_finish(region_file_id).await;
329    }
330
331    pub(crate) async fn handle_index_build_failed(
332        &mut self,
333        region_id: RegionId,
334        request: IndexBuildFailed,
335    ) {
336        error!(request.err; "Index build failed for region: {}", region_id);
337        self.index_build_scheduler
338            .on_failure(region_id, request.err.clone())
339            .await;
340    }
341
342    pub(crate) async fn handle_index_build_stopped(
343        &mut self,
344        region_id: RegionId,
345        request: IndexBuildStopped,
346    ) {
347        self.index_build_scheduler
348            .on_task_stopped(region_id, request.file_id);
349    }
350}