diff --git a/src/flow/src/batching_mode/table_creator.rs b/src/flow/src/batching_mode/table_creator.rs index 73a17f94967..778ee939a8e 100644 --- a/src/flow/src/batching_mode/table_creator.rs +++ b/src/flow/src/batching_mode/table_creator.rs @@ -254,15 +254,22 @@ mod test { use std::sync::Arc; use api::v1::column_def::try_as_column_schema; + use catalog::RegisterTableRequest; + use catalog::memory::new_memory_catalog_manager; + use common_catalog::consts::{DEFAULT_CATALOG_NAME, DEFAULT_SCHEMA_NAME}; use datafusion::arrow::datatypes::{ DataType as ArrowDataType, Field, Schema as ArrowSchema, TimeUnit, }; use datafusion_common::DFSchema; use datafusion_expr::logical_plan::EmptyRelation; use datatypes::prelude::ConcreteDataType; - use datatypes::schema::ColumnSchema; + use datatypes::schema::{ColumnSchema, Schema}; use pretty_assertions::assert_eq; + use query::options::QueryOptions; + use query::{QueryEngineFactory, QueryEngineRef}; use session::context::QueryContext; + use table::metadata::{TableInfoBuilder, TableMetaBuilder}; + use table::test_util::EmptyTable; use super::*; use crate::adapter::{AUTO_CREATED_PLACEHOLDER_TS_COL, AUTO_CREATED_UPDATE_AT_TS_COL}; @@ -308,6 +315,120 @@ mod test { assert!(columns[2].is_time_index()); } + /// Creates a query engine holding a Prometheus shaped table `http_requests`: tags + /// (`host`, `idc`), a single f64 value column (`val`) and a time index (`ts`), so that + /// TQL queries can be planned against it. + fn create_tql_test_query_engine() -> QueryEngineRef { + let catalog_list = new_memory_catalog_manager().unwrap(); + let table_meta = TableMetaBuilder::empty() + .schema(Arc::new(Schema::new(vec![ + ColumnSchema::new( + "ts", + ConcreteDataType::timestamp_millisecond_datatype(), + false, + ) + .with_time_index(true), + ColumnSchema::new("host", ConcreteDataType::string_datatype(), true), + ColumnSchema::new("idc", ConcreteDataType::string_datatype(), true), + ColumnSchema::new("val", ConcreteDataType::float64_datatype(), true), + ]))) + // `host` and `idc` are tags, `val` is the value column + .primary_key_indices(vec![1, 2]) + .value_indices(vec![3]) + .engine("mito".to_string()) + .next_column_id(1026) + .build() + .unwrap(); + let table_info = TableInfoBuilder::default() + .name("http_requests".to_string()) + .meta(table_meta) + .build() + .unwrap(); + assert!( + catalog_list + .register_table_sync(RegisterTableRequest { + catalog: DEFAULT_CATALOG_NAME.to_string(), + schema: DEFAULT_SCHEMA_NAME.to_string(), + table_name: "http_requests".to_string(), + table_id: 1026, + table: EmptyTable::from_table_info(&table_info), + }) + .is_ok() + ); + + QueryEngineFactory::new( + catalog_list, + None, + None, + None, + None, + false, + QueryOptions::default(), + ) + .query_engine() + } + + /// A TQL `count_values` flow must keep the generated label as a primary key of the auto + /// created sink table. `count_values("status_code", http_requests)` groups by the sample + /// value column and projects the label as a unary scalar expression of that column + /// (`prom_float_to_string(val) AS status_code`), so the label only derives from a group by + /// column instead of referencing it directly. + /// + /// The plan is built with the same path a flow task uses (including the DataFusion + /// optimizers, which may rewrite the shape of the alias), and the assertion covers both. + #[tokio::test] + async fn test_tql_count_values_generated_label_is_primary_key() { + let query_engine = create_tql_test_query_engine(); + let ctx = QueryContext::arc(); + + for optimize in [false, true] { + let plan = sql_to_df_plan( + ctx.clone(), + query_engine.clone(), + r#"TQL EVAL (0, 15, '5s') count_values("status_code", http_requests)"#, + optimize, + ) + .await + .unwrap(); + let plan_display = plan.display_indent_schema().to_string(); + let expr = create_table_with_expr( + &plan, + &[ + "greptime".to_string(), + "public".to_string(), + "sink".to_string(), + ], + &QueryType::Tql, + ) + .unwrap(); + let columns = expr + .column_defs + .iter() + .map(|column| try_as_column_schema(column).unwrap()) + .collect::>(); + + assert_eq!( + vec!["status_code".to_string()], + expr.primary_keys, + "optimize={optimize}, plan:\n{plan_display}" + ); + assert_eq!( + "ts", expr.time_index, + "optimize={optimize}, plan:\n{plan_display}" + ); + // the aggregation output is a value column, the generated label is a tag column + assert_eq!( + "count(http_requests.val)", columns[0].name, + "optimize={optimize}, plan:\n{plan_display}" + ); + assert_eq!(ConcreteDataType::float64_datatype(), columns[0].data_type); + assert_eq!("ts", columns[1].name); + assert!(columns[1].is_time_index()); + assert_eq!("status_code", columns[2].name); + assert_eq!(ConcreteDataType::string_datatype(), columns[2].data_type); + } + } + #[tokio::test] async fn test_gen_create_table_sql() { let query_engine = create_test_query_engine(); diff --git a/src/flow/src/batching_mode/utils.rs b/src/flow/src/batching_mode/utils.rs index 0220f4ef48d..267fa170cdd 100644 --- a/src/flow/src/batching_mode/utils.rs +++ b/src/flow/src/batching_mode/utils.rs @@ -1164,6 +1164,16 @@ pub fn df_plan_to_sql(plan: &LogicalPlan) -> Result { Ok(sql.to_string()) } +/// Returns whether `expr` directly renames the group by expression `group_expr`. +/// +/// Renames are identified only by expression name. Derived unary expressions are not group +/// keys, even when they reference a single group key column. A PromQL `count_values("label", +/// metric)` plan lands here too: it groups by the formatted sample value +/// (`prom_float_to_string(value)`) and projects that group key column as the generated label. +fn is_alias_of_group_expr(group_expr: &Expr, expr: &Expr) -> DfResult { + Ok(group_expr.name_for_alias()? == expr.name_for_alias()?) +} + /// Helper to find the innermost group by expr in schema, return None if no group by expr #[derive(Debug, Clone, Default)] pub struct FindGroupByFinalName { @@ -1220,24 +1230,28 @@ impl TreeNodeVisitor<'_> for FindGroupByFinalName { Ok(TreeNodeRecursion::Continue) } - /// deal with projection when going up with group exprs + /// Applies only direct, name-matched aliases when propagating group expressions upward. fn f_up(&mut self, node: &Self::Node) -> datafusion_common::Result { if let LogicalPlan::Projection(projection) = node { + // Only direct aliases of group expressions rename group keys. Derived unary + // expressions remain ordinary projected columns. for expr in &projection.expr { let Some(group_exprs) = &mut self.group_exprs else { return Ok(TreeNodeRecursion::Continue); }; if let datafusion_expr::Expr::Alias(alias) = expr { // if a alias exist, replace with the new alias - let mut new_group_exprs = group_exprs.clone(); + let mut renamed_group_expr = None; for group_expr in group_exprs.iter() { - if group_expr.name_for_alias()? == alias.expr.name_for_alias()? { - new_group_exprs.remove(group_expr); - new_group_exprs.insert(expr.clone()); + if is_alias_of_group_expr(group_expr, alias.expr.as_ref())? { + renamed_group_expr = Some(group_expr.clone()); break; } } - *group_exprs = new_group_exprs; + if let Some(group_expr) = renamed_group_expr { + group_exprs.remove(&group_expr); + group_exprs.insert(expr.clone()); + } } } } diff --git a/src/flow/src/batching_mode/utils/test.rs b/src/flow/src/batching_mode/utils/test.rs index 6315ae3d679..1f390d1594c 100644 --- a/src/flow/src/batching_mode/utils/test.rs +++ b/src/flow/src/batching_mode/utils/test.rs @@ -1033,6 +1033,34 @@ async fn test_find_group_by_exprs() { } } +#[tokio::test] +async fn test_find_group_by_exprs_does_not_replace_group_key_with_derived_expr() { + let query_engine = create_test_query_engine(); + let ctx = QueryContext::arc(); + let sql = "SELECT host, lower(host) AS host_lc, ts, SUM(val) AS total \ + FROM (SELECT CAST(number AS STRING) AS host, ts, number AS val FROM numbers_with_ts) \ + GROUP BY host, ts"; + + for optimize in [false, true] { + let plan = sql_to_df_plan(ctx.clone(), query_engine.clone(), sql, optimize) + .await + .unwrap(); + let plan_display = plan.display_indent_schema().to_string(); + let mut group_finder = FindGroupByFinalName::default(); + plan.visit(&mut group_finder).unwrap(); + let group_names = group_finder.get_group_expr_names().unwrap_or_default(); + + assert!( + group_names.contains("host"), + "optimize={optimize}, group keys {group_names:?} must keep host:\n{plan_display}" + ); + assert!( + !group_names.contains("host_lc"), + "optimize={optimize}, derived host_lc must not replace host:\n{plan_display}" + ); + } +} + #[tokio::test] async fn test_analyze_incremental_aggregate_plan() { let query_engine = create_test_query_engine(); @@ -1243,6 +1271,64 @@ async fn test_rewrite_incremental_aggregate_allows_alias_wrapped_scan() { assert_eq!(rewritten_fields, analysis.output_field_names); } +#[tokio::test] +async fn test_count_values_flow_plan_keeps_generated_label_as_group_key() { + // PromQL `count_values("label", metric)` groups by the *formatted* sample value + // (`prom_float_to_string(value)`), so such a flow carries a scalar UDF inside its GROUP + // BY instead of the raw value column: + // + // Projection: count(val), ts, prom_float_to_string(val) AS v + // Aggregate: groupBy=[[ts, prom_float_to_string(CAST(val AS Float64))]], aggr=[[count(val)]] + // + // Two things must keep working for that shape: the TQL flow transport (a Substrait + // encoded insert plan) and the group-key resolution that derives the sink schema, which + // is what keeps the generated label a primary key of the auto-created sink table. + let query_engine = create_test_query_engine(); + let ctx = QueryContext::arc(); + for optimize in [false, true] { + let plan = sql_to_df_plan( + ctx.clone(), + query_engine.clone(), + r#"TQL EVAL (0, 15, '5s') count_values("v", numbers_with_ts)"#, + optimize, + ) + .await + .unwrap(); + let plan_display = plan.display_indent_schema().to_string(); + + DFLogicalSubstraitConvertor {} + .encode(&plan, DefaultSerializer) + .unwrap_or_else(|err| panic!("optimize={optimize}: {err}\n{plan_display}")); + + let mut group_finder = FindGroupByFinalName::default(); + plan.visit(&mut group_finder).unwrap(); + let group_names = group_finder.get_group_expr_names().unwrap_or_default(); + assert!( + group_names.contains("v"), + "optimize={optimize}, group keys {group_names:?} must keep the generated label:\n{plan_display}" + ); + assert!( + group_names.contains("ts"), + "optimize={optimize}, group keys {group_names:?}:\n{plan_display}" + ); + + // The flow plan keeps the generated label as a group key name, and the plan shape + // (PromQL extension nodes above the aggregate) stays outside the incremental + // aggregate rewrite: `prepare_plan_for_incremental` only rewrites `QueryType::Sql` + // flows, so a `count_values` flow keeps running as a full snapshot. + let analysis = analyze_incremental_aggregate_plan(&plan).unwrap().unwrap(); + assert!( + analysis.group_key_names.contains(&"v".to_string()), + "optimize={optimize}, analysis {analysis:?}:\n{plan_display}" + ); + assert!( + !analysis.unsupported_exprs.is_empty(), + "optimize={optimize}, a count_values flow plan must not become an incremental \ + aggregate rewrite candidate: {analysis:?}\n{plan_display}" + ); + } +} + #[tokio::test] async fn test_analyze_incremental_aggregate_plan_rejects_having_filter() { let sql = diff --git a/src/query/src/promql/planner.rs b/src/query/src/promql/planner.rs index 2142ee82371..73332c237c8 100644 --- a/src/query/src/promql/planner.rs +++ b/src/query/src/promql/planner.rs @@ -713,13 +713,76 @@ impl PromPlanner { let label = Self::get_param_value_as_str(*op, param)?; // `count_values` must be grouped by fields, // and project the fields to the new label. - let count_value_exprs = prev_field_exprs.iter().map(|expr| { - match expr { - DfExpr::Column(column) => DfExpr::Column(column.clone()), - _ => DfExpr::Column(Column::from_name(expr.schema_name().to_string())), - } - .alias(label) - }); + // + // The generated label is a real label (column) of the output, so it must be + // registered in `ctx.tag_columns` below. Otherwise enclosing expressions + // rebuild their projections from `ctx.tag_columns` and silently drop it. + // + // PromQL sets the generated label *before* the grouping key is built, so an + // input label with the same name is overwritten by the sample value and must + // not remain a grouping key either: samples are grouped by the generated + // label only. Dropping it from the projected tag columns is required as well: + // projecting both would emit two columns with the same name (rejected as an + // ambiguous reference). + self.ctx.tag_columns.retain(|tag| tag != label); + group_exprs.retain( + |expr| !matches!(expr, DfExpr::Column(column) if column.name == label), + ); + // The tag columns projected below are unqualified `Column` references, so they + // end up qualified with whatever qualifier they carry in the input plan. Give + // the generated label the same qualifier, otherwise qualified references to it + // (e.g. from an enclosing binary expression or a vector join) fail to resolve. + let label_qualifier = self + .ctx + .time_index_column + .as_deref() + .and_then(|time_index| { + builder + .schema() + .qualified_field_with_unqualified_name(time_index) + .ok() + }) + .and_then(|(qualifier, _)| qualifier.cloned()); + // The generated label carries the *sample value*, so it must be materialized + // as a string using Prometheus' format + // (`strconv.FormatFloat(value, 'f', -1, 64)`: shortest decimal form without + // an exponent, `1.0` becomes "1"). Labels are inferred from string columns + // when a result is converted into the Prometheus HTTP API JSON format, so a + // numeric label column would be mistaken for the sample value of the series. + let count_value_exprs = prev_field_exprs + .iter() + .map(|expr| { + let value = match expr { + DfExpr::Column(column) => { + let value = DfExpr::Column(column.clone()); + // `prom_float_to_string` formats exactly like Prometheus, + // while arrow's `Float64 -> Utf8` cast would render `1.0`. + let value = if matches!( + builder.schema().field_with_unqualified_name(&column.name), + Ok(field) if field.data_type() == &ArrowDataType::Float64 + ) { + value + } else { + DfExpr::Cast(Cast::new( + Box::new(value), + ArrowDataType::Float64, + )) + }; + DfExpr::ScalarFunction(ScalarFunction { + func: Arc::new(PromqlFloatToString::scalar_udf()), + args: vec![value], + }) + } + // The value is already formatted (e.g. a native histogram is + // converted to its string form by the aggregation), and the + // aggregate output names it by its schema name. + _ => DfExpr::Column(Column::from_name( + expr.schema_name().to_string(), + )), + }; + DfExpr::Alias(Alias::new(value, label_qualifier.clone(), label)) + }) + .collect::>(); let aggregate_group_exprs = group_exprs .iter() .cloned() @@ -733,11 +796,14 @@ impl PromPlanner { .chain(Some(self.create_time_index_column_expr()?)) .chain(count_value_exprs); - builder + let builder = builder .aggregate(aggregate_group_exprs, aggr_exprs) .context(DataFusionPlanningSnafu)? .project(project_fields) - .context(DataFusionPlanningSnafu)? + .context(DataFusionPlanningSnafu)?; + // The label only exists in the output schema from here on. + self.ctx.tag_columns.push(label.to_string()); + builder } else { builder .aggregate(group_exprs.clone(), aggr_exprs) @@ -4445,9 +4511,10 @@ impl PromPlanner { /// /// Returns a tuple of `(aggregate_expressions, previous_field_expressions)` where: /// - `aggregate_expressions`: Expressions that apply the aggregate function to the original fields - /// - `previous_field_expressions`: Original field expressions before aggregation. This is non-empty - /// only when the operation is `count_values`, as this operation requires preserving the original - /// values for grouping. + /// - `previous_field_expressions`: Field expressions naming the pre-aggregation values. This is + /// non-empty only when the operation is `count_values`, which groups by the sample value and + /// projects it as the generated label, so these expressions are passed through the same + /// formatting as that label (`prom_float_to_string`). /// fn create_aggregate_exprs( &mut self, @@ -4491,14 +4558,13 @@ impl PromPlanner { .collect::>>()?; // if the aggregator is `count_values`, it must be grouped by current fields. + // + // The grouping key is the *formatted* sample value, i.e. the same expression that + // produces the generated label below: PromQL groups by the value, and the label is + // that value in Prometheus' textual form (`strconv.FormatFloat(value, 'f', -1, 64)`), + // so grouping by the raw value would split samples that render to one label into + // several groups, each emitting the same label set for one timestamp. let prev_field_exprs = if op.id() == token::T_COUNT_VALUES { - let prev_field_exprs: Vec<_> = self - .ctx - .field_columns - .iter() - .map(|col| DfExpr::Column(Column::from_name(col))) - .collect(); - ensure!( self.ctx.field_columns.len() == 1, UnsupportedExprSnafu { @@ -4506,7 +4572,28 @@ impl PromPlanner { } ); - prev_field_exprs + self.ctx + .field_columns + .iter() + .map(|col| { + let value = DfExpr::Column(Column::from_name(col)); + // Normalize non `Float64` inputs the same way the label projection does, + // so both sides agree on the formatted value: `prom_float_to_string` + // formats exactly like Prometheus, while arrow's `Float64 -> Utf8` cast + // would render `1.0`. + let value = if Self::field_column_type(input_plan.schema(), col) + == Some(&ArrowDataType::Float64) + { + value + } else { + DfExpr::Cast(Cast::new(Box::new(value), ArrowDataType::Float64)) + }; + DfExpr::ScalarFunction(ScalarFunction { + func: Arc::new(PromqlFloatToString::scalar_udf()), + args: vec![value], + }) + }) + .collect() } else { vec![] }; @@ -7633,7 +7720,7 @@ mod test { use common_query::test_util::DummyDecoder; use common_recordbatch::RecordBatch as GreptimeRecordBatch; use datafusion::arrow::array::{ - Array, Float64Array, Int64Array, StringArray, TimestampMillisecondArray, + Array, ArrayRef, Float64Array, Int64Array, StringArray, TimestampMillisecondArray, }; use datafusion::arrow::datatypes::{Field, Schema as ArrowSchema}; use datafusion::arrow::record_batch::RecordBatch; @@ -13081,9 +13168,9 @@ Filter: up.field_0 IS NOT NULL [timestamp:Timestamp(ms), field_0:Float64;N, foo: PromPlanner::stmt_to_plan(table_provider, &eval_stmt, &build_query_engine_state()) .await .unwrap(); - let expected = "Sort: prometheus_tsdb_head_series.ip ASC NULLS LAST, prometheus_tsdb_head_series.greptime_timestamp ASC NULLS LAST, series ASC NULLS LAST [count(prometheus_tsdb_head_series.greptime_value):Int64, ip:Utf8, greptime_timestamp:Timestamp(ms), series:Float64;N]\ - \n Projection: count(prometheus_tsdb_head_series.greptime_value), prometheus_tsdb_head_series.ip, prometheus_tsdb_head_series.greptime_timestamp, prometheus_tsdb_head_series.greptime_value AS series [count(prometheus_tsdb_head_series.greptime_value):Int64, ip:Utf8, greptime_timestamp:Timestamp(ms), series:Float64;N]\ - \n Aggregate: groupBy=[[prometheus_tsdb_head_series.ip, prometheus_tsdb_head_series.greptime_timestamp, prometheus_tsdb_head_series.greptime_value]], aggr=[[count(prometheus_tsdb_head_series.greptime_value)]] [ip:Utf8, greptime_timestamp:Timestamp(ms), greptime_value:Float64;N, count(prometheus_tsdb_head_series.greptime_value):Int64]\ + let expected = "Sort: prometheus_tsdb_head_series.ip ASC NULLS LAST, prometheus_tsdb_head_series.greptime_timestamp ASC NULLS LAST, prometheus_tsdb_head_series.series ASC NULLS LAST [count(prometheus_tsdb_head_series.greptime_value):Int64, ip:Utf8, greptime_timestamp:Timestamp(ms), series:Utf8;N]\ + \n Projection: count(prometheus_tsdb_head_series.greptime_value), prometheus_tsdb_head_series.ip, prometheus_tsdb_head_series.greptime_timestamp, prom_float_to_string(prometheus_tsdb_head_series.greptime_value) AS series [count(prometheus_tsdb_head_series.greptime_value):Int64, ip:Utf8, greptime_timestamp:Timestamp(ms), series:Utf8;N]\ + \n Aggregate: groupBy=[[prometheus_tsdb_head_series.ip, prometheus_tsdb_head_series.greptime_timestamp, prom_float_to_string(prometheus_tsdb_head_series.greptime_value)]], aggr=[[count(prometheus_tsdb_head_series.greptime_value)]] [ip:Utf8, greptime_timestamp:Timestamp(ms), prom_float_to_string(prometheus_tsdb_head_series.greptime_value):Utf8;N, count(prometheus_tsdb_head_series.greptime_value):Int64]\ \n PromInstantManipulate: range=[0..100000000], lookback=[1000], interval=[5000], time index=[greptime_timestamp] [ip:Utf8, greptime_timestamp:Timestamp(ms), greptime_value:Float64;N]\ \n PromSeriesDivide: tags=[\"ip\"] [ip:Utf8, greptime_timestamp:Timestamp(ms), greptime_value:Float64;N]\ \n Sort: prometheus_tsdb_head_series.ip ASC NULLS FIRST, prometheus_tsdb_head_series.greptime_timestamp ASC NULLS FIRST [ip:Utf8, greptime_timestamp:Timestamp(ms), greptime_value:Float64;N]\ @@ -13129,10 +13216,10 @@ Filter: up.field_0 IS NOT NULL [timestamp:Timestamp(ms), field_0:Float64;N, foo: .await .unwrap(); let expected = r#" -Projection: count(prometheus_tsdb_head_series.greptime_value) AS my_series, prometheus_tsdb_head_series.ip, prometheus_tsdb_head_series.greptime_timestamp [my_series:Int64, ip:Utf8, greptime_timestamp:Timestamp(ms)] - Sort: prometheus_tsdb_head_series.ip ASC NULLS LAST, prometheus_tsdb_head_series.greptime_timestamp ASC NULLS LAST, series ASC NULLS LAST [count(prometheus_tsdb_head_series.greptime_value):Int64, ip:Utf8, greptime_timestamp:Timestamp(ms), series:Float64;N] - Projection: count(prometheus_tsdb_head_series.greptime_value), prometheus_tsdb_head_series.ip, prometheus_tsdb_head_series.greptime_timestamp, prometheus_tsdb_head_series.greptime_value AS series [count(prometheus_tsdb_head_series.greptime_value):Int64, ip:Utf8, greptime_timestamp:Timestamp(ms), series:Float64;N] - Aggregate: groupBy=[[prometheus_tsdb_head_series.ip, prometheus_tsdb_head_series.greptime_timestamp, prometheus_tsdb_head_series.greptime_value]], aggr=[[count(prometheus_tsdb_head_series.greptime_value)]] [ip:Utf8, greptime_timestamp:Timestamp(ms), greptime_value:Float64;N, count(prometheus_tsdb_head_series.greptime_value):Int64] +Projection: count(prometheus_tsdb_head_series.greptime_value) AS my_series, prometheus_tsdb_head_series.ip, prometheus_tsdb_head_series.series, prometheus_tsdb_head_series.greptime_timestamp [my_series:Int64, ip:Utf8, series:Utf8;N, greptime_timestamp:Timestamp(ms)] + Sort: prometheus_tsdb_head_series.ip ASC NULLS LAST, prometheus_tsdb_head_series.greptime_timestamp ASC NULLS LAST, prometheus_tsdb_head_series.series ASC NULLS LAST [count(prometheus_tsdb_head_series.greptime_value):Int64, ip:Utf8, greptime_timestamp:Timestamp(ms), series:Utf8;N] + Projection: count(prometheus_tsdb_head_series.greptime_value), prometheus_tsdb_head_series.ip, prometheus_tsdb_head_series.greptime_timestamp, prom_float_to_string(prometheus_tsdb_head_series.greptime_value) AS series [count(prometheus_tsdb_head_series.greptime_value):Int64, ip:Utf8, greptime_timestamp:Timestamp(ms), series:Utf8;N] + Aggregate: groupBy=[[prometheus_tsdb_head_series.ip, prometheus_tsdb_head_series.greptime_timestamp, prom_float_to_string(prometheus_tsdb_head_series.greptime_value)]], aggr=[[count(prometheus_tsdb_head_series.greptime_value)]] [ip:Utf8, greptime_timestamp:Timestamp(ms), prom_float_to_string(prometheus_tsdb_head_series.greptime_value):Utf8;N, count(prometheus_tsdb_head_series.greptime_value):Int64] PromInstantManipulate: range=[0..100000000], lookback=[1000], interval=[5000], time index=[greptime_timestamp] [ip:Utf8, greptime_timestamp:Timestamp(ms), greptime_value:Float64;N] PromSeriesDivide: tags=["ip"] [ip:Utf8, greptime_timestamp:Timestamp(ms), greptime_value:Float64;N] Sort: prometheus_tsdb_head_series.ip ASC NULLS FIRST, prometheus_tsdb_head_series.greptime_timestamp ASC NULLS FIRST [ip:Utf8, greptime_timestamp:Timestamp(ms), greptime_value:Float64;N] @@ -14884,4 +14971,395 @@ Projection: count(prometheus_tsdb_head_series.greptime_value) AS my_series, prom let sample_count = batches.iter().map(RecordBatch::num_rows).sum::(); assert_eq!(sample_count, 2); } + + /// Table provider with a single metric `cv_metric` holding three series at `ts=1000`: + /// `k="a"` and `k="c"` carry the value 1.0, `k="b"` carries 2.0. + async fn build_count_values_table_provider() -> DfTableSourceProvider { + build_count_values_table_provider_with_values(&[1.0, 2.0, 1.0]).await + } + + /// Table provider with a single metric `cv_metric` holding one series per given value at + /// `ts=1000` (`k` names the series, the sample value is the given value). + async fn build_count_values_table_provider_with_values( + values: &[f64], + ) -> DfTableSourceProvider { + build_count_values_table_provider_with_value_array(Arc::new(Float64Array::from( + values.to_vec(), + ))) + .await + } + + /// Like [`build_count_values_table_provider_with_values`], but with a caller provided value + /// column, so tests can cover value columns that are not `Float64` (e.g. `BIGINT`). + async fn build_count_values_table_provider_with_value_array( + values: ArrayRef, + ) -> DfTableSourceProvider { + let value_data_type = ConcreteDataType::from_arrow_type(values.data_type()); + let catalog_list = MemoryCatalogManager::with_default_setup(); + let columns = vec![ + ColumnSchema::new("k".to_string(), ConcreteDataType::string_datatype(), false), + ColumnSchema::new( + "timestamp".to_string(), + ConcreteDataType::timestamp_millisecond_datatype(), + false, + ) + .with_time_index(true), + ColumnSchema::new(greptime_value().to_string(), value_data_type, true), + ]; + let schema = Arc::new(Schema::new(columns)); + let table_meta = TableMetaBuilder::empty() + .schema(schema.clone()) + .primary_key_indices(vec![0]) + .value_indices(vec![2]) + .next_column_id(1024) + .build() + .unwrap(); + let table_info = Arc::new( + TableInfoBuilder::default() + .table_id(3_001) + .name("cv_metric") + .meta(table_meta) + .build() + .unwrap(), + ); + let batch = RecordBatch::try_new( + schema.arrow_schema().clone(), + vec![ + Arc::new(StringArray::from( + (0..values.len()) + .map(|index| format!("k{index}")) + .collect::>(), + )), + Arc::new(TimestampMillisecondArray::from(vec![1_000; values.len()])), + values.clone(), + ], + ) + .unwrap(); + let backing = GreptimeMemTable::new_with_catalog( + "cv_metric", + GreptimeRecordBatch::from_df_record_batch(schema, batch), + 3_001, + DEFAULT_CATALOG_NAME.to_string(), + DEFAULT_SCHEMA_NAME.to_string(), + ); + let table = Arc::new(Table::new( + table_info, + FilterPushDownType::Unsupported, + backing.data_source(), + )); + + assert!( + catalog_list + .register_table_sync(RegisterTableRequest { + catalog: DEFAULT_CATALOG_NAME.to_string(), + schema: DEFAULT_SCHEMA_NAME.to_string(), + table_name: "cv_metric".to_string(), + table_id: 3_001, + table, + }) + .is_ok() + ); + + DfTableSourceProvider::new( + catalog_list, + false, + QueryContext::arc(), + DummyDecoder::arc(), + false, + ) + } + + /// Collects `(label, value)` pairs of a `count_values` result, where `label` is the + /// PromQL label generated by `count_values` and `value` is the aggregated sample value. + /// + /// The generated label holds the original sample value in PromQL's textual form, so it + /// is asserted as a string: comparing it as a number would not catch formatting bugs + /// (`1.0` instead of `1`, scientific notation, ...). + fn count_values_rows<'a>(batches: &'a [RecordBatch], label: &str) -> Vec<(&'a str, f64)> { + let mut rows = batches + .iter() + .flat_map(|batch| { + // The aggregated value is the only numeric column that is not the generated label. + let value_index = batch + .schema() + .fields() + .iter() + .position(|field| { + field.name() != label + && matches!( + field.data_type(), + ArrowDataType::Float64 + | ArrowDataType::Int64 + | ArrowDataType::UInt64 + ) + }) + .expect("no aggregated value column"); + let labels = batch + .column_by_name(label) + .expect("no generated label column") + .as_any() + .downcast_ref::() + .expect("the generated label must be a string column"); + let values = datafusion::arrow::compute::cast( + batch.column(value_index), + &ArrowDataType::Float64, + ) + .unwrap(); + let values = values.as_any().downcast_ref::().unwrap(); + labels + .iter() + .zip(values.iter()) + .map(|(label, value)| (label.unwrap(), value.unwrap())) + .collect::>() + }) + .collect::>(); + rows.sort_by(|left, right| left.0.cmp(right.0).then(left.1.total_cmp(&right.1))); + rows + } + + /// Asserts that a `count_values` result holds one sample per label set and evaluation + /// timestamp: Prometheus groups by the generated label, so a timestamp must never repeat + /// a label set (that would mean the samples were still grouped by the overwritten input + /// label). + fn assert_unique_label_set_per_timestamp(batches: &[RecordBatch], label: &str) { + let mut seen = HashMap::>::new(); + for batch in batches { + let timestamp_index = batch + .schema() + .fields() + .iter() + .position(|field| matches!(field.data_type(), ArrowDataType::Timestamp(..))) + .expect("no timestamp column"); + let timestamps = batch + .column(timestamp_index) + .as_any() + .downcast_ref::() + .expect("timestamp column is not a millisecond timestamp"); + let labels = batch + .column_by_name(label) + .expect("no generated label column") + .as_any() + .downcast_ref::() + .expect("the generated label must be a string column"); + for (timestamp, label) in timestamps.iter().zip(labels.iter()) { + let timestamp = timestamp.unwrap(); + let label = label.unwrap(); + assert!( + seen.entry(timestamp).or_default().insert(label.to_string()), + "duplicated label set `{label}` at timestamp {timestamp}" + ); + } + } + } + + #[tokio::test] + async fn test_count_values_generated_label_survives_enclosing_expr() { + // https://github.com/GreptimeTeam/greptimedb/issues/9181 + for (case, label) in [ + (r#"count_values("v", prometheus_tsdb_head_series)"#, "v"), + ( + r#"abs(count_values("v", prometheus_tsdb_head_series))"#, + "v", + ), + ( + r#"round(count_values("v", prometheus_tsdb_head_series))"#, + "v", + ), + (r#"count_values("v", prometheus_tsdb_head_series) + 1"#, "v"), + ( + r#"topk(1, count_values("v", prometheus_tsdb_head_series))"#, + "v", + ), + ( + r#"sum by (v) (count_values("v", prometheus_tsdb_head_series))"#, + "v", + ), + ( + r#"label_replace(count_values("v", prometheus_tsdb_head_series), "vcopy", "$1", "v", "(.*)")"#, + "v", + ), + ( + r#"count_values("v", prometheus_tsdb_head_series) by (ip) + 1"#, + "v", + ), + // The generated label overwrites an input label with the same name. + ( + r#"count_values("ip", prometheus_tsdb_head_series) by (ip)"#, + "ip", + ), + ( + r#"count_values("ip", prometheus_tsdb_head_series) by (ip) + 1"#, + "ip", + ), + ] { + let plan = PromPlanner::stmt_to_plan( + build_test_table_provider_with_fields( + &[( + DEFAULT_SCHEMA_NAME.to_string(), + "prometheus_tsdb_head_series".to_string(), + )], + &["ip"], + ) + .await, + &build_eval_stmt(case), + &build_query_engine_state(), + ) + .await + .unwrap(); + + let label_columns = plan + .schema() + .fields() + .iter() + .filter(|field| field.name() == label) + .count(); + assert_eq!( + label_columns, + 1, + "{case}: the `{label}` label must survive: {}", + plan.display_indent() + ); + } + } + + #[tokio::test] + async fn test_count_values_generated_label_in_enclosing_expr_execute() { + // https://github.com/GreptimeTeam/greptimedb/issues/9181 + let state = build_query_engine_state(); + // (query, generated label, expected `(label value, aggregated value)` pairs) + for (query, label, expected) in [ + ( + r#"count_values("v", cv_metric)"#, + "v", + vec![("1", 2.0), ("2", 1.0)], + ), + ( + r#"abs(count_values("v", cv_metric))"#, + "v", + vec![("1", 2.0), ("2", 1.0)], + ), + ( + r#"round(count_values("v", cv_metric))"#, + "v", + vec![("1", 2.0), ("2", 1.0)], + ), + ( + r#"count_values("v", cv_metric) + 1"#, + "v", + vec![("1", 3.0), ("2", 2.0)], + ), + ( + r#"sum by (v) (count_values("v", cv_metric))"#, + "v", + vec![("1", 2.0), ("2", 1.0)], + ), + ( + r#"topk(10, count_values("v", cv_metric))"#, + "v", + vec![("1", 2.0), ("2", 1.0)], + ), + // The generated label overwrites the input label with the same name, and the + // samples are grouped by the generated label only: `{k="1"}` holds the two + // samples of value `1.0` instead of one row per (overwritten label, value). + ( + r#"count_values("k", cv_metric) by (k)"#, + "k", + vec![("1", 2.0), ("2", 1.0)], + ), + ] { + let plan = PromPlanner::stmt_to_plan( + build_count_values_table_provider().await, + &operator_eval_stmt(query), + &state, + ) + .await + .unwrap_or_else(|err| panic!("{query}: {err}")); + + assert_eq!( + plan.schema() + .fields() + .iter() + .filter(|field| field.name() == label) + .count(), + 1, + "{query}: {}", + plan.display_indent() + ); + + let (_, batches) = execute(plan, &state).await; + assert_eq!(count_values_rows(&batches, label), expected, "{query}"); + assert_unique_label_set_per_timestamp(&batches, label); + } + } + + #[tokio::test] + async fn test_count_values_label_is_prometheus_formatted_value() { + // PromQL materializes the generated label with `strconv.FormatFloat(value, 'f', -1, 64)`: + // the shortest decimal form of the sample value without an exponent. The label is a + // label, so it must be a string column holding exactly that text: arrow's + // `Float64 -> Utf8` cast would render `1`/`200`/`1e21` as `1.0`/`200.0`/`1e21`. + let state = build_query_engine_state(); + // Samples are grouped by that formatted text, exactly like Prometheus groups by the + // generated label: `-0.0` and `0.0` are two series (`-0` and `0`), while values that + // round to the same text share one group. `0.0` formats as "0". + let plan = PromPlanner::stmt_to_plan( + build_count_values_table_provider_with_values(&[ + -0.0, 0.0, 1.0, 0.5, 200.0, 1e21, 1e-7, 2.5, + ]) + .await, + &operator_eval_stmt(r#"count_values("v", cv_metric)"#), + &state, + ) + .await + .unwrap(); + + let (_, batches) = execute(plan, &state).await; + assert_unique_label_set_per_timestamp(&batches, "v"); + let mut labels = count_values_rows(&batches, "v") + .into_iter() + .map(|(label, _)| label) + .collect::>(); + labels.sort(); + assert_eq!( + labels, + vec![ + "-0", + "0", + "0.0000001", + "0.5", + "1", + "1000000000000000000000", + "2.5", + "200", + ] + ); + } + + #[tokio::test] + async fn test_count_values_groups_by_formatted_value_for_bigint_input() { + // The grouping key of `count_values` is the formatted sample value, not the raw input + // value. Two `BIGINT` values that differ below the `Float64` precision (`2^53` and + // `2^53 + 1`) cast and format to the same label, so they must share one group and one + // count, exactly like Prometheus, which groups by the generated label text. + let state = build_query_engine_state(); + let plan = PromPlanner::stmt_to_plan( + build_count_values_table_provider_with_value_array(Arc::new(Int64Array::from(vec![ + 9_007_199_254_740_992_i64, + 9_007_199_254_740_993_i64, + ]))) + .await, + &operator_eval_stmt(r#"count_values("v", cv_metric)"#), + &state, + ) + .await + .unwrap(); + + let (_, batches) = execute(plan, &state).await; + // One timestamp must never carry the same label set twice. + assert_unique_label_set_per_timestamp(&batches, "v"); + assert_eq!( + count_values_rows(&batches, "v"), + vec![("9007199254740992", 2.0)] + ); + } } diff --git a/src/servers/tests/http/mod.rs b/src/servers/tests/http/mod.rs index cca2dfd7873..ffd76b3753b 100644 --- a/src/servers/tests/http/mod.rs +++ b/src/servers/tests/http/mod.rs @@ -16,4 +16,5 @@ mod authorize; mod http_handler_test; mod influxdb_test; mod opentsdb_test; +mod prom_count_values_test; mod prom_store_test; diff --git a/src/servers/tests/http/prom_count_values_test.rs b/src/servers/tests/http/prom_count_values_test.rs new file mode 100644 index 00000000000..ab17ef8791f --- /dev/null +++ b/src/servers/tests/http/prom_count_values_test.rs @@ -0,0 +1,192 @@ +// Copyright 2023 Greptime Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +//! Regression tests for the Prometheus HTTP API conversion of `count_values` results. +//! +//! `count_values` turns a sample value into a label, so the generated column must be a +//! *string* column holding Prometheus' textual form of the value: the conversion infers +//! labels from string columns and the sample value from the first numeric column. A numeric +//! label column is reported as the sample value of the series instead of a label. + +use std::sync::Arc; + +use catalog::RegisterTableRequest; +use catalog::memory::MemoryCatalogManager; +use common_catalog::consts::{DEFAULT_CATALOG_NAME, DEFAULT_SCHEMA_NAME}; +use common_query::prelude::greptime_value; +use common_recordbatch::RecordBatch; +use datatypes::prelude::ConcreteDataType; +use datatypes::schema::{ColumnSchema, Schema}; +use datatypes::vectors::{Float64Vector, StringVector, TimestampMillisecondVector}; +use promql_parser::parser::value::ValueType; +use query::options::QueryOptions; +use query::parser::{PromQuery, QueryLanguageParser}; +use query::query_engine::QueryEngineFactory; +use servers::http::prometheus::{ + PromData, PromQueryResult, PrometheusJsonResponse, PrometheusResponse, +}; +use session::context::QueryContext; +use table::Table; +use table::metadata::{FilterPushDownType, TableInfoBuilder, TableMetaBuilder}; +use table::test_util::MemTable; + +/// Catalog with a single metric `cv_metric` holding one series per value at `ts=1000ms`. +fn metric_catalog_manager(values: &[f64]) -> Arc { + let catalog_list = MemoryCatalogManager::with_default_setup(); + let columns = vec![ + ColumnSchema::new("k".to_string(), ConcreteDataType::string_datatype(), false), + ColumnSchema::new( + "timestamp".to_string(), + ConcreteDataType::timestamp_millisecond_datatype(), + false, + ) + .with_time_index(true), + ColumnSchema::new( + greptime_value().to_string(), + ConcreteDataType::float64_datatype(), + true, + ), + ]; + let schema = Arc::new(Schema::new(columns)); + let table_meta = TableMetaBuilder::empty() + .schema(schema.clone()) + .primary_key_indices(vec![0]) + .value_indices(vec![2]) + .next_column_id(1024) + .build() + .unwrap(); + let table_info = Arc::new( + TableInfoBuilder::default() + .table_id(3_001) + .name("cv_metric") + .meta(table_meta) + .build() + .unwrap(), + ); + let batch = RecordBatch::new( + schema.clone(), + vec![ + Arc::new(StringVector::from( + (0..values.len()) + .map(|index| format!("k{index}")) + .collect::>(), + )) as _, + Arc::new(TimestampMillisecondVector::from_vec(vec![ + 1_000; + values.len() + ])) as _, + Arc::new(Float64Vector::from_vec(values.to_vec())) as _, + ], + ) + .unwrap(); + let backing = MemTable::new_with_catalog( + "cv_metric", + batch, + 3_001, + DEFAULT_CATALOG_NAME.to_string(), + DEFAULT_SCHEMA_NAME.to_string(), + ); + let table = Arc::new(Table::new( + table_info, + FilterPushDownType::Unsupported, + backing.data_source(), + )); + + assert!( + catalog_list + .register_table_sync(RegisterTableRequest { + catalog: DEFAULT_CATALOG_NAME.to_string(), + schema: DEFAULT_SCHEMA_NAME.to_string(), + table_name: "cv_metric".to_string(), + table_id: 3_001, + table, + }) + .is_ok() + ); + + catalog_list +} + +#[tokio::test] +async fn count_values_generated_label_is_reported_as_a_prometheus_label() { + // Two samples carry the value 5.0, one carries 0.5 and one carries 200.0. The query + // evaluates at `ts=1s`, the only timestamp holding samples. + let catalog_list = metric_catalog_manager(&[5.0, 5.0, 0.5, 200.0]); + let query_engine = QueryEngineFactory::new( + catalog_list, + None, + None, + None, + None, + false, + QueryOptions::default(), + ) + .query_engine(); + + let query_ctx = QueryContext::arc(); + let prom_query = PromQuery { + query: r#"count_values("v", cv_metric) + 1"#.to_string(), + start: "1".to_string(), + end: "1".to_string(), + step: "5s".to_string(), + lookback: "5m".to_string(), + alias: None, + }; + let statement = QueryLanguageParser::parse_promql(&prom_query, &query_ctx).unwrap(); + let plan = query_engine + .planner() + .plan(&statement, query_ctx.clone()) + .await + .unwrap(); + let output = query_engine.execute(plan, query_ctx).await.unwrap(); + + let response = + PrometheusJsonResponse::from_query_result(Ok(output), None, ValueType::Vector, None).await; + let PrometheusResponse::PromData(PromData { + result: PromQueryResult::Vector(series), + .. + }) = response.data + else { + panic!("expected a vector response"); + }; + + // `v` is the generated label, so it must be reported as a label holding Prometheus' + // textual form of the sample value (`5`, not `5.0`), and the series carries the + // aggregation result (`count_values("v", cv_metric) + 1`) as its sample value. + let mut actual = series + .into_iter() + .map(|series| { + assert_eq!( + series.metric.keys().collect::>(), + vec!["v"], + "`v` must be the only label of the series" + ); + let label = series.metric["v"].clone(); + let (timestamp, value) = series + .value + .expect("every series of a vector result carries a sample"); + (label, timestamp, value) + }) + .collect::>(); + actual.sort_by(|left, right| left.0.cmp(&right.0)); + + assert_eq!( + actual, + vec![ + ("0.5".to_string(), 1.0, "2.0".to_string()), + ("200".to_string(), 1.0, "2.0".to_string()), + ("5".to_string(), 1.0, "3.0".to_string()), + ] + ); +} diff --git a/tests/cases/distributed/flow-tql/flow_tql.result b/tests/cases/distributed/flow-tql/flow_tql.result index 925fff3f242..30158a176f8 100644 --- a/tests/cases/distributed/flow-tql/flow_tql.result +++ b/tests/cases/distributed/flow-tql/flow_tql.result @@ -228,14 +228,14 @@ SELECT * FROM cnt_reqs ORDER BY ts, status_code; +--------------------------+---------------------+-------------+ | count(http_requests.val) | ts | status_code | +--------------------------+---------------------+-------------+ -| 3.0 | 1970-01-01T00:00:00 | 200.0 | -| 1.0 | 1970-01-01T00:00:00 | 401.0 | -| 1.0 | 1970-01-01T00:00:05 | 401.0 | -| 2.0 | 1970-01-01T00:00:05 | 404.0 | -| 1.0 | 1970-01-01T00:00:05 | 500.0 | -| 2.0 | 1970-01-01T00:00:10 | 200.0 | -| 2.0 | 1970-01-01T00:00:10 | 201.0 | -| 4.0 | 1970-01-01T00:00:15 | 500.0 | +| 3.0 | 1970-01-01T00:00:00 | 200 | +| 1.0 | 1970-01-01T00:00:00 | 401 | +| 1.0 | 1970-01-01T00:00:05 | 401 | +| 2.0 | 1970-01-01T00:00:05 | 404 | +| 1.0 | 1970-01-01T00:00:05 | 500 | +| 2.0 | 1970-01-01T00:00:10 | 200 | +| 2.0 | 1970-01-01T00:00:10 | 201 | +| 4.0 | 1970-01-01T00:00:15 | 500 | +--------------------------+---------------------+-------------+ DROP FLOW calc_reqs; diff --git a/tests/cases/standalone/common/promql/count_values.result b/tests/cases/standalone/common/promql/count_values.result index e2f65eae342..886c3a143b7 100644 --- a/tests/cases/standalone/common/promql/count_values.result +++ b/tests/cases/standalone/common/promql/count_values.result @@ -61,6 +61,145 @@ TQL EVAL (0, 15, '5s') count_values("status_code", http_requests) by (idc); | 2 | idc2 | 1970-01-01T00:00:15 | 500 | +--------------------------+------+---------------------+-------------+ +-- https://github.com/GreptimeTeam/greptimedb/issues/9181 +-- The label generated by `count_values` must survive an enclosing expression, +-- e.g. a per-sample function, an arithmetic operation or an enclosing aggregation. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (0, 15, '5s') abs(count_values("v", http_requests)); + ++---------------------+-------------------------------+-----+ +| ts | abs(count(http_requests.val)) | v | ++---------------------+-------------------------------+-----+ +| 1970-01-01T00:00:00 | 1 | 401 | +| 1970-01-01T00:00:00 | 3 | 200 | +| 1970-01-01T00:00:05 | 1 | 401 | +| 1970-01-01T00:00:05 | 1 | 500 | +| 1970-01-01T00:00:05 | 2 | 404 | +| 1970-01-01T00:00:10 | 2 | 200 | +| 1970-01-01T00:00:10 | 2 | 201 | +| 1970-01-01T00:00:15 | 4 | 500 | ++---------------------+-------------------------------+-----+ + +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (0, 15, '5s') round(count_values("v", http_requests)); + ++---------------------+-------------------------------------------------+-----+ +| ts | prom_round(count(http_requests.val),Float64(0)) | v | ++---------------------+-------------------------------------------------+-----+ +| 1970-01-01T00:00:00 | 1.0 | 401 | +| 1970-01-01T00:00:00 | 3.0 | 200 | +| 1970-01-01T00:00:05 | 1.0 | 401 | +| 1970-01-01T00:00:05 | 1.0 | 500 | +| 1970-01-01T00:00:05 | 2.0 | 404 | +| 1970-01-01T00:00:10 | 2.0 | 200 | +| 1970-01-01T00:00:10 | 2.0 | 201 | +| 1970-01-01T00:00:15 | 4.0 | 500 | ++---------------------+-------------------------------------------------+-----+ + +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (0, 15, '5s') count_values("v", http_requests) + 1; + ++-----+---------------------+---------------------------------------+ +| v | ts | count(http_requests.val) + Float64(1) | ++-----+---------------------+---------------------------------------+ +| 200 | 1970-01-01T00:00:00 | 4.0 | +| 200 | 1970-01-01T00:00:10 | 3.0 | +| 201 | 1970-01-01T00:00:10 | 3.0 | +| 401 | 1970-01-01T00:00:00 | 2.0 | +| 401 | 1970-01-01T00:00:05 | 2.0 | +| 404 | 1970-01-01T00:00:05 | 3.0 | +| 500 | 1970-01-01T00:00:05 | 2.0 | +| 500 | 1970-01-01T00:00:15 | 5.0 | ++-----+---------------------+---------------------------------------+ + +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (0, 15, '5s') sum by (v) (count_values("v", http_requests)); + ++-----+---------------------+-------------------------------+ +| v | ts | sum(count(http_requests.val)) | ++-----+---------------------+-------------------------------+ +| 200 | 1970-01-01T00:00:00 | 3 | +| 200 | 1970-01-01T00:00:10 | 2 | +| 201 | 1970-01-01T00:00:10 | 2 | +| 401 | 1970-01-01T00:00:00 | 1 | +| 401 | 1970-01-01T00:00:05 | 1 | +| 404 | 1970-01-01T00:00:05 | 2 | +| 500 | 1970-01-01T00:00:05 | 1 | +| 500 | 1970-01-01T00:00:15 | 4 | ++-----+---------------------+-------------------------------+ + +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (0, 15, '5s') topk(10, count_values("v", http_requests)); + ++--------------------------+-----+---------------------+ +| count(http_requests.val) | v | ts | ++--------------------------+-----+---------------------+ +| 1 | 401 | 1970-01-01T00:00:00 | +| 1 | 401 | 1970-01-01T00:00:05 | +| 1 | 500 | 1970-01-01T00:00:05 | +| 2 | 200 | 1970-01-01T00:00:10 | +| 2 | 201 | 1970-01-01T00:00:10 | +| 2 | 404 | 1970-01-01T00:00:05 | +| 3 | 200 | 1970-01-01T00:00:00 | +| 4 | 500 | 1970-01-01T00:00:15 | ++--------------------------+-----+---------------------+ + +-- The generated label overwrites an input label with the same name (PromQL semantics), +-- and must not turn the generated column into an ambiguous reference. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (0, 15, '5s') count_values("idc", http_requests) by (idc); + ++--------------------------+---------------------+-----+ +| count(http_requests.val) | ts | idc | ++--------------------------+---------------------+-----+ +| 1 | 1970-01-01T00:00:00 | 401 | +| 1 | 1970-01-01T00:00:05 | 401 | +| 1 | 1970-01-01T00:00:05 | 500 | +| 2 | 1970-01-01T00:00:05 | 404 | +| 2 | 1970-01-01T00:00:10 | 200 | +| 2 | 1970-01-01T00:00:10 | 201 | +| 3 | 1970-01-01T00:00:00 | 200 | +| 4 | 1970-01-01T00:00:15 | 500 | ++--------------------------+---------------------+-----+ + +-- The generated label is the sample value in PromQL's textual form +-- (`strconv.FormatFloat(value, 'f', -1, 64)`: the shortest decimal form without an +-- exponent), so it is a string: `200.0` becomes the label `200` and `1e21` becomes +-- `1000000000000000000000`. +CREATE TABLE http_requests_double ( + ts timestamp(3) time index, + host STRING, + idc STRING, + val DOUBLE, + PRIMARY KEY(host, idc), +); + +Affected Rows: 0 + +INSERT INTO TABLE http_requests_double VALUES + (0, 'host1', "idc1", 200.0), + (0, 'host2', "idc1", 0.5), + (5000, 'host1', "idc1", 200.0), + (5000, 'host2', "idc1", 1e21); + +Affected Rows: 4 + +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (0, 5, '5s') count_values("val_label", http_requests_double); + ++---------------------------------+---------------------+------------------------+ +| count(http_requests_double.val) | ts | val_label | ++---------------------------------+---------------------+------------------------+ +| 1 | 1970-01-01T00:00:00 | 0.5 | +| 1 | 1970-01-01T00:00:00 | 200 | +| 1 | 1970-01-01T00:00:05 | 1000000000000000000000 | +| 1 | 1970-01-01T00:00:05 | 200 | ++---------------------------------+---------------------+------------------------+ + +DROP TABLE http_requests_double; + +Affected Rows: 0 + DROP TABLE http_requests; Affected Rows: 0 diff --git a/tests/cases/standalone/common/promql/count_values.sql b/tests/cases/standalone/common/promql/count_values.sql index 09c2b136b83..f23cdec5fbe 100644 --- a/tests/cases/standalone/common/promql/count_values.sql +++ b/tests/cases/standalone/common/promql/count_values.sql @@ -28,4 +28,50 @@ TQL EVAL (0, 15, '5s') count_values("status_code", http_requests); TQL EVAL (0, 15, '5s') count_values("status_code", http_requests) by (idc); +-- https://github.com/GreptimeTeam/greptimedb/issues/9181 +-- The label generated by `count_values` must survive an enclosing expression, +-- e.g. a per-sample function, an arithmetic operation or an enclosing aggregation. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (0, 15, '5s') abs(count_values("v", http_requests)); + +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (0, 15, '5s') round(count_values("v", http_requests)); + +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (0, 15, '5s') count_values("v", http_requests) + 1; + +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (0, 15, '5s') sum by (v) (count_values("v", http_requests)); + +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (0, 15, '5s') topk(10, count_values("v", http_requests)); + +-- The generated label overwrites an input label with the same name (PromQL semantics), +-- and must not turn the generated column into an ambiguous reference. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (0, 15, '5s') count_values("idc", http_requests) by (idc); + +-- The generated label is the sample value in PromQL's textual form +-- (`strconv.FormatFloat(value, 'f', -1, 64)`: the shortest decimal form without an +-- exponent), so it is a string: `200.0` becomes the label `200` and `1e21` becomes +-- `1000000000000000000000`. +CREATE TABLE http_requests_double ( + ts timestamp(3) time index, + host STRING, + idc STRING, + val DOUBLE, + PRIMARY KEY(host, idc), +); + +INSERT INTO TABLE http_requests_double VALUES + (0, 'host1', "idc1", 200.0), + (0, 'host2', "idc1", 0.5), + (5000, 'host1', "idc1", 200.0), + (5000, 'host2', "idc1", 1e21); + +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (0, 5, '5s') count_values("val_label", http_requests_double); + +DROP TABLE http_requests_double; + DROP TABLE http_requests; diff --git a/tests/cases/standalone/flow-tql/flow_tql.result b/tests/cases/standalone/flow-tql/flow_tql.result index de93fb7c518..13239bf4e1f 100644 --- a/tests/cases/standalone/flow-tql/flow_tql.result +++ b/tests/cases/standalone/flow-tql/flow_tql.result @@ -228,14 +228,14 @@ SELECT * FROM cnt_reqs ORDER BY ts, status_code; +--------------------------+---------------------+-------------+ | count(http_requests.val) | ts | status_code | +--------------------------+---------------------+-------------+ -| 3.0 | 1970-01-01T00:00:00 | 200.0 | -| 1.0 | 1970-01-01T00:00:00 | 401.0 | -| 1.0 | 1970-01-01T00:00:05 | 401.0 | -| 2.0 | 1970-01-01T00:00:05 | 404.0 | -| 1.0 | 1970-01-01T00:00:05 | 500.0 | -| 2.0 | 1970-01-01T00:00:10 | 200.0 | -| 2.0 | 1970-01-01T00:00:10 | 201.0 | -| 4.0 | 1970-01-01T00:00:15 | 500.0 | +| 3.0 | 1970-01-01T00:00:00 | 200 | +| 1.0 | 1970-01-01T00:00:00 | 401 | +| 1.0 | 1970-01-01T00:00:05 | 401 | +| 2.0 | 1970-01-01T00:00:05 | 404 | +| 1.0 | 1970-01-01T00:00:05 | 500 | +| 2.0 | 1970-01-01T00:00:10 | 200 | +| 2.0 | 1970-01-01T00:00:10 | 201 | +| 4.0 | 1970-01-01T00:00:15 | 500 | +--------------------------+---------------------+-------------+ DROP FLOW calc_reqs;