Skip to main content

metric_engine/engine/
alter.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
15mod extract_new_columns;
16mod validate;
17
18use std::collections::{BTreeSet, HashMap, HashSet};
19
20use api::v1::SemanticType;
21use common_query::native_histogram::is_native_histogram_value_type;
22use extract_new_columns::extract_new_columns;
23use snafu::{OptionExt, ResultExt, ensure};
24use store_api::metadata::ColumnMetadata;
25use store_api::metric_engine_consts::ALTER_PHYSICAL_EXTENSION_KEY;
26use store_api::region_request::{AffectedRows, AlterKind, RegionAlterRequest};
27use store_api::storage::RegionId;
28use validate::validate_alter_region_requests;
29
30use crate::engine::MetricEngineInner;
31use crate::error::{
32    AddingFieldColumnSnafu, LogicalRegionNotFoundSnafu, PhysicalRegionNotFoundSnafu, Result,
33    SerializeColumnMetadataSnafu, UnexpectedRequestSnafu,
34};
35use crate::utils::{append_manifest_info, encode_manifest_info_to_extensions, to_data_region_id};
36
37impl MetricEngineInner {
38    pub async fn alter_regions(
39        &self,
40        mut requests: Vec<(RegionId, RegionAlterRequest)>,
41        extension_return_value: &mut HashMap<String, Vec<u8>>,
42    ) -> Result<AffectedRows> {
43        if requests.is_empty() {
44            return Ok(0);
45        }
46
47        let first_region_id = &requests.first().unwrap().0;
48        if self.is_physical_region(*first_region_id) {
49            ensure!(
50                requests.len() == 1,
51                UnexpectedRequestSnafu {
52                    reason: "Physical table must be altered with single request".to_string(),
53                }
54            );
55            let (region_id, request) = requests.pop().unwrap();
56            self.alter_physical_region(region_id, request).await?;
57        } else {
58            // Fast path for single logical region alter request
59            if requests.len() == 1 {
60                // Safety: requests is not empty
61                let region_id = requests.first().unwrap().0;
62                let physical_region_id = self
63                    .state
64                    .read()
65                    .unwrap()
66                    .get_physical_region_id(region_id)
67                    .with_context(|| LogicalRegionNotFoundSnafu { region_id })?;
68                let mut manifest_infos = Vec::with_capacity(1);
69                self.alter_logical_regions(physical_region_id, requests, extension_return_value)
70                    .await?;
71                append_manifest_info(&self.mito, physical_region_id, &mut manifest_infos);
72                encode_manifest_info_to_extensions(&manifest_infos, extension_return_value)?;
73            } else {
74                let grouped_requests =
75                    self.group_logical_region_requests_by_physical_region_id(requests)?;
76                let mut manifest_infos = Vec::with_capacity(grouped_requests.len());
77                for (physical_region_id, requests) in grouped_requests {
78                    self.alter_logical_regions(
79                        physical_region_id,
80                        requests,
81                        extension_return_value,
82                    )
83                    .await?;
84                    append_manifest_info(&self.mito, physical_region_id, &mut manifest_infos);
85                }
86                encode_manifest_info_to_extensions(&manifest_infos, extension_return_value)?;
87            }
88        }
89        Ok(0)
90    }
91
92    /// Groups the alter logical region requests by physical region id.
93    fn group_logical_region_requests_by_physical_region_id(
94        &self,
95        requests: Vec<(RegionId, RegionAlterRequest)>,
96    ) -> Result<HashMap<RegionId, Vec<(RegionId, RegionAlterRequest)>>> {
97        let mut result = HashMap::with_capacity(requests.len());
98        let state = self.state.read().unwrap();
99
100        for (region_id, request) in requests {
101            let physical_region_id = state
102                .get_physical_region_id(region_id)
103                .with_context(|| LogicalRegionNotFoundSnafu { region_id })?;
104            result
105                .entry(physical_region_id)
106                .or_insert_with(Vec::new)
107                .push((region_id, request));
108        }
109
110        Ok(result)
111    }
112
113    /// Alter multiple logical regions on the same physical region.
114    pub async fn alter_logical_regions(
115        &self,
116        physical_region_id: RegionId,
117        requests: Vec<(RegionId, RegionAlterRequest)>,
118        extension_return_value: &mut HashMap<String, Vec<u8>>,
119    ) -> Result<AffectedRows> {
120        // Checks all alter requests are add columns.
121        validate_alter_region_requests(&requests)?;
122        self.validate_logical_field_alters(physical_region_id, &requests)
123            .await?;
124
125        // Finds new columns to add
126        let mut new_column_names = HashSet::new();
127        let mut new_columns_to_add = vec![];
128
129        let index_options = {
130            let state = &self.state.read().unwrap();
131            let region_state = state
132                .physical_region_states()
133                .get(&physical_region_id)
134                .with_context(|| PhysicalRegionNotFoundSnafu {
135                    region_id: physical_region_id,
136                })?;
137            let physical_columns = region_state.physical_columns();
138
139            extract_new_columns(
140                &requests,
141                physical_columns,
142                &mut new_column_names,
143                &mut new_columns_to_add,
144            )?;
145
146            region_state.options().index
147        };
148        let data_region_id = to_data_region_id(physical_region_id);
149
150        // Acquire logical region locks in a deterministic order to avoid deadlocks when multiple
151        // alter operations target overlapping regions concurrently.
152        let region_ids = requests
153            .iter()
154            .map(|(region_id, _)| *region_id)
155            .collect::<BTreeSet<_>>();
156
157        let mut write_guards = Vec::with_capacity(region_ids.len());
158        for region_id in region_ids {
159            write_guards.push(
160                self.metadata_region
161                    .write_lock_logical_region(region_id)
162                    .await?,
163            );
164        }
165
166        self.data_region
167            .add_columns(data_region_id, new_columns_to_add, index_options)
168            .await?;
169
170        let physical_columns = self.data_region.physical_columns(data_region_id).await?;
171        let physical_schema_map = physical_columns
172            .iter()
173            .map(|metadata| (metadata.column_schema.name.as_str(), metadata))
174            .collect::<HashMap<_, _>>();
175
176        let logical_region_columns = requests.iter().map(|(region_id, request)| {
177            let AlterKind::AddColumns { columns } = &request.kind else {
178                unreachable!()
179            };
180            (
181                *region_id,
182                columns
183                    .iter()
184                    .map(|col| {
185                        let column_name = col.column_metadata.column_schema.name.as_str();
186                        let column_metadata = *physical_schema_map.get(column_name).unwrap();
187                        (column_name, column_metadata)
188                    })
189                    .collect::<HashMap<_, _>>(),
190            )
191        });
192
193        let new_add_columns = new_column_names.iter().map(|name| {
194            // Safety: previous steps ensure the physical region exist
195            let column_metadata = *physical_schema_map.get(name).unwrap();
196            (name.to_string(), column_metadata.clone())
197        });
198
199        // Writes logical regions metadata to metadata region
200        self.metadata_region
201            .add_logical_regions(physical_region_id, false, logical_region_columns)
202            .await?;
203
204        extension_return_value.insert(
205            ALTER_PHYSICAL_EXTENSION_KEY.to_string(),
206            ColumnMetadata::encode_list(&physical_columns).context(SerializeColumnMetadataSnafu)?,
207        );
208
209        let mut state = self.state.write().unwrap();
210        state.add_physical_columns(data_region_id, new_add_columns);
211        state.invalid_logical_regions_cache(requests.iter().map(|(region_id, _)| *region_id));
212
213        Ok(0)
214    }
215
216    async fn validate_logical_field_alters(
217        &self,
218        physical_region_id: RegionId,
219        requests: &[(RegionId, RegionAlterRequest)],
220    ) -> Result<()> {
221        // Logical metric tables have one field column. Native histograms are a
222        // special struct field, so field alters must leave exactly that field.
223        for (region_id, request) in requests {
224            let AlterKind::AddColumns { columns } = &request.kind else {
225                unreachable!()
226            };
227            let added_fields = columns
228                .iter()
229                .filter(|col| col.column_metadata.semantic_type == SemanticType::Field)
230                .collect::<Vec<_>>();
231            let Some(&first_added_field) = added_fields.first() else {
232                continue;
233            };
234
235            let mut fields = self
236                .load_logical_columns(physical_region_id, *region_id)
237                .await?
238                .into_iter()
239                .filter(|col| col.semantic_type == SemanticType::Field)
240                .collect::<Vec<_>>();
241            fields.extend(
242                added_fields
243                    .into_iter()
244                    .map(|col| col.column_metadata.clone()),
245            );
246
247            ensure!(
248                fields.len() == 1
249                    && is_native_histogram_value_type(&fields[0].column_schema.data_type),
250                AddingFieldColumnSnafu {
251                    name: first_added_field.column_metadata.column_schema.name.clone(),
252                }
253            );
254        }
255
256        Ok(())
257    }
258
259    async fn alter_physical_region(
260        &self,
261        region_id: RegionId,
262        request: RegionAlterRequest,
263    ) -> Result<()> {
264        self.data_region
265            .alter_region_options(region_id, request)
266            .await?;
267        Ok(())
268    }
269}
270
271#[cfg(test)]
272mod test {
273    use api::v1::SemanticType;
274    use common_meta::ddl::test_util::assert_column_name_and_id;
275    use common_meta::ddl::utils::{parse_column_metadatas, parse_manifest_infos_from_extensions};
276    use common_query::prelude::{greptime_timestamp, greptime_value};
277    use store_api::metric_engine_consts::ALTER_PHYSICAL_EXTENSION_KEY;
278    use store_api::region_engine::RegionEngine;
279    use store_api::region_request::{
280        AlterKind, BatchRegionDdlRequest, RegionAlterRequest, SetRegionOption,
281    };
282    use store_api::storage::RegionId;
283    use store_api::storage::consts::ReservedColumnId;
284
285    use crate::test_util::{TestEnv, alter_logical_region_request, create_logical_region_request};
286
287    #[tokio::test]
288    async fn test_alter_region() {
289        let env = TestEnv::new().await;
290        env.init_metric_region().await;
291        let engine = env.metric();
292        let engine_inner = engine.inner;
293
294        // alter physical region
295        let physical_region_id = env.default_physical_region_id();
296        let request = alter_logical_region_request(&["tag1"]);
297
298        let result = engine_inner
299            .alter_physical_region(physical_region_id, request.clone())
300            .await;
301        assert!(result.is_err());
302        assert_eq!(
303            result.unwrap_err().to_string(),
304            "Alter request to physical region is forbidden".to_string()
305        );
306
307        // skip WAL on the physical region should be forwarded to the data region
308        let alter_region_option_request = RegionAlterRequest {
309            kind: AlterKind::SetRegionOptions {
310                options: vec![SetRegionOption::SkipWal],
311            },
312        };
313        engine_inner
314            .alter_physical_region(physical_region_id, alter_region_option_request)
315            .await
316            .unwrap();
317
318        // alter logical region
319        let metadata_region = env.metadata_region();
320        let logical_region_id = env.default_logical_region_id();
321        let is_column_exist = metadata_region
322            .column_semantic_type(physical_region_id, logical_region_id, "tag1")
323            .await
324            .unwrap()
325            .is_some();
326        assert!(!is_column_exist);
327
328        let region_id = env.default_logical_region_id();
329        let response = env
330            .metric()
331            .handle_batch_ddl_requests(BatchRegionDdlRequest::Alter(vec![(
332                region_id,
333                request.clone(),
334            )]))
335            .await
336            .unwrap();
337        let manifest_infos = parse_manifest_infos_from_extensions(&response.extensions).unwrap();
338        assert_eq!(manifest_infos[0].0, physical_region_id);
339        assert!(manifest_infos[0].1.is_metric());
340
341        let semantic_type = metadata_region
342            .column_semantic_type(physical_region_id, logical_region_id, "tag1")
343            .await
344            .unwrap()
345            .unwrap();
346        assert_eq!(semantic_type, SemanticType::Tag);
347        let timestamp_index = metadata_region
348            .column_semantic_type(physical_region_id, logical_region_id, greptime_timestamp())
349            .await
350            .unwrap()
351            .unwrap();
352        assert_eq!(timestamp_index, SemanticType::Timestamp);
353        let column_metadatas =
354            parse_column_metadatas(&response.extensions, ALTER_PHYSICAL_EXTENSION_KEY).unwrap();
355        assert_column_name_and_id(
356            &column_metadatas,
357            &[
358                (greptime_timestamp(), 0),
359                (greptime_value(), 1),
360                ("__table_id", ReservedColumnId::table_id()),
361                ("__tsid", ReservedColumnId::tsid()),
362                ("job", 2),
363                ("tag1", 3),
364            ],
365        );
366    }
367
368    #[tokio::test]
369    async fn test_alter_logical_regions() {
370        let env = TestEnv::new().await;
371        let engine = env.metric();
372        let physical_region_id1 = RegionId::new(1024, 0);
373        let physical_region_id2 = RegionId::new(1024, 1);
374        let logical_region_id1 = RegionId::new(1025, 0);
375        let logical_region_id2 = RegionId::new(1025, 1);
376        env.create_physical_region(physical_region_id1, "/test_dir1", vec![])
377            .await;
378        env.create_physical_region(physical_region_id2, "/test_dir2", vec![])
379            .await;
380
381        let region_create_request1 = crate::test_util::create_logical_region_request(
382            &["job"],
383            physical_region_id1,
384            "logical1",
385        );
386        let region_create_request2 =
387            create_logical_region_request(&["job"], physical_region_id2, "logical2");
388        engine
389            .handle_batch_ddl_requests(BatchRegionDdlRequest::Create(vec![
390                (logical_region_id1, region_create_request1),
391                (logical_region_id2, region_create_request2),
392            ]))
393            .await
394            .unwrap();
395
396        let region_alter_request1 = alter_logical_region_request(&["tag1"]);
397        let region_alter_request2 = alter_logical_region_request(&["tag1"]);
398        let response = engine
399            .handle_batch_ddl_requests(BatchRegionDdlRequest::Alter(vec![
400                (logical_region_id1, region_alter_request1),
401                (logical_region_id2, region_alter_request2),
402            ]))
403            .await
404            .unwrap();
405
406        let manifest_infos = parse_manifest_infos_from_extensions(&response.extensions).unwrap();
407        assert_eq!(manifest_infos.len(), 2);
408        let region_ids = manifest_infos.into_iter().map(|i| i.0).collect::<Vec<_>>();
409        assert!(region_ids.contains(&physical_region_id1));
410        assert!(region_ids.contains(&physical_region_id2));
411
412        let column_metadatas =
413            parse_column_metadatas(&response.extensions, ALTER_PHYSICAL_EXTENSION_KEY).unwrap();
414        assert_column_name_and_id(
415            &column_metadatas,
416            &[
417                (greptime_timestamp(), 0),
418                (greptime_value(), 1),
419                ("__table_id", ReservedColumnId::table_id()),
420                ("__tsid", ReservedColumnId::tsid()),
421                ("job", 2),
422                ("tag1", 3),
423            ],
424        );
425    }
426}