Skip to main content

mito2/sst/index/indexer/
update.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::warn;
16use datatypes::arrow::record_batch::RecordBatch;
17
18use crate::read::Batch;
19use crate::sst::index::Indexer;
20
21impl Indexer {
22    pub(crate) async fn do_update(&mut self, batch: &mut Batch) {
23        if batch.is_empty() {
24            return;
25        }
26
27        if !self.do_prepare_primary_key(batch) {
28            self.do_abort().await;
29            return;
30        }
31
32        if !self.do_update_inverted_index(batch).await {
33            self.do_abort().await;
34        }
35        if !self.do_update_fulltext_index(batch).await {
36            self.do_abort().await;
37        }
38        if !self.do_update_bloom_filter(batch).await {
39            self.do_abort().await;
40        }
41    }
42
43    /// Handles decode errors before entering asynchronous cleanup, following the
44    /// creators' update policy without carrying a decode error across an await.
45    fn do_prepare_primary_key(&self, batch: &mut Batch) -> bool {
46        let Err(err) = self.prepare_primary_key(batch) else {
47            return true;
48        };
49        if cfg!(any(test, feature = "test")) {
50            panic!(
51                "Failed to decode primary key for indexes, region_id: {}, file_id: {}, err: {:?}",
52                self.region_id, self.file_id, err
53            );
54        } else {
55            warn!(err; "Failed to decode primary key for indexes, region_id: {}, file_id: {}", self.region_id, self.file_id);
56        }
57        false
58    }
59
60    /// Returns false if the update failed.
61    async fn do_update_inverted_index(&mut self, batch: &mut Batch) -> bool {
62        let Some(creator) = self.inverted_indexer.as_mut() else {
63            return true;
64        };
65
66        let Err(err) = creator.update(batch).await else {
67            return true;
68        };
69
70        if cfg!(any(test, feature = "test")) {
71            panic!(
72                "Failed to update inverted index, region_id: {}, file_id: {}, err: {:?}",
73                self.region_id, self.file_id, err
74            );
75        } else {
76            warn!(
77                err; "Failed to update inverted index, region_id: {}, file_id: {}",
78                self.region_id, self.file_id,
79            );
80        }
81
82        false
83    }
84
85    /// Returns false if the update failed.
86    async fn do_update_fulltext_index(&mut self, batch: &mut Batch) -> bool {
87        let Some(creator) = self.fulltext_indexer.as_mut() else {
88            return true;
89        };
90
91        let Err(err) = creator.update(batch).await else {
92            return true;
93        };
94
95        if cfg!(any(test, feature = "test")) {
96            panic!(
97                "Failed to update full-text index, region_id: {}, file_id: {}, err: {:?}",
98                self.region_id, self.file_id, err
99            );
100        } else {
101            warn!(
102                err; "Failed to update full-text index, region_id: {}, file_id: {}",
103                self.region_id, self.file_id,
104            );
105        }
106
107        false
108    }
109
110    /// Returns false if the update failed.
111    async fn do_update_bloom_filter(&mut self, batch: &mut Batch) -> bool {
112        let Some(creator) = self.bloom_filter_indexer.as_mut() else {
113            return true;
114        };
115
116        let Err(err) = creator.update(batch).await else {
117            return true;
118        };
119
120        if cfg!(any(test, feature = "test")) {
121            panic!(
122                "Failed to update bloom filter, region_id: {}, file_id: {}, err: {:?}",
123                self.region_id, self.file_id, err
124            );
125        } else {
126            warn!(
127                err; "Failed to update bloom filter, region_id: {}, file_id: {}",
128                self.region_id, self.file_id,
129            );
130        }
131
132        false
133    }
134
135    pub(crate) async fn do_update_flat(&mut self, batch: &RecordBatch) {
136        if batch.num_rows() == 0 {
137            return;
138        }
139
140        if !self.do_update_flat_inverted_index(batch).await {
141            self.do_abort().await;
142        }
143        if !self.do_update_flat_fulltext_index(batch).await {
144            self.do_abort().await;
145        }
146        if !self.do_update_flat_bloom_filter(batch).await {
147            self.do_abort().await;
148        }
149    }
150
151    /// Returns false if the update failed.
152    async fn do_update_flat_inverted_index(&mut self, batch: &RecordBatch) -> bool {
153        let Some(creator) = self.inverted_indexer.as_mut() else {
154            return true;
155        };
156
157        let Err(err) = creator.update_flat(batch).await else {
158            return true;
159        };
160
161        if cfg!(any(test, feature = "test")) {
162            panic!(
163                "Failed to update inverted index with flat format, region_id: {}, file_id: {}, err: {:?}",
164                self.region_id, self.file_id, err
165            );
166        } else {
167            warn!(
168                err; "Failed to update inverted index with flat format, region_id: {}, file_id: {}",
169                self.region_id, self.file_id,
170            );
171        }
172
173        false
174    }
175
176    /// Returns false if the update failed.
177    async fn do_update_flat_fulltext_index(&mut self, batch: &RecordBatch) -> bool {
178        let Some(creator) = self.fulltext_indexer.as_mut() else {
179            return true;
180        };
181
182        let Err(err) = creator.update_flat(batch).await else {
183            return true;
184        };
185
186        if cfg!(any(test, feature = "test")) {
187            panic!(
188                "Failed to update full-text index with flat format, region_id: {}, file_id: {}, err: {:?}",
189                self.region_id, self.file_id, err
190            );
191        } else {
192            warn!(
193                err; "Failed to update full-text index with flat format, region_id: {}, file_id: {}",
194                self.region_id, self.file_id,
195            );
196        }
197
198        false
199    }
200
201    /// Returns false if the update failed.
202    async fn do_update_flat_bloom_filter(&mut self, batch: &RecordBatch) -> bool {
203        let Some(creator) = self.bloom_filter_indexer.as_mut() else {
204            return true;
205        };
206
207        let Err(err) = creator.update_flat(batch).await else {
208            return true;
209        };
210
211        if cfg!(any(test, feature = "test")) {
212            panic!(
213                "Failed to update bloom filter with flat format, region_id: {}, file_id: {}, err: {:?}",
214                self.region_id, self.file_id, err
215            );
216        } else {
217            warn!(
218                err; "Failed to update bloom filter with flat format, region_id: {}, file_id: {}",
219                self.region_id, self.file_id,
220            );
221        }
222
223        false
224    }
225}