1use std::collections::HashMap;
18
19use arrow_schema::extension::ExtensionType;
20use common_meta::SchemaOptions;
21use datatypes::extension::json::{Json2ExtensionType, parse_legacy_json2_settings};
22use datatypes::schema::{
23 COLUMN_FULLTEXT_OPT_KEY_ANALYZER, COLUMN_FULLTEXT_OPT_KEY_BACKEND,
24 COLUMN_FULLTEXT_OPT_KEY_CASE_SENSITIVE, COLUMN_FULLTEXT_OPT_KEY_FALSE_POSITIVE_RATE,
25 COLUMN_FULLTEXT_OPT_KEY_GRANULARITY, COLUMN_SKIPPING_INDEX_OPT_KEY_FALSE_POSITIVE_RATE,
26 COLUMN_SKIPPING_INDEX_OPT_KEY_GRANULARITY, COLUMN_SKIPPING_INDEX_OPT_KEY_TYPE, COMMENT_KEY,
27 ColumnDefaultConstraint, ColumnSchema, FulltextBackend, SchemaRef,
28};
29use datatypes::types::JsonFormat;
30use snafu::ResultExt;
31use sql::ast::{ColumnDef, ColumnOption, ColumnOptionDef, DataType, Expr, Ident, ObjectName};
32use sql::dialect::GreptimeDbDialect;
33use sql::parser::ParserContext;
34use sql::statements::create::{Column, ColumnExtensions, CreateTable, TableConstraint};
35use sql::statements::{self, OptionMap, concrete_data_type_to_sql_data_type};
36use store_api::metric_engine_consts::{is_metric_engine, is_metric_engine_internal_column};
37use table::metadata::{TableInfoRef, TableMeta};
38use table::requests::{
39 COMMENT_KEY as TABLE_COMMENT_KEY, FILE_TABLE_META_KEY, SKIP_WAL_KEY, TTL_KEY,
40 WRITE_BUFFER_SIZE_KEY,
41};
42
43use crate::error::{
44 ConvertSqlTypeSnafu, ConvertSqlValueSnafu, GetFulltextOptionsSnafu,
45 GetSkippingIndexOptionsSnafu, Result, SqlSnafu,
46};
47
48fn create_sql_options(table_meta: &TableMeta, schema_options: Option<SchemaOptions>) -> OptionMap {
50 let table_opts = &table_meta.options;
51 let mut options = OptionMap::default();
52 if let Some(write_buffer_size) = table_opts.write_buffer_size {
53 options.insert(
54 WRITE_BUFFER_SIZE_KEY.to_string(),
55 write_buffer_size.to_string(),
56 );
57 }
58 if let Some(ttl) = table_opts.ttl.map(|t| t.to_string()) {
59 options.insert(TTL_KEY.to_string(), ttl);
60 } else if let Some(database_ttl) = schema_options
61 .as_ref()
62 .and_then(|o| o.ttl)
63 .map(|ttl| ttl.to_string())
64 {
65 options.insert(TTL_KEY.to_string(), database_ttl);
66 };
67 for (k, v) in table_opts
68 .extra_options
69 .iter()
70 .filter(|(k, _)| k != &FILE_TABLE_META_KEY)
71 {
72 options.insert(k.clone(), v.clone());
73 }
74 if table_opts.skip_wal {
75 options.insert(SKIP_WAL_KEY.to_string(), true.to_string());
76 }
77 options
78}
79
80#[inline]
81fn column_option_def(option: ColumnOption) -> ColumnOptionDef {
82 ColumnOptionDef { name: None, option }
83}
84
85fn create_column(column_schema: &ColumnSchema, quote_style: char) -> Result<Column> {
86 let name = &column_schema.name;
87 let mut options = Vec::with_capacity(2);
88 let mut extensions = ColumnExtensions::default();
89
90 if column_schema.is_nullable() {
91 options.push(column_option_def(ColumnOption::Null));
92 } else {
93 options.push(column_option_def(ColumnOption::NotNull));
94 }
95
96 if let Some(c) = column_schema.default_constraint() {
97 let expr = match c {
98 ColumnDefaultConstraint::Value(v) => Expr::Value(
99 statements::value_to_sql_value(v)
100 .with_context(|_| ConvertSqlValueSnafu { value: v.clone() })?
101 .into(),
102 ),
103 ColumnDefaultConstraint::Function(expr) => {
104 ParserContext::parse_function(expr, &GreptimeDbDialect {}).context(SqlSnafu)?
105 }
106 };
107
108 options.push(column_option_def(ColumnOption::Default(expr)));
109 }
110
111 if let Some(c) = column_schema.metadata().get(COMMENT_KEY) {
112 options.push(column_option_def(ColumnOption::Comment(c.clone())));
113 }
114
115 if let Some(opt) = column_schema
116 .fulltext_options()
117 .context(GetFulltextOptionsSnafu)?
118 && opt.enable
119 {
120 let mut map = HashMap::from([
121 (
122 COLUMN_FULLTEXT_OPT_KEY_ANALYZER.to_string(),
123 opt.analyzer.to_string(),
124 ),
125 (
126 COLUMN_FULLTEXT_OPT_KEY_CASE_SENSITIVE.to_string(),
127 opt.case_sensitive.to_string(),
128 ),
129 (
130 COLUMN_FULLTEXT_OPT_KEY_BACKEND.to_string(),
131 opt.backend.to_string(),
132 ),
133 ]);
134 if opt.backend == FulltextBackend::Bloom {
135 map.insert(
136 COLUMN_FULLTEXT_OPT_KEY_GRANULARITY.to_string(),
137 opt.granularity.to_string(),
138 );
139 map.insert(
140 COLUMN_FULLTEXT_OPT_KEY_FALSE_POSITIVE_RATE.to_string(),
141 opt.false_positive_rate().to_string(),
142 );
143 }
144 extensions.fulltext_index_options = Some(map.into());
145 }
146
147 if let Some(opt) = column_schema
148 .skipping_index_options()
149 .context(GetSkippingIndexOptionsSnafu)?
150 {
151 let map = HashMap::from([
152 (
153 COLUMN_SKIPPING_INDEX_OPT_KEY_GRANULARITY.to_string(),
154 opt.granularity.to_string(),
155 ),
156 (
157 COLUMN_SKIPPING_INDEX_OPT_KEY_FALSE_POSITIVE_RATE.to_string(),
158 opt.false_positive_rate().to_string(),
159 ),
160 (
161 COLUMN_SKIPPING_INDEX_OPT_KEY_TYPE.to_string(),
162 opt.index_type.to_string(),
163 ),
164 ]);
165 extensions.skipping_index_options = Some(map.into());
166 }
167
168 if column_schema.is_inverted_indexed() {
169 extensions.inverted_index_options = Some(HashMap::new().into());
170 }
171
172 let mut data_type = concrete_data_type_to_sql_data_type(&column_schema.data_type)
173 .with_context(|_| ConvertSqlTypeSnafu {
174 datatype: column_schema.data_type.clone(),
175 })?;
176
177 if matches!(
178 &column_schema.data_type,
179 datatypes::data_type::ConcreteDataType::Json(json_type)
180 if matches!(json_type.format, JsonFormat::Json2(_))
181 ) {
182 data_type = DataType::Custom(ObjectName::from(vec![Ident::new("JSON2")]), vec![]);
183 }
184
185 let settings = if let Some(extension) = column_schema.extension_type::<Json2ExtensionType>()? {
186 Some(extension.metadata().json_settings().clone())
187 } else {
188 parse_legacy_json2_settings(column_schema.metadata())?
189 };
190 if let Some(settings) = settings {
191 extensions.set_json_settings(settings).context(SqlSnafu)?;
192 }
193
194 Ok(Column {
195 column_def: ColumnDef {
196 name: Ident::with_quote(quote_style, name),
197 data_type,
198 options,
199 },
200 extensions,
201 })
202}
203
204fn primary_key_columns_for_show_create<'a>(
208 table_meta: &'a TableMeta,
209 engine: &str,
210) -> Vec<&'a String> {
211 let is_metric_engine = is_metric_engine(engine);
212 if is_metric_engine {
213 table_meta
214 .row_key_column_names()
215 .filter(|name| !is_metric_engine_internal_column(name))
216 .collect()
217 } else {
218 table_meta.row_key_column_names().collect()
219 }
220}
221
222fn create_table_constraints(
223 engine: &str,
224 schema: &SchemaRef,
225 table_meta: &TableMeta,
226 quote_style: char,
227) -> Vec<TableConstraint> {
228 let mut constraints = Vec::with_capacity(2);
229 if let Some(timestamp_column) = schema.timestamp_column() {
230 let column_name = ×tamp_column.name;
231 constraints.push(TableConstraint::TimeIndex {
232 column: Ident::with_quote(quote_style, column_name),
233 });
234 }
235 if !table_meta.primary_key_indices.is_empty() {
236 let columns = primary_key_columns_for_show_create(table_meta, engine)
237 .into_iter()
238 .map(|name| Ident::with_quote(quote_style, name))
239 .collect();
240 constraints.push(TableConstraint::PrimaryKey { columns });
241 }
242
243 constraints
244}
245
246pub fn create_table_stmt(
248 table_info: &TableInfoRef,
249 schema_options: Option<SchemaOptions>,
250 quote_style: char,
251) -> Result<CreateTable> {
252 let table_meta = &table_info.meta;
253 let table_name = &table_info.name;
254 let schema = &table_info.meta.schema;
255 let is_metric_engine = is_metric_engine(&table_meta.engine);
256 let columns = schema
257 .column_schemas()
258 .iter()
259 .filter_map(|c| {
260 if is_metric_engine && is_metric_engine_internal_column(&c.name) {
261 None
262 } else {
263 Some(create_column(c, quote_style))
264 }
265 })
266 .collect::<Result<Vec<_>>>()?;
267
268 let constraints = create_table_constraints(&table_meta.engine, schema, table_meta, quote_style);
269
270 let mut options = create_sql_options(table_meta, schema_options);
271 if let Some(comment) = &table_info.desc
272 && options.get(TABLE_COMMENT_KEY).is_none()
273 {
274 options.insert(format!("'{TABLE_COMMENT_KEY}'"), comment.clone());
275 }
276
277 Ok(CreateTable {
278 if_not_exists: true,
279 table_id: table_info.ident.table_id,
280 name: ObjectName::from(vec![Ident::with_quote(quote_style, table_name)]),
281 columns,
282 engine: table_meta.engine.clone(),
283 constraints,
284 options,
285 partitions: None,
286 })
287}
288
289#[cfg(test)]
290mod tests {
291 use std::sync::Arc;
292 use std::time::Duration;
293
294 use common_time::timestamp::TimeUnit;
295 use datatypes::extension::json::JsonExtensionType;
296 use datatypes::prelude::ConcreteDataType;
297 use datatypes::schema::{FulltextOptions, Schema, SchemaRef, SkippingIndexOptions};
298 use table::metadata::*;
299 use table::requests::{
300 FILE_TABLE_FORMAT_KEY, FILE_TABLE_LOCATION_KEY, FILE_TABLE_META_KEY, TableOptions,
301 };
302
303 use super::*;
304
305 #[test]
306 fn test_show_create_table_sql() {
307 let schema = vec![
308 ColumnSchema::new("id", ConcreteDataType::uint32_datatype(), true)
309 .with_skipping_options(SkippingIndexOptions {
310 granularity: 4096,
311 ..Default::default()
312 })
313 .unwrap(),
314 ColumnSchema::new("host", ConcreteDataType::string_datatype(), true)
315 .with_inverted_index(true),
316 ColumnSchema::new("cpu", ConcreteDataType::float64_datatype(), true),
317 ColumnSchema::new("disk", ConcreteDataType::float32_datatype(), true),
318 ColumnSchema::new("msg", ConcreteDataType::string_datatype(), true)
319 .with_fulltext_options(FulltextOptions {
320 enable: true,
321 ..Default::default()
322 })
323 .unwrap(),
324 ColumnSchema::new("embedding", ConcreteDataType::vector_datatype(4), true),
325 ColumnSchema::new(
326 "ts",
327 ConcreteDataType::timestamp_datatype(TimeUnit::Millisecond),
328 false,
329 )
330 .with_default_constraint(Some(ColumnDefaultConstraint::Function(String::from(
331 "current_timestamp()",
332 ))))
333 .unwrap()
334 .with_time_index(true),
335 ];
336
337 let table_schema = SchemaRef::new(Schema::new(schema));
338 let table_name = "system_metrics";
339 let schema_name = "public".to_string();
340 let catalog_name = "greptime".to_string();
341
342 let mut options = table::requests::TableOptions {
343 ttl: Some(Duration::from_secs(30).into()),
344 skip_wal: true,
345 ..Default::default()
346 };
347
348 let _ = options
349 .extra_options
350 .insert("compaction.type".to_string(), "twcs".to_string());
351
352 let meta = TableMetaBuilder::empty()
353 .schema(table_schema)
354 .primary_key_indices(vec![0, 1])
355 .value_indices(vec![2, 3])
356 .engine("mito".to_string())
357 .next_column_id(0)
358 .options(options)
359 .created_on(Default::default())
360 .build()
361 .unwrap();
362
363 let info = Arc::new(
364 TableInfoBuilder::default()
365 .table_id(1024)
366 .table_version(0 as TableVersion)
367 .name(table_name)
368 .schema_name(schema_name)
369 .catalog_name(catalog_name)
370 .desc(None)
371 .table_type(TableType::Base)
372 .meta(meta)
373 .build()
374 .unwrap(),
375 );
376
377 let stmt = create_table_stmt(&info, None, '"').unwrap();
378
379 let sql = format!("\n{}", stmt);
380 assert_eq!(
381 r#"
382CREATE TABLE IF NOT EXISTS "system_metrics" (
383 "id" INT UNSIGNED NULL SKIPPING INDEX WITH(false_positive_rate = '0.01', granularity = '4096', type = 'BLOOM'),
384 "host" STRING NULL INVERTED INDEX,
385 "cpu" DOUBLE NULL,
386 "disk" FLOAT NULL,
387 "msg" STRING NULL FULLTEXT INDEX WITH(analyzer = 'English', backend = 'bloom', case_sensitive = 'false', false_positive_rate = '0.01', granularity = '10240'),
388 "embedding" VECTOR(4) NULL,
389 "ts" TIMESTAMP(3) NOT NULL DEFAULT current_timestamp(),
390 TIME INDEX ("ts"),
391 PRIMARY KEY ("id", "host")
392)
393ENGINE=mito
394WITH(
395 'compaction.type' = 'twcs',
396 skip_wal = 'true',
397 ttl = '30s'
398)"#,
399 sql
400 );
401
402 let mut table_meta = info.meta.clone();
403 table_meta.options.skip_wal = false;
404 table_meta
405 .options
406 .extra_options
407 .insert(SKIP_WAL_KEY.to_string(), false.to_string());
408 assert_eq!(
409 Some("false"),
410 create_sql_options(&table_meta, None).get(SKIP_WAL_KEY)
411 );
412
413 let mut schema_options = SchemaOptions::default();
414 schema_options
415 .extra_options
416 .insert(SKIP_WAL_KEY.to_string(), true.to_string());
417 assert_eq!(
418 Some("false"),
419 create_sql_options(&table_meta, Some(schema_options)).get(SKIP_WAL_KEY)
420 );
421 }
422
423 #[test]
424 fn test_show_create_legacy_json_with_json_extension() {
425 let mut json_column = ColumnSchema::new("j", ConcreteDataType::json_datatype(), true);
426 json_column.with_extension_type(&JsonExtensionType);
427
428 let table_schema = SchemaRef::new(Schema::new(vec![
429 json_column,
430 ColumnSchema::new(
431 "ts",
432 ConcreteDataType::timestamp_datatype(TimeUnit::Millisecond),
433 false,
434 )
435 .with_time_index(true),
436 ]));
437 let table_name = "legacy_json";
438 let meta = TableMetaBuilder::empty()
439 .schema(table_schema)
440 .primary_key_indices(vec![])
441 .value_indices(vec![0])
442 .engine("mito".to_string())
443 .next_column_id(0)
444 .options(Default::default())
445 .created_on(Default::default())
446 .build()
447 .unwrap();
448
449 let info = Arc::new(
450 TableInfoBuilder::default()
451 .table_id(1024)
452 .table_version(0 as TableVersion)
453 .name(table_name)
454 .schema_name("public")
455 .catalog_name("greptime")
456 .desc(None)
457 .table_type(TableType::Base)
458 .meta(meta)
459 .build()
460 .unwrap(),
461 );
462
463 let stmt = create_table_stmt(&info, None, '"').unwrap();
464 let sql = format!("\n{}", stmt);
465 assert_eq!(
466 r#"
467CREATE TABLE IF NOT EXISTS "legacy_json" (
468 "j" JSON NULL,
469 "ts" TIMESTAMP(3) NOT NULL,
470 TIME INDEX ("ts")
471)
472ENGINE=mito
473"#,
474 sql
475 );
476 }
477
478 #[test]
479 fn test_show_create_external_table_sql() {
480 let schema = vec![
481 ColumnSchema::new("host", ConcreteDataType::string_datatype(), true),
482 ColumnSchema::new("cpu", ConcreteDataType::float64_datatype(), true),
483 ];
484 let table_schema = SchemaRef::new(Schema::new(schema));
485 let table_name = "system_metrics";
486 let schema_name = "public".to_string();
487 let catalog_name = "greptime".to_string();
488 let mut options: TableOptions = Default::default();
489 let _ = options
490 .extra_options
491 .insert(FILE_TABLE_LOCATION_KEY.to_string(), "foo.csv".to_string());
492 let _ = options.extra_options.insert(
493 FILE_TABLE_META_KEY.to_string(),
494 "{{\"files\":[\"foo.csv\"]}}".to_string(),
495 );
496 let _ = options
497 .extra_options
498 .insert(FILE_TABLE_FORMAT_KEY.to_string(), "csv".to_string());
499 let meta = TableMetaBuilder::empty()
500 .schema(table_schema)
501 .primary_key_indices(vec![])
502 .engine("file".to_string())
503 .next_column_id(0)
504 .options(options)
505 .created_on(Default::default())
506 .build()
507 .unwrap();
508
509 let info = Arc::new(
510 TableInfoBuilder::default()
511 .table_id(1024)
512 .table_version(0 as TableVersion)
513 .name(table_name)
514 .schema_name(schema_name)
515 .catalog_name(catalog_name)
516 .desc(None)
517 .table_type(TableType::Base)
518 .meta(meta)
519 .build()
520 .unwrap(),
521 );
522
523 let stmt = create_table_stmt(&info, None, '"').unwrap();
524
525 let sql = format!("\n{}", stmt);
526 assert_eq!(
527 r#"
528CREATE EXTERNAL TABLE IF NOT EXISTS "system_metrics" (
529 "host" STRING NULL,
530 "cpu" DOUBLE NULL,
531
532)
533ENGINE=file
534WITH(
535 format = 'csv',
536 location = 'foo.csv'
537)"#,
538 sql
539 );
540 }
541}