Skip to main content

mito2/sst/index/indexer/
finish.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 common_telemetry::{debug, warn};
16use puffin::puffin_manager::{PuffinManager, PuffinWriter};
17use store_api::storage::ColumnId;
18
19use crate::sst::file::{RegionFileId, RegionIndexId};
20use crate::sst::index::puffin_manager::SstPuffinWriter;
21use crate::sst::index::statistics::{ByteCount, RowCount};
22use crate::sst::index::{
23    BloomFilterOutput, FulltextIndexOutput, IndexOutput, Indexer, InvertedIndexOutput,
24};
25
26impl Indexer {
27    pub(crate) async fn do_finish(&mut self) -> IndexOutput {
28        self.dense_pk_decoder = None;
29        let mut output = IndexOutput::default();
30
31        let Some(mut writer) = self.build_puffin_writer().await else {
32            self.do_abort().await;
33            return output;
34        };
35
36        let success = self
37            .do_finish_inverted_index(&mut writer, &mut output)
38            .await;
39        if !success {
40            self.do_abort().await;
41            return IndexOutput::default();
42        }
43
44        let success = self
45            .do_finish_fulltext_index(&mut writer, &mut output)
46            .await;
47        if !success {
48            self.do_abort().await;
49            return IndexOutput::default();
50        }
51
52        let success = self.do_finish_bloom_filter(&mut writer, &mut output).await;
53        if !success {
54            self.do_abort().await;
55            return IndexOutput::default();
56        }
57
58        self.do_prune_intm_sst_dir().await;
59        output.file_size = self.do_finish_puffin_writer(writer).await;
60        output.version = self.index_version;
61        output
62    }
63
64    async fn build_puffin_writer(&mut self) -> Option<SstPuffinWriter> {
65        let puffin_manager = self.puffin_manager.clone()?;
66
67        let err = match puffin_manager
68            .writer(&RegionIndexId::new(
69                RegionFileId::new(self.physical_region_id, self.file_id),
70                self.index_version,
71            ))
72            .await
73        {
74            Ok(writer) => return Some(writer),
75            Err(err) => err,
76        };
77
78        if cfg!(any(test, feature = "test")) {
79            panic!(
80                "Failed to create puffin writer, region_id: {}, file_id: {}, err: {:?}",
81                self.region_id, self.file_id, err
82            );
83        } else {
84            warn!(
85                err; "Failed to create puffin writer, region_id: {}, file_id: {}",
86                self.region_id, self.file_id,
87            );
88        }
89
90        None
91    }
92
93    async fn do_finish_puffin_writer(&mut self, writer: SstPuffinWriter) -> ByteCount {
94        let err = match writer.finish().await {
95            Ok(size) => return size,
96            Err(err) => err,
97        };
98
99        if cfg!(any(test, feature = "test")) {
100            panic!(
101                "Failed to finish puffin writer, region_id: {}, file_id: {}, err: {:?}",
102                self.region_id, self.file_id, err
103            );
104        } else {
105            warn!(
106                err; "Failed to finish puffin writer, region_id: {}, file_id: {}",
107                self.region_id, self.file_id,
108            );
109        }
110
111        0
112    }
113
114    /// Returns false if the finish failed.
115    async fn do_finish_inverted_index(
116        &mut self,
117        puffin_writer: &mut SstPuffinWriter,
118        index_output: &mut IndexOutput,
119    ) -> bool {
120        let Some(mut indexer) = self.inverted_indexer.take() else {
121            return true;
122        };
123
124        let column_ids = indexer.column_ids().collect();
125        let err = match indexer.finish(puffin_writer).await {
126            Ok((row_count, byte_count)) => {
127                self.fill_inverted_index_output(
128                    &mut index_output.inverted_index,
129                    row_count,
130                    byte_count,
131                    column_ids,
132                );
133                return true;
134            }
135            Err(err) => err,
136        };
137
138        if cfg!(any(test, feature = "test")) {
139            panic!(
140                "Failed to finish inverted index, region_id: {}, file_id: {}, err: {:?}",
141                self.region_id, self.file_id, err
142            );
143        } else {
144            warn!(
145                err; "Failed to finish inverted index, region_id: {}, file_id: {}",
146                self.region_id, self.file_id,
147            );
148        }
149
150        false
151    }
152
153    async fn do_finish_fulltext_index(
154        &mut self,
155        puffin_writer: &mut SstPuffinWriter,
156        index_output: &mut IndexOutput,
157    ) -> bool {
158        let Some(mut indexer) = self.fulltext_indexer.take() else {
159            return true;
160        };
161
162        let column_ids = indexer.column_ids().collect();
163        let err = match indexer.finish(puffin_writer).await {
164            Ok((row_count, byte_count)) => {
165                self.fill_fulltext_index_output(
166                    &mut index_output.fulltext_index,
167                    row_count,
168                    byte_count,
169                    column_ids,
170                );
171                return true;
172            }
173            Err(err) => err,
174        };
175
176        if cfg!(any(test, feature = "test")) {
177            panic!(
178                "Failed to finish full-text index, region_id: {}, file_id: {}, err: {:?}",
179                self.region_id, self.file_id, err
180            );
181        } else {
182            warn!(
183                err; "Failed to finish full-text index, region_id: {}, file_id: {}",
184                self.region_id, self.file_id,
185            );
186        }
187
188        false
189    }
190
191    async fn do_finish_bloom_filter(
192        &mut self,
193        puffin_writer: &mut SstPuffinWriter,
194        index_output: &mut IndexOutput,
195    ) -> bool {
196        let Some(mut indexer) = self.bloom_filter_indexer.take() else {
197            return true;
198        };
199
200        let column_ids = indexer.column_ids().collect();
201        let err = match indexer.finish(puffin_writer).await {
202            Ok((row_count, byte_count)) => {
203                self.fill_bloom_filter_output(
204                    &mut index_output.bloom_filter,
205                    row_count,
206                    byte_count,
207                    column_ids,
208                );
209                return true;
210            }
211            Err(err) => err,
212        };
213
214        if cfg!(any(test, feature = "test")) {
215            panic!(
216                "Failed to finish bloom filter, region_id: {}, file_id: {}, err: {:?}",
217                self.region_id, self.file_id, err
218            );
219        } else {
220            warn!(
221                err; "Failed to finish bloom filter, region_id: {}, file_id: {}",
222                self.region_id, self.file_id,
223            );
224        }
225
226        false
227    }
228
229    fn fill_inverted_index_output(
230        &mut self,
231        output: &mut InvertedIndexOutput,
232        row_count: RowCount,
233        byte_count: ByteCount,
234        column_ids: Vec<ColumnId>,
235    ) {
236        debug!(
237            "Inverted index created, region_id: {}, file_id: {}, written_bytes: {}, written_rows: {}, columns: {:?}",
238            self.region_id, self.file_id, byte_count, row_count, column_ids
239        );
240
241        output.index_size = byte_count;
242        output.row_count = row_count;
243        output.columns = column_ids;
244    }
245
246    fn fill_fulltext_index_output(
247        &mut self,
248        output: &mut FulltextIndexOutput,
249        row_count: RowCount,
250        byte_count: ByteCount,
251        column_ids: Vec<ColumnId>,
252    ) {
253        debug!(
254            "Full-text index created, region_id: {}, file_id: {}, written_bytes: {}, written_rows: {}, columns: {:?}",
255            self.region_id, self.file_id, byte_count, row_count, column_ids
256        );
257
258        output.index_size = byte_count;
259        output.row_count = row_count;
260        output.columns = column_ids;
261    }
262
263    fn fill_bloom_filter_output(
264        &mut self,
265        output: &mut BloomFilterOutput,
266        row_count: RowCount,
267        byte_count: ByteCount,
268        column_ids: Vec<ColumnId>,
269    ) {
270        debug!(
271            "Bloom filter created, region_id: {}, file_id: {}, written_bytes: {}, written_rows: {}, columns: {:?}",
272            self.region_id, self.file_id, byte_count, row_count, column_ids
273        );
274
275        output.index_size = byte_count;
276        output.row_count = row_count;
277        output.columns = column_ids;
278    }
279
280    pub(crate) async fn do_prune_intm_sst_dir(&mut self) {
281        if let Some(manager) = self.intermediate_manager.take()
282            && let Err(e) = manager.prune_sst_dir(&self.region_id, &self.file_id).await
283        {
284            warn!(e; "Failed to prune intermediate SST directory, region_id: {}, file_id: {}", self.region_id, self.file_id);
285        }
286    }
287}