Skip to main content

flow/batching_mode/
table_creator.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 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    /// query is a tql query
32    Tql,
33    /// query is a sql query
34    Sql,
35}
36
37// auto created table have a auto added column `update_at`, and optional have a `AUTO_CREATED_PLACEHOLDER_TS_COL` column for time index placeholder if no timestamp column is specified
38// TODO(discord9): for now no default value is set for auto added column for compatibility reason with streaming mode, but this might change in favor of simpler code?
39pub(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            // first try build from aggr, then from tql schema because tql query might not have aggr node
54            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                // if is val column, need to rename as val DOUBLE NULL
82                // if is tag column, need to cast type as STRING NULL
83                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                    // time index column
95                    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
141/// simply build by schema, return first timestamp column and no primary key
142fn 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
156/// Return first timestamp column found in output schema and all string columns
157fn 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
191/// Return first timestamp column which is in group by clause and other columns which are also in group by clause
192///
193/// # Returns
194///
195/// * `Option<String>` - first timestamp column which is in group by clause
196/// * `Vec<String>` - other columns which are also in group by clause
197///
198/// if no aggregation found, return None
199fn 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    // if no group by clause, return empty with first timestamp column found in output schema
209    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    // Auto-created tables use the first timestamp column in the group-by keys
229    // as the time index. It is possible that timestamp columns appear only as
230    // aggregate outputs (for example `max(ts)`) and are not group-by keys; in
231    // that case `first_time_stamp` stays `None` and the caller falls back to a
232    // placeholder time index column.
233    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    /// Creates a query engine holding a Prometheus shaped table `http_requests`: tags
319    /// (`host`, `idc`), a single f64 value column (`val`) and a time index (`ts`), so that
320    /// TQL queries can be planned against it.
321    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            // `host` and `idc` are tags, `val` is the value column
336            .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    /// A TQL `count_values` flow must keep the generated label as a primary key of the auto
372    /// created sink table. `count_values("status_code", http_requests)` groups by the sample
373    /// value column and projects the label as a unary scalar expression of that column
374    /// (`prom_float_to_string(val) AS status_code`), so the label only derives from a group by
375    /// column instead of referencing it directly.
376    ///
377    /// The plan is built with the same path a flow task uses (including the DataFusion
378    /// optimizers, which may rewrite the shape of the alias), and the assertion covers both.
379    #[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            // the aggregation output is a value column, the generated label is a tag column
420            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            // TODO(discord9): assert expr
543            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}