1mod 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 if requests.len() == 1 {
60 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 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 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 validate_alter_region_requests(&requests)?;
122 self.validate_logical_field_alters(physical_region_id, &requests)
123 .await?;
124
125 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 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 let column_metadata = *physical_schema_map.get(name).unwrap();
196 (name.to_string(), column_metadata.clone())
197 });
198
199 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 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 async fn check_alter_physical_skip_wal(
288 engine: &super::MetricEngineInner,
289 region_id: RegionId,
290 skip_wal: bool,
291 ) {
292 let request = RegionAlterRequest {
293 kind: AlterKind::SetRegionOptions {
294 options: vec![SetRegionOption::SkipWal(skip_wal)],
295 },
296 };
297 engine
298 .alter_physical_region(region_id, request)
299 .await
300 .unwrap();
301 }
302
303 #[tokio::test]
304 async fn test_alter_region() {
305 let env = TestEnv::new().await;
306 env.init_metric_region().await;
307 let engine = env.metric();
308 let engine_inner = engine.inner;
309
310 let physical_region_id = env.default_physical_region_id();
312 let request = alter_logical_region_request(&["tag1"]);
313
314 let result = engine_inner
315 .alter_physical_region(physical_region_id, request.clone())
316 .await;
317 assert!(result.is_err());
318 assert_eq!(
319 result.unwrap_err().to_string(),
320 "Alter request to physical region is forbidden".to_string()
321 );
322
323 check_alter_physical_skip_wal(&engine_inner, physical_region_id, true).await;
325 check_alter_physical_skip_wal(&engine_inner, physical_region_id, false).await;
326
327 let metadata_region = env.metadata_region();
329 let logical_region_id = env.default_logical_region_id();
330 let is_column_exist = metadata_region
331 .column_semantic_type(physical_region_id, logical_region_id, "tag1")
332 .await
333 .unwrap()
334 .is_some();
335 assert!(!is_column_exist);
336
337 let region_id = env.default_logical_region_id();
338 let response = env
339 .metric()
340 .handle_batch_ddl_requests(BatchRegionDdlRequest::Alter(vec![(
341 region_id,
342 request.clone(),
343 )]))
344 .await
345 .unwrap();
346 let manifest_infos = parse_manifest_infos_from_extensions(&response.extensions).unwrap();
347 assert_eq!(manifest_infos[0].0, physical_region_id);
348 assert!(manifest_infos[0].1.is_metric());
349
350 let semantic_type = metadata_region
351 .column_semantic_type(physical_region_id, logical_region_id, "tag1")
352 .await
353 .unwrap()
354 .unwrap();
355 assert_eq!(semantic_type, SemanticType::Tag);
356 let timestamp_index = metadata_region
357 .column_semantic_type(physical_region_id, logical_region_id, greptime_timestamp())
358 .await
359 .unwrap()
360 .unwrap();
361 assert_eq!(timestamp_index, SemanticType::Timestamp);
362 let column_metadatas =
363 parse_column_metadatas(&response.extensions, ALTER_PHYSICAL_EXTENSION_KEY).unwrap();
364 assert_column_name_and_id(
365 &column_metadatas,
366 &[
367 (greptime_timestamp(), 0),
368 (greptime_value(), 1),
369 ("__table_id", ReservedColumnId::table_id()),
370 ("__tsid", ReservedColumnId::tsid()),
371 ("job", 2),
372 ("tag1", 3),
373 ],
374 );
375 }
376
377 #[tokio::test]
378 async fn test_alter_logical_regions() {
379 let env = TestEnv::new().await;
380 let engine = env.metric();
381 let physical_region_id1 = RegionId::new(1024, 0);
382 let physical_region_id2 = RegionId::new(1024, 1);
383 let logical_region_id1 = RegionId::new(1025, 0);
384 let logical_region_id2 = RegionId::new(1025, 1);
385 env.create_physical_region(physical_region_id1, "/test_dir1", vec![])
386 .await;
387 env.create_physical_region(physical_region_id2, "/test_dir2", vec![])
388 .await;
389
390 let region_create_request1 = crate::test_util::create_logical_region_request(
391 &["job"],
392 physical_region_id1,
393 "logical1",
394 );
395 let region_create_request2 =
396 create_logical_region_request(&["job"], physical_region_id2, "logical2");
397 engine
398 .handle_batch_ddl_requests(BatchRegionDdlRequest::Create(vec![
399 (logical_region_id1, region_create_request1),
400 (logical_region_id2, region_create_request2),
401 ]))
402 .await
403 .unwrap();
404
405 let region_alter_request1 = alter_logical_region_request(&["tag1"]);
406 let region_alter_request2 = alter_logical_region_request(&["tag1"]);
407 let response = engine
408 .handle_batch_ddl_requests(BatchRegionDdlRequest::Alter(vec![
409 (logical_region_id1, region_alter_request1),
410 (logical_region_id2, region_alter_request2),
411 ]))
412 .await
413 .unwrap();
414
415 let manifest_infos = parse_manifest_infos_from_extensions(&response.extensions).unwrap();
416 assert_eq!(manifest_infos.len(), 2);
417 let region_ids = manifest_infos.into_iter().map(|i| i.0).collect::<Vec<_>>();
418 assert!(region_ids.contains(&physical_region_id1));
419 assert!(region_ids.contains(&physical_region_id2));
420
421 let column_metadatas =
422 parse_column_metadatas(&response.extensions, ALTER_PHYSICAL_EXTENSION_KEY).unwrap();
423 assert_column_name_and_id(
424 &column_metadatas,
425 &[
426 (greptime_timestamp(), 0),
427 (greptime_value(), 1),
428 ("__table_id", ReservedColumnId::table_id()),
429 ("__tsid", ReservedColumnId::tsid()),
430 ("job", 2),
431 ("tag1", 3),
432 ],
433 );
434 }
435}