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