1use api::v1::CreateTableExpr;
16use common_recordbatch::map_dictionary_to_values_data_type;
17use datafusion_common::tree_node::TreeNode;
18use datafusion_expr::LogicalPlan;
19use datatypes::prelude::ConcreteDataType;
20use datatypes::schema::ColumnSchema;
21use operator::expr_helper::column_schemas_to_defs;
22use snafu::ResultExt;
23
24use crate::Error;
25use crate::adapter::{AUTO_CREATED_PLACEHOLDER_TS_COL, AUTO_CREATED_UPDATE_AT_TS_COL};
26use crate::batching_mode::utils::FindGroupByFinalName;
27use crate::error::{ConvertColumnSchemaSnafu, DatafusionSnafu};
28
29#[derive(Debug, Clone, PartialEq, Eq)]
30pub enum QueryType {
31 Tql,
33 Sql,
35}
36
37pub(super) fn create_table_with_expr(
40 plan: &LogicalPlan,
41 sink_table_name: &[String; 3],
42 query_type: &QueryType,
43) -> Result<CreateTableExpr, Error> {
44 let table_def = match query_type {
45 &QueryType::Sql => {
46 if let Some(def) = build_pk_from_aggr(plan)? {
47 def
48 } else {
49 build_by_sql_schema(plan)?
50 }
51 }
52 QueryType::Tql => {
53 if let Some(table_def) = build_pk_from_aggr(plan)? {
55 table_def
56 } else {
57 build_by_tql_schema(plan)?
58 }
59 }
60 };
61 let first_time_stamp = table_def.ts_col;
62 let primary_keys = table_def.pks;
63
64 let mut column_schemas = Vec::new();
65 for field in plan.schema().fields() {
66 let name = field.name();
67 let ty = map_dictionary_to_values_data_type(&ConcreteDataType::from_arrow_type(
68 field.data_type(),
69 ));
70 let col_schema = if first_time_stamp == Some(name.clone()) {
71 ColumnSchema::new(name, ty, false).with_time_index(true)
72 } else {
73 ColumnSchema::new(name, ty, true)
74 };
75
76 match query_type {
77 QueryType::Sql => {
78 column_schemas.push(col_schema);
79 }
80 QueryType::Tql => {
81 let is_tag_column = primary_keys.contains(name);
84 let is_val_column = !is_tag_column && first_time_stamp.as_ref() != Some(name);
85 if is_val_column {
86 let col_schema =
87 ColumnSchema::new(name, ConcreteDataType::float64_datatype(), true);
88 column_schemas.push(col_schema);
89 } else if is_tag_column {
90 let col_schema =
91 ColumnSchema::new(name, ConcreteDataType::string_datatype(), true);
92 column_schemas.push(col_schema);
93 } else {
94 column_schemas.push(col_schema);
96 }
97 }
98 }
99 }
100
101 if query_type == &QueryType::Sql {
102 let update_at_schema = ColumnSchema::new(
103 AUTO_CREATED_UPDATE_AT_TS_COL,
104 ConcreteDataType::timestamp_millisecond_datatype(),
105 true,
106 );
107 column_schemas.push(update_at_schema);
108 }
109
110 let time_index = if let Some(time_index) = first_time_stamp {
111 time_index
112 } else {
113 column_schemas.push(
114 ColumnSchema::new(
115 AUTO_CREATED_PLACEHOLDER_TS_COL,
116 ConcreteDataType::timestamp_millisecond_datatype(),
117 false,
118 )
119 .with_time_index(true),
120 );
121 AUTO_CREATED_PLACEHOLDER_TS_COL.to_string()
122 };
123
124 let column_defs =
125 column_schemas_to_defs(column_schemas, &primary_keys).context(ConvertColumnSchemaSnafu)?;
126 Ok(CreateTableExpr {
127 catalog_name: sink_table_name[0].clone(),
128 schema_name: sink_table_name[1].clone(),
129 table_name: sink_table_name[2].clone(),
130 desc: "Auto created table by flow engine".to_string(),
131 column_defs,
132 time_index,
133 primary_keys,
134 create_if_not_exists: true,
135 table_options: Default::default(),
136 table_id: None,
137 engine: "mito".to_string(),
138 })
139}
140
141fn build_by_sql_schema(plan: &LogicalPlan) -> Result<TableDef, Error> {
143 let first_time_stamp = plan.schema().fields().iter().find_map(|f| {
144 if ConcreteDataType::from_arrow_type(f.data_type()).is_timestamp() {
145 Some(f.name().clone())
146 } else {
147 None
148 }
149 });
150 Ok(TableDef {
151 ts_col: first_time_stamp,
152 pks: vec![],
153 })
154}
155
156fn build_by_tql_schema(plan: &LogicalPlan) -> Result<TableDef, Error> {
158 let first_time_stamp = plan.schema().fields().iter().find_map(|f| {
159 if ConcreteDataType::from_arrow_type(f.data_type()).is_timestamp() {
160 Some(f.name().clone())
161 } else {
162 None
163 }
164 });
165 let string_columns = plan
166 .schema()
167 .fields()
168 .iter()
169 .filter_map(|f| {
170 if map_dictionary_to_values_data_type(&ConcreteDataType::from_arrow_type(f.data_type()))
171 .is_string()
172 {
173 Some(f.name().clone())
174 } else {
175 None
176 }
177 })
178 .collect::<Vec<_>>();
179
180 Ok(TableDef {
181 ts_col: first_time_stamp,
182 pks: string_columns,
183 })
184}
185
186struct TableDef {
187 ts_col: Option<String>,
188 pks: Vec<String>,
189}
190
191fn build_pk_from_aggr(plan: &LogicalPlan) -> Result<Option<TableDef>, Error> {
200 let fields = plan.schema().fields();
201 let mut pk_names = FindGroupByFinalName::default();
202
203 plan.visit(&mut pk_names)
204 .with_context(|_| DatafusionSnafu {
205 context: format!("Can't find aggr expr in plan {plan:?}"),
206 })?;
207
208 let Some(pk_final_names) = pk_names.get_group_expr_names() else {
210 return Ok(None);
211 };
212 if pk_final_names.is_empty() {
213 let first_ts_col = fields
214 .iter()
215 .find(|f| ConcreteDataType::from_arrow_type(f.data_type()).is_timestamp())
216 .map(|f| f.name().clone());
217 return Ok(Some(TableDef {
218 ts_col: first_ts_col,
219 pks: vec![],
220 }));
221 }
222
223 let all_pk_cols: Vec<_> = fields
224 .iter()
225 .filter(|f| pk_final_names.contains(f.name()))
226 .map(|f| f.name().clone())
227 .collect();
228 let first_time_stamp = fields
234 .iter()
235 .find(|f| {
236 all_pk_cols.contains(&f.name().clone())
237 && ConcreteDataType::from_arrow_type(f.data_type()).is_timestamp()
238 })
239 .map(|f| f.name().clone());
240
241 let all_pk_cols: Vec<_> = all_pk_cols
242 .into_iter()
243 .filter(|col| first_time_stamp.as_ref() != Some(col))
244 .collect();
245
246 Ok(Some(TableDef {
247 ts_col: first_time_stamp,
248 pks: all_pk_cols,
249 }))
250}
251
252#[cfg(test)]
253mod test {
254 use std::sync::Arc;
255
256 use api::v1::column_def::try_as_column_schema;
257 use catalog::RegisterTableRequest;
258 use catalog::memory::new_memory_catalog_manager;
259 use common_catalog::consts::{DEFAULT_CATALOG_NAME, DEFAULT_SCHEMA_NAME};
260 use datafusion::arrow::datatypes::{
261 DataType as ArrowDataType, Field, Schema as ArrowSchema, TimeUnit,
262 };
263 use datafusion_common::DFSchema;
264 use datafusion_expr::logical_plan::EmptyRelation;
265 use datatypes::prelude::ConcreteDataType;
266 use datatypes::schema::{ColumnSchema, Schema};
267 use pretty_assertions::assert_eq;
268 use query::options::QueryOptions;
269 use query::{QueryEngineFactory, QueryEngineRef};
270 use session::context::QueryContext;
271 use table::metadata::{TableInfoBuilder, TableMetaBuilder};
272 use table::test_util::EmptyTable;
273
274 use super::*;
275 use crate::adapter::{AUTO_CREATED_PLACEHOLDER_TS_COL, AUTO_CREATED_UPDATE_AT_TS_COL};
276 use crate::batching_mode::utils::sql_to_df_plan;
277 use crate::test_utils::create_test_query_engine;
278
279 #[test]
280 fn test_tql_dictionary_string_is_label() {
281 let arrow_schema = Arc::new(ArrowSchema::new(vec![
282 Field::new_dictionary("host", ArrowDataType::UInt32, ArrowDataType::Utf8, true),
283 Field::new("value", ArrowDataType::Float64, true),
284 Field::new(
285 "ts",
286 ArrowDataType::Timestamp(TimeUnit::Millisecond, None),
287 false,
288 ),
289 ]));
290 let plan = LogicalPlan::EmptyRelation(EmptyRelation {
291 produce_one_row: false,
292 schema: Arc::new(DFSchema::try_from(arrow_schema).unwrap()),
293 });
294
295 let expr = create_table_with_expr(
296 &plan,
297 &[
298 "greptime".to_string(),
299 "public".to_string(),
300 "sink".to_string(),
301 ],
302 &QueryType::Tql,
303 )
304 .unwrap();
305 let columns = expr
306 .column_defs
307 .iter()
308 .map(|column| try_as_column_schema(column).unwrap())
309 .collect::<Vec<_>>();
310
311 assert_eq!(vec!["host".to_string()], expr.primary_keys);
312 assert_eq!("ts", expr.time_index);
313 assert_eq!(ConcreteDataType::string_datatype(), columns[0].data_type);
314 assert_eq!(ConcreteDataType::float64_datatype(), columns[1].data_type);
315 assert!(columns[2].is_time_index());
316 }
317
318 fn create_tql_test_query_engine() -> QueryEngineRef {
322 let catalog_list = new_memory_catalog_manager().unwrap();
323 let table_meta = TableMetaBuilder::empty()
324 .schema(Arc::new(Schema::new(vec![
325 ColumnSchema::new(
326 "ts",
327 ConcreteDataType::timestamp_millisecond_datatype(),
328 false,
329 )
330 .with_time_index(true),
331 ColumnSchema::new("host", ConcreteDataType::string_datatype(), true),
332 ColumnSchema::new("idc", ConcreteDataType::string_datatype(), true),
333 ColumnSchema::new("val", ConcreteDataType::float64_datatype(), true),
334 ])))
335 .primary_key_indices(vec![1, 2])
337 .value_indices(vec![3])
338 .engine("mito".to_string())
339 .next_column_id(1026)
340 .build()
341 .unwrap();
342 let table_info = TableInfoBuilder::default()
343 .name("http_requests".to_string())
344 .meta(table_meta)
345 .build()
346 .unwrap();
347 assert!(
348 catalog_list
349 .register_table_sync(RegisterTableRequest {
350 catalog: DEFAULT_CATALOG_NAME.to_string(),
351 schema: DEFAULT_SCHEMA_NAME.to_string(),
352 table_name: "http_requests".to_string(),
353 table_id: 1026,
354 table: EmptyTable::from_table_info(&table_info),
355 })
356 .is_ok()
357 );
358
359 QueryEngineFactory::new(
360 catalog_list,
361 None,
362 None,
363 None,
364 None,
365 false,
366 QueryOptions::default(),
367 )
368 .query_engine()
369 }
370
371 #[tokio::test]
380 async fn test_tql_count_values_generated_label_is_primary_key() {
381 let query_engine = create_tql_test_query_engine();
382 let ctx = QueryContext::arc();
383
384 for optimize in [false, true] {
385 let plan = sql_to_df_plan(
386 ctx.clone(),
387 query_engine.clone(),
388 r#"TQL EVAL (0, 15, '5s') count_values("status_code", http_requests)"#,
389 optimize,
390 )
391 .await
392 .unwrap();
393 let plan_display = plan.display_indent_schema().to_string();
394 let expr = create_table_with_expr(
395 &plan,
396 &[
397 "greptime".to_string(),
398 "public".to_string(),
399 "sink".to_string(),
400 ],
401 &QueryType::Tql,
402 )
403 .unwrap();
404 let columns = expr
405 .column_defs
406 .iter()
407 .map(|column| try_as_column_schema(column).unwrap())
408 .collect::<Vec<_>>();
409
410 assert_eq!(
411 vec!["status_code".to_string()],
412 expr.primary_keys,
413 "optimize={optimize}, plan:\n{plan_display}"
414 );
415 assert_eq!(
416 "ts", expr.time_index,
417 "optimize={optimize}, plan:\n{plan_display}"
418 );
419 assert_eq!(
421 "count(http_requests.val)", columns[0].name,
422 "optimize={optimize}, plan:\n{plan_display}"
423 );
424 assert_eq!(ConcreteDataType::float64_datatype(), columns[0].data_type);
425 assert_eq!("ts", columns[1].name);
426 assert!(columns[1].is_time_index());
427 assert_eq!("status_code", columns[2].name);
428 assert_eq!(ConcreteDataType::string_datatype(), columns[2].data_type);
429 }
430 }
431
432 #[tokio::test]
433 async fn test_gen_create_table_sql() {
434 let query_engine = create_test_query_engine();
435 let ctx = QueryContext::arc();
436 struct TestCase {
437 sql: String,
438 sink_table_name: String,
439 column_schemas: Vec<ColumnSchema>,
440 primary_keys: Vec<String>,
441 time_index: String,
442 }
443
444 let update_at_schema = ColumnSchema::new(
445 AUTO_CREATED_UPDATE_AT_TS_COL,
446 ConcreteDataType::timestamp_millisecond_datatype(),
447 true,
448 );
449
450 let ts_placeholder_schema = ColumnSchema::new(
451 AUTO_CREATED_PLACEHOLDER_TS_COL,
452 ConcreteDataType::timestamp_millisecond_datatype(),
453 false,
454 )
455 .with_time_index(true);
456
457 let testcases = vec![
458 TestCase {
459 sql: "SELECT number, ts FROM numbers_with_ts".to_string(),
460 sink_table_name: "new_table".to_string(),
461 column_schemas: vec![
462 ColumnSchema::new("number", ConcreteDataType::uint32_datatype(), true),
463 ColumnSchema::new(
464 "ts",
465 ConcreteDataType::timestamp_millisecond_datatype(),
466 false,
467 )
468 .with_time_index(true),
469 update_at_schema.clone(),
470 ],
471 primary_keys: vec![],
472 time_index: "ts".to_string(),
473 },
474 TestCase {
475 sql: "SELECT number, max(ts) FROM numbers_with_ts GROUP BY number".to_string(),
476 sink_table_name: "new_table".to_string(),
477 column_schemas: vec![
478 ColumnSchema::new("number", ConcreteDataType::uint32_datatype(), true),
479 ColumnSchema::new(
480 "max(numbers_with_ts.ts)",
481 ConcreteDataType::timestamp_millisecond_datatype(),
482 true,
483 ),
484 update_at_schema.clone(),
485 ts_placeholder_schema.clone(),
486 ],
487 primary_keys: vec!["number".to_string()],
488 time_index: AUTO_CREATED_PLACEHOLDER_TS_COL.to_string(),
489 },
490 TestCase {
491 sql: "SELECT max(number), ts FROM numbers_with_ts GROUP BY ts".to_string(),
492 sink_table_name: "new_table".to_string(),
493 column_schemas: vec![
494 ColumnSchema::new(
495 "max(numbers_with_ts.number)",
496 ConcreteDataType::uint32_datatype(),
497 true,
498 ),
499 ColumnSchema::new(
500 "ts",
501 ConcreteDataType::timestamp_millisecond_datatype(),
502 false,
503 )
504 .with_time_index(true),
505 update_at_schema.clone(),
506 ],
507 primary_keys: vec![],
508 time_index: "ts".to_string(),
509 },
510 TestCase {
511 sql: "SELECT number, ts FROM numbers_with_ts GROUP BY ts, number".to_string(),
512 sink_table_name: "new_table".to_string(),
513 column_schemas: vec![
514 ColumnSchema::new("number", ConcreteDataType::uint32_datatype(), true),
515 ColumnSchema::new(
516 "ts",
517 ConcreteDataType::timestamp_millisecond_datatype(),
518 false,
519 )
520 .with_time_index(true),
521 update_at_schema.clone(),
522 ],
523 primary_keys: vec!["number".to_string()],
524 time_index: "ts".to_string(),
525 },
526 ];
527
528 for tc in testcases {
529 let plan = sql_to_df_plan(ctx.clone(), query_engine.clone(), &tc.sql, true)
530 .await
531 .unwrap();
532 let expr = create_table_with_expr(
533 &plan,
534 &[
535 "greptime".to_string(),
536 "public".to_string(),
537 tc.sink_table_name.clone(),
538 ],
539 &QueryType::Sql,
540 )
541 .unwrap();
542 let column_schemas = expr
544 .column_defs
545 .iter()
546 .map(|c| try_as_column_schema(c).unwrap())
547 .collect::<Vec<_>>();
548 assert_eq!(tc.column_schemas, column_schemas, "{:?}", tc.sql);
549 assert_eq!(tc.primary_keys, expr.primary_keys, "{:?}", tc.sql);
550 assert_eq!(tc.time_index, expr.time_index, "{:?}", tc.sql);
551 }
552 }
553}