refactor(promql): split planner into focused modules (#9354)

* refactor(promql): move planner tests into separate module

Signed-off-by: discord9 <55937128+discord9@users.noreply.github.com>

* refactor(promql): move function-specific planner methods into module

Signed-off-by: discord9 <55937128+discord9@users.noreply.github.com>

* refactor(promql): move OR planner method into module

Signed-off-by: discord9 <55937128+discord9@users.noreply.github.com>

* refactor(promql): move binary island planner into module

Signed-off-by: discord9 <55937128+discord9@users.noreply.github.com>

* refactor(promql): keep binary result labels in planner root

Signed-off-by: discord9 <55937128+discord9@users.noreply.github.com>

* refactor(promql): limit planner helper visibility to parent module

Signed-off-by: discord9 <55937128+discord9@users.noreply.github.com>

* refactor(promql): group set operator planning and localize imports

Signed-off-by: discord9 <55937128+discord9@users.noreply.github.com>

* fix(promql): use crate-rooted planner imports

Signed-off-by: discord9 <55937128+discord9@users.noreply.github.com>

---------

Signed-off-by: discord9 <55937128+discord9@users.noreply.github.com>
This commit is contained in:
discord9
2026-09-28 05:30:33 +00:00
committed by GitHub
parent 825402dc70
commit f956f4da30
5 changed files with 9451 additions and 9390 deletions
File diff suppressed because it is too large Load Diff
@@ -0,0 +1,473 @@
// 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.
//! Dedicated helper plans for PromQL functions that cannot be expressed as a
//! plain scalar function over the input plan: the histogram helpers, `vector()`,
//! `scalar()`, and `absent()`.
use std::sync::Arc;
use common_query::prelude::greptime_value;
use datafusion::functions_aggregate::expr_fn::first_value;
use datafusion::logical_expr::expr::ScalarFunction;
use datafusion::logical_expr::{Extension, LogicalPlan, LogicalPlanBuilder};
use datafusion::optimizer::simplify_expressions::ExprSimplifier;
use datafusion::prelude::{Column, Expr as DfExpr};
use datafusion::scalar::ScalarValue;
use datafusion_common::DFSchema;
use datafusion_expr::expr_fn::when;
use datafusion_expr::simplify::SimplifyContext;
use datafusion_expr::{col, lit};
use datafusion_functions::core::coalesce;
use datatypes::arrow::datatypes::DataType as ArrowDataType;
use promql::extension_plan::{
Absent, EmptyMetric, HistogramFold, HistogramFoldOperation, ScalarCalculate,
};
use promql::functions::{NativeHistogramDrop, NativeHistogramFraction, NativeHistogramQuantile};
use promql_parser::label::MatchOp;
use promql_parser::parser::{Expr as PromExpr, FunctionArgs as PromFunctionArgs};
use snafu::{OptionExt, ResultExt, ensure};
use crate::promql::error::{
DataFusionPlanningSnafu, FunctionInvalidArgumentSnafu, MultiFieldsNotSupportedSnafu,
PromqlPlanNodeSnafu, Result, TimeIndexNotFoundSnafu, ValueNotFoundSnafu,
};
use crate::promql::planner::{
LE_COLUMN_NAME, PromPlanner, SCALAR_FUNCTION, SPECIAL_ABSENT_FUNCTION,
SPECIAL_HISTOGRAM_FRACTION, SPECIAL_HISTOGRAM_QUANTILE, SPECIAL_TIME_FUNCTION,
SPECIAL_VECTOR_FUNCTION,
};
use crate::query_engine::QueryEngineState;
impl PromPlanner {
/// Create a classic, native, or mixed histogram helper plan.
pub(super) async fn create_histogram_plan(
&mut self,
function_name: &str,
args: &PromFunctionArgs,
query_engine_state: &QueryEngineState,
) -> Result<LogicalPlan> {
let float_literal = |param: &PromExpr| -> Result<f64> {
let value = (|| {
let expr = Self::get_param_as_literal_expr(
Some(param),
None,
Some(ArrowDataType::Float64),
)
.ok()?;
let simplifier = ExprSimplifier::new(SimplifyContext::default());
let expr = simplifier.coerce(expr, &DFSchema::empty()).ok()?;
let DfExpr::Literal(value, _) = simplifier.simplify(expr).ok()? else {
return None;
};
let ScalarValue::Float64(Some(value)) =
value.cast_to(&ArrowDataType::Float64).ok()?
else {
return None;
};
Some(value)
})()
.with_context(|| FunctionInvalidArgumentSnafu {
fn_name: function_name.to_string(),
})?;
Ok(value)
};
let (function, input) = match (function_name, args.args.as_slice()) {
(SPECIAL_HISTOGRAM_QUANTILE, [quantile, input]) => (
HistogramFoldOperation::Quantile(float_literal(quantile)?.into()),
input.as_ref().clone(),
),
(SPECIAL_HISTOGRAM_FRACTION, [lower, upper, input]) => (
HistogramFoldOperation::Fraction {
lower: float_literal(lower)?.into(),
upper: float_literal(upper)?.into(),
},
input.as_ref().clone(),
),
_ => {
return FunctionInvalidArgumentSnafu {
fn_name: function_name.to_string(),
}
.fail();
}
};
let input_plan = self.prom_expr_to_plan(&input, query_engine_state).await?;
// Histogram helpers fold buckets across `le`, so `__tsid` (which includes `le`) is not a
// stable series identifier anymore. HistogramFold must not treat it as a label column.
let input_plan = self.strip_tsid_column(input_plan)?;
self.ctx.use_tsid = false;
if let Some((float_field, histogram_field)) =
Self::alternative_sample_columns(input_plan.schema(), &self.ctx.field_columns)
.map(|(float, histogram)| (float.to_string(), histogram.to_string()))
{
if self.ctx.has_le_tag() {
return self.create_mixed_histogram_plan(
function,
input_plan,
float_field,
histogram_field,
);
}
self.ctx.field_columns = vec![histogram_field];
}
if self.all_field_columns_are_native_histograms(input_plan.schema()) {
return self.create_native_histogram_plan(function, input_plan);
}
if !self.ctx.has_le_tag() {
// Return empty result instead of error when 'le' column is not found
// This handles the case when histogram metrics don't exist
return Ok(LogicalPlan::EmptyRelation(
datafusion::logical_expr::EmptyRelation {
produce_one_row: false,
schema: input_plan.schema().clone(),
},
));
}
let time_index_column =
self.ctx
.time_index_column
.clone()
.with_context(|| TimeIndexNotFoundSnafu {
table: self.ctx.table_name.clone().unwrap_or_default(),
})?;
// FIXME(ruihang): support multi fields
let field_column = self
.ctx
.field_columns
.first()
.with_context(|| FunctionInvalidArgumentSnafu {
fn_name: function.function_name().to_string(),
})?
.clone();
// remove le column from tag columns
self.ctx.tag_columns.retain(|col| col != LE_COLUMN_NAME);
let fold = HistogramFold::new_with_operation(
LE_COLUMN_NAME.to_string(),
field_column,
time_index_column,
function,
None,
input_plan,
)
.context(DataFusionPlanningSnafu)?;
Ok(LogicalPlan::Extension(Extension {
node: Arc::new(fold),
}))
}
fn create_native_histogram_expr(
&self,
function: HistogramFoldOperation,
field_column: &str,
) -> DfExpr {
let field = DfExpr::Column(Column::from_name(field_column));
let (func, args) = match function {
HistogramFoldOperation::Quantile(quantile) => (
Arc::new(NativeHistogramQuantile::scalar_udf_with_collector(
self.promql_annotations.clone(),
)),
vec![field, lit(f64::from(quantile))],
),
HistogramFoldOperation::Fraction { lower, upper } => (
Arc::new(NativeHistogramFraction::scalar_udf_with_collector(
self.promql_annotations.clone(),
)),
vec![field, lit(f64::from(lower)), lit(f64::from(upper))],
),
};
DfExpr::ScalarFunction(ScalarFunction { func, args })
}
fn create_native_histogram_plan(
&mut self,
function: HistogramFoldOperation,
input_plan: LogicalPlan,
) -> Result<LogicalPlan> {
ensure!(
self.ctx.field_columns.len() == 1,
MultiFieldsNotSupportedSnafu {
operator: function.function_name()
},
);
let field_column = self.ctx.field_columns[0].clone();
let function_expr = self.create_native_histogram_expr(function, &field_column);
let display_name = function_expr.schema_name().to_string();
self.ctx.field_columns = vec![display_name.clone()];
let project_exprs = std::iter::once(self.create_time_index_column_expr()?)
.chain(std::iter::once(function_expr.alias(display_name)))
.chain(self.create_tag_column_exprs()?)
.collect::<Vec<_>>();
LogicalPlanBuilder::from(input_plan)
.project(project_exprs)
.context(DataFusionPlanningSnafu)?
.filter(self.create_empty_values_filter_expr(false)?)
.context(DataFusionPlanningSnafu)?
.build()
.context(DataFusionPlanningSnafu)
}
fn create_mixed_histogram_plan(
&mut self,
function: HistogramFoldOperation,
input_plan: LogicalPlan,
float_field: String,
histogram_field: String,
) -> Result<LogicalPlan> {
let time_index_column =
self.ctx
.time_index_column
.clone()
.with_context(|| TimeIndexNotFoundSnafu {
table: self.ctx.table_name.clone().unwrap_or_default(),
})?;
let tag_columns = self.ctx.tag_columns.clone();
let folded = HistogramFold::new_with_operation(
LE_COLUMN_NAME.to_string(),
float_field.clone(),
time_index_column.clone(),
function,
Some(histogram_field.clone()),
input_plan,
)
.context(DataFusionPlanningSnafu)?;
let record_collision = DfExpr::ScalarFunction(ScalarFunction {
func: Arc::new(NativeHistogramDrop::warning_bool_false_udf(
"vector contains a mix of classic and native histograms".to_string(),
self.promql_annotations.clone(),
)),
args: vec![col(&float_field), col(&histogram_field)],
});
let keep = when(
col(&float_field)
.is_not_null()
.and(col(&histogram_field).is_not_null()),
record_collision,
)
.otherwise(lit(true))
.context(DataFusionPlanningSnafu)?;
let native_expr = self.create_native_histogram_expr(function, &histogram_field);
let output_field = native_expr.schema_name().to_string();
let value = DfExpr::ScalarFunction(ScalarFunction {
func: coalesce(),
args: vec![col(&float_field), native_expr],
});
self.ctx.field_columns = vec![output_field.clone()];
LogicalPlanBuilder::from(LogicalPlan::Extension(Extension {
node: Arc::new(folded),
}))
.filter(keep)
.context(DataFusionPlanningSnafu)?
.project(
std::iter::once(col(&time_index_column))
.chain(std::iter::once(value.alias(output_field)))
.chain(tag_columns.iter().map(col)),
)
.context(DataFusionPlanningSnafu)?
.build()
.context(DataFusionPlanningSnafu)
}
/// Create a [SPECIAL_VECTOR_FUNCTION] plan
pub(super) async fn create_vector_plan(
&mut self,
args: &PromFunctionArgs,
) -> Result<LogicalPlan> {
if args.args.len() != 1 {
return FunctionInvalidArgumentSnafu {
fn_name: SPECIAL_VECTOR_FUNCTION.to_string(),
}
.fail();
}
let lit = Self::get_param_as_literal_expr(Some(args.args[0].as_ref()), None, None)?;
// reuse `SPECIAL_TIME_FUNCTION` as name of time index column
self.ctx.time_index_column = Some(SPECIAL_TIME_FUNCTION.to_string());
self.ctx.reset_table_name_and_schema();
self.ctx.tag_columns = vec![];
self.ctx.aggregation_field_labels.clear();
self.ctx.field_columns = vec![greptime_value().to_string()];
Ok(LogicalPlan::Extension(Extension {
node: Arc::new(
EmptyMetric::new(
self.ctx.start,
self.ctx.end,
self.ctx.interval,
SPECIAL_TIME_FUNCTION.to_string(),
greptime_value().to_string(),
Some(lit),
)
.context(DataFusionPlanningSnafu)?,
),
}))
}
/// Create a [SCALAR_FUNCTION] plan
pub(super) async fn create_scalar_plan(
&mut self,
args: &PromFunctionArgs,
query_engine_state: &QueryEngineState,
) -> Result<LogicalPlan> {
ensure!(
args.len() == 1,
FunctionInvalidArgumentSnafu {
fn_name: SCALAR_FUNCTION
}
);
let input = self
.prom_expr_to_plan(&args.args[0], query_engine_state)
.await?;
let input_schema = input.schema().clone();
let alternative_samples =
Self::field_columns_are_alternative_samples(&input_schema, &self.ctx.field_columns);
let histogram_fields = self
.ctx
.field_columns
.iter()
.filter(|field| Self::field_column_is_native_histogram(&input_schema, field))
.count();
ensure!(
self.ctx.field_columns.len() == 1 || alternative_samples,
MultiFieldsNotSupportedSnafu {
operator: SCALAR_FUNCTION
},
);
let scalar_field = self
.ctx
.field_columns
.iter()
.find(|field| !Self::field_column_is_native_histogram(&input_schema, field))
.or_else(|| self.ctx.field_columns.first())
.cloned()
.with_context(|| FunctionInvalidArgumentSnafu {
fn_name: SCALAR_FUNCTION,
})?;
let input = if histogram_fields == self.ctx.field_columns.len() {
// scalar() ignores histogram samples. An empty input makes ScalarCalculate emit NaN
// for every evaluation timestamp without attempting a Struct-to-Float64 cast.
LogicalPlanBuilder::from(input)
.filter(lit(false))
.context(DataFusionPlanningSnafu)?
.build()
.context(DataFusionPlanningSnafu)?
} else if histogram_fields > 0 {
// A mixed vector contributes only its float samples to scalar().
LogicalPlanBuilder::from(input)
.filter(DfExpr::Column(Column::from_name(&scalar_field)).is_not_null())
.context(DataFusionPlanningSnafu)?
.build()
.context(DataFusionPlanningSnafu)?
} else {
input
};
let scalar_plan = LogicalPlan::Extension(Extension {
node: Arc::new(
ScalarCalculate::new(
self.ctx.start,
self.ctx.end,
self.ctx.interval,
input,
self.ctx.time_index_column.as_ref().unwrap(),
&self.ctx.tag_columns,
&scalar_field,
self.ctx.table_name.as_deref(),
)
.context(PromqlPlanNodeSnafu)?,
),
});
// scalar plan have no tag columns
self.ctx.tag_columns.clear();
self.ctx.aggregation_field_labels.clear();
self.ctx.field_columns.clear();
self.ctx
.field_columns
.push(scalar_plan.schema().field(1).name().clone());
Ok(scalar_plan)
}
/// Create a [SPECIAL_ABSENT_FUNCTION] plan
pub(super) async fn create_absent_plan(
&mut self,
args: &PromFunctionArgs,
query_engine_state: &QueryEngineState,
) -> Result<LogicalPlan> {
if args.args.len() != 1 {
return FunctionInvalidArgumentSnafu {
fn_name: SPECIAL_ABSENT_FUNCTION.to_string(),
}
.fail();
}
let input = self
.prom_expr_to_plan(&args.args[0], query_engine_state)
.await?;
let time_index_expr = self.create_time_index_column_expr()?;
let first_field_expr =
self.create_field_column_exprs()?
.pop()
.with_context(|| ValueNotFoundSnafu {
table: self.ctx.table_name.clone().unwrap_or_default(),
})?;
let first_value_expr = first_value(first_field_expr, vec![]);
let ordered_aggregated_input = LogicalPlanBuilder::from(input)
.aggregate(
vec![time_index_expr.clone()],
vec![first_value_expr.clone()],
)
.context(DataFusionPlanningSnafu)?
.sort(vec![time_index_expr.sort(true, false)])
.context(DataFusionPlanningSnafu)?
.build()
.context(DataFusionPlanningSnafu)?;
let fake_labels = self
.ctx
.selector_matcher
.iter()
.filter_map(|matcher| match matcher.op {
MatchOp::Equal => Some((matcher.name.clone(), matcher.value.clone())),
_ => None,
})
.collect::<Vec<_>>();
// Create the absent plan
let absent_plan = LogicalPlan::Extension(Extension {
node: Arc::new(
Absent::try_new(
self.ctx.start,
self.ctx.end,
self.ctx.interval,
self.ctx.time_index_column.as_ref().unwrap().clone(),
self.ctx.field_columns[0].clone(),
fake_labels,
ordered_aggregated_input,
)
.context(DataFusionPlanningSnafu)?,
),
});
// The absent series carries the equality matchers as labels, not the input's
// tags or value fields, so the input's field grouping labels no longer apply.
self.ctx.aggregation_field_labels.clear();
Ok(absent_plan)
}
}
+536
View File
@@ -0,0 +1,536 @@
// 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.
//! The binary-island fast path of the PromQL planner.
use std::collections::{BTreeSet, HashMap};
use common_query::prelude::OTLP_AGGREGATION_TEMPORALITY_LABEL;
use datafusion::common::DFSchemaRef;
use datafusion::logical_expr::expr::Alias;
use datafusion::logical_expr::{LogicalPlan, LogicalPlanBuilder};
use datafusion::prelude::{Column, Expr as DfExpr, JoinType};
use datafusion_common::{NullEquality, TableReference};
use datafusion_expr::lit;
use promql_parser::label::{METRIC_NAME, MatchOp, Matcher};
use promql_parser::parser::token::{self, TokenType};
use promql_parser::parser::{
BinaryExpr as PromBinaryExpr, Expr as PromExpr, Offset, ParenExpr, UnaryExpr,
VectorMatchCardinality, VectorSelector,
};
use snafu::ResultExt;
use crate::promql::error::{DataFusionPlanningSnafu, Result};
use crate::promql::planner::{PromPlanner, PromPlannerContext};
/// Prefix for generated binary island leaf aliases.
const BINARY_ISLAND_LEAF_ALIAS_PREFIX: &str = "__prom_v";
#[derive(Debug, Clone, PartialEq, Eq, Hash)]
struct VectorLeafKey {
metric_name: String,
matchers: Vec<(String, String, String)>,
or_matchers: Vec<Vec<(String, String, String)>>,
offset_ms: i128,
at: String,
}
#[derive(Debug, Clone)]
struct IslandLeaf {
selector: VectorSelector,
display_table: String,
}
#[derive(Debug, Clone)]
enum IslandExpr {
VectorLeaf(usize),
Scalar(DfExpr),
Unary {
input: Box<IslandExpr>,
},
Binary {
op: TokenType,
lhs: Box<IslandExpr>,
rhs: Box<IslandExpr>,
},
}
impl IslandExpr {
fn try_new(expr: &PromExpr, env: &mut IslandCollectEnv) -> Option<Self> {
if let Some(expr) = PromPlanner::try_build_literal_expr(expr) {
return Some(Self::Scalar(expr));
}
match expr {
PromExpr::Paren(ParenExpr { expr }) => Self::try_new(expr, env),
PromExpr::VectorSelector(selector) => {
let leaf = env.intern_leaf(selector)?;
Some(Self::VectorLeaf(leaf))
}
PromExpr::Unary(UnaryExpr { expr }) => {
let input = Self::try_new(expr, env)?;
Some(Self::Unary {
input: Box::new(input),
})
}
PromExpr::Binary(PromBinaryExpr {
lhs,
rhs,
op,
modifier,
}) if matches!(
op.id(),
token::T_ADD
| token::T_SUB
| token::T_MUL
| token::T_DIV
| token::T_MOD
| token::T_POW
| token::T_ATAN2
) && modifier.as_ref().is_none_or(|modifier| {
!modifier.return_bool
&& modifier.matching.is_none()
&& matches!(modifier.card, VectorMatchCardinality::OneToOne)
&& modifier.fill_values.lhs.is_none()
&& modifier.fill_values.rhs.is_none()
}) =>
{
let lhs = Self::try_new(lhs, env)?;
let rhs = Self::try_new(rhs, env)?;
Some(Self::Binary {
op: *op,
lhs: Box::new(lhs),
rhs: Box::new(rhs),
})
}
_ => None,
}
}
}
#[derive(Debug, Default)]
struct IslandCollectEnv {
leaf_by_key: HashMap<VectorLeafKey, usize>,
leaves: Vec<IslandLeaf>,
vector_occurrences: usize,
}
#[derive(Debug)]
struct PlannedIslandLeaf {
plan: LogicalPlan,
ctx: PromPlannerContext,
alias: TableReference,
display_table: String,
}
#[derive(Debug)]
struct IslandFieldExprs {
exprs: Vec<DfExpr>,
names: Vec<String>,
scalar: bool,
}
impl VectorLeafKey {
fn from_selector(selector: &VectorSelector) -> Option<Self> {
let mut metric_name = selector.name.clone();
let mut matchers = Vec::with_capacity(selector.matchers.matchers.len());
let matcher_key = |matcher: &Matcher| {
(
matcher.name.clone(),
matcher.op.to_string(),
matcher.value.clone(),
)
};
for matcher in &selector.matchers.matchers {
if matcher.name == METRIC_NAME {
if matcher.op != MatchOp::Equal || metric_name.is_some() {
return None;
}
metric_name = Some(matcher.value.clone());
} else {
matchers.push(matcher_key(matcher));
}
}
matchers.sort();
let mut or_matchers = selector
.matchers
.or_matchers
.iter()
.map(|group| {
let mut group = group.iter().map(matcher_key).collect::<Vec<_>>();
group.sort();
group
})
.collect::<Vec<_>>();
or_matchers.sort();
Some(Self {
metric_name: metric_name?,
matchers,
or_matchers,
offset_ms: match &selector.offset {
Some(Offset::Pos(duration)) => duration.as_millis() as i128,
Some(Offset::Neg(duration)) => -(duration.as_millis() as i128),
None => 0,
},
at: format!("{:?}", selector.at),
})
}
}
impl IslandCollectEnv {
fn intern_leaf(&mut self, selector: &VectorSelector) -> Option<usize> {
self.vector_occurrences += 1;
let key = VectorLeafKey::from_selector(selector)?;
if let Some(id) = self.leaf_by_key.get(&key) {
return Some(*id);
}
let id = self.leaves.len();
self.leaves.push(IslandLeaf {
selector: selector.clone(),
display_table: key.metric_name.clone(),
});
self.leaf_by_key.insert(key, id);
Some(id)
}
}
impl PromPlanner {
pub(super) async fn try_plan_binary_island(
&mut self,
binary_expr: &PromBinaryExpr,
) -> Result<Option<LogicalPlan>> {
let original_ctx = self.ctx.clone();
let mut collect_env = IslandCollectEnv::default();
let Some(island_expr) =
IslandExpr::try_new(&PromExpr::Binary(binary_expr.clone()), &mut collect_env)
else {
return Ok(None);
};
if collect_env.leaves.is_empty()
|| collect_env.vector_occurrences <= collect_env.leaves.len()
{
return Ok(None);
}
let mut planned_leaves = Vec::with_capacity(collect_env.leaves.len());
for (idx, leaf) in collect_env.leaves.iter().enumerate() {
let plan = self
.prom_vector_selector_to_plan(&leaf.selector, false)
.await?;
let ctx = self.ctx.clone();
let alias = TableReference::bare(format!("{BINARY_ISLAND_LEAF_ALIAS_PREFIX}{idx}"));
let plan = LogicalPlanBuilder::from(plan)
.alias(alias.clone())
.context(DataFusionPlanningSnafu)?
.build()
.context(DataFusionPlanningSnafu)?;
planned_leaves.push(PlannedIslandLeaf {
plan,
ctx,
alias,
display_table: leaf.display_table.clone(),
});
}
if planned_leaves.iter().any(|leaf| {
Self::field_columns_contain_native_histogram(
leaf.plan.schema(),
&leaf.ctx.field_columns,
)
}) {
self.ctx = original_ctx;
return Ok(None);
}
if !Self::binary_island_join_contexts_supported(&planned_leaves) {
self.ctx = original_ctx;
return Ok(None);
}
let mut input = planned_leaves[0].plan.clone();
for right_idx in 1..planned_leaves.len() {
input = self.join_binary_island_leaf(
input,
&planned_leaves[0],
&planned_leaves[right_idx],
)?;
}
let field_exprs =
Self::build_binary_island_field_exprs(&island_expr, &planned_leaves, input.schema())?;
if field_exprs.scalar || field_exprs.exprs.is_empty() {
self.ctx = original_ctx;
return Ok(None);
}
let plan = self.project_binary_island(
input,
&planned_leaves[0].alias,
&planned_leaves[0].ctx,
field_exprs,
)?;
Ok(Some(plan))
}
fn binary_island_join_contexts_supported(leaves: &[PlannedIslandLeaf]) -> bool {
if leaves
.iter()
.any(|leaf| leaf.ctx.time_index_column.is_none())
{
return false;
}
if leaves.len() <= 1 {
return true;
}
let first_tags = leaves[0].ctx.tag_columns.iter().collect::<BTreeSet<_>>();
leaves.iter().skip(1).all(|leaf| {
(Self::plan_has_tsid_column(&leaves[0].plan) && Self::plan_has_tsid_column(&leaf.plan))
|| leaf.ctx.tag_columns.iter().collect::<BTreeSet<_>>() == first_tags
})
}
fn join_binary_island_leaf(
&self,
left: LogicalPlan,
first_leaf: &PlannedIslandLeaf,
right_leaf: &PlannedIslandLeaf,
) -> Result<LogicalPlan> {
let only_join_time_index = (first_leaf.ctx.tag_columns.is_empty()
|| right_leaf.ctx.tag_columns.is_empty())
&& !first_leaf
.ctx
.tag_columns
.iter()
.chain(&right_leaf.ctx.tag_columns)
.any(|tag| tag == OTLP_AGGREGATION_TEMPORALITY_LABEL);
let (mut left_keys, mut right_keys, force_empty_join) = self.binary_join_key_columns(
left.schema(),
right_leaf.plan.schema(),
&first_leaf.ctx,
&right_leaf.ctx,
only_join_time_index,
&None,
)?;
if let (Some(left_time_index_column), Some(right_time_index_column)) = (
first_leaf.ctx.time_index_column.clone(),
right_leaf.ctx.time_index_column.clone(),
) {
left_keys.insert(left_time_index_column);
right_keys.insert(right_time_index_column);
}
LogicalPlanBuilder::from(left)
.join_detailed(
right_leaf.plan.clone(),
JoinType::Inner,
(
left_keys
.into_iter()
.map(|name| Column::new(Some(first_leaf.alias.clone()), name))
.collect::<Vec<_>>(),
right_keys
.into_iter()
.map(|name| Column::new(Some(right_leaf.alias.clone()), name))
.collect::<Vec<_>>(),
),
force_empty_join.then_some(lit(false)),
NullEquality::NullEqualsNull,
)
.context(DataFusionPlanningSnafu)?
.build()
.context(DataFusionPlanningSnafu)
}
fn build_binary_island_field_exprs(
expr: &IslandExpr,
leaves: &[PlannedIslandLeaf],
schema: &DFSchemaRef,
) -> Result<IslandFieldExprs> {
match expr {
IslandExpr::VectorLeaf(id) => {
let leaf = &leaves[*id];
let exprs = leaf
.ctx
.field_columns
.iter()
.map(|field| {
schema
.qualified_field_with_name(Some(&leaf.alias), field)
.context(DataFusionPlanningSnafu)
.map(|field| DfExpr::Column(field.into()))
})
.collect::<Result<Vec<_>>>()?;
let names = leaf
.ctx
.field_columns
.iter()
.map(|field| format!("{}.{}", leaf.display_table, field))
.collect();
Ok(IslandFieldExprs {
exprs,
names,
scalar: false,
})
}
IslandExpr::Scalar(expr) => Ok(IslandFieldExprs {
exprs: vec![expr.clone()],
names: vec![expr.schema_name().to_string()],
scalar: true,
}),
IslandExpr::Unary { input } => {
let input = Self::build_binary_island_field_exprs(input, leaves, schema)?;
let mut exprs = Vec::with_capacity(input.exprs.len());
let mut names = Vec::with_capacity(input.names.len());
for (expr, name) in input.exprs.into_iter().zip(input.names) {
exprs.push(DfExpr::Negative(Box::new(expr)));
names.push(format!("-{name}"));
}
Ok(IslandFieldExprs {
exprs,
names,
scalar: input.scalar,
})
}
IslandExpr::Binary { op, lhs, rhs } => {
let same_leaf = match (&**lhs, &**rhs) {
(IslandExpr::VectorLeaf(left), IslandExpr::VectorLeaf(right))
if left == right =>
{
Some(*left)
}
_ => None,
};
let lhs = Self::build_binary_island_field_exprs(lhs, leaves, schema)?;
let rhs = Self::build_binary_island_field_exprs(rhs, leaves, schema)?;
let expr_builder = Self::prom_token_to_binary_expr_builder(*op)?;
let scalar = lhs.scalar && rhs.scalar;
let op = op.to_string();
let (exprs, names) = match (lhs.scalar, rhs.scalar) {
(true, true) => {
let expr = expr_builder(lhs.exprs[0].clone(), rhs.exprs[0].clone())?;
let name = format!("{} {op} {}", lhs.names[0], rhs.names[0]);
(vec![expr], vec![name])
}
(true, false) => {
let mut exprs = Vec::with_capacity(rhs.exprs.len());
let mut names = Vec::with_capacity(rhs.names.len());
for (rhs_expr, rhs_name) in rhs.exprs.into_iter().zip(rhs.names) {
exprs.push(expr_builder(lhs.exprs[0].clone(), rhs_expr)?);
names.push(format!("{} {op} {rhs_name}", lhs.names[0]));
}
(exprs, names)
}
(false, true) => {
let mut exprs = Vec::with_capacity(lhs.exprs.len());
let mut names = Vec::with_capacity(lhs.names.len());
for (lhs_expr, lhs_name) in lhs.exprs.into_iter().zip(lhs.names) {
exprs.push(expr_builder(lhs_expr, rhs.exprs[0].clone())?);
names.push(format!("{lhs_name} {op} {}", rhs.names[0]));
}
(exprs, names)
}
(false, false) => {
let mut exprs = Vec::new();
let mut names = Vec::new();
for (idx, ((lhs_expr, rhs_expr), (mut lhs_name, mut rhs_name))) in lhs
.exprs
.into_iter()
.zip(rhs.exprs)
.zip(lhs.names.into_iter().zip(rhs.names))
.enumerate()
{
if let Some(leaf) = same_leaf {
let field = leaves[leaf]
.ctx
.field_columns
.get(idx)
.cloned()
.unwrap_or_else(|| lhs_name.clone());
lhs_name = format!("lhs.{field}");
rhs_name = format!("rhs.{field}");
}
exprs.push(expr_builder(lhs_expr, rhs_expr)?);
names.push(format!("{lhs_name} {op} {rhs_name}"));
}
(exprs, names)
}
};
Ok(IslandFieldExprs {
exprs,
names,
scalar,
})
}
}
}
fn project_binary_island(
&mut self,
input: LogicalPlan,
base_alias: &TableReference,
base_ctx: &PromPlannerContext,
field_exprs: IslandFieldExprs,
) -> Result<LogicalPlan> {
self.ctx = base_ctx.clone();
let schema = input.schema();
let non_field_exprs = base_ctx
.tag_columns
.iter()
.chain(base_ctx.time_index_column.iter())
.map(|column| {
schema
.qualified_field_with_name(Some(base_alias), column)
.context(DataFusionPlanningSnafu)
.map(|field| DfExpr::Column(field.into()))
});
let tsid_expr = Self::optional_tsid_projection(schema, Some(base_alias), base_ctx.use_tsid)
.into_iter()
.map(Ok);
self.ctx.field_columns = field_exprs.names;
let field_exprs = field_exprs
.exprs
.into_iter()
.zip(self.ctx.field_columns.iter())
.map(|(expr, name)| Ok(DfExpr::Alias(Alias::new(expr, None::<String>, name))));
let project_exprs = non_field_exprs
.chain(tsid_expr)
.chain(field_exprs)
.collect::<Result<Vec<_>>>()?;
let plan = LogicalPlanBuilder::from(input)
.project(project_exprs)
.context(DataFusionPlanningSnafu)?
.build()
.context(DataFusionPlanningSnafu)?;
self.ctx.table_name = None;
self.ctx.schema_name = None;
Ok(plan)
}
}
@@ -0,0 +1,832 @@
// 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.
//! Planning of the PromQL set operators (`and`, `unless`, and `or`).
use std::collections::{BTreeSet, HashMap, HashSet};
use std::sync::Arc;
use datafusion::logical_expr::{Cast, Extension, LogicalPlan, LogicalPlanBuilder};
use datafusion::prelude::{Column, Expr as DfExpr, JoinType};
use datafusion::scalar::ScalarValue;
use datafusion_common::{NullEquality, TableReference};
use datatypes::arrow::datatypes::DataType as ArrowDataType;
use promql::extension_plan::UnionDistinctOn;
use promql_parser::parser::token::{self, TokenType};
use promql_parser::parser::{BinModifier, LabelModifier, VectorMatchCardinality};
use snafu::{OptionExt, ResultExt, ensure};
use store_api::metric_engine_consts::DATA_SCHEMA_TSID_COLUMN_NAME;
use crate::promql::error::{
ColumnNotFoundSnafu, CombineTableColumnMismatchSnafu, DataFusionPlanningSnafu,
MultiFieldsNotSupportedSnafu, Result, TimeIndexNotFoundSnafu, UnexpectedPlanExprSnafu,
UnexpectedTokenSnafu, UnsupportedVectorMatchSnafu,
};
use crate::promql::planner::{
OR_FLOAT_FIELD_PREFIX, OR_HISTOGRAM_FIELD_PREFIX, PromPlanner, PromPlannerContext,
};
impl PromPlanner {
// TODO(ruihang): change function name
#[allow(clippy::too_many_arguments)]
pub(super) fn or_operator(
&mut self,
left: LogicalPlan,
right: LogicalPlan,
left_tag_cols_set: HashSet<String>,
right_tag_cols_set: HashSet<String>,
left_context: PromPlannerContext,
right_context: PromPlannerContext,
modifier: &Option<BinModifier>,
) -> Result<LogicalPlan> {
let left_is_empty = Self::is_zero_row_empty_relation(&left);
let right_is_empty = Self::is_zero_row_empty_relation(&right);
match (left_is_empty, right_is_empty) {
(true, false) => {
self.ctx = right_context;
return Ok(right);
}
(false, true) => {
self.ctx = left_context;
return Ok(left);
}
(true, true) => {
self.ctx = left_context;
return Ok(left);
}
(false, false) => {}
}
ensure!(
!left.schema().fields().is_empty() && !right.schema().fields().is_empty(),
UnexpectedPlanExprSnafu {
desc: "OR operator input has zero columns",
}
);
let left_has_alternative_samples =
Self::field_columns_are_alternative_samples(left.schema(), &left_context.field_columns);
let right_has_alternative_samples = Self::field_columns_are_alternative_samples(
right.schema(),
&right_context.field_columns,
);
ensure!(
left_context.field_columns.len() == 1 || left_has_alternative_samples,
MultiFieldsNotSupportedSnafu {
operator: "OR operator"
}
);
ensure!(
right_context.field_columns.len() == 1 || right_has_alternative_samples,
MultiFieldsNotSupportedSnafu {
operator: "OR operator"
}
);
// prepare hash sets
let all_tags = left_tag_cols_set
.union(&right_tag_cols_set)
.cloned()
.collect::<HashSet<_>>();
let left_qualifier = left.schema().qualified_field(0).0.cloned();
let right_qualifier = right.schema().qualified_field(0).0.cloned();
let left_qualifier_string = left_qualifier
.as_ref()
.map(|l| l.to_string())
.unwrap_or_default();
let right_qualifier_string = right_qualifier
.as_ref()
.map(|r| r.to_string())
.unwrap_or_default();
let left_time_index_column =
left_context
.time_index_column
.clone()
.with_context(|| TimeIndexNotFoundSnafu {
table: left_qualifier_string.clone(),
})?;
let right_time_index_column =
right_context
.time_index_column
.clone()
.with_context(|| TimeIndexNotFoundSnafu {
table: right_qualifier_string.clone(),
})?;
let native_histogram_type = Self::native_histogram_arrow_type();
let is_numeric = |data_type: &ArrowDataType| {
matches!(
data_type,
ArrowDataType::Int8
| ArrowDataType::Int16
| ArrowDataType::Int32
| ArrowDataType::Int64
| ArrowDataType::UInt8
| ArrowDataType::UInt16
| ArrowDataType::UInt32
| ArrowDataType::UInt64
| ArrowDataType::Float32
| ArrowDataType::Float64
)
};
let left_fields = left_context
.field_columns
.iter()
.map(|name| {
left.schema()
.iter()
.find(|(_, field)| field.name() == name)
.map(|(qualifier, field)| {
(name.clone(), qualifier.cloned(), field.data_type().clone())
})
.with_context(|| ColumnNotFoundSnafu { col: name.clone() })
})
.collect::<Result<Vec<_>>>()?;
let right_fields = right_context
.field_columns
.iter()
.map(|name| {
right
.schema()
.iter()
.find(|(_, field)| field.name() == name)
.map(|(qualifier, field)| {
(name.clone(), qualifier.cloned(), field.data_type().clone())
})
.with_context(|| ColumnNotFoundSnafu { col: name.clone() })
})
.collect::<Result<Vec<_>>>()?;
let left_field = &left_fields[0];
let right_field = &right_fields[0];
let left_field_col = &left_field.0;
let right_field_col = &right_field.0;
let fields_are_samples = |fields: &[(String, Option<TableReference>, ArrowDataType)]| {
fields.iter().all(|(_, _, data_type)| {
is_numeric(data_type) || data_type == &native_histogram_type
})
};
let mixed_sample_types = if left_has_alternative_samples || right_has_alternative_samples {
if !fields_are_samples(&left_fields) || !fields_are_samples(&right_fields) {
return UnexpectedPlanExprSnafu {
desc: format!(
"OR value fields have incompatible types: {:?} and {:?}",
left_fields
.iter()
.map(|(_, _, data_type)| data_type)
.collect::<Vec<_>>(),
right_fields
.iter()
.map(|(_, _, data_type)| data_type)
.collect::<Vec<_>>()
),
}
.fail();
}
true
} else {
(left_field.2 == native_histogram_type && is_numeric(&right_field.2))
|| (right_field.2 == native_histogram_type && is_numeric(&left_field.2))
};
let target_field_type = if mixed_sample_types {
// Mixed vectors use the existing response representation: one nullable float column
// and one nullable native-histogram column.
ArrowDataType::Float64
} else if left_field.2 == right_field.2 {
left_field.2.clone()
} else if is_numeric(&left_field.2) && is_numeric(&right_field.2) {
ArrowDataType::Float64
} else {
return UnexpectedPlanExprSnafu {
desc: format!(
"OR value fields have incompatible types: {:?} and {:?}",
left_field.2, right_field.2
),
}
.fail();
};
let (mixed_float_field_col, mixed_histogram_field_col) = if mixed_sample_types {
let mut reserved_names = left
.schema()
.fields()
.iter()
.chain(right.schema().fields().iter())
.map(|field| field.name().clone())
.collect::<HashSet<_>>();
for (name, _, _) in left_fields.iter().chain(&right_fields) {
reserved_names.remove(name);
}
reserved_names.extend(all_tags.iter().cloned());
let unique_name = |prefix: &str, reserved_names: &mut HashSet<String>| {
let mut index = 0;
loop {
let name = format!("{prefix}{index}");
index += 1;
if reserved_names.insert(name.clone()) {
break name;
}
}
};
let float_field = unique_name(OR_FLOAT_FIELD_PREFIX, &mut reserved_names);
let histogram_field = unique_name(OR_HISTOGRAM_FIELD_PREFIX, &mut reserved_names);
(float_field, histogram_field)
} else {
(left_field_col.clone(), String::new())
};
let left_tag_types = left_tag_cols_set
.iter()
.map(|label| {
left.schema()
.fields()
.iter()
.find(|field| field.name() == label)
.map(|field| (label.clone(), field.data_type().clone()))
.with_context(|| ColumnNotFoundSnafu { col: label.clone() })
})
.collect::<Result<HashMap<_, _>>>()?;
let right_tag_types = right_tag_cols_set
.iter()
.map(|label| {
right
.schema()
.fields()
.iter()
.find(|field| field.name() == label)
.map(|field| (label.clone(), field.data_type().clone()))
.with_context(|| ColumnNotFoundSnafu { col: label.clone() })
})
.collect::<Result<HashMap<_, _>>>()?;
let mut target_tag_types = HashMap::with_capacity(all_tags.len());
for label in &all_tags {
let Some(data_type) =
Self::common_label_data_type(left_tag_types.get(label), right_tag_types.get(label))
else {
return UnexpectedPlanExprSnafu {
desc: format!(
"OR label {label} has incompatible types: {:?} and {:?}",
left_tag_types.get(label),
right_tag_types.get(label)
),
}
.fail();
};
target_tag_types.insert(label.clone(), data_type);
}
let left_has_tsid = left
.schema()
.fields()
.iter()
.any(|field| field.name() == DATA_SCHEMA_TSID_COLUMN_NAME);
let right_has_tsid = right
.schema()
.fields()
.iter()
.any(|field| field.name() == DATA_SCHEMA_TSID_COLUMN_NAME);
// step 0: fill all columns in output schema
let mut all_columns_set = left
.schema()
.fields()
.iter()
.chain(right.schema().fields().iter())
.map(|field| field.name().clone())
.collect::<HashSet<_>>();
// Keep `__tsid` only when both sides contain it, otherwise it may break schema alignment
// (e.g. `unknown_metric or some_metric`).
if !(left_has_tsid && right_has_tsid) {
all_columns_set.remove(DATA_SCHEMA_TSID_COLUMN_NAME);
}
// remove time index column
all_columns_set.remove(&left_time_index_column);
all_columns_set.remove(&right_time_index_column);
if mixed_sample_types {
for (name, _, _) in left_fields.iter().chain(&right_fields) {
all_columns_set.remove(name);
}
all_columns_set.extend(all_tags.iter().cloned());
all_columns_set.insert(mixed_float_field_col.clone());
all_columns_set.insert(mixed_histogram_field_col.clone());
} else if left_field_col != right_field_col {
// remove field column in the right
all_columns_set.remove(right_field_col);
}
let mut all_columns = all_columns_set.into_iter().collect::<Vec<_>>();
// sort to ensure the generated schema is not volatile
all_columns.sort_unstable();
// use left time index column name as the result time index column name
all_columns.insert(0, left_time_index_column.clone());
let mut occupied_column_names = left
.schema()
.fields()
.iter()
.chain(right.schema().fields().iter())
.map(|field| field.name().clone())
.collect::<HashSet<_>>();
// step 1: align schema using project, fill non-exist columns with null
let aligned_label_expr = |col: &String, source_types: &HashMap<String, ArrowDataType>| {
let target_type = &target_tag_types[col];
if let Some(source_type) = source_types.get(col) {
let expr = DfExpr::Column(Column::new(None::<String>, col));
if source_type == target_type {
expr
} else {
DfExpr::Cast(Cast::new(Box::new(expr), target_type.clone())).alias(col.clone())
}
} else {
DfExpr::Literal(
Self::string_scalar_value(target_type, None)
.expect("target label type is a string"),
None,
)
.alias(col.clone())
}
};
let null_histogram =
ScalarValue::try_new_null(&native_histogram_type).context(DataFusionPlanningSnafu)?;
let mixed_value_expr = |fields: &[(String, Option<TableReference>, ArrowDataType)],
output_col: &String| {
if output_col == &mixed_float_field_col {
if let Some((name, qualifier, data_type)) = fields
.iter()
.find(|(_, _, data_type)| is_numeric(data_type))
{
let expr = DfExpr::Column(Column::new(qualifier.clone(), name));
if data_type == &ArrowDataType::Float64 {
expr.alias(output_col)
} else {
DfExpr::Cast(Cast::new(Box::new(expr), ArrowDataType::Float64))
.alias(output_col)
}
} else {
DfExpr::Literal(ScalarValue::Float64(None), None).alias(output_col)
}
} else {
fields
.iter()
.find(|(_, _, data_type)| data_type == &native_histogram_type)
.map(|(name, qualifier, _)| {
DfExpr::Column(Column::new(qualifier.clone(), name)).alias(output_col)
})
.unwrap_or_else(|| {
DfExpr::Literal(null_histogram.clone(), None).alias(output_col)
})
}
};
let left_proj_exprs = all_columns.iter().map(|col| {
if mixed_sample_types
&& (col == &mixed_float_field_col || col == &mixed_histogram_field_col)
{
mixed_value_expr(&left_fields, col)
} else if !mixed_sample_types
&& col == left_field_col
&& left_field.2 != target_field_type
{
DfExpr::Cast(Cast::new(
Box::new(DfExpr::Column(Column::new(
left_field.1.clone(),
left_field_col,
))),
target_field_type.clone(),
))
.alias(left_field_col.clone())
} else if target_tag_types.contains_key(col) {
aligned_label_expr(col, &left_tag_types)
} else {
DfExpr::Column(Column::new(None::<String>, col))
}
});
let right_time_index_expr = DfExpr::Column(Column::new(
right_qualifier.clone(),
right_time_index_column,
))
.alias(left_time_index_column.clone());
// The field column in right side may not have qualifier (it may be removed by join operation),
// so we need to find it from the schema.
// `skip(1)` to skip the time index column
let right_proj_exprs_without_time_index = all_columns.iter().skip(1).map(|col| {
// expr
if mixed_sample_types
&& (col == &mixed_float_field_col || col == &mixed_histogram_field_col)
{
mixed_value_expr(&right_fields, col)
} else if !mixed_sample_types && col == left_field_col {
let expr = DfExpr::Column(Column::new(right_field.1.clone(), right_field_col));
if right_field.2 != target_field_type {
DfExpr::Cast(Cast::new(Box::new(expr), target_field_type.clone()))
.alias(left_field_col.clone())
} else if left_field_col != right_field_col {
expr.alias(left_field_col.clone())
} else {
expr
}
} else if target_tag_types.contains_key(col) {
aligned_label_expr(col, &right_tag_types)
} else {
DfExpr::Column(Column::new(None::<String>, col))
}
});
let right_proj_exprs = [right_time_index_expr]
.into_iter()
.chain(right_proj_exprs_without_time_index);
let left_projected = LogicalPlanBuilder::from(left)
.project(left_proj_exprs)
.context(DataFusionPlanningSnafu)?
.alias(left_qualifier_string.clone())
.context(DataFusionPlanningSnafu)?
.build()
.context(DataFusionPlanningSnafu)?;
let right_projected = LogicalPlanBuilder::from(right)
.project(right_proj_exprs)
.context(DataFusionPlanningSnafu)?
.alias(right_qualifier_string.clone())
.context(DataFusionPlanningSnafu)?
.build()
.context(DataFusionPlanningSnafu)?;
// step 2: compute match columns
let mut match_columns = if let Some(modifier) = modifier
&& let Some(matching) = &modifier.matching
{
match matching {
// keeps columns mentioned in `on`
LabelModifier::Include(on) => on.labels.clone(),
// removes columns memtioned in `ignoring`
LabelModifier::Exclude(ignoring) => {
let ignoring = ignoring.labels.iter().cloned().collect::<HashSet<_>>();
all_tags.difference(&ignoring).cloned().collect()
}
}
} else {
all_tags.iter().cloned().collect()
};
// sort to ensure the generated plan is not volatile
match_columns.sort_unstable();
match_columns.dedup();
occupied_column_names.extend(
left_projected
.schema()
.fields()
.iter()
.chain(right_projected.schema().fields().iter())
.map(|field| field.name().clone()),
);
let visible_schema = left_projected.schema().clone();
let visible_left_exprs = left_projected
.schema()
.iter()
.map(|(qualifier, field)| {
DfExpr::Column(Column::new(qualifier.cloned(), field.name().clone()))
})
.collect::<Vec<_>>();
let visible_right_exprs = right_projected
.schema()
.iter()
.map(|(qualifier, field)| {
DfExpr::Column(Column::new(qualifier.cloned(), field.name().clone()))
})
.collect::<Vec<_>>();
let mut left_match_exprs = Vec::with_capacity(match_columns.len());
let mut right_match_exprs = Vec::with_capacity(match_columns.len());
let mut next_internal_column = 0;
for label in &match_columns {
let left_field = if left_tag_cols_set.contains(label) {
Some(
left_projected
.schema()
.iter()
.find(|(_, field)| field.name() == label)
.map(|(qualifier, field)| (qualifier.cloned(), field.data_type().clone()))
.with_context(|| ColumnNotFoundSnafu { col: label.clone() })?,
)
} else {
None
};
let right_field = if right_tag_cols_set.contains(label) {
Some(
right_projected
.schema()
.iter()
.find(|(_, field)| field.name() == label)
.map(|(qualifier, field)| (qualifier.cloned(), field.data_type().clone()))
.with_context(|| ColumnNotFoundSnafu { col: label.clone() })?,
)
} else {
None
};
let data_type = match (left_field.as_ref(), right_field.as_ref()) {
(Some((_, left_type)), Some((_, right_type))) if left_type == right_type => {
left_type.clone()
}
(Some((_, left_type)), Some((_, right_type))) => {
return UnexpectedPlanExprSnafu {
desc: format!(
"OR match label {label} has incompatible types: {left_type:?} and {right_type:?}"
),
}
.fail();
}
(Some((_, data_type)), None) | (None, Some((_, data_type))) => data_type.clone(),
(None, None) => ArrowDataType::Utf8,
};
let Some(value_type) = Self::string_value_data_type(&data_type).cloned() else {
return UnexpectedPlanExprSnafu {
desc: format!("OR match label {label} must be a string"),
}
.fail();
};
let internal_name = loop {
let name = format!("__promql_or_match_{next_internal_column}");
next_internal_column += 1;
if occupied_column_names.insert(name.clone()) {
break name;
}
};
left_match_exprs.push(Self::normalized_match_key_expr(
label,
left_field,
&value_type,
&internal_name,
));
right_match_exprs.push(Self::normalized_match_key_expr(
label,
right_field,
&value_type,
&internal_name,
));
}
let left_augmented = LogicalPlanBuilder::from(left_projected)
.project(visible_left_exprs.into_iter().chain(left_match_exprs))
.context(DataFusionPlanningSnafu)?
.build()
.context(DataFusionPlanningSnafu)?;
let right_augmented = LogicalPlanBuilder::from(right_projected)
.project(visible_right_exprs.into_iter().chain(right_match_exprs))
.context(DataFusionPlanningSnafu)?
.build()
.context(DataFusionPlanningSnafu)?;
// step 3: build `UnionDistinctOn` with normalized internal match keys.
let visible_field_count = visible_schema.fields().len();
let compare_key_indices =
(visible_field_count..visible_field_count + match_columns.len()).collect::<Vec<_>>();
let (time_qualifier, _) = visible_schema
.iter()
.find(|(_, field)| field.name() == &left_time_index_column)
.with_context(|| TimeIndexNotFoundSnafu {
table: left_qualifier_string.clone(),
})?;
let ts_col_idx = left_augmented
.schema()
.iter()
.position(|(qualifier, field)| {
qualifier == time_qualifier && field.name() == &left_time_index_column
})
.with_context(|| TimeIndexNotFoundSnafu {
table: left_qualifier_string.clone(),
})?;
let union_distinct_on = UnionDistinctOn::try_new(
left_augmented,
right_augmented,
compare_key_indices,
ts_col_idx,
)
.context(DataFusionPlanningSnafu)?;
let augmented_result = LogicalPlan::Extension(Extension {
node: Arc::new(union_distinct_on),
});
let result = LogicalPlanBuilder::from(augmented_result)
.project(visible_schema.iter().map(|(qualifier, field)| {
DfExpr::Column(Column::new(qualifier.cloned(), field.name().clone()))
}))
.context(DataFusionPlanningSnafu)?
.build()
.context(DataFusionPlanningSnafu)?;
// step 4: update context
let output_field_col = left_field_col.clone();
let mut output_context = left_context;
let mut visible_tags = all_tags.into_iter().collect::<Vec<_>>();
visible_tags.sort_unstable();
output_context.time_index_column = Some(left_time_index_column);
output_context.tag_columns = visible_tags;
output_context.field_columns = if mixed_sample_types {
vec![mixed_float_field_col, mixed_histogram_field_col]
} else {
vec![output_field_col]
};
output_context.use_tsid = left_has_tsid && right_has_tsid;
self.ctx = output_context;
Ok(result)
}
/// Build a set operator (AND/OR/UNLESS)
pub(super) fn set_op_on_non_field_columns(
&mut self,
mut left: LogicalPlan,
mut right: LogicalPlan,
left_context: PromPlannerContext,
right_context: PromPlannerContext,
op: TokenType,
modifier: &Option<BinModifier>,
) -> Result<LogicalPlan> {
let left_tag_col_set = left_context
.tag_columns
.iter()
.cloned()
.collect::<HashSet<_>>();
let right_tag_col_set = right_context
.tag_columns
.iter()
.cloned()
.collect::<HashSet<_>>();
if matches!(op.id(), token::T_LOR) {
return self.or_operator(
left,
right,
left_tag_col_set,
right_tag_col_set,
left_context,
right_context,
modifier,
);
}
if let Some(modifier) = modifier {
ensure!(
matches!(
modifier.card,
VectorMatchCardinality::OneToOne | VectorMatchCardinality::ManyToMany
),
UnsupportedVectorMatchSnafu {
name: modifier.card.clone(),
},
);
}
let output_context = left_context.clone();
let visible_left_schema = left.schema().clone();
let mut left_context = left_context;
let mut right_context = right_context;
let added_marker_to_left = if Self::only_temporality_match_label_mismatches(
&left_context,
&right_context,
modifier,
) {
let aligned = Self::align_temporality_match_column(
left,
right,
&mut left_context,
&mut right_context,
)?;
left = aligned.0;
right = aligned.1;
aligned.2
} else {
false
};
let mut left_tag_col_set = left_context
.tag_columns
.iter()
.cloned()
.collect::<BTreeSet<_>>();
let mut right_tag_col_set = right_context
.tag_columns
.iter()
.cloned()
.collect::<BTreeSet<_>>();
if let Some(matching) = modifier
.as_ref()
.and_then(|modifier| modifier.matching.as_ref())
{
match matching {
LabelModifier::Include(on) => {
let mask = on.labels.iter().cloned().collect::<BTreeSet<_>>();
left_tag_col_set = left_tag_col_set.intersection(&mask).cloned().collect();
right_tag_col_set = right_tag_col_set.intersection(&mask).cloned().collect();
}
LabelModifier::Exclude(ignoring) => {
for label in &ignoring.labels {
let _ = left_tag_col_set.remove(label);
let _ = right_tag_col_set.remove(label);
}
}
}
}
ensure!(
left_tag_col_set == right_tag_col_set,
CombineTableColumnMismatchSnafu {
left: left_tag_col_set.iter().cloned().collect::<Vec<_>>(),
right: right_tag_col_set.iter().cloned().collect::<Vec<_>>(),
}
);
let left_time_index = left_context.time_index_column.clone().unwrap();
let right_time_index = right_context.time_index_column.clone().unwrap();
// alias right time index column if necessary
if left_context.time_index_column != right_context.time_index_column {
let right_project_exprs = right
.schema()
.fields()
.iter()
.map(|field| {
if field.name() == &right_time_index {
DfExpr::Column(Column::from_name(&right_time_index)).alias(&left_time_index)
} else {
DfExpr::Column(Column::from_name(field.name()))
}
})
.collect::<Vec<_>>();
right = LogicalPlanBuilder::from(right)
.project(right_project_exprs)
.context(DataFusionPlanningSnafu)?
.build()
.context(DataFusionPlanningSnafu)?;
}
let join_keys = left_tag_col_set
.into_iter()
.chain([left_time_index])
.collect::<Vec<_>>();
ensure!(
left_context.field_columns.len() == 1
|| Self::field_columns_are_alternative_samples(
left.schema(),
&left_context.field_columns,
),
MultiFieldsNotSupportedSnafu {
operator: "AND/UNLESS operator"
}
);
// Generate join plan.
// All set operations in PromQL are "distinct"
let result = match op.id() {
token::T_LAND => LogicalPlanBuilder::from(left)
.distinct()
.context(DataFusionPlanningSnafu)?
.join_detailed(
right,
JoinType::LeftSemi,
(join_keys.clone(), join_keys),
None,
NullEquality::NullEqualsNull,
)
.context(DataFusionPlanningSnafu)?
.build()
.context(DataFusionPlanningSnafu),
token::T_LUNLESS => LogicalPlanBuilder::from(left)
.distinct()
.context(DataFusionPlanningSnafu)?
.join_detailed(
right,
JoinType::LeftAnti,
(join_keys.clone(), join_keys),
None,
NullEquality::NullEqualsNull,
)
.context(DataFusionPlanningSnafu)?
.build()
.context(DataFusionPlanningSnafu),
token::T_LOR => {
// OR is handled at the beginning of this function, as it cannot
// be expressed using JOIN like AND and UNLESS.
unreachable!()
}
_ => UnexpectedTokenSnafu { token: op }.fail(),
}?;
let result = if added_marker_to_left {
LogicalPlanBuilder::from(result)
.project(visible_left_schema.iter().map(|(qualifier, field)| {
DfExpr::Column(Column::new(qualifier.cloned(), field.name().clone()))
}))
.context(DataFusionPlanningSnafu)?
.build()
.context(DataFusionPlanningSnafu)?
} else {
result
};
// AND/UNLESS preserve the complete left operand's visible columns and values; encoded
// markers are decoded.
self.ctx = output_context;
Ok(result)
}
}
File diff suppressed because it is too large Load Diff