From 88197f4019eb76f4d68b293b1f2b5c026a901a49 Mon Sep 17 00:00:00 2001 From: dennis zhuang Date: Wed, 23 Sep 2026 11:08:18 +0000 Subject: [PATCH] fix(promql): derive vector matching result labels and reject ambiguous matchings (#9306) * fix(promql): derive vector matching result labels and reject ambiguous matchings A vector-vector binary operation projected one operand's whole tag set and inner-joined without any cardinality check, so `on()`/`ignoring()` did not reduce the result labels, `group_left`/`group_right` changed nothing, and a non-unique match group produced a cross product that PromQL cannot represent. Result labels now follow Prometheus `resultMetric`: `on(...)` keeps the matching labels, `ignoring(...)` drops them, and a group modifier keeps the many side's labels plus the `group_x(...)` labels taken from the one side. A label the one side does not carry is deleted from the result. The reduced label set no longer identifies the operand series, so `__tsid` is dropped from the context on this path. Cardinality is enforced with a `count(1) OVER (PARTITION BY match keys, ts)` window and a scalar UDF that fails the query on a repeated group: on the one side before the join, and on the result labels after it, matching where Prometheus raises each of its three errors. Series are unique by their whole tag set, so the window is only planted when the match keys drop a tag; plain arithmetic and `on()` plan exactly as before. Closes #9207, closes #9208, closes #9209. Signed-off-by: Dennis Zhuang * perf(promql): keep __tsid when the result labels are an operand's whole tag set Deriving the result labels dropped `__tsid` from the context unconditionally, so an enclosing operation fell back to joining on the tag columns even where the column still identified the result series. Keep it when every result label comes from one operand and covers that operand's whole tag set: no other operand value reaches the labels, and the matching gives each of its rows a single partner, so its `__tsid` is still one per result series. That is the common `on()` and bare `group_left` shape; a matching that actually drops a tag still clears it. The column is re-qualified as the result's own, which is how the enclosing expression and the context look it up. Signed-off-by: Dennis Zhuang * fix(promql): keep the match group count column unambiguous The cardinality check aliased its row count to a fixed `__promql_match_group_count`. An operand carrying a label of that name made the window output two fields with the same name, and planning failed with "Schema contains qualified field name collide_right.__promql_match_group_count and unqualified field name __promql_match_group_count which would be ambiguous". Pick a name the operand does not already have, the way the `or` operator allocates its match key columns. Signed-off-by: Dennis Zhuang * test(promql): cover a match group spread over several regions The cardinality check runs above the merge of the region scans, so it counts a match group globally. Nothing pinned that: every table in these cases holds a single region, and a check evaluated per region would pass them all. Partition the operand on a column outside the match keys, which puts the two series of one match group in different regions, and assert both the pre-join and the post-join check still reject it. The case runs in the distributed environment too, where the regions sit on different datanodes. Signed-off-by: Dennis Zhuang * fix(promql): group an outer aggregate on the labels the operand kept `by`/`without` planning topped up a missing grouping column by walking down to the table scan and re-projecting it. That is right for a column the scan pruned for efficiency, but the labels a matching modifier deletes are also absent from the operand's output, and they were restored the same way: sum without(host) (a / on(host) b) `on(host)` leaves the operand with `host` alone, so the sum covers everything and Prometheus answers `{} 10`. Instead `device` came back from the scan under `a` and split the result into `{device="d1"} 5` and `{device="d2"} 5`. Same for `sum by(device)` of that operand, which has no `device` to group on at all. Take the grouping labels from the operand's own label set rather than from the row keys of the scan beneath it. A label pruned from the plan is still in that set and still gets restored; a label the operand dropped is not. Signed-off-by: Dennis Zhuang * refactor(promql): drop the unreachable aggregation tag top-up `by`/`without` planning could restore a grouping column that the plan no longer carried by rewriting the scan underneath it. Once the grouping labels come from the operand's own label set, there is nothing left for it to restore: a scan projects every label of `ctx.tag_columns` (`scan_tag_columns` only ever adds matcher columns to that set), so a label in the set is always in the schema. Stubbing the rewriter to a no-op passed the whole sqlness suite, in both the standalone and the distributed environment. Signed-off-by: Dennis Zhuang * test(promql): drop cases that guarded the removed tag top-up Three plain selector aggregates were there to show that restoring a pruned grouping column still worked. With the restore gone they only repeat what the aggregate cases already cover. Also fix a comment that still said the metric engine scan prunes tag columns: it projects every label of the operand. Signed-off-by: Dennis Zhuang * docs(promql): note the duplicate a propagated matcher hides A matcher copied onto the one-side operand removes groups without a partner before the cardinality check sees them, so a duplicate in such a group is not reported. Prometheus checks every group of the one side and fails the query. Signed-off-by: Dennis Zhuang --------- Signed-off-by: Dennis Zhuang --- src/promql/src/functions.rs | 2 + src/promql/src/functions/vector_matching.rs | 208 +++++ src/query/src/promql/planner.rs | 742 ++++++++++++------ .../src/promql/planner/matching_filters.rs | 7 +- .../common/promql/matching_groups.result | 108 ++- .../common/promql/matching_groups.sql | 8 +- .../common/promql/set_operation.result | 32 +- .../promql/tsid_binary_join_regression.result | 51 +- .../promql/tsid_binary_join_regression.sql | 5 + .../common/promql/vector_matching.result | 406 ++++++++++ .../common/promql/vector_matching.sql | 221 ++++++ .../common/tql/matching_filter.result | 79 +- .../standalone/common/tql/matching_filter.sql | 4 +- 13 files changed, 1484 insertions(+), 389 deletions(-) create mode 100644 src/promql/src/functions/vector_matching.rs create mode 100644 tests/cases/standalone/common/promql/vector_matching.result create mode 100644 tests/cases/standalone/common/promql/vector_matching.sql diff --git a/src/promql/src/functions.rs b/src/promql/src/functions.rs index 22cd64b9147..34618f0a953 100644 --- a/src/promql/src/functions.rs +++ b/src/promql/src/functions.rs @@ -27,6 +27,7 @@ mod resets; mod round; #[cfg(test)] mod test_util; +mod vector_matching; pub use aggr_over_time::{ AbsentOverTime, AvgOverTime, CountOverTime, LastOverTime, MaxOverTime, MinOverTime, @@ -61,6 +62,7 @@ pub use quantile::QuantileOverTime; pub use quantile_aggr::{QUANTILE_NAME, quantile_udaf}; pub use resets::Resets; pub use round::Round; +pub use vector_matching::{MatchGroupViolation, UniqueMatchGroup}; use crate::range_array::RangeArray; diff --git a/src/promql/src/functions/vector_matching.rs b/src/promql/src/functions/vector_matching.rs new file mode 100644 index 00000000000..fdef83f6f05 --- /dev/null +++ b/src/promql/src/functions/vector_matching.rs @@ -0,0 +1,208 @@ +// 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. + +//! Cardinality checks for PromQL vector matching. + +use datafusion::arrow::array::{Array, Int64Array}; +use datafusion::arrow::util::display::array_value_to_string; +use datafusion::common::{DataFusionError, Result as DfResult}; +use datafusion::logical_expr::{ScalarUDF, Volatility}; +use datafusion::physical_plan::ColumnarValue; +use datafusion_common::ScalarValue; +use datafusion_expr::{ScalarFunctionArgs, ScalarUDFImpl, Signature}; +use datatypes::arrow::datatypes::DataType; + +use crate::functions::extract_array; + +/// The rule a repeated match group violates. +#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)] +pub enum MatchGroupViolation { + /// Several series share a match group on the side that must hold one series per group. + DuplicateOnOneSide { one_side_is_left: bool }, + /// One-to-one matching found several matches for the same group. + ImplicitManyToOne, + /// A group modifier left several matches with the same result label set. + AmbiguousGroupLabels, +} + +impl MatchGroupViolation { + fn message(&self, group: &str) -> String { + match self { + Self::DuplicateOnOneSide { one_side_is_left } => { + let side = if *one_side_is_left { "left" } else { "right" }; + format!( + "found duplicate series for the match group {group} on the {side} hand-side \ + of the operation; many-to-many matching not allowed: matching labels must \ + be unique on one side" + ) + } + Self::ImplicitManyToOne => format!( + "multiple matches for labels {group}: many-to-one matching must be explicit \ + (group_left/group_right)" + ), + Self::AmbiguousGroupLabels => format!( + "multiple matches for labels {group}: grouping labels must ensure unique matches" + ), + } + } +} + +/// Rejects a vector matching whose match groups are not unique. +/// +/// Takes the per-group row count as first argument and the label columns that form the group +/// as the remaining ones. Returns `true` when every group holds a single row, and fails the +/// query otherwise; PromQL has no way to express the duplicate series it would produce. +pub struct UniqueMatchGroup; + +impl UniqueMatchGroup { + pub const fn name() -> &'static str { + "prom_assert_unique_match_group" + } + + pub fn scalar_udf(labels: Vec, violation: MatchGroupViolation) -> ScalarUDF { + ScalarUDF::new_from_impl(AssertUniqueMatchGroup { + signature: Signature::variadic_any(Volatility::Volatile), + labels, + violation, + }) + } +} + +#[derive(Debug, Clone, PartialEq, Eq, Hash)] +struct AssertUniqueMatchGroup { + signature: Signature, + /// Label names of the group, in the same order as the label arguments. + labels: Vec, + violation: MatchGroupViolation, +} + +impl AssertUniqueMatchGroup { + fn render_group(&self, args: &[ColumnarValue], row: usize) -> DfResult { + let mut rendered = Vec::with_capacity(self.labels.len()); + for (label, arg) in self.labels.iter().zip(args) { + let array = extract_array(arg)?; + // A scalar argument was expanded to a single row above. + let row = if array.len() == 1 { 0 } else { row }; + if row >= array.len() || array.is_null(row) { + continue; + } + let value = array_value_to_string(&array, row) + .map_err(|e| DataFusionError::ArrowError(Box::new(e), None))?; + if value.is_empty() { + continue; + } + rendered.push(format!("{label}=\"{value}\"")); + } + Ok(format!("{{{}}}", rendered.join(", "))) + } +} + +impl ScalarUDFImpl for AssertUniqueMatchGroup { + fn name(&self) -> &str { + UniqueMatchGroup::name() + } + + fn signature(&self) -> &Signature { + &self.signature + } + + fn return_type(&self, _arg_types: &[DataType]) -> DfResult { + Ok(DataType::Boolean) + } + + fn invoke_with_args(&self, args: ScalarFunctionArgs) -> DfResult { + let Some((counts, labels)) = args.args.split_first() else { + return Err(DataFusionError::Execution(format!( + "{} expects the match group row count as first argument", + UniqueMatchGroup::name() + ))); + }; + let counts = extract_array(counts)?; + let counts = counts + .as_any() + .downcast_ref::() + .ok_or_else(|| { + DataFusionError::Execution(format!( + "{} expects an Int64 match group row count, found {}", + UniqueMatchGroup::name(), + counts.data_type() + )) + })?; + + if let Some(row) = + (0..counts.len()).find(|row| !counts.is_null(*row) && counts.value(*row) > 1) + { + let group = self.render_group(labels, row)?; + return Err(DataFusionError::Execution(self.violation.message(&group))); + } + + Ok(ColumnarValue::Scalar(ScalarValue::Boolean(Some(true)))) + } +} + +#[cfg(test)] +mod tests { + use std::sync::Arc; + + use datafusion::arrow::array::StringArray; + use datafusion::arrow::datatypes::Field; + + use super::*; + + fn invoke(counts: Vec, hosts: Vec>) -> DfResult { + let udf = UniqueMatchGroup::scalar_udf( + vec!["host".to_string()], + MatchGroupViolation::ImplicitManyToOne, + ); + let number_rows = counts.len(); + udf.invoke_with_args(ScalarFunctionArgs { + args: vec![ + ColumnarValue::Array(Arc::new(Int64Array::from(counts))), + ColumnarValue::Array(Arc::new(StringArray::from(hosts))), + ], + arg_fields: vec![ + Arc::new(Field::new("count", DataType::Int64, true)), + Arc::new(Field::new("host", DataType::Utf8, true)), + ], + number_rows, + return_field: Arc::new(Field::new("assert", DataType::Boolean, false)), + config_options: Arc::new(Default::default()), + }) + } + + #[test] + fn unique_groups_pass() { + let result = invoke(vec![1, 1], vec![Some("a"), Some("b")]).unwrap(); + assert!(matches!( + result, + ColumnarValue::Scalar(ScalarValue::Boolean(Some(true))) + )); + } + + #[test] + fn duplicate_group_reports_its_labels() { + let err = invoke(vec![1, 2], vec![Some("a"), Some("b")]).unwrap_err(); + assert!( + err.to_string() + .contains("multiple matches for labels {host=\"b\"}"), + "{err}" + ); + } + + #[test] + fn null_label_is_omitted_from_the_group() { + let err = invoke(vec![2], vec![None]).unwrap_err(); + assert!(err.to_string().contains("labels {}"), "{err}"); + } +} diff --git a/src/query/src/promql/planner.rs b/src/query/src/promql/planner.rs index 44211b8d2b7..2142ee82371 100644 --- a/src/query/src/promql/planner.rs +++ b/src/query/src/promql/planner.rs @@ -50,15 +50,12 @@ use datafusion::optimizer::simplify_expressions::ExprSimplifier; use datafusion::prelude as df_prelude; use datafusion::prelude::{Column, Expr as DfExpr, JoinType}; use datafusion::scalar::ScalarValue; -use datafusion_common::tree_node::{Transformed, TreeNode, TreeNodeRewriter}; use datafusion_common::{DFSchema, NullEquality, TableReference}; use datafusion_expr::expr::WindowFunctionParams; use datafusion_expr::expr_fn::when; use datafusion_expr::simplify::SimplifyContext; use datafusion_expr::utils::{conjunction, disjunction}; -use datafusion_expr::{ - ExprSchemable, Literal, Projection, SortExpr, TableScanBuilder, TableSource, col, lit, -}; +use datafusion_expr::{ExprSchemable, Literal, SortExpr, TableSource, col, lit}; use datafusion_functions::core::coalesce; use datatypes::arrow::datatypes::{DataType as ArrowDataType, TimeUnit as ArrowTimeUnit}; use datatypes::data_type::{ConcreteDataType, DataType as GreptimeDataType}; @@ -71,7 +68,7 @@ use promql::extension_plan::{ }; use promql::functions::{ AbsentOverTime, AvgOverTime, Changes, CountOverTime, Delta, Deriv, DoubleExponentialSmoothing, - IDelta, Increase, LastOverTime, MaxOverTime, MinOverTime, MixedRange, + IDelta, Increase, LastOverTime, MatchGroupViolation, MaxOverTime, MinOverTime, MixedRange, NativeHistogramAbsentOverTime, NativeHistogramAdd, NativeHistogramAggAvg, NativeHistogramAggSum, NativeHistogramAvg, NativeHistogramAvgOverTime, NativeHistogramChanges, NativeHistogramCount, NativeHistogramCountOverTime, NativeHistogramDelta, @@ -82,7 +79,8 @@ use promql::functions::{ NativeHistogramRate, NativeHistogramResets, NativeHistogramScalarMul, NativeHistogramStddev, NativeHistogramStdvar, NativeHistogramSub, NativeHistogramSum, NativeHistogramSumOverTime, NativeHistogramToString, PredictLinear, PresentOverTime, PromqlFloatToString, QuantileOverTime, - Rate, Resets, Round, StddevOverTime, StdvarOverTime, SumOverTime, quantile_udaf, + Rate, Resets, Round, StddevOverTime, StdvarOverTime, SumOverTime, UniqueMatchGroup, + quantile_udaf, }; use promql_parser::label::{METRIC_NAME, MatchOp, Matcher, Matchers}; use promql_parser::parser::token::TokenType; @@ -154,6 +152,8 @@ const BINARY_ISLAND_LEAF_ALIAS_PREFIX: &str = "__prom_v"; const OR_FLOAT_FIELD_PREFIX: &str = "__promql_or_float_"; const OR_HISTOGRAM_FIELD_PREFIX: &str = "__promql_or_histogram_"; const TIMESTAMP_VALUE_PREFIX: &str = "__promql_timestamp_value_"; +/// Per-match-group row count the vector matching cardinality check is built on. +const MATCH_GROUP_COUNT_COLUMN: &str = "__promql_match_group_count"; /// Threshold for scatter scan mode const MAX_SCATTER_POINTS: i64 = 400; @@ -295,6 +295,35 @@ struct PlannedIslandLeaf { display_table: String, } +/// Result labels a vector-vector binary operation derives from its matching modifier, projected +/// from the operand each label belongs to. +#[derive(Debug)] +struct BinaryResultLabels { + exprs: Vec, + names: Vec, + aggregation_field_labels: Vec, + /// `__tsid` column of the operand the labels come from, when that operand contributes its + /// whole tag set and the column still identifies the result series. + tsid: Option, +} + +impl BinaryResultLabels { + fn apply(&self, ctx: &mut PromPlannerContext) { + ctx.tag_columns = self.names.clone(); + ctx.aggregation_field_labels = self.aggregation_field_labels.clone(); + ctx.use_tsid = self.tsid.is_some(); + } + + /// `__tsid` is projected from whichever operand the labels come from, so it carries that + /// operand's qualifier. Re-qualify it as the result's own, which is what the context names + /// and what the enclosing expression looks the column up by. + fn tsid_projection(&self, table_ref: Option) -> Option { + self.tsid + .clone() + .map(|tsid| tsid.alias_qualified(table_ref, DATA_SCHEMA_TSID_COLUMN_NAME)) + } +} + #[derive(Debug)] struct IslandFieldExprs { exprs: Vec, @@ -622,40 +651,12 @@ impl PromPlanner { param, } = aggr_expr; - let mut input = self.prom_expr_to_plan(expr, query_engine_state).await?; + let input = self.prom_expr_to_plan(expr, query_engine_state).await?; let input_has_tsid = input.schema().fields().iter().any(|field| { field.name() == DATA_SCHEMA_TSID_COLUMN_NAME && field.data_type() == &ArrowDataType::UInt64 }); - // `__tsid` based scan projection may prune tag columns. Ensure tags referenced in - // aggregation modifiers (`by`/`without`) are available before planning group keys. - let required_group_tags = match modifier { - None => BTreeSet::new(), - Some(LabelModifier::Include(labels)) => labels - .labels - .iter() - .filter(|label| !is_metric_engine_internal_column(label.as_str())) - .cloned() - .collect(), - Some(LabelModifier::Exclude(labels)) => { - let mut all_tags = self.collect_row_key_tag_columns_from_plan(&input)?; - for label in &labels.labels { - let _ = all_tags.remove(label); - } - all_tags - } - }; - - if !required_group_tags.is_empty() - && required_group_tags - .iter() - .any(|tag| Self::find_case_sensitive_column(input.schema(), tag.as_str()).is_none()) - { - input = self.ensure_tag_columns_available(input, &required_group_tags)?; - self.refresh_tag_columns_from_schema(input.schema()); - } - match (*op).id() { token::T_TOPK | token::T_BOTTOMK => { self.prom_topk_bottomk_to_plan(aggr_expr, input).await @@ -1634,6 +1635,16 @@ impl PromPlanner { // `vector()` uses EmptyMetric and keeps GreptimeDB's timestamp broadcast. let has_empty_metric_operand = left_is_empty_metric || right_is_empty_metric; + let only_join_time_index = lhs.value_type() == ValueType::Scalar + || rhs.value_type() == ValueType::Scalar + || has_empty_metric_operand + || ((left_context.tag_columns.is_empty() + || right_context.tag_columns.is_empty()) + && !left_context + .tag_columns + .iter() + .chain(&right_context.tag_columns) + .any(|tag| tag == OTLP_AGGREGATION_TEMPORALITY_LABEL)); let join_plan = self.join_on_non_field_columns( left_input, right_input, @@ -1641,21 +1652,34 @@ impl PromPlanner { right_table_ref.clone(), left_time_index_column, right_time_index_column, - lhs.value_type() == ValueType::Scalar - || rhs.value_type() == ValueType::Scalar - || has_empty_metric_operand - || ((left_context.tag_columns.is_empty() - || right_context.tag_columns.is_empty()) - && !left_context - .tag_columns - .iter() - .chain(&right_context.tag_columns) - .any(|tag| tag == OTLP_AGGREGATION_TEMPORALITY_LABEL)), + only_join_time_index, modifier, &left_context, &right_context, )?; let join_plan_schema = join_plan.schema().clone(); + // The matching modifier derives the result labels instead of the operands keeping + // their own tag set. An operand that is broadcast rather than matched (a scalar, + // or a vector without tags) keeps the other side's labels as before. + let result_labels = if only_join_time_index { + None + } else { + Self::binary_result_labels(&left_context, &right_context, modifier) + .map(|labels| { + Self::binary_result_label_projection( + &join_plan_schema, + &left_table_ref, + &right_table_ref, + &left_context, + &right_context, + labels, + ) + }) + .transpose()? + }; + if let Some(labels) = &result_labels { + labels.apply(&mut self.ctx); + } let promql_annotations = self.promql_annotations.clone(); // These predicates always pass; they only evaluate otherwise-discarded pairs // while collecting annotations. @@ -1776,10 +1800,18 @@ impl PromPlanner { ), }; project_context.field_columns = project_field_columns; - self.project_binary_join_side(filtered, project_table_ref, &project_context) + self.project_binary_join_side( + filtered, + project_table_ref, + &project_context, + result_labels.as_ref(), + ) } else { - let projected = - self.projection_for_each_field_column(join_plan, bin_expr_builder)?; + let projected = self.projection_for_each_field_column_with_labels( + join_plan, + result_labels.as_ref(), + bin_expr_builder, + )?; let preserve_any_value = Self::field_columns_are_alternative_samples( projected.schema(), &self.ctx.field_columns, @@ -1867,6 +1899,7 @@ impl PromPlanner { input: LogicalPlan, table_ref: &TableReference, context: &PromPlannerContext, + result_labels: Option<&BinaryResultLabels>, ) -> Result { let schema = input.schema(); @@ -1891,20 +1924,27 @@ impl PromPlanner { project_exprs.push(DfExpr::Column(field_col)); } - // Project tag columns from the chosen side. - for tag_column in &context.tag_columns { - let tag_col = schema - .qualified_field_with_name(Some(table_ref), tag_column) - .context(DataFusionPlanningSnafu)? - .into(); - project_exprs.push(DfExpr::Column(tag_col)); + // Project tag columns: the labels the matching modifier derived, or the chosen side's. + match result_labels { + Some(labels) => project_exprs.extend(labels.exprs.iter().cloned()), + None => { + for tag_column in &context.tag_columns { + let tag_col = schema + .qualified_field_with_name(Some(table_ref), tag_column) + .context(DataFusionPlanningSnafu)? + .into(); + project_exprs.push(DfExpr::Column(tag_col)); + } + } } // Preserve `__tsid` if present, so it can still be used internally downstream. It's // stripped from the final output anyway. - if let Some(tsid_col) = - Self::optional_tsid_projection(schema, Some(table_ref), context.use_tsid) - { + let tsid_col = match result_labels { + Some(labels) => labels.tsid_projection(Some(table_ref.clone())), + None => Self::optional_tsid_projection(schema, Some(table_ref), context.use_tsid), + }; + if let Some(tsid_col) = tsid_col { project_exprs.push(tsid_col); } @@ -1919,6 +1959,9 @@ impl PromPlanner { self.ctx = context.clone(); self.ctx.table_name = None; self.ctx.schema_name = None; + if let Some(labels) = result_labels { + labels.apply(&mut self.ctx); + } Ok(plan) } @@ -3214,130 +3257,6 @@ impl PromPlanner { Ok(out) } - fn ensure_tag_columns_available( - &self, - plan: LogicalPlan, - required_tags: &BTreeSet, - ) -> Result { - if required_tags.is_empty() { - return Ok(plan); - } - - struct Rewriter { - required_tags: BTreeSet, - } - - impl TreeNodeRewriter for Rewriter { - type Node = LogicalPlan; - - fn f_up( - &mut self, - node: Self::Node, - ) -> datafusion_common::Result> { - match node { - LogicalPlan::TableScan(scan) => { - let schema = scan.source.schema(); - let mut projection = match scan.projection.clone() { - Some(p) => p, - None => { - // Scanning all columns already covers required tags. - return Ok(Transformed::no(LogicalPlan::TableScan(scan))); - } - }; - - let mut changed = false; - for tag in &self.required_tags { - if let Some((idx, _)) = schema - .fields() - .iter() - .enumerate() - .find(|(_, field)| field.name() == tag) - && !projection.contains(&idx) - { - projection.push(idx); - changed = true; - } - } - - if !changed { - return Ok(Transformed::no(LogicalPlan::TableScan(scan))); - } - - projection.sort_unstable(); - projection.dedup(); - - let new_scan = - TableScanBuilder::new(scan.table_name.clone(), scan.source.clone()) - .with_projection(Some(projection)) - .with_filters(scan.filters) - .with_fetch(scan.fetch) - .build()?; - Ok(Transformed::yes(LogicalPlan::TableScan(new_scan))) - } - LogicalPlan::Projection(proj) => { - let input_schema = proj.input.schema(); - - let existing = proj - .schema - .fields() - .iter() - .map(|f| f.name().as_str()) - .collect::>(); - - let mut expr = proj.expr.clone(); - let mut has_changed = false; - for tag in &self.required_tags { - if existing.contains(tag.as_str()) { - continue; - } - - if let Some(idx) = input_schema.index_of_column_by_name(None, tag) { - expr.push(DfExpr::Column(Column::from( - input_schema.qualified_field(idx), - ))); - has_changed = true; - } - } - - if !has_changed { - return Ok(Transformed::no(LogicalPlan::Projection(proj))); - } - - let new_proj = Projection::try_new(expr, proj.input)?; - Ok(Transformed::yes(LogicalPlan::Projection(new_proj))) - } - other => Ok(Transformed::no(other)), - } - } - } - - let mut rewriter = Rewriter { - required_tags: required_tags.clone(), - }; - let rewritten = plan - .rewrite(&mut rewriter) - .context(DataFusionPlanningSnafu)?; - Ok(rewritten.data) - } - - fn refresh_tag_columns_from_schema(&mut self, schema: &DFSchemaRef) { - let time_index = self.ctx.time_index_column.as_deref(); - let field_columns = self.ctx.field_columns.iter().collect::>(); - - let mut tags = schema - .fields() - .iter() - .map(|f| f.name()) - .filter(|name| Some(name.as_str()) != time_index) - .filter(|name| !field_columns.contains(name)) - .filter(|name| !is_metric_engine_internal_column(name)) - .cloned() - .collect::>(); - tags.sort_unstable(); - tags.dedup(); - self.ctx.tag_columns = tags; - } - /// Setup [PromPlannerContext]'s state fields. /// /// Returns a logical plan for an empty metric. @@ -5956,6 +5875,225 @@ impl PromPlanner { Ok((left_tag_columns, right_tag_columns, force_empty_join)) } + /// Result labels of a vector-vector binary operation, following Prometheus `resultMetric`: + /// `on(...)` keeps only the matching labels, `ignoring(...)` drops them, and a group modifier + /// keeps the "many" side's labels plus the `group_x(...)` labels taken from the "one" side. + /// + /// The flag of each entry tells which operand the label is projected from. `None` means the + /// operation keeps a whole operand tag set, which the default projection already does. + fn binary_result_labels( + left_context: &PromPlannerContext, + right_context: &PromPlannerContext, + modifier: &Option, + ) -> Option> { + let modifier = modifier.as_ref()?; + let (many_is_left, include) = match &modifier.card { + VectorMatchCardinality::OneToOne => (true, None), + VectorMatchCardinality::ManyToOne(labels) => (true, Some(labels)), + VectorMatchCardinality::OneToMany(labels) => (false, Some(labels)), + // Set operators keep their operands' labels and don't reach this path. + VectorMatchCardinality::ManyToMany => return None, + }; + + let Some(include) = include else { + let matching = modifier.matching.as_ref()?; + let reduced = |keep: bool, labels: &BTreeSet<&String>| { + left_context + .tag_columns + .iter() + .filter(|tag| labels.contains(tag) == keep) + .map(|tag| (true, tag.clone())) + .collect() + }; + return Some(match matching { + LabelModifier::Include(on) => reduced(true, &on.labels.iter().collect()), + LabelModifier::Exclude(ignoring) => { + reduced(false, &ignoring.labels.iter().collect()) + } + }); + }; + + let (many_context, one_context) = if many_is_left { + (left_context, right_context) + } else { + (right_context, left_context) + }; + let include = include.labels.iter().collect::>(); + // An included label the "one" side doesn't carry is deleted from the result, so it is + // dropped from the "many" side as well. + let mut labels = many_context + .tag_columns + .iter() + .filter(|tag| !include.contains(tag)) + .map(|tag| (many_is_left, tag.clone())) + .collect::>(); + labels.extend( + include + .into_iter() + .filter(|label| one_context.tag_columns.contains(label)) + .map(|label| (!many_is_left, label.clone())), + ); + Some(labels) + } + + /// Resolve [`Self::binary_result_labels`] against the join output. + fn binary_result_label_projection( + schema: &DFSchemaRef, + left_table_ref: &TableReference, + right_table_ref: &TableReference, + left_context: &PromPlannerContext, + right_context: &PromPlannerContext, + labels: Vec<(bool, String)>, + ) -> Result { + let mut exprs = Vec::with_capacity(labels.len()); + let mut names = Vec::with_capacity(labels.len()); + let mut aggregation_field_labels = Vec::new(); + let mut sources = HashSet::new(); + for (from_left, label) in labels { + let (table_ref, context) = if from_left { + (left_table_ref, left_context) + } else { + (right_table_ref, right_context) + }; + let field = schema + .qualified_field_with_name(Some(table_ref), &label) + .context(DataFusionPlanningSnafu)?; + exprs.push(DfExpr::Column(field.into())); + if context.aggregation_field_labels.contains(&label) { + aggregation_field_labels.push(label.clone()); + } + sources.insert(from_left); + names.push(label); + } + + // One operand contributing every one of its tags as the whole result label set keeps the + // one-to-one correspondence between its `__tsid` and a result series: no other operand + // value reaches the labels, and the matching gives each of its rows a single partner. + let source_context = match sources.into_iter().collect::>().as_slice() { + [true] => Some((left_table_ref, left_context)), + [false] => Some((right_table_ref, right_context)), + _ => None, + }; + let tsid = source_context + .filter(|(_, context)| { + context.use_tsid + && context.tag_columns.len() == names.len() + && context.tag_columns.iter().all(|tag| names.contains(tag)) + }) + .and_then(|(table_ref, _)| { + Self::optional_tsid_projection(schema, Some(table_ref), true) + }); + + Ok(BinaryResultLabels { + exprs, + names, + aggregation_field_labels, + tsid, + }) + } + + /// Whether the result of a binary operation can hold two series with the same labels, which + /// only [`Self::binary_result_labels`] can introduce: one-to-one matching on a subset of the + /// tags, or a group modifier that overwrites a label of the "many" side. + fn binary_result_labels_may_repeat( + left_context: &PromPlannerContext, + right_context: &PromPlannerContext, + modifier: &Option, + ) -> bool { + let Some(modifier) = modifier else { + return false; + }; + match &modifier.card { + VectorMatchCardinality::OneToOne => match &modifier.matching { + None => false, + Some(LabelModifier::Include(on)) => { + let on = on.labels.iter().collect::>(); + !left_context.tag_columns.iter().all(|tag| on.contains(tag)) + } + Some(LabelModifier::Exclude(ignoring)) => ignoring + .labels + .iter() + .any(|label| left_context.tag_columns.contains(label)), + }, + VectorMatchCardinality::ManyToOne(include) => include + .labels + .iter() + .any(|label| left_context.tag_columns.contains(label)), + VectorMatchCardinality::OneToMany(include) => include + .labels + .iter() + .any(|label| right_context.tag_columns.contains(label)), + VectorMatchCardinality::ManyToMany => false, + } + } + + /// Wrap `plan` in a check that fails the query when a match group holds more than one row at + /// a timestamp. `group_exprs` are the label columns of the group, resolved against `plan`. + fn assert_unique_match_group( + plan: LogicalPlan, + group_exprs: Vec, + group_labels: Vec, + time_index_expr: DfExpr, + violation: MatchGroupViolation, + ) -> Result { + let mut partition_by = group_exprs.clone(); + partition_by.push(time_index_expr); + // A label may carry the generated name, and the count column has to stay unambiguous + // against every column the operand already has. + let occupied_column_names = plan + .schema() + .fields() + .iter() + .map(|field| field.name().as_str()) + .collect::>(); + let mut next_suffix = 0; + let count_column = loop { + let name = match next_suffix { + 0 => MATCH_GROUP_COUNT_COLUMN.to_string(), + suffix => format!("{MATCH_GROUP_COUNT_COLUMN}_{suffix}"), + }; + next_suffix += 1; + if !occupied_column_names.contains(name.as_str()) { + break name; + } + }; + let count = DfExpr::WindowFunction(Box::new(WindowFunction { + fun: WindowFunctionDefinition::AggregateUDF(count_udaf()), + params: WindowFunctionParams { + args: vec![lit(1i64)], + partition_by, + order_by: vec![], + window_frame: WindowFrame::new(None), + null_treatment: None, + distinct: false, + filter: None, + }, + })) + .alias(count_column.as_str()); + + let output_exprs = plan + .schema() + .iter() + .map(|(qualifier, field)| DfExpr::Column(Column::new(qualifier.cloned(), field.name()))) + .collect::>(); + let assert_expr = DfExpr::ScalarFunction(ScalarFunction { + func: Arc::new(UniqueMatchGroup::scalar_udf(group_labels, violation)), + args: std::iter::once(col(count_column.as_str())) + .chain(group_exprs) + .collect(), + }); + + LogicalPlanBuilder::from(plan) + .window(vec![count]) + .context(DataFusionPlanningSnafu)? + .filter(assert_expr) + .context(DataFusionPlanningSnafu)? + .project(output_exprs) + .context(DataFusionPlanningSnafu)? + .build() + .context(DataFusionPlanningSnafu) + } + fn binary_modifier_preserves_tsid_join_key( &self, left_context: &PromPlannerContext, @@ -6023,7 +6161,7 @@ impl PromPlanner { && !force_empty_join && left_tag_columns == BTreeSet::from([DATA_SCHEMA_TSID_COLUMN_NAME.to_string()]) && right_tag_columns == BTreeSet::from([DATA_SCHEMA_TSID_COLUMN_NAME.to_string()]); - let (left, right) = if !only_join_time_index + let (left, right, left_matched_tags, right_matched_tags) = if !only_join_time_index && !use_tsid_join && Self::only_temporality_match_label_mismatches(left_context, right_context, modifier) { @@ -6044,28 +6182,79 @@ impl PromPlanner { false, modifier, )?; - (left, right) + ( + left, + right, + aligned_left_context.tag_columns, + aligned_right_context.tag_columns, + ) } else { + ( + left, + right, + left_context.tag_columns.clone(), + right_context.tag_columns.clone(), + ) + }; + + // A join key that covers the whole tag set of a side already makes that side's match + // groups unique, and so does a join on `__tsid`. Only the reduced keys need the check. + let checked_matching = !only_join_time_index + && !force_empty_join + && !use_tsid_join + && !matches!( + modifier.as_ref().map(|modifier| &modifier.card), + Some(VectorMatchCardinality::ManyToMany) + ); + // The right operand is the "one" side of the matching, unless `group_right` swaps them. + let one_side_is_left = matches!( + modifier.as_ref().map(|modifier| &modifier.card), + Some(VectorMatchCardinality::OneToMany(_)) + ); + let (left, right) = if !checked_matching { (left, right) + } else if one_side_is_left { + ( + Self::assert_unique_one_side( + left, + &left_tag_columns, + &left_matched_tags, + left_time_index_column.as_deref(), + true, + )?, + right, + ) + } else { + ( + left, + Self::assert_unique_one_side( + right, + &right_tag_columns, + &right_matched_tags, + right_time_index_column.as_deref(), + false, + )?, + ) }; // push time index column if it exists - if let (Some(left_time_index_column), Some(right_time_index_column)) = - (left_time_index_column, right_time_index_column) - { + if let (Some(left_time_index_column), Some(right_time_index_column)) = ( + left_time_index_column.clone(), + right_time_index_column.clone(), + ) { left_tag_columns.insert(left_time_index_column); right_tag_columns.insert(right_time_index_column); } let right = LogicalPlanBuilder::from(right) - .alias(right_table_ref) + .alias(right_table_ref.clone()) .context(DataFusionPlanningSnafu)? .build() .context(DataFusionPlanningSnafu)?; // Inner Join on time index column to concat two operator - LogicalPlanBuilder::from(left) - .alias(left_table_ref) + let join_plan = LogicalPlanBuilder::from(left) + .alias(left_table_ref.clone()) .context(DataFusionPlanningSnafu)? .join_detailed( right, @@ -6085,7 +6274,90 @@ impl PromPlanner { ) .context(DataFusionPlanningSnafu)? .build() - .context(DataFusionPlanningSnafu) + .context(DataFusionPlanningSnafu)?; + + // The result labels no longer identify a series on their own: the "many" side can hold + // several series per result label set, which PromQL cannot represent. + let labels = checked_matching + .then(|| Self::binary_result_labels(left_context, right_context, modifier)) + .flatten() + .filter(|_| { + Self::binary_result_labels_may_repeat(left_context, right_context, modifier) + }); + let Some(labels) = labels else { + return Ok(join_plan); + }; + let schema = join_plan.schema().clone(); + let group_exprs = labels + .iter() + .map(|(from_left, label)| { + let table_ref = if *from_left { + &left_table_ref + } else { + &right_table_ref + }; + schema + .qualified_field_with_name(Some(table_ref), label) + .map(|field| DfExpr::Column(field.into())) + .context(DataFusionPlanningSnafu) + }) + .collect::>>()?; + let time_index_expr = left_time_index_column + .as_deref() + .map(|column| { + schema + .qualified_field_with_name(Some(&left_table_ref), column) + .map(|field| DfExpr::Column(field.into())) + .context(DataFusionPlanningSnafu) + }) + .transpose()? + .with_context(|| UnexpectedPlanExprSnafu { + desc: "vector matching on a plan without a time index", + })?; + let violation = if matches!( + modifier.as_ref().map(|modifier| &modifier.card), + Some(VectorMatchCardinality::OneToOne) + ) { + MatchGroupViolation::ImplicitManyToOne + } else { + MatchGroupViolation::AmbiguousGroupLabels + }; + Self::assert_unique_match_group( + join_plan, + group_exprs, + labels.into_iter().map(|(_, label)| label).collect(), + time_index_expr, + violation, + ) + } + + /// Guard the side of a vector matching that must hold one series per match group. + fn assert_unique_one_side( + plan: LogicalPlan, + join_keys: &BTreeSet, + tag_columns: &[String], + time_index_column: Option<&str>, + one_side_is_left: bool, + ) -> Result { + let Some(time_index_column) = time_index_column else { + return Ok(plan); + }; + if tag_columns.iter().all(|tag| join_keys.contains(tag)) { + return Ok(plan); + } + + let group_labels = join_keys.iter().cloned().collect::>(); + let group_exprs = group_labels + .iter() + .map(|label| DfExpr::Column(Column::from_name(label))) + .collect(); + Self::assert_unique_match_group( + plan, + group_exprs, + group_labels, + DfExpr::Column(Column::from_name(time_index_column)), + MatchGroupViolation::DuplicateOnOneSide { one_side_is_left }, + ) } fn selected_binary_match_labels( @@ -7118,6 +7390,20 @@ impl PromPlanner { input: LogicalPlan, name_to_expr: F, ) -> Result + where + F: FnMut(&String) -> Result, + { + self.projection_for_each_field_column_with_labels(input, None, name_to_expr) + } + + /// Like [`Self::projection_for_each_field_column`], but projects `result_labels` instead of + /// the context tag columns when a binary operation derived its own result label set. + fn projection_for_each_field_column_with_labels( + &mut self, + input: LogicalPlan, + result_labels: Option<&BinaryResultLabels>, + name_to_expr: F, + ) -> Result where F: FnMut(&String) -> Result, { @@ -7128,22 +7414,30 @@ impl PromPlanner { let table_ref = self.ctx.table_name.clone().map(TableReference::bare); // Derived labels can be unqualified even when the context still names the source table. let input_schema = input.schema().clone(); - let non_field_columns_iter = self - .ctx - .tag_columns - .iter() - .chain(self.ctx.time_index_column.iter()) - .map(|col| { - input_schema - .qualified_field_with_name(table_ref.as_ref(), col) - .or_else(|_| input_schema.qualified_field_with_unqualified_name(col)) - .map(|field| DfExpr::Column(field.into())) - .context(DataFusionPlanningSnafu) - }); - let tsid_iter = - Self::optional_tsid_projection(input.schema(), table_ref.as_ref(), self.ctx.use_tsid) - .into_iter() - .map(Ok); + let lookup = |col: &String| { + input_schema + .qualified_field_with_name(table_ref.as_ref(), col) + .or_else(|_| input_schema.qualified_field_with_unqualified_name(col)) + .map(|field| DfExpr::Column(field.into())) + .context(DataFusionPlanningSnafu) + }; + let tag_columns_iter = match result_labels { + Some(labels) => labels.exprs.iter().cloned().map(Ok).collect::>(), + None => self.ctx.tag_columns.iter().map(lookup).collect::>(), + }; + let non_field_columns_iter = tag_columns_iter + .into_iter() + .chain(self.ctx.time_index_column.iter().map(lookup)); + let tsid_iter = match result_labels { + Some(labels) => labels.tsid_projection(table_ref.clone()), + None => Self::optional_tsid_projection( + input.schema(), + table_ref.as_ref(), + self.ctx.use_tsid, + ), + } + .into_iter() + .map(Ok); // build computation exprs let result_field_columns = self @@ -11848,20 +12142,26 @@ mod test { PromPlanner::stmt_to_plan(table_provider, &eval_stmt, &build_query_engine_state()) .await .unwrap(); - let expected = "Projection: http_server_requests_seconds_count.uri, http_server_requests_seconds_count.kubernetes_namespace, http_server_requests_seconds_count.kubernetes_pod_name, http_server_requests_seconds_count.greptime_timestamp, CAST(http_server_requests_seconds_sum.greptime_value AS Float64) / CAST(http_server_requests_seconds_count.greptime_value AS Float64) AS http_server_requests_seconds_sum.greptime_value / http_server_requests_seconds_count.greptime_value\ - \n Inner Join: http_server_requests_seconds_sum.greptime_timestamp = http_server_requests_seconds_count.greptime_timestamp, http_server_requests_seconds_sum.uri = http_server_requests_seconds_count.uri\ - \n SubqueryAlias: http_server_requests_seconds_sum\ - \n PromInstantManipulate: range=[0..100000000], lookback=[1000], interval=[5000], time index=[greptime_timestamp]\ - \n PromSeriesDivide: tags=[\"uri\", \"kubernetes_namespace\", \"kubernetes_pod_name\"]\ - \n Sort: http_server_requests_seconds_sum.uri ASC NULLS FIRST, http_server_requests_seconds_sum.kubernetes_namespace ASC NULLS FIRST, http_server_requests_seconds_sum.kubernetes_pod_name ASC NULLS FIRST, http_server_requests_seconds_sum.greptime_timestamp ASC NULLS FIRST\ - \n Filter: http_server_requests_seconds_sum.uri = Utf8(\"/accounts/login\") AND http_server_requests_seconds_sum.greptime_timestamp >= TimestampMillisecond(-999, None) AND http_server_requests_seconds_sum.greptime_timestamp <= TimestampMillisecond(100000000, None)\ - \n TableScan: http_server_requests_seconds_sum\ - \n SubqueryAlias: http_server_requests_seconds_count\ - \n PromInstantManipulate: range=[0..100000000], lookback=[1000], interval=[5000], time index=[greptime_timestamp]\ - \n PromSeriesDivide: tags=[\"uri\", \"kubernetes_namespace\", \"kubernetes_pod_name\"]\ - \n Sort: http_server_requests_seconds_count.uri ASC NULLS FIRST, http_server_requests_seconds_count.kubernetes_namespace ASC NULLS FIRST, http_server_requests_seconds_count.kubernetes_pod_name ASC NULLS FIRST, http_server_requests_seconds_count.greptime_timestamp ASC NULLS FIRST\ - \n Filter: http_server_requests_seconds_count.uri = Utf8(\"/accounts/login\") AND http_server_requests_seconds_count.greptime_timestamp >= TimestampMillisecond(-999, None) AND http_server_requests_seconds_count.greptime_timestamp <= TimestampMillisecond(100000000, None)\ - \n TableScan: http_server_requests_seconds_count"; + let expected = "Projection: http_server_requests_seconds_sum.uri, http_server_requests_seconds_count.greptime_timestamp, CAST(http_server_requests_seconds_sum.greptime_value AS Float64) / CAST(http_server_requests_seconds_count.greptime_value AS Float64) AS http_server_requests_seconds_sum.greptime_value / http_server_requests_seconds_count.greptime_value\ + \n Projection: http_server_requests_seconds_sum.uri, http_server_requests_seconds_sum.kubernetes_namespace, http_server_requests_seconds_sum.kubernetes_pod_name, http_server_requests_seconds_sum.greptime_timestamp, http_server_requests_seconds_sum.greptime_value, http_server_requests_seconds_count.uri, http_server_requests_seconds_count.kubernetes_namespace, http_server_requests_seconds_count.kubernetes_pod_name, http_server_requests_seconds_count.greptime_timestamp, http_server_requests_seconds_count.greptime_value\ + \n Filter: prom_assert_unique_match_group(__promql_match_group_count, http_server_requests_seconds_sum.uri)\ + \n WindowAggr: windowExpr=[[count(Int64(1)) PARTITION BY [http_server_requests_seconds_sum.uri, http_server_requests_seconds_sum.greptime_timestamp] ROWS BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING AS __promql_match_group_count]]\ + \n Inner Join: http_server_requests_seconds_sum.greptime_timestamp = http_server_requests_seconds_count.greptime_timestamp, http_server_requests_seconds_sum.uri = http_server_requests_seconds_count.uri\ + \n SubqueryAlias: http_server_requests_seconds_sum\ + \n PromInstantManipulate: range=[0..100000000], lookback=[1000], interval=[5000], time index=[greptime_timestamp]\ + \n PromSeriesDivide: tags=[\"uri\", \"kubernetes_namespace\", \"kubernetes_pod_name\"]\ + \n Sort: http_server_requests_seconds_sum.uri ASC NULLS FIRST, http_server_requests_seconds_sum.kubernetes_namespace ASC NULLS FIRST, http_server_requests_seconds_sum.kubernetes_pod_name ASC NULLS FIRST, http_server_requests_seconds_sum.greptime_timestamp ASC NULLS FIRST\ + \n Filter: http_server_requests_seconds_sum.uri = Utf8(\"/accounts/login\") AND http_server_requests_seconds_sum.greptime_timestamp >= TimestampMillisecond(-999, None) AND http_server_requests_seconds_sum.greptime_timestamp <= TimestampMillisecond(100000000, None)\ + \n TableScan: http_server_requests_seconds_sum\ + \n SubqueryAlias: http_server_requests_seconds_count\ + \n Projection: http_server_requests_seconds_count.uri, http_server_requests_seconds_count.kubernetes_namespace, http_server_requests_seconds_count.kubernetes_pod_name, http_server_requests_seconds_count.greptime_timestamp, http_server_requests_seconds_count.greptime_value\ + \n Filter: prom_assert_unique_match_group(__promql_match_group_count, http_server_requests_seconds_count.uri)\ + \n WindowAggr: windowExpr=[[count(Int64(1)) PARTITION BY [http_server_requests_seconds_count.uri, http_server_requests_seconds_count.greptime_timestamp] ROWS BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING AS __promql_match_group_count]]\ + \n PromInstantManipulate: range=[0..100000000], lookback=[1000], interval=[5000], time index=[greptime_timestamp]\ + \n PromSeriesDivide: tags=[\"uri\", \"kubernetes_namespace\", \"kubernetes_pod_name\"]\ + \n Sort: http_server_requests_seconds_count.uri ASC NULLS FIRST, http_server_requests_seconds_count.kubernetes_namespace ASC NULLS FIRST, http_server_requests_seconds_count.kubernetes_pod_name ASC NULLS FIRST, http_server_requests_seconds_count.greptime_timestamp ASC NULLS FIRST\ + \n Filter: http_server_requests_seconds_count.uri = Utf8(\"/accounts/login\") AND http_server_requests_seconds_count.greptime_timestamp >= TimestampMillisecond(-999, None) AND http_server_requests_seconds_count.greptime_timestamp <= TimestampMillisecond(100000000, None)\ + \n TableScan: http_server_requests_seconds_count"; assert_eq!(plan.to_string(), expected); } diff --git a/src/query/src/promql/planner/matching_filters.rs b/src/query/src/promql/planner/matching_filters.rs index e38d48044da..d70e1bc4192 100644 --- a/src/query/src/promql/planner/matching_filters.rs +++ b/src/query/src/promql/planner/matching_filters.rs @@ -56,9 +56,10 @@ const LABEL_PRESERVING_RANGE_FUNCTIONS: [&str; 20] = [ /// equality, so a matcher on a matching label already holds for every surviving pair. Copying it /// to the other operand can only drop rows that had no partner, whatever the matcher kind, and /// `NULL` pairs only with `NULL`. Nor can it split a match group, whose series agree on every -/// matching label, so it does not affect a one-to-one cardinality check such as #9209. -/// Grouped matching is conservatively restricted to a provably unique one-side operand until -/// its cardinality checks are implemented (#9208). +/// matching label: a propagated matcher removes whole groups, so the cardinality check still +/// sees every group that takes part in the matching. A duplicate in a group without a partner +/// goes unreported as a result, where Prometheus fails the query on it. +/// Grouped matching is conservatively restricted to a provably unique one-side operand. /// /// `left_tags` and `right_tags` are the tag columns of the planned operands: a matcher may just /// as well constrain a value field, which the parsed expression does not distinguish from a diff --git a/tests/cases/standalone/common/promql/matching_groups.result b/tests/cases/standalone/common/promql/matching_groups.result index 183ae6d218c..04b1588b292 100644 --- a/tests/cases/standalone/common/promql/matching_groups.result +++ b/tests/cases/standalone/common/promql/matching_groups.result @@ -30,19 +30,16 @@ INSERT INTO matching_groups_right VALUES Affected Rows: 5 --- The result labels are the right operand's tag set rather than what the modifier implies --- (#9207), so `zone` is absent wherever the right side aggregates it away and rows differing --- only by it are indistinguishable. -- Filtering whole ranking partitions preserves the aggregate and scalar arithmetic. -- SQLNESS SORT_RESULT 3 1 TQL EVAL (0, 0, '1s') (8 * matching_groups_left{host="x"}) / on(host) group_left topk by(host)(1, max by(host)(matching_groups_right)); -+------+---------------------+--------------------------------------------------------------------------------------------------------------------+ -| host | ts | matching_groups_left.Float64(8) * greptime_value / matching_groups_right.max(matching_groups_right.greptime_value) | -+------+---------------------+--------------------------------------------------------------------------------------------------------------------+ -| x | 1970-01-01T00:00:00 | 24.0 | -| x | 1970-01-01T00:00:00 | 48.0 | -+------+---------------------+--------------------------------------------------------------------------------------------------------------------+ ++------+------+---------------------+--------------------------------------------------------------------------------------------------------------------+ +| host | zone | ts | matching_groups_left.Float64(8) * greptime_value / matching_groups_right.max(matching_groups_right.greptime_value) | ++------+------+---------------------+--------------------------------------------------------------------------------------------------------------------+ +| x | a | 1970-01-01T00:00:00 | 24.0 | +| x | b | 1970-01-01T00:00:00 | 48.0 | ++------+------+---------------------+--------------------------------------------------------------------------------------------------------------------+ -- SQLNESS SORT_RESULT 3 1 TQL EVAL (0, 0, '1s') sum by(host)(matching_groups_left) / on(host) group_right() (matching_groups_right{host="x"} * 2); @@ -57,45 +54,45 @@ TQL EVAL (0, 0, '1s') sum by(host)(matching_groups_left) / on(host) group_right( -- SQLNESS SORT_RESULT 3 1 TQL EVAL (0, 0, '1s') matching_groups_left{host!="y"} / on(host) group_left max by(host)(matching_groups_right); -+------+---------------------+-------------------------------------------------------------------------------------------------------+ -| host | ts | matching_groups_left.greptime_value / matching_groups_right.max(matching_groups_right.greptime_value) | -+------+---------------------+-------------------------------------------------------------------------------------------------------+ -| | 1970-01-01T00:00:00 | 6.0 | -| | 1970-01-01T00:00:00 | 6.0 | -| x | 1970-01-01T00:00:00 | 3.0 | -| x | 1970-01-01T00:00:00 | 6.0 | -+------+---------------------+-------------------------------------------------------------------------------------------------------+ ++------+------+---------------------+-------------------------------------------------------------------------------------------------------+ +| host | zone | ts | matching_groups_left.greptime_value / matching_groups_right.max(matching_groups_right.greptime_value) | ++------+------+---------------------+-------------------------------------------------------------------------------------------------------+ +| | a | 1970-01-01T00:00:00 | 6.0 | +| | a | 1970-01-01T00:00:00 | 6.0 | +| x | a | 1970-01-01T00:00:00 | 3.0 | +| x | b | 1970-01-01T00:00:00 | 6.0 | ++------+------+---------------------+-------------------------------------------------------------------------------------------------------+ -- SQLNESS SORT_RESULT 3 1 TQL EVAL (0, 0, '1s') matching_groups_left{host=~"x|y"} / on(host) group_left max by(host)(matching_groups_right); -+------+---------------------+-------------------------------------------------------------------------------------------------------+ -| host | ts | matching_groups_left.greptime_value / matching_groups_right.max(matching_groups_right.greptime_value) | -+------+---------------------+-------------------------------------------------------------------------------------------------------+ -| x | 1970-01-01T00:00:00 | 3.0 | -| x | 1970-01-01T00:00:00 | 6.0 | -| y | 1970-01-01T00:00:00 | 0.36 | -+------+---------------------+-------------------------------------------------------------------------------------------------------+ ++------+------+---------------------+-------------------------------------------------------------------------------------------------------+ +| host | zone | ts | matching_groups_left.greptime_value / matching_groups_right.max(matching_groups_right.greptime_value) | ++------+------+---------------------+-------------------------------------------------------------------------------------------------------+ +| x | a | 1970-01-01T00:00:00 | 3.0 | +| x | b | 1970-01-01T00:00:00 | 6.0 | +| y | a | 1970-01-01T00:00:00 | 0.36 | ++------+------+---------------------+-------------------------------------------------------------------------------------------------------+ -- SQLNESS SORT_RESULT 3 1 TQL EVAL (0, 0, '1s') matching_groups_left{host=""} / on(host) group_left max by(host)(matching_groups_right); -+------+---------------------+-------------------------------------------------------------------------------------------------------+ -| host | ts | matching_groups_left.greptime_value / matching_groups_right.max(matching_groups_right.greptime_value) | -+------+---------------------+-------------------------------------------------------------------------------------------------------+ -| | 1970-01-01T00:00:00 | 6.0 | -| | 1970-01-01T00:00:00 | 6.0 | -+------+---------------------+-------------------------------------------------------------------------------------------------------+ ++------+------+---------------------+-------------------------------------------------------------------------------------------------------+ +| host | zone | ts | matching_groups_left.greptime_value / matching_groups_right.max(matching_groups_right.greptime_value) | ++------+------+---------------------+-------------------------------------------------------------------------------------------------------+ +| | a | 1970-01-01T00:00:00 | 6.0 | +| | a | 1970-01-01T00:00:00 | 6.0 | ++------+------+---------------------+-------------------------------------------------------------------------------------------------------+ -- SQLNESS SORT_RESULT 3 1 TQL EVAL (0, 0, '1s') matching_groups_left{host!~"x|y"} / on(host) group_left max by(host)(matching_groups_right); -+------+---------------------+-------------------------------------------------------------------------------------------------------+ -| host | ts | matching_groups_left.greptime_value / matching_groups_right.max(matching_groups_right.greptime_value) | -+------+---------------------+-------------------------------------------------------------------------------------------------------+ -| | 1970-01-01T00:00:00 | 6.0 | -| | 1970-01-01T00:00:00 | 6.0 | -+------+---------------------+-------------------------------------------------------------------------------------------------------+ ++------+------+---------------------+-------------------------------------------------------------------------------------------------------+ +| host | zone | ts | matching_groups_left.greptime_value / matching_groups_right.max(matching_groups_right.greptime_value) | ++------+------+---------------------+-------------------------------------------------------------------------------------------------------+ +| | a | 1970-01-01T00:00:00 | 6.0 | +| | a | 1970-01-01T00:00:00 | 6.0 | ++------+------+---------------------+-------------------------------------------------------------------------------------------------------+ -- bottomk partitions by host when zone is excluded from grouping. -- SQLNESS SORT_RESULT 3 1 @@ -130,46 +127,29 @@ TQL EVAL (0, 0, '1s') matching_groups_left / on(host,zone) bottomk(1, matching_g -- SQLNESS SORT_RESULT 3 1 TQL EVAL (0, 0, '1s') matching_groups_left{host="x"} / on(host) group_left topk(1, max by(host)(matching_groups_right)); -+------+----+-------------------------------------------------------------------------------------------------------+ -| host | ts | matching_groups_left.greptime_value / matching_groups_right.max(matching_groups_right.greptime_value) | -+------+----+-------------------------------------------------------------------------------------------------------+ -+------+----+-------------------------------------------------------------------------------------------------------+ ++------+------+----+-------------------------------------------------------------------------------------------------------+ +| host | zone | ts | matching_groups_left.greptime_value / matching_groups_right.max(matching_groups_right.greptime_value) | ++------+------+----+-------------------------------------------------------------------------------------------------------+ ++------+------+----+-------------------------------------------------------------------------------------------------------+ -- SQLNESS SORT_RESULT 3 1 TQL EVAL (0, 0, '1s') matching_groups_left{host="x"} / on(host) group_left topk by(host)(1, topk(1, max by(host)(matching_groups_right))); -+------+----+-------------------------------------------------------------------------------------------------------+ -| host | ts | matching_groups_left.greptime_value / matching_groups_right.max(matching_groups_right.greptime_value) | -+------+----+-------------------------------------------------------------------------------------------------------+ -+------+----+-------------------------------------------------------------------------------------------------------+ ++------+------+----+-------------------------------------------------------------------------------------------------------+ +| host | zone | ts | matching_groups_left.greptime_value / matching_groups_right.max(matching_groups_right.greptime_value) | ++------+------+----+-------------------------------------------------------------------------------------------------------+ ++------+------+----+-------------------------------------------------------------------------------------------------------+ --- Both one-sides carry two series per `host`. Prometheus rejects this with "many-to-many --- matching not allowed: matching labels must be unique on one side"; GreptimeDB has no --- cardinality check (#9209) and returns the cross product. The rewrite leaves these operands --- alone, so the recorded output is the unchanged baseline. +-- Both one-sides carry two series per `host`, which the matching cannot resolve. -- SQLNESS SORT_RESULT 3 1 TQL EVAL (0, 0, '1s') matching_groups_left{host="x"} / on(host) group_left matching_groups_right; -+------+------+---------------------+----------------------------------------------------------------------------+ -| host | zone | ts | matching_groups_left.greptime_value / matching_groups_right.greptime_value | -+------+------+---------------------+----------------------------------------------------------------------------+ -| x | a | 1970-01-01T00:00:00 | 12.0 | -| x | a | 1970-01-01T00:00:00 | 6.0 | -| x | b | 1970-01-01T00:00:00 | 3.0 | -| x | b | 1970-01-01T00:00:00 | 6.0 | -+------+------+---------------------+----------------------------------------------------------------------------+ +Error: 3001(EngineExecuteQuery), Execution error: found duplicate series for the match group {host="x"} on the right hand-side of the operation; many-to-many matching not allowed: matching labels must be unique on one side -- SQLNESS SORT_RESULT 3 1 TQL EVAL (0, 0, '1s') matching_groups_left{host="x"} / on(host) group_left max by(host,zone)(matching_groups_right); -+------+------+---------------------+-------------------------------------------------------------------------------------------------------+ -| host | zone | ts | matching_groups_left.greptime_value / matching_groups_right.max(matching_groups_right.greptime_value) | -+------+------+---------------------+-------------------------------------------------------------------------------------------------------+ -| x | a | 1970-01-01T00:00:00 | 12.0 | -| x | a | 1970-01-01T00:00:00 | 6.0 | -| x | b | 1970-01-01T00:00:00 | 3.0 | -| x | b | 1970-01-01T00:00:00 | 6.0 | -+------+------+---------------------+-------------------------------------------------------------------------------------------------------+ +Error: 3001(EngineExecuteQuery), Execution error: found duplicate series for the match group {host="x"} on the right hand-side of the operation; many-to-many matching not allowed: matching labels must be unique on one side DROP TABLE matching_groups_left; diff --git a/tests/cases/standalone/common/promql/matching_groups.sql b/tests/cases/standalone/common/promql/matching_groups.sql index 97e46da0437..03bdf8eaaa7 100644 --- a/tests/cases/standalone/common/promql/matching_groups.sql +++ b/tests/cases/standalone/common/promql/matching_groups.sql @@ -19,9 +19,6 @@ INSERT INTO matching_groups_right VALUES ('x', 'a', 0, 2), ('x', 'b', 0, 4), ('y', 'a', 0, 100), ('', 'a', 0, 8), (NULL, 'a', 0, 10); --- The result labels are the right operand's tag set rather than what the modifier implies --- (#9207), so `zone` is absent wherever the right side aggregates it away and rows differing --- only by it are indistinguishable. -- Filtering whole ranking partitions preserves the aggregate and scalar arithmetic. -- SQLNESS SORT_RESULT 3 1 TQL EVAL (0, 0, '1s') (8 * matching_groups_left{host="x"}) / on(host) group_left topk by(host)(1, max by(host)(matching_groups_right)); @@ -52,10 +49,7 @@ TQL EVAL (0, 0, '1s') matching_groups_left{host="x"} / on(host) group_left topk( -- SQLNESS SORT_RESULT 3 1 TQL EVAL (0, 0, '1s') matching_groups_left{host="x"} / on(host) group_left topk by(host)(1, topk(1, max by(host)(matching_groups_right))); --- Both one-sides carry two series per `host`. Prometheus rejects this with "many-to-many --- matching not allowed: matching labels must be unique on one side"; GreptimeDB has no --- cardinality check (#9209) and returns the cross product. The rewrite leaves these operands --- alone, so the recorded output is the unchanged baseline. +-- Both one-sides carry two series per `host`, which the matching cannot resolve. -- SQLNESS SORT_RESULT 3 1 TQL EVAL (0, 0, '1s') matching_groups_left{host="x"} / on(host) group_left matching_groups_right; -- SQLNESS SORT_RESULT 3 1 diff --git a/tests/cases/standalone/common/promql/set_operation.result b/tests/cases/standalone/common/promql/set_operation.result index 2f0041f1a36..db97fed7e6b 100644 --- a/tests/cases/standalone/common/promql/set_operation.result +++ b/tests/cases/standalone/common/promql/set_operation.result @@ -688,26 +688,26 @@ tql eval (3, 4, '1s') cache_hit_with_null_label / (cache_miss_with_null_label + -- SQLNESS SORT_RESULT 3 1 tql eval (3, 4, '1s') cache_hit_with_null_label / ignoring(null_label) (cache_miss_with_null_label + ignoring(null_label) cache_hit_with_null_label); -+-------+------------+---------------------+---------------------------------------------------------------------------------------------------------------+ -| job | null_label | ts | lhs.greptime_value / rhs.cache_miss_with_null_label.greptime_value + cache_hit_with_null_label.greptime_value | -+-------+------------+---------------------+---------------------------------------------------------------------------------------------------------------+ -| read | | 1970-01-01T00:00:03 | 0.5 | -| read | | 1970-01-01T00:00:04 | 0.75 | -| write | | 1970-01-01T00:00:03 | 0.5 | -| write | | 1970-01-01T00:00:04 | 0.6666666666666666 | -+-------+------------+---------------------+---------------------------------------------------------------------------------------------------------------+ ++-------+---------------------+---------------------------------------------------------------------------------------------------------------+ +| job | ts | lhs.greptime_value / rhs.cache_miss_with_null_label.greptime_value + cache_hit_with_null_label.greptime_value | ++-------+---------------------+---------------------------------------------------------------------------------------------------------------+ +| read | 1970-01-01T00:00:03 | 0.5 | +| read | 1970-01-01T00:00:04 | 0.75 | +| write | 1970-01-01T00:00:03 | 0.5 | +| write | 1970-01-01T00:00:04 | 0.6666666666666666 | ++-------+---------------------+---------------------------------------------------------------------------------------------------------------+ -- SQLNESS SORT_RESULT 3 1 tql eval (3, 4, '1s') cache_hit_with_null_label / on(job) (cache_miss_with_null_label + on(job) cache_hit_with_null_label); -+-------+------------+---------------------+---------------------------------------------------------------------------------------------------------------+ -| job | null_label | ts | lhs.greptime_value / rhs.cache_miss_with_null_label.greptime_value + cache_hit_with_null_label.greptime_value | -+-------+------------+---------------------+---------------------------------------------------------------------------------------------------------------+ -| read | | 1970-01-01T00:00:03 | 0.5 | -| read | | 1970-01-01T00:00:04 | 0.75 | -| write | | 1970-01-01T00:00:03 | 0.5 | -| write | | 1970-01-01T00:00:04 | 0.6666666666666666 | -+-------+------------+---------------------+---------------------------------------------------------------------------------------------------------------+ ++-------+---------------------+---------------------------------------------------------------------------------------------------------------+ +| job | ts | lhs.greptime_value / rhs.cache_miss_with_null_label.greptime_value + cache_hit_with_null_label.greptime_value | ++-------+---------------------+---------------------------------------------------------------------------------------------------------------+ +| read | 1970-01-01T00:00:03 | 0.5 | +| read | 1970-01-01T00:00:04 | 0.75 | +| write | 1970-01-01T00:00:03 | 0.5 | +| write | 1970-01-01T00:00:04 | 0.6666666666666666 | ++-------+---------------------+---------------------------------------------------------------------------------------------------------------+ drop table cache_hit_with_null_label; diff --git a/tests/cases/standalone/common/promql/tsid_binary_join_regression.result b/tests/cases/standalone/common/promql/tsid_binary_join_regression.result index d80d29ce000..54fd5ad5e26 100644 --- a/tests/cases/standalone/common/promql/tsid_binary_join_regression.result +++ b/tests/cases/standalone/common/promql/tsid_binary_join_regression.result @@ -243,12 +243,19 @@ TQL ANALYZE (0, 5, '5s') tsid_binary_join_left / ignoring(host) tsid_binary_join +-+-+-+ | stage | node | plan_| +-+-+-+ -| 0_| 0_|_ProjectionExec: expr=[host@0 as host, job@1 as job, ts@2 as ts, __tsid@3 as __tsid, greptime_value@4 / greptime_value@5 as tsid_binary_join_left.greptime_value / tsid_binary_join_right.greptime_value] REDACTED -|_|_|_HashJoinExec: mode=Partitioned, join_type=Inner, on=[(job@1, job@2), (ts@2, ts@4)], projection=[host@4, job@5, ts@7, __tsid@6, greptime_value@0, greptime_value@3], NullsEqual: true REDACTED +| 0_| 0_|_ProjectionExec: expr=[job@0 as job, ts@1 as ts, greptime_value@2 / greptime_value@3 as tsid_binary_join_left.greptime_value / tsid_binary_join_right.greptime_value] REDACTED +|_|_|_FilterExec: prom_assert_unique_match_group(__promql_match_group_count@4, job@1), projection=[job@1, ts@3, greptime_value@0, greptime_value@2] REDACTED +|_|_|_ProjectionExec: expr=[greptime_value@0 as greptime_value, job@1 as job, greptime_value@3 as greptime_value, ts@4 as ts, __promql_match_group_count@5 as __promql_match_group_count] REDACTED +|_|_|_WindowAggExec: wdw=[__promql_match_group_count: Ok(Field { name: "__promql_match_group_count", data_type: Int64 }), frame: WindowFrame { units: Rows, start_bound: Preceding(UInt64(NULL)), end_bound: Following(UInt64(NULL)), is_causal: false }] REDACTED +|_|_|_HashJoinExec: mode=Partitioned, join_type=Inner, on=[(job@1, job@1), (ts@2, ts@2)], projection=[greptime_value@0, job@1, ts@2, greptime_value@3, ts@5], NullsEqual: true REDACTED |_|_|_RepartitionExec: partitioning=Hash([job@1, ts@2],REDACTED |_|_|_ProjectionExec: expr=[greptime_value@0 as greptime_value, job@2 as job, ts@4 as ts] REDACTED |_|_|_MergeScanExec: REDACTED -|_|_|_RepartitionExec: partitioning=Hash([job@2, ts@4],REDACTED +|_|_|_FilterExec: prom_assert_unique_match_group(__promql_match_group_count@3, job@1), projection=[greptime_value@0, job@1, ts@2] REDACTED +|_|_|_WindowAggExec: wdw=[__promql_match_group_count: Ok(Field { name: "__promql_match_group_count", data_type: Int64 }), frame: WindowFrame { units: Rows, start_bound: Preceding(UInt64(NULL)), end_bound: Following(UInt64(NULL)), is_causal: false }] REDACTED +|_|_|_SortExec: expr=[job@1 ASC NULLS LAST, ts@2 ASC NULLS LAST], preserve_partitioning=[true] REDACTED +|_|_|_RepartitionExec: partitioning=Hash([job@1, ts@2],REDACTED +|_|_|_ProjectionExec: expr=[greptime_value@0 as greptime_value, job@2 as job, ts@4 as ts] REDACTED |_|_|_MergeScanExec: REDACTED |_|_|_| | 1_| 0_|_PromInstantManipulateExec: range=[0..5000], lookback=[300000], interval=[5000], time index=[ts] REDACTED @@ -282,12 +289,19 @@ TQL ANALYZE (0, 5, '5s') tsid_binary_join_left / on(job) tsid_binary_join_right_ +-+-+-+ | stage | node | plan_| +-+-+-+ -| 0_| 0_|_ProjectionExec: expr=[job@0 as job, ts@1 as ts, __tsid@2 as __tsid, greptime_value@3 / greptime_value@4 as tsid_binary_join_left.greptime_value / tsid_binary_join_right_by_job.greptime_value] REDACTED -|_|_|_HashJoinExec: mode=Partitioned, join_type=Inner, on=[(job@1, job@1), (ts@2, ts@3)], projection=[job@4, ts@6, __tsid@5, greptime_value@0, greptime_value@3], NullsEqual: true REDACTED +| 0_| 0_|_ProjectionExec: expr=[job@0 as job, ts@1 as ts, greptime_value@2 / greptime_value@3 as tsid_binary_join_left.greptime_value / tsid_binary_join_right_by_job.greptime_value] REDACTED +|_|_|_FilterExec: prom_assert_unique_match_group(__promql_match_group_count@4, job@1), projection=[job@1, ts@3, greptime_value@0, greptime_value@2] REDACTED +|_|_|_ProjectionExec: expr=[greptime_value@0 as greptime_value, job@1 as job, greptime_value@3 as greptime_value, ts@4 as ts, __promql_match_group_count@5 as __promql_match_group_count] REDACTED +|_|_|_WindowAggExec: wdw=[__promql_match_group_count: Ok(Field { name: "__promql_match_group_count", data_type: Int64 }), frame: WindowFrame { units: Rows, start_bound: Preceding(UInt64(NULL)), end_bound: Following(UInt64(NULL)), is_causal: false }] REDACTED +|_|_|_SortExec: expr=[job@1 ASC NULLS LAST, ts@2 ASC NULLS LAST], preserve_partitioning=[true] REDACTED +|_|_|_RepartitionExec: partitioning=Hash([job@1, ts@2],REDACTED +|_|_|_CoalescePartitionsExec REDACTED +|_|_|_HashJoinExec: mode=Partitioned, join_type=Inner, on=[(job@1, job@1), (ts@2, ts@2)], projection=[greptime_value@0, job@1, ts@2, greptime_value@3, ts@5], NullsEqual: true REDACTED |_|_|_RepartitionExec: partitioning=Hash([job@1, ts@2],REDACTED |_|_|_ProjectionExec: expr=[greptime_value@0 as greptime_value, job@2 as job, ts@4 as ts] REDACTED |_|_|_MergeScanExec: REDACTED -|_|_|_RepartitionExec: partitioning=Hash([job@1, ts@3],REDACTED +|_|_|_RepartitionExec: partitioning=Hash([job@1, ts@2],REDACTED +|_|_|_ProjectionExec: expr=[greptime_value@0 as greptime_value, job@1 as job, ts@3 as ts] REDACTED |_|_|_MergeScanExec: REDACTED |_|_|_| | 1_| 0_|_PromInstantManipulateExec: range=[0..5000], lookback=[300000], interval=[5000], time index=[ts] REDACTED @@ -566,13 +580,16 @@ TQL ANALYZE (0, 5, '5s') (tsid_binary_join_left / ignoring(host) group_left tsid |_|_|_HashJoinExec: mode=CollectLeft, join_type=Inner, on=[(__tsid@1, __tsid@3), (ts@0, ts@4)], projection=[host@4, job@5, ts@7, __tsid@6, tsid_binary_join_left.greptime_value / tsid_binary_join_right.greptime_value@2, greptime_value@3], NullsEqual: true REDACTED |_|_|_CoalescePartitionsExec REDACTED |_|_|_ProjectionExec: expr=[ts@0 as ts, __tsid@1 as __tsid, greptime_value@2 / greptime_value@3 as tsid_binary_join_left.greptime_value / tsid_binary_join_right.greptime_value] REDACTED -|_|_|_HashJoinExec: mode=Partitioned, join_type=Inner, on=[(job@1, job@1), (ts@2, ts@3)], projection=[ts@6, __tsid@5, greptime_value@0, greptime_value@3], NullsEqual: true REDACTED -|_|_|_RepartitionExec: partitioning=Hash([REDACTED -|_|_|_ProjectionExec: expr=[greptime_value@0 as greptime_value, job@2 as job, ts@4 as ts] REDACTED -|_|_|_MergeScanExec: REDACTED +|_|_|_HashJoinExec: mode=Partitioned, join_type=Inner, on=[(job@1, job@1), (ts@3, ts@2)], projection=[ts@6, __tsid@2, greptime_value@0, greptime_value@4], NullsEqual: true REDACTED |_|_|_RepartitionExec: partitioning=Hash([REDACTED |_|_|_ProjectionExec: expr=[greptime_value@0 as greptime_value, job@2 as job, __tsid@3 as __tsid, ts@4 as ts] REDACTED |_|_|_MergeScanExec: REDACTED +|_|_|_FilterExec: prom_assert_unique_match_group(__promql_match_group_count@3, job@1), projection=[greptime_value@0, job@1, ts@2] REDACTED +|_|_|_WindowAggExec: wdw=[__promql_match_group_count: Ok(Field { name: "__promql_match_group_count", data_type: Int64 }), frame: WindowFrame { units: Rows, start_bound: Preceding(UInt64(NULL)), end_bound: Following(UInt64(NULL)), is_causal: false }] REDACTED +|_|_|_SortExec: expr=[job@1 ASC NULLS LAST, ts@2 ASC NULLS LAST], preserve_partitioning=[true] REDACTED +|_|_|_RepartitionExec: partitioning=Hash([REDACTED +|_|_|_ProjectionExec: expr=[greptime_value@0 as greptime_value, job@2 as job, ts@4 as ts] REDACTED +|_|_|_MergeScanExec: REDACTED |_|_|_RepartitionExec: partitioning=REDACTED |_|_|_MergeScanExec: REDACTED |_|_|_| @@ -597,6 +614,20 @@ TQL ANALYZE (0, 5, '5s') (tsid_binary_join_left / ignoring(host) group_left tsid |_|_| Total rows: 4_| +-+-+-+ +-- A group modifier that keeps the many side's whole tag set keeps its `__tsid` too, so the +-- enclosing operation matches on it. Same rows as the label-based plan. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (0, 5, '5s') (tsid_binary_join_left / ignoring(host) group_left tsid_binary_join_right) / tsid_binary_join_left; + ++-------+------+---------------------+--------------------------------------------------------------------------------------------------------------------------------------------+ +| host | job | ts | tsid_binary_join_right.tsid_binary_join_left.greptime_value / tsid_binary_join_right.greptime_value / tsid_binary_join_left.greptime_value | ++-------+------+---------------------+--------------------------------------------------------------------------------------------------------------------------------------------+ +| host1 | job1 | 1970-01-01T00:00:00 | 0.3333333333333333 | +| host1 | job1 | 1970-01-01T00:00:05 | 0.2 | +| host2 | job2 | 1970-01-01T00:00:00 | 0.16666666666666666 | +| host2 | job2 | 1970-01-01T00:00:05 | 0.14285714285714285 | ++-------+------+---------------------+--------------------------------------------------------------------------------------------------------------------------------------------+ + -- SQLNESS SORT_RESULT 3 1 TQL EVAL (0, 5, '5s') tsid_binary_join_left / tsid_binary_join_right; diff --git a/tests/cases/standalone/common/promql/tsid_binary_join_regression.sql b/tests/cases/standalone/common/promql/tsid_binary_join_regression.sql index 7a3bc897c75..627680330aa 100644 --- a/tests/cases/standalone/common/promql/tsid_binary_join_regression.sql +++ b/tests/cases/standalone/common/promql/tsid_binary_join_regression.sql @@ -226,6 +226,11 @@ TQL ANALYZE (0, 5, '5s') (tsid_binary_join_left or tsid_binary_join_right) / tsi -- SQLNESS REPLACE region=\d+\(\d+,\s+\d+\) region=REDACTED TQL ANALYZE (0, 5, '5s') (tsid_binary_join_left / ignoring(host) group_left tsid_binary_join_right) / tsid_binary_join_left; +-- A group modifier that keeps the many side's whole tag set keeps its `__tsid` too, so the +-- enclosing operation matches on it. Same rows as the label-based plan. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (0, 5, '5s') (tsid_binary_join_left / ignoring(host) group_left tsid_binary_join_right) / tsid_binary_join_left; + -- SQLNESS SORT_RESULT 3 1 TQL EVAL (0, 5, '5s') tsid_binary_join_left / tsid_binary_join_right; diff --git a/tests/cases/standalone/common/promql/vector_matching.result b/tests/cases/standalone/common/promql/vector_matching.result new file mode 100644 index 00000000000..5382e94967c --- /dev/null +++ b/tests/cases/standalone/common/promql/vector_matching.result @@ -0,0 +1,406 @@ +CREATE TABLE counter_metric ( + host STRING NULL, + device STRING NULL, + ts TIMESTAMP(3) TIME INDEX, + greptime_value DOUBLE, + PRIMARY KEY(host, device) +); + +Affected Rows: 0 + +CREATE TABLE gauge_metric ( + host STRING NULL, + device STRING NULL, + ts TIMESTAMP(3) TIME INDEX, + greptime_value DOUBLE, + PRIMARY KEY(host, device) +); + +Affected Rows: 0 + +INSERT INTO counter_metric VALUES + ('host1', 'eth0', 0, 10), ('host1', 'eth1', 0, 20), ('host2', 'eth0', 0, 30); + +Affected Rows: 3 + +INSERT INTO gauge_metric VALUES + ('host1', 'eth0', 0, 2), ('host1', 'eth1', 0, 4), ('host2', 'eth0', 0, 5); + +Affected Rows: 3 + +-- Matching on the whole tag set keeps it. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (0, 0, '5s') counter_metric / gauge_metric; + ++-------+--------+---------------------+-------------------------------------------------------------+ +| host | device | ts | counter_metric.greptime_value / gauge_metric.greptime_value | ++-------+--------+---------------------+-------------------------------------------------------------+ +| host1 | eth0 | 1970-01-01T00:00:00 | 5.0 | +| host1 | eth1 | 1970-01-01T00:00:00 | 5.0 | +| host2 | eth0 | 1970-01-01T00:00:00 | 6.0 | ++-------+--------+---------------------+-------------------------------------------------------------+ + +-- `on(...)` reduces the result to the matching labels. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (0, 0, '5s') counter_metric{device="eth0"} / on(host) gauge_metric{device="eth0"}; + ++-------+---------------------+-------------------------------------------------------------+ +| host | ts | counter_metric.greptime_value / gauge_metric.greptime_value | ++-------+---------------------+-------------------------------------------------------------+ +| host1 | 1970-01-01T00:00:00 | 5.0 | +| host2 | 1970-01-01T00:00:00 | 6.0 | ++-------+---------------------+-------------------------------------------------------------+ + +-- `ignoring(...)` drops the ignored labels, here from operands that don't even share them. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (0, 0, '5s') counter_metric{device="eth0"} / ignoring(device) gauge_metric{device="eth1"}; + ++-------+---------------------+-------------------------------------------------------------+ +| host | ts | counter_metric.greptime_value / gauge_metric.greptime_value | ++-------+---------------------+-------------------------------------------------------------+ +| host1 | 1970-01-01T00:00:00 | 2.5 | ++-------+---------------------+-------------------------------------------------------------+ + +-- `on()` matches everything against everything, so both sides must hold a single series. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (0, 0, '5s') counter_metric{host="host1",device="eth0"} / on() gauge_metric{host="host1",device="eth0"}; + ++---------------------+-------------------------------------------------------------+ +| ts | counter_metric.greptime_value / gauge_metric.greptime_value | ++---------------------+-------------------------------------------------------------+ +| 1970-01-01T00:00:00 | 5.0 | ++---------------------+-------------------------------------------------------------+ + +-- A filtering comparison keeps the left sample values and the derived labels. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (0, 0, '5s') counter_metric{device="eth0"} > on(host) gauge_metric{device="eth0"}; + ++---------------------+----------------+-------+ +| ts | greptime_value | host | ++---------------------+----------------+-------+ +| 1970-01-01T00:00:00 | 10.0 | host1 | +| 1970-01-01T00:00:00 | 30.0 | host2 | ++---------------------+----------------+-------+ + +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (0, 0, '5s') counter_metric{device="eth0"} > bool on(host) gauge_metric{device="eth0"}; + ++-------+---------------------+-------------------------------------------------------------+ +| host | ts | counter_metric.greptime_value > gauge_metric.greptime_value | ++-------+---------------------+-------------------------------------------------------------+ +| host1 | 1970-01-01T00:00:00 | 1.0 | +| host2 | 1970-01-01T00:00:00 | 1.0 | ++-------+---------------------+-------------------------------------------------------------+ + +-- `host1` carries two devices on both sides: the match group repeats on the one side. +TQL EVAL (0, 0, '5s') counter_metric / on(host) gauge_metric; + +Error: 3001(EngineExecuteQuery), Execution error: found duplicate series for the match group {host="host1"} on the right hand-side of the operation; many-to-many matching not allowed: matching labels must be unique on one side + +-- The one side is unique here, the many side is not, and no group modifier makes it explicit. +TQL EVAL (0, 0, '5s') counter_metric / on(host) gauge_metric{device="eth0"}; + +Error: 3001(EngineExecuteQuery), Execution error: multiple matches for labels {host="host1"}: many-to-one matching must be explicit (group_left/group_right) + +-- `group_left` takes the result labels from the many side. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (0, 0, '5s') counter_metric / on(host) group_left gauge_metric{device="eth0"}; + ++-------+--------+---------------------+-------------------------------------------------------------+ +| host | device | ts | counter_metric.greptime_value / gauge_metric.greptime_value | ++-------+--------+---------------------+-------------------------------------------------------------+ +| host1 | eth0 | 1970-01-01T00:00:00 | 5.0 | +| host1 | eth1 | 1970-01-01T00:00:00 | 10.0 | +| host2 | eth0 | 1970-01-01T00:00:00 | 6.0 | ++-------+--------+---------------------+-------------------------------------------------------------+ + +-- `group_right` swaps the sides. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (0, 0, '5s') counter_metric{device="eth0"} / on(host) group_right gauge_metric; + ++-------+--------+---------------------+-------------------------------------------------------------+ +| host | device | ts | counter_metric.greptime_value / gauge_metric.greptime_value | ++-------+--------+---------------------+-------------------------------------------------------------+ +| host1 | eth0 | 1970-01-01T00:00:00 | 5.0 | +| host1 | eth1 | 1970-01-01T00:00:00 | 2.5 | +| host2 | eth0 | 1970-01-01T00:00:00 | 6.0 | ++-------+--------+---------------------+-------------------------------------------------------------+ + +-- `group_right` makes the left operand the one side, which `host1` breaks. +TQL EVAL (0, 0, '5s') counter_metric / on(host) group_right gauge_metric{device="eth0"}; + +Error: 3001(EngineExecuteQuery), Execution error: found duplicate series for the match group {host="host1"} on the left hand-side of the operation; many-to-many matching not allowed: matching labels must be unique on one side + +-- `group_left(device)` copies `device` from the one side, so both `host1` rows end up with the +-- same labels. +TQL EVAL (0, 0, '5s') counter_metric / on(host) group_left(device) gauge_metric{device="eth0"}; + +Error: 3001(EngineExecuteQuery), Execution error: multiple matches for labels {host="host1", device="eth0"}: grouping labels must ensure unique matches + +-- An included label the one side doesn't carry is dropped from the result. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (0, 0, '5s') counter_metric{device="eth0"} / on(host) group_left(device) sum by(host)(gauge_metric); + ++-------+---------------------+-------------------------------------------------------------------------------+ +| host | ts | counter_metric.greptime_value / gauge_metric.sum(gauge_metric.greptime_value) | ++-------+---------------------+-------------------------------------------------------------------------------+ +| host1 | 1970-01-01T00:00:00 | 1.6666666666666667 | +| host2 | 1970-01-01T00:00:00 | 6.0 | ++-------+---------------------+-------------------------------------------------------------------------------+ + +DROP TABLE counter_metric; + +Affected Rows: 0 + +DROP TABLE gauge_metric; + +Affected Rows: 0 + +-- A label can carry the name the cardinality check generates for its row count. +CREATE TABLE collide_left ( + host STRING NULL, + __promql_match_group_count STRING NULL, + ts TIMESTAMP(3) TIME INDEX, + greptime_value DOUBLE, + PRIMARY KEY(host, __promql_match_group_count) +); + +Affected Rows: 0 + +CREATE TABLE collide_right ( + host STRING NULL, + __promql_match_group_count STRING NULL, + ts TIMESTAMP(3) TIME INDEX, + greptime_value DOUBLE, + PRIMARY KEY(host, __promql_match_group_count) +); + +Affected Rows: 0 + +INSERT INTO collide_left VALUES ('host1', 'a', 0, 10); + +Affected Rows: 1 + +INSERT INTO collide_right VALUES ('host1', 'b', 0, 2); + +Affected Rows: 1 + +TQL EVAL (0, 0, '5s') collide_left / on(host) collide_right; + ++-------+---------------------+------------------------------------------------------------+ +| host | ts | collide_left.greptime_value / collide_right.greptime_value | ++-------+---------------------+------------------------------------------------------------+ +| host1 | 1970-01-01T00:00:00 | 5.0 | ++-------+---------------------+------------------------------------------------------------+ + +DROP TABLE collide_left; + +Affected Rows: 0 + +DROP TABLE collide_right; + +Affected Rows: 0 + +-- Partitioning on a column outside the match keys puts one match group in several regions. The +-- cardinality check counts the merged data, so it must still see the duplicate. +CREATE TABLE spread_left ( + host STRING NULL, + device STRING NULL, + ts TIMESTAMP(3) TIME INDEX, + greptime_value DOUBLE, + PRIMARY KEY(host, device) +) +PARTITION ON COLUMNS (device) ( + device < 'eth1', + device >= 'eth1' +); + +Affected Rows: 0 + +CREATE TABLE spread_right ( + host STRING NULL, + ts TIMESTAMP(3) TIME INDEX, + greptime_value DOUBLE, + PRIMARY KEY(host) +); + +Affected Rows: 0 + +INSERT INTO spread_left VALUES ('host1', 'eth0', 0, 10), ('host1', 'eth1', 0, 20); + +Affected Rows: 2 + +INSERT INTO spread_right VALUES ('host1', 0, 2); + +Affected Rows: 1 + +TQL EVAL (0, 0, '5s') spread_left / on(host) spread_right; + +Error: 3001(EngineExecuteQuery), Execution error: multiple matches for labels {host="host1"}: many-to-one matching must be explicit (group_left/group_right) + +TQL EVAL (0, 0, '5s') spread_right / on(host) group_left spread_left; + +Error: 3001(EngineExecuteQuery), Execution error: found duplicate series for the match group {host="host1"} on the right hand-side of the operation; many-to-many matching not allowed: matching labels must be unique on one side + +DROP TABLE spread_left; + +Affected Rows: 0 + +DROP TABLE spread_right; + +Affected Rows: 0 + +-- An outer aggregate groups on the labels the operand has left, not on the ones the scan +-- underneath could still produce. +CREATE TABLE outer_agg_a ( + host STRING NULL, + device STRING NULL, + ts TIMESTAMP(3) TIME INDEX, + greptime_value DOUBLE, + PRIMARY KEY(host, device) +); + +Affected Rows: 0 + +CREATE TABLE outer_agg_b ( + host STRING NULL, + ts TIMESTAMP(3) TIME INDEX, + greptime_value DOUBLE, + PRIMARY KEY(host) +); + +Affected Rows: 0 + +INSERT INTO outer_agg_a VALUES ('h1', 'd1', 0, 10), ('h2', 'd2', 0, 20); + +Affected Rows: 2 + +INSERT INTO outer_agg_b VALUES ('h1', 0, 2), ('h2', 0, 4); + +Affected Rows: 2 + +-- `on(host)` leaves only `host`, so `without(host)` groups everything together. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (0, 0, '5s') sum without(host) (outer_agg_a / on(host) outer_agg_b); + ++---------------------+--------------------------------------------------------------+ +| ts | sum(outer_agg_a.greptime_value / outer_agg_b.greptime_value) | ++---------------------+--------------------------------------------------------------+ +| 1970-01-01T00:00:00 | 10.0 | ++---------------------+--------------------------------------------------------------+ + +-- `device` is not a label of the operand any more, so it groups everything together too. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (0, 0, '5s') sum by(device) (outer_agg_a / on(host) outer_agg_b); + ++---------------------+--------------------------------------------------------------+ +| ts | sum(outer_agg_a.greptime_value / outer_agg_b.greptime_value) | ++---------------------+--------------------------------------------------------------+ +| 1970-01-01T00:00:00 | 10.0 | ++---------------------+--------------------------------------------------------------+ + +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (0, 0, '5s') sum by(host) (outer_agg_a / on(host) outer_agg_b); + ++------+---------------------+--------------------------------------------------------------+ +| host | ts | sum(outer_agg_a.greptime_value / outer_agg_b.greptime_value) | ++------+---------------------+--------------------------------------------------------------+ +| h1 | 1970-01-01T00:00:00 | 5.0 | +| h2 | 1970-01-01T00:00:00 | 5.0 | ++------+---------------------+--------------------------------------------------------------+ + +DROP TABLE outer_agg_a; + +Affected Rows: 0 + +DROP TABLE outer_agg_b; + +Affected Rows: 0 + +-- Same over the metric engine, whose operands carry `__tsid`. +CREATE TABLE outer_agg_physical ( + ts TIMESTAMP(3) TIME INDEX, + greptime_value DOUBLE, +) ENGINE = metric WITH ("physical_metric_table" = ""); + +Affected Rows: 0 + +CREATE TABLE outer_agg_metric_a ( + host STRING NULL, + device STRING NULL, + ts TIMESTAMP(3) NOT NULL, + greptime_value DOUBLE NULL, + TIME INDEX (ts), + PRIMARY KEY(host, device), +) +ENGINE = metric +WITH( + on_physical_table = 'outer_agg_physical' +); + +Affected Rows: 0 + +CREATE TABLE outer_agg_metric_b ( + host STRING NULL, + ts TIMESTAMP(3) NOT NULL, + greptime_value DOUBLE NULL, + TIME INDEX (ts), + PRIMARY KEY(host), +) +ENGINE = metric +WITH( + on_physical_table = 'outer_agg_physical' +); + +Affected Rows: 0 + +INSERT INTO outer_agg_metric_a (host, device, ts, greptime_value) VALUES + ('h1', 'd1', 0, 10), ('h2', 'd2', 0, 20); + +Affected Rows: 2 + +INSERT INTO outer_agg_metric_b (host, ts, greptime_value) VALUES + ('h1', 0, 2), ('h2', 0, 4); + +Affected Rows: 2 + +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (0, 0, '5s') sum without(host) (outer_agg_metric_a / on(host) outer_agg_metric_b); + ++---------------------+----------------------------------------------------------------------------+ +| ts | sum(outer_agg_metric_a.greptime_value / outer_agg_metric_b.greptime_value) | ++---------------------+----------------------------------------------------------------------------+ +| 1970-01-01T00:00:00 | 10.0 | ++---------------------+----------------------------------------------------------------------------+ + +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (0, 0, '5s') sum by(device) (outer_agg_metric_a / on(host) outer_agg_metric_b); + ++---------------------+----------------------------------------------------------------------------+ +| ts | sum(outer_agg_metric_a.greptime_value / outer_agg_metric_b.greptime_value) | ++---------------------+----------------------------------------------------------------------------+ +| 1970-01-01T00:00:00 | 10.0 | ++---------------------+----------------------------------------------------------------------------+ + +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (0, 0, '5s') sum by(host) (outer_agg_metric_a / on(host) outer_agg_metric_b); + ++------+---------------------+----------------------------------------------------------------------------+ +| host | ts | sum(outer_agg_metric_a.greptime_value / outer_agg_metric_b.greptime_value) | ++------+---------------------+----------------------------------------------------------------------------+ +| h1 | 1970-01-01T00:00:00 | 5.0 | +| h2 | 1970-01-01T00:00:00 | 5.0 | ++------+---------------------+----------------------------------------------------------------------------+ + +DROP TABLE outer_agg_metric_a; + +Affected Rows: 0 + +DROP TABLE outer_agg_metric_b; + +Affected Rows: 0 + +DROP TABLE outer_agg_physical; + +Affected Rows: 0 + diff --git a/tests/cases/standalone/common/promql/vector_matching.sql b/tests/cases/standalone/common/promql/vector_matching.sql new file mode 100644 index 00000000000..1fd0091c901 --- /dev/null +++ b/tests/cases/standalone/common/promql/vector_matching.sql @@ -0,0 +1,221 @@ +CREATE TABLE counter_metric ( + host STRING NULL, + device STRING NULL, + ts TIMESTAMP(3) TIME INDEX, + greptime_value DOUBLE, + PRIMARY KEY(host, device) +); + +CREATE TABLE gauge_metric ( + host STRING NULL, + device STRING NULL, + ts TIMESTAMP(3) TIME INDEX, + greptime_value DOUBLE, + PRIMARY KEY(host, device) +); + +INSERT INTO counter_metric VALUES + ('host1', 'eth0', 0, 10), ('host1', 'eth1', 0, 20), ('host2', 'eth0', 0, 30); + +INSERT INTO gauge_metric VALUES + ('host1', 'eth0', 0, 2), ('host1', 'eth1', 0, 4), ('host2', 'eth0', 0, 5); + +-- Matching on the whole tag set keeps it. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (0, 0, '5s') counter_metric / gauge_metric; + +-- `on(...)` reduces the result to the matching labels. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (0, 0, '5s') counter_metric{device="eth0"} / on(host) gauge_metric{device="eth0"}; + +-- `ignoring(...)` drops the ignored labels, here from operands that don't even share them. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (0, 0, '5s') counter_metric{device="eth0"} / ignoring(device) gauge_metric{device="eth1"}; + +-- `on()` matches everything against everything, so both sides must hold a single series. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (0, 0, '5s') counter_metric{host="host1",device="eth0"} / on() gauge_metric{host="host1",device="eth0"}; + +-- A filtering comparison keeps the left sample values and the derived labels. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (0, 0, '5s') counter_metric{device="eth0"} > on(host) gauge_metric{device="eth0"}; + +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (0, 0, '5s') counter_metric{device="eth0"} > bool on(host) gauge_metric{device="eth0"}; + +-- `host1` carries two devices on both sides: the match group repeats on the one side. +TQL EVAL (0, 0, '5s') counter_metric / on(host) gauge_metric; + +-- The one side is unique here, the many side is not, and no group modifier makes it explicit. +TQL EVAL (0, 0, '5s') counter_metric / on(host) gauge_metric{device="eth0"}; + +-- `group_left` takes the result labels from the many side. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (0, 0, '5s') counter_metric / on(host) group_left gauge_metric{device="eth0"}; + +-- `group_right` swaps the sides. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (0, 0, '5s') counter_metric{device="eth0"} / on(host) group_right gauge_metric; + +-- `group_right` makes the left operand the one side, which `host1` breaks. +TQL EVAL (0, 0, '5s') counter_metric / on(host) group_right gauge_metric{device="eth0"}; + +-- `group_left(device)` copies `device` from the one side, so both `host1` rows end up with the +-- same labels. +TQL EVAL (0, 0, '5s') counter_metric / on(host) group_left(device) gauge_metric{device="eth0"}; + +-- An included label the one side doesn't carry is dropped from the result. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (0, 0, '5s') counter_metric{device="eth0"} / on(host) group_left(device) sum by(host)(gauge_metric); + +DROP TABLE counter_metric; + +DROP TABLE gauge_metric; + +-- A label can carry the name the cardinality check generates for its row count. +CREATE TABLE collide_left ( + host STRING NULL, + __promql_match_group_count STRING NULL, + ts TIMESTAMP(3) TIME INDEX, + greptime_value DOUBLE, + PRIMARY KEY(host, __promql_match_group_count) +); + +CREATE TABLE collide_right ( + host STRING NULL, + __promql_match_group_count STRING NULL, + ts TIMESTAMP(3) TIME INDEX, + greptime_value DOUBLE, + PRIMARY KEY(host, __promql_match_group_count) +); + +INSERT INTO collide_left VALUES ('host1', 'a', 0, 10); + +INSERT INTO collide_right VALUES ('host1', 'b', 0, 2); + +TQL EVAL (0, 0, '5s') collide_left / on(host) collide_right; + +DROP TABLE collide_left; + +DROP TABLE collide_right; + +-- Partitioning on a column outside the match keys puts one match group in several regions. The +-- cardinality check counts the merged data, so it must still see the duplicate. +CREATE TABLE spread_left ( + host STRING NULL, + device STRING NULL, + ts TIMESTAMP(3) TIME INDEX, + greptime_value DOUBLE, + PRIMARY KEY(host, device) +) +PARTITION ON COLUMNS (device) ( + device < 'eth1', + device >= 'eth1' +); + +CREATE TABLE spread_right ( + host STRING NULL, + ts TIMESTAMP(3) TIME INDEX, + greptime_value DOUBLE, + PRIMARY KEY(host) +); + +INSERT INTO spread_left VALUES ('host1', 'eth0', 0, 10), ('host1', 'eth1', 0, 20); + +INSERT INTO spread_right VALUES ('host1', 0, 2); + +TQL EVAL (0, 0, '5s') spread_left / on(host) spread_right; + +TQL EVAL (0, 0, '5s') spread_right / on(host) group_left spread_left; + +DROP TABLE spread_left; + +DROP TABLE spread_right; + +-- An outer aggregate groups on the labels the operand has left, not on the ones the scan +-- underneath could still produce. +CREATE TABLE outer_agg_a ( + host STRING NULL, + device STRING NULL, + ts TIMESTAMP(3) TIME INDEX, + greptime_value DOUBLE, + PRIMARY KEY(host, device) +); + +CREATE TABLE outer_agg_b ( + host STRING NULL, + ts TIMESTAMP(3) TIME INDEX, + greptime_value DOUBLE, + PRIMARY KEY(host) +); + +INSERT INTO outer_agg_a VALUES ('h1', 'd1', 0, 10), ('h2', 'd2', 0, 20); + +INSERT INTO outer_agg_b VALUES ('h1', 0, 2), ('h2', 0, 4); + +-- `on(host)` leaves only `host`, so `without(host)` groups everything together. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (0, 0, '5s') sum without(host) (outer_agg_a / on(host) outer_agg_b); + +-- `device` is not a label of the operand any more, so it groups everything together too. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (0, 0, '5s') sum by(device) (outer_agg_a / on(host) outer_agg_b); + +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (0, 0, '5s') sum by(host) (outer_agg_a / on(host) outer_agg_b); + +DROP TABLE outer_agg_a; + +DROP TABLE outer_agg_b; + +-- Same over the metric engine, whose operands carry `__tsid`. +CREATE TABLE outer_agg_physical ( + ts TIMESTAMP(3) TIME INDEX, + greptime_value DOUBLE, +) ENGINE = metric WITH ("physical_metric_table" = ""); + +CREATE TABLE outer_agg_metric_a ( + host STRING NULL, + device STRING NULL, + ts TIMESTAMP(3) NOT NULL, + greptime_value DOUBLE NULL, + TIME INDEX (ts), + PRIMARY KEY(host, device), +) +ENGINE = metric +WITH( + on_physical_table = 'outer_agg_physical' +); + +CREATE TABLE outer_agg_metric_b ( + host STRING NULL, + ts TIMESTAMP(3) NOT NULL, + greptime_value DOUBLE NULL, + TIME INDEX (ts), + PRIMARY KEY(host), +) +ENGINE = metric +WITH( + on_physical_table = 'outer_agg_physical' +); + +INSERT INTO outer_agg_metric_a (host, device, ts, greptime_value) VALUES + ('h1', 'd1', 0, 10), ('h2', 'd2', 0, 20); + +INSERT INTO outer_agg_metric_b (host, ts, greptime_value) VALUES + ('h1', 0, 2), ('h2', 0, 4); + +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (0, 0, '5s') sum without(host) (outer_agg_metric_a / on(host) outer_agg_metric_b); + +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (0, 0, '5s') sum by(device) (outer_agg_metric_a / on(host) outer_agg_metric_b); + +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (0, 0, '5s') sum by(host) (outer_agg_metric_a / on(host) outer_agg_metric_b); + +DROP TABLE outer_agg_metric_a; + +DROP TABLE outer_agg_metric_b; + +DROP TABLE outer_agg_physical; diff --git a/tests/cases/standalone/common/tql/matching_filter.result b/tests/cases/standalone/common/tql/matching_filter.result index 8f5bebbe8c3..c0df22a738e 100644 --- a/tests/cases/standalone/common/tql/matching_filter.result +++ b/tests/cases/standalone/common/tql/matching_filter.result @@ -179,24 +179,11 @@ tql eval(0, 5, '5s') counter_metric offset 5s / on(host, device) gauge_metric{ho | host1 | eth1 | 1970-01-01T00:00:05 | 10.0 | +-------+--------+---------------------+---------------------------------------+ --- Matching on a subset of the labels. Prometheus rejects this as many-to-many; GreptimeDB --- returns the cross product instead (#9209). Recorded here to show the rewrite does not --- change it, not to endorse it. +-- Matching on a subset of the labels: both operands hold several devices per host. -- SQLNESS SORT_RESULT 3 1 tql eval(0, 5, '5s') counter_metric / on(host) gauge_metric{host="host1"}; -+-------+--------+---------------------+---------------------------------------+ -| host | device | ts | counter_metric.val / gauge_metric.val | -+-------+--------+---------------------+---------------------------------------+ -| host1 | eth0 | 1970-01-01T00:00:00 | 10.0 | -| host1 | eth0 | 1970-01-01T00:00:00 | 5.0 | -| host1 | eth0 | 1970-01-01T00:00:05 | 10.0 | -| host1 | eth0 | 1970-01-01T00:00:05 | 5.0 | -| host1 | eth1 | 1970-01-01T00:00:00 | 10.0 | -| host1 | eth1 | 1970-01-01T00:00:00 | 5.0 | -| host1 | eth1 | 1970-01-01T00:00:05 | 10.0 | -| host1 | eth1 | 1970-01-01T00:00:05 | 5.0 | -+-------+--------+---------------------+---------------------------------------+ +Error: 3001(EngineExecuteQuery), Execution error: found duplicate series for the match group {host="host1"} on the right hand-side of the operation; many-to-many matching not allowed: matching labels must be unique on one side -- `ignoring(...)`: every label that is not ignored is a matching label. -- SQLNESS SORT_RESULT 3 1 @@ -214,48 +201,19 @@ tql eval(0, 5, '5s') counter_metric / ignoring(missing_label) gauge_metric{host= -- SQLNESS SORT_RESULT 3 1 tql eval(0, 5, '5s') counter_metric / ignoring(device) gauge_metric{host="host1"}; -+-------+--------+---------------------+---------------------------------------+ -| host | device | ts | counter_metric.val / gauge_metric.val | -+-------+--------+---------------------+---------------------------------------+ -| host1 | eth0 | 1970-01-01T00:00:00 | 10.0 | -| host1 | eth0 | 1970-01-01T00:00:00 | 5.0 | -| host1 | eth0 | 1970-01-01T00:00:05 | 10.0 | -| host1 | eth0 | 1970-01-01T00:00:05 | 5.0 | -| host1 | eth1 | 1970-01-01T00:00:00 | 10.0 | -| host1 | eth1 | 1970-01-01T00:00:00 | 5.0 | -| host1 | eth1 | 1970-01-01T00:00:05 | 10.0 | -| host1 | eth1 | 1970-01-01T00:00:05 | 5.0 | -+-------+--------+---------------------+---------------------------------------+ +Error: 3001(EngineExecuteQuery), Execution error: found duplicate series for the match group {host="host1"} on the right hand-side of the operation; many-to-many matching not allowed: matching labels must be unique on one side -- An ignored label is not a matching label: the left operand keeps its `eth1` series. -- SQLNESS SORT_RESULT 3 1 tql eval(0, 5, '5s') counter_metric / ignoring(device) gauge_metric{device="eth0"}; -+-------+--------+---------------------+---------------------------------------+ -| host | device | ts | counter_metric.val / gauge_metric.val | -+-------+--------+---------------------+---------------------------------------+ -| host1 | eth0 | 1970-01-01T00:00:00 | 10.0 | -| host1 | eth0 | 1970-01-01T00:00:00 | 5.0 | -| host1 | eth0 | 1970-01-01T00:00:05 | 10.0 | -| host1 | eth0 | 1970-01-01T00:00:05 | 5.0 | -| host2 | eth0 | 1970-01-01T00:00:00 | 20.0 | -| host2 | eth0 | 1970-01-01T00:00:05 | 20.0 | -+-------+--------+---------------------+---------------------------------------+ +Error: 3001(EngineExecuteQuery), Execution error: multiple matches for labels {host="host1"}: many-to-one matching must be explicit (group_left/group_right) -- Same for a label left out of `on(...)`. -- SQLNESS SORT_RESULT 3 1 tql eval(0, 5, '5s') counter_metric / on(host) gauge_metric{device="eth0"}; -+-------+--------+---------------------+---------------------------------------+ -| host | device | ts | counter_metric.val / gauge_metric.val | -+-------+--------+---------------------+---------------------------------------+ -| host1 | eth0 | 1970-01-01T00:00:00 | 10.0 | -| host1 | eth0 | 1970-01-01T00:00:00 | 5.0 | -| host1 | eth0 | 1970-01-01T00:00:05 | 10.0 | -| host1 | eth0 | 1970-01-01T00:00:05 | 5.0 | -| host2 | eth0 | 1970-01-01T00:00:00 | 20.0 | -| host2 | eth0 | 1970-01-01T00:00:05 | 20.0 | -+-------+--------+---------------------+---------------------------------------+ +Error: 3001(EngineExecuteQuery), Execution error: multiple matches for labels {host="host1"}: many-to-one matching must be explicit (group_left/group_right) -- Matcher kinds other than equality are enforced by the join just the same. -- SQLNESS SORT_RESULT 3 1 @@ -321,16 +279,7 @@ tql eval(0, 5, '5s') sum by(host)(counter_metric{device="eth0"}) / on(host) sum -- SQLNESS SORT_RESULT 3 1 tql eval(0, 5, '5s') counter_metric / on(host) sum by(host)(gauge_metric{device="eth0"}); -+-------+---------------------+---------------------------------------------------------+ -| host | ts | counter_metric.val / gauge_metric.sum(gauge_metric.val) | -+-------+---------------------+---------------------------------------------------------+ -| host1 | 1970-01-01T00:00:00 | 10.0 | -| host1 | 1970-01-01T00:00:00 | 5.0 | -| host1 | 1970-01-01T00:00:05 | 10.0 | -| host1 | 1970-01-01T00:00:05 | 5.0 | -| host2 | 1970-01-01T00:00:00 | 20.0 | -| host2 | 1970-01-01T00:00:05 | 20.0 | -+-------+---------------------+---------------------------------------------------------+ +Error: 3001(EngineExecuteQuery), Execution error: multiple matches for labels {host="host1"}: many-to-one matching must be explicit (group_left/group_right) -- A grouping label that is a value field, not a tag: its value varies between samples of one -- series, so a matcher on it must not filter the scan before sample selection (#9242). @@ -411,14 +360,14 @@ tql eval(0, 5, '5s') counter_metric > on(host, device) gauge_metric{host="host1" -- SQLNESS SORT_RESULT 3 1 tql eval(0, 5, '5s') counter_metric / on(host) group_left sum by(host)(gauge_metric{host="host1"}); -+-------+---------------------+---------------------------------------------------------+ -| host | ts | counter_metric.val / gauge_metric.sum(gauge_metric.val) | -+-------+---------------------+---------------------------------------------------------+ -| host1 | 1970-01-01T00:00:00 | 2.5 | -| host1 | 1970-01-01T00:00:00 | 5.0 | -| host1 | 1970-01-01T00:00:05 | 2.5 | -| host1 | 1970-01-01T00:00:05 | 5.0 | -+-------+---------------------+---------------------------------------------------------+ ++-------+--------+---------------------+---------------------------------------------------------+ +| host | device | ts | counter_metric.val / gauge_metric.sum(gauge_metric.val) | ++-------+--------+---------------------+---------------------------------------------------------+ +| host1 | eth0 | 1970-01-01T00:00:00 | 2.5 | +| host1 | eth0 | 1970-01-01T00:00:05 | 2.5 | +| host1 | eth1 | 1970-01-01T00:00:00 | 5.0 | +| host1 | eth1 | 1970-01-01T00:00:05 | 5.0 | ++-------+--------+---------------------+---------------------------------------------------------+ -- `absent()` reports the inner selector's equality matchers as labels, so a copied matcher -- must not reach the context an enclosing expression reads. diff --git a/tests/cases/standalone/common/tql/matching_filter.sql b/tests/cases/standalone/common/tql/matching_filter.sql index 72da3da70a0..1b8492c4351 100644 --- a/tests/cases/standalone/common/tql/matching_filter.sql +++ b/tests/cases/standalone/common/tql/matching_filter.sql @@ -73,9 +73,7 @@ tql eval(0, 5, '5s') count_over_time(counter_metric[10s]) / on(host, device) cou -- SQLNESS SORT_RESULT 3 1 tql eval(0, 5, '5s') counter_metric offset 5s / on(host, device) gauge_metric{host="host1"}; --- Matching on a subset of the labels. Prometheus rejects this as many-to-many; GreptimeDB --- returns the cross product instead (#9209). Recorded here to show the rewrite does not --- change it, not to endorse it. +-- Matching on a subset of the labels: both operands hold several devices per host. -- SQLNESS SORT_RESULT 3 1 tql eval(0, 5, '5s') counter_metric / on(host) gauge_metric{host="host1"};