From 97cbf79fb6dfe9216e75ab2f22bbca32b61f6f67 Mon Sep 17 00:00:00 2001 From: shuiyisong <113876041+shuiyisong@users.noreply.github.com> Date: Tue, 11 Aug 2026 16:01:06 +0800 Subject: [PATCH] feat(promql): support native histogram vector operators (#8798) * feat(promql): support native histogram vector operators Signed-off-by: shuiyisong * fix: CR issue Signed-off-by: shuiyisong * chore: rebase main Signed-off-by: shuiyisong * fix: cr issue Signed-off-by: shuiyisong * fix: cr issue Signed-off-by: shuiyisong * fix: cr issue Signed-off-by: shuiyisong --------- Signed-off-by: shuiyisong --- .../src/extension_plan/instant_manipulate.rs | 72 +- src/promql/src/functions/native_histogram.rs | 20 +- src/query/src/promql/planner.rs | 1741 +++++++++++++++-- 3 files changed, 1668 insertions(+), 165 deletions(-) diff --git a/src/promql/src/extension_plan/instant_manipulate.rs b/src/promql/src/extension_plan/instant_manipulate.rs index f92a4e4490..4851916a01 100644 --- a/src/promql/src/extension_plan/instant_manipulate.rs +++ b/src/promql/src/extension_plan/instant_manipulate.rs @@ -18,6 +18,7 @@ use std::pin::Pin; use std::sync::Arc; use std::task::{Context, Poll}; +use common_query::prelude::{greptime_native_histogram, greptime_value}; use datafusion::arrow::array::{Array, TimestampMillisecondArray, UInt64Array}; use datafusion::arrow::datatypes::{DataType, SchemaRef}; use datafusion::arrow::record_batch::RecordBatch; @@ -52,6 +53,15 @@ use crate::metrics::PROMQL_SERIES_COUNT; const MAX_INSTANT_MANIPULATE_OUTPUT_POINTS: usize = 1_000_000; +fn mixed_sample_fields(field: Option<&str>) -> [Option<&str>; 2] { + let companion = match field { + Some(field) if field == greptime_value() => Some(greptime_native_histogram()), + Some(field) if field == greptime_native_histogram() => Some(greptime_value()), + _ => None, + }; + [field, companion] +} + /// Manipulate the input record batch to make it suitable for Instant Operator. /// /// This plan will try to align the input time series, for every timestamp between @@ -66,7 +76,7 @@ pub struct InstantManipulate { time_index_column: String, // Planner-provided tag-column hint for execution fast paths. tag_columns: Vec, - /// A optional column for validating staleness + /// Primary sample column used to derive the columns checked for staleness. field_column: Option, input: LogicalPlan, unfix: Option, @@ -97,9 +107,7 @@ impl UserDefinedLogicalNodeCore for InstantManipulate { } let mut exprs = vec![col(&self.time_index_column)]; - if let Some(field) = &self.field_column { - exprs.push(col(field)); - } + exprs.extend(self.staleness_field_columns().map(col)); exprs } @@ -116,7 +124,7 @@ impl UserDefinedLogicalNodeCore for InstantManipulate { let mut required = output_columns.to_vec(); required.push(input_schema.index_of_column_by_name(None, &self.time_index_column)?); - if let Some(field) = &self.field_column { + for field in self.staleness_field_columns() { required.push(input_schema.index_of_column_by_name(None, field)?); } @@ -223,6 +231,21 @@ impl InstantManipulate { "InstantManipulate" } + fn staleness_field_columns(&self) -> impl Iterator { + let [field, companion] = mixed_sample_fields(self.field_column.as_deref()); + [ + field, + companion.filter(|companion| { + self.input + .schema() + .index_of_column_by_name(None, companion) + .is_some() + }), + ] + .into_iter() + .flatten() + } + fn resolve_tag_columns(input: &LogicalPlan, tag_columns: &[String]) -> Vec { if !tag_columns.is_empty() { return tag_columns.to_vec(); @@ -385,11 +408,9 @@ impl ExecutionPlan for InstantManipulateExec { .column_with_name(&self.time_index_column) .expect("time index column not found") .0; - let field_index = self - .field_column - .as_ref() - .and_then(|name| schema.column_with_name(name)) - .map(|x| x.0); + let field_indices = mixed_sample_fields(self.field_column.as_deref()).map(|field| { + field.and_then(|field| schema.column_with_name(field).map(|(index, _)| index)) + }); let tsid_index = schema .column_with_name("__tsid") .filter(|(_, field)| field.data_type() == &DataType::UInt64) @@ -400,7 +421,7 @@ impl ExecutionPlan for InstantManipulateExec { lookback_delta: self.lookback_delta, interval: self.interval, time_index, - field_index, + field_indices, tsid_index, reuse_tsid_column: self.reuse_tsid_column && tsid_index.is_some(), schema, @@ -467,7 +488,7 @@ pub struct InstantManipulateStream { interval: Millisecond, // Column index of TIME INDEX column's position in schema time_index: usize, - field_index: Option, + field_indices: [Option; 2], tsid_index: Option, reuse_tsid_column: bool, @@ -532,11 +553,16 @@ impl InstantManipulateStream { return Ok(input); } - // Field column for staleness checks, classified once per batch. - let stale_sample_column = self - .field_index - .map(|index| input.column(index).as_ref()) - .and_then(prometheus_stale_sample_column); + // Field columns for staleness checks, classified once per batch. + let stale_sample_columns = self.field_indices.map(|index| { + index.and_then(|index| prometheus_stale_sample_column(input.column(index).as_ref())) + }); + let is_stale = |row| { + stale_sample_columns + .iter() + .flatten() + .any(|column| is_prometheus_stale_sample(*column, row)) + }; // Optimize iteration range based on actual data bounds let first_ts = ts_column.value(0); @@ -581,9 +607,7 @@ impl InstantManipulateStream { let curr = ts_column.value(cursor); match curr.cmp(&expected_ts) { Ordering::Equal => { - if stale_sample_column - .is_some_and(|column| is_prometheus_stale_sample(column, cursor)) - { + if is_stale(cursor) { // Ignore the stale marker. } else { take_indices.push(cursor as u64); @@ -615,9 +639,7 @@ impl InstantManipulateStream { let prev_ts = ts_column.value(prev_cursor); if prev_ts + self.lookback_delta > expected_ts { // only use the point in the time range - if stale_sample_column - .is_some_and(|column| is_prometheus_stale_sample(column, prev_cursor)) - { + if is_stale(prev_cursor) { // Do not use a stale marker as the newest value. continue; } @@ -626,9 +648,7 @@ impl InstantManipulateStream { aligned_ts.push(expected_ts); } } - } else if stale_sample_column - .is_some_and(|column| is_prometheus_stale_sample(column, cursor)) - { + } else if is_stale(cursor) { // Do not use a stale marker as the newest value. } else { // use this point diff --git a/src/promql/src/functions/native_histogram.rs b/src/promql/src/functions/native_histogram.rs index 1863e663b8..e6c9e84d59 100644 --- a/src/promql/src/functions/native_histogram.rs +++ b/src/promql/src/functions/native_histogram.rs @@ -173,12 +173,20 @@ impl ScalarUDFImpl for NativeHistogramAnnotationUdf { } fn invoke_with_args(&self, args: ScalarFunctionArgs) -> DfResult { - if let Some(collector) = args - .config_options - .extensions - .get::() - .cloned() - .or_else(|| self.collector.clone()) + let has_dropped_sample = !args.args.is_empty() + && (0..args.number_rows).any(|row| { + args.args.iter().all(|arg| match arg { + ColumnarValue::Array(array) => array.is_valid(row), + ColumnarValue::Scalar(value) => !value.is_null(), + }) + }); + if has_dropped_sample + && let Some(collector) = args + .config_options + .extensions + .get::() + .cloned() + .or_else(|| self.collector.clone()) { collector.record_info(self.message.clone()); } diff --git a/src/query/src/promql/planner.rs b/src/query/src/promql/planner.rs index d21a90bf12..10462ea182 100644 --- a/src/query/src/promql/planner.rs +++ b/src/query/src/promql/planner.rs @@ -64,14 +64,17 @@ use promql::extension_plan::{ use promql::functions::{ AbsentOverTime, AvgOverTime, Changes, CountOverTime, Delta, Deriv, DoubleExponentialSmoothing, IDelta, Increase, LastOverTime, MaxOverTime, MinOverTime, MixedRange, - NativeHistogramAbsentOverTime, NativeHistogramAvg, NativeHistogramAvgOverTime, - NativeHistogramChanges, NativeHistogramCount, NativeHistogramCountOverTime, - NativeHistogramDelta, NativeHistogramDrop, NativeHistogramFraction, NativeHistogramIDelta, + NativeHistogramAbsentOverTime, NativeHistogramAdd, NativeHistogramAvg, + NativeHistogramAvgOverTime, NativeHistogramChanges, NativeHistogramCount, + NativeHistogramCountOverTime, NativeHistogramDelta, NativeHistogramDivScalar, + NativeHistogramDrop, NativeHistogramEq, NativeHistogramFraction, NativeHistogramIDelta, NativeHistogramIRate, NativeHistogramIncrease, NativeHistogramLastOverTime, + NativeHistogramMulScalar, NativeHistogramNeg, NativeHistogramNotEq, NativeHistogramPresentOverTime, NativeHistogramQuantile, NativeHistogramRate, - NativeHistogramResets, NativeHistogramStddev, NativeHistogramStdvar, NativeHistogramSum, - NativeHistogramSumOverTime, PredictLinear, PresentOverTime, QuantileOverTime, Rate, Resets, - Round, StddevOverTime, StdvarOverTime, SumOverTime, quantile_udaf, + NativeHistogramResets, NativeHistogramScalarMul, NativeHistogramStddev, NativeHistogramStdvar, + NativeHistogramSub, NativeHistogramSum, NativeHistogramSumOverTime, PredictLinear, + PresentOverTime, QuantileOverTime, Rate, Resets, Round, StddevOverTime, StdvarOverTime, + SumOverTime, quantile_udaf, }; use promql_parser::label::{METRIC_NAME, MatchOp, Matcher, Matchers}; use promql_parser::parser::token::TokenType; @@ -392,6 +395,8 @@ pub struct PromPlanner { promql_annotations: Option, } +type BinaryFieldPair<'a> = (&'a String, &'a String); + impl PromPlanner { pub async fn stmt_to_plan( table_provider: DfTableSourceProvider, @@ -823,8 +828,20 @@ impl PromPlanner { let UnaryExpr { expr } = unary_expr; // Unary Expr in PromQL implys the `-` operator let input = self.prom_expr_to_plan(expr, query_engine_state).await?; + self.negate_field_columns(input) + } + + fn negate_field_columns(&mut self, input: LogicalPlan) -> Result { + let input_schema = input.schema().clone(); self.projection_for_each_field_column(input, |col| { - Ok(DfExpr::Negative(Box::new(DfExpr::Column(col.into())))) + if Self::field_column_is_native_histogram(&input_schema, col) { + Ok(DfExpr::ScalarFunction(ScalarFunction { + func: Arc::new(NativeHistogramNeg::scalar_udf()), + args: vec![DfExpr::Column(col.into())], + })) + } else { + Ok(DfExpr::Negative(Box::new(DfExpr::Column(col.into())))) + } }) } @@ -866,6 +883,16 @@ impl PromPlanner { }); } + 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); @@ -1214,10 +1241,46 @@ impl PromPlanner { if let Some(time_expr) = self.try_build_special_time_expr_with_context(lhs) { expr = time_expr } + let input_schema = input.schema().clone(); + let preserve_any_value = Self::field_columns_are_alternative_samples( + &input_schema, + &self.ctx.field_columns, + ); + let has_native_histogram = Self::field_columns_contain_native_histogram( + &input_schema, + &self.ctx.field_columns, + ); + let retain_field_columns = self + .ctx + .field_columns + .iter() + .map(|col| { + Self::binary_result_is_histogram( + *op, + false, + Self::field_column_is_native_histogram(&input_schema, col), + ) + .is_some() + }) + .collect(); + let promql_annotations = self.promql_annotations.clone(); let bin_expr_builder = |col: &String| { let binary_expr_builder = Self::prom_token_to_binary_expr_builder(*op)?; - let mut binary_expr = - binary_expr_builder(expr.clone(), DfExpr::Column(col.into()))?; + let rhs_is_histogram = + Self::field_column_is_native_histogram(&input_schema, col); + let rhs = DfExpr::Column(col.into()); + let mut binary_expr = match Self::native_histogram_binary_expr( + *op, + expr.clone(), + false, + rhs.clone(), + rhs_is_histogram, + is_comparison_op && !should_return_bool, + promql_annotations.clone(), + )? { + Some(expr) => expr, + None => binary_expr_builder(expr.clone(), rhs)?, + }; if is_comparison_op && should_return_bool { binary_expr = DfExpr::Cast(Cast { @@ -1230,7 +1293,14 @@ impl PromPlanner { if is_comparison_op && !should_return_bool { self.filter_on_field_column(input, bin_expr_builder) } else { - self.projection_for_each_field_column(input, bin_expr_builder) + let projected = + self.projection_for_each_field_column(input, bin_expr_builder)?; + self.filter_binary_projection( + projected, + has_native_histogram, + preserve_any_value, + retain_field_columns, + ) } } // lhs is a column, rhs is a literal @@ -1240,10 +1310,46 @@ impl PromPlanner { if let Some(time_expr) = self.try_build_special_time_expr_with_context(rhs) { expr = time_expr } + let input_schema = input.schema().clone(); + let preserve_any_value = Self::field_columns_are_alternative_samples( + &input_schema, + &self.ctx.field_columns, + ); + let has_native_histogram = Self::field_columns_contain_native_histogram( + &input_schema, + &self.ctx.field_columns, + ); + let retain_field_columns = self + .ctx + .field_columns + .iter() + .map(|col| { + Self::binary_result_is_histogram( + *op, + Self::field_column_is_native_histogram(&input_schema, col), + false, + ) + .is_some() + }) + .collect(); + let promql_annotations = self.promql_annotations.clone(); let bin_expr_builder = |col: &String| { let binary_expr_builder = Self::prom_token_to_binary_expr_builder(*op)?; - let mut binary_expr = - binary_expr_builder(DfExpr::Column(col.into()), expr.clone())?; + let lhs_is_histogram = + Self::field_column_is_native_histogram(&input_schema, col); + let lhs = DfExpr::Column(col.into()); + let mut binary_expr = match Self::native_histogram_binary_expr( + *op, + lhs.clone(), + lhs_is_histogram, + expr.clone(), + false, + is_comparison_op && !should_return_bool, + promql_annotations.clone(), + )? { + Some(expr) => expr, + None => binary_expr_builder(lhs, expr.clone())?, + }; if is_comparison_op && should_return_bool { binary_expr = DfExpr::Cast(Cast { @@ -1256,7 +1362,14 @@ impl PromPlanner { if is_comparison_op && !should_return_bool { self.filter_on_field_column(input, bin_expr_builder) } else { - self.projection_for_each_field_column(input, bin_expr_builder) + let projected = + self.projection_for_each_field_column(input, bin_expr_builder)?; + self.filter_binary_projection( + projected, + has_native_histogram, + preserve_any_value, + retain_field_columns, + ) } } // both are columns. join them on time index @@ -1291,6 +1404,14 @@ impl PromPlanner { ); } + let has_native_histogram = Self::field_columns_contain_native_histogram( + left_input.schema(), + &left_field_columns, + ) || Self::field_columns_contain_native_histogram( + right_input.schema(), + &right_field_columns, + ); + // normal join if left_table_ref == right_table_ref { // rename table references to avoid ambiguity @@ -1308,21 +1429,41 @@ impl PromPlanner { self.ctx.table_name = Some("rhs".to_string()); } } - let (output_field_columns, field_columns) = - Self::align_binary_field_columns(&left_field_columns, &right_field_columns); - let left_aligned_field_columns = field_columns + // Computed scalars reach this join path instead of the literal projection paths. + // Broadcast them for arithmetic in the same way as literal scalars. + let broadcast_scalar = !is_comparison_op; + let (field_groups, invalid_field_pairs) = Self::align_binary_field_columns( + left_input.schema(), + right_input.schema(), + &left_field_columns, + &right_field_columns, + *op, + broadcast_scalar && lhs.value_type() == ValueType::Scalar, + broadcast_scalar && rhs.value_type() == ValueType::Scalar, + ); + let left_aligned_field_columns = field_groups .iter() - .map(|(left_col_name, _)| (*left_col_name).clone()) + .flat_map(|(_, pairs)| { + pairs + .iter() + .map(|(left_col_name, _)| (*left_col_name).clone()) + }) .collect::>(); - let right_aligned_field_columns = field_columns + let right_aligned_field_columns = field_groups .iter() - .map(|(_, right_col_name)| (*right_col_name).clone()) + .flat_map(|(_, pairs)| { + pairs + .iter() + .map(|(_, right_col_name)| (*right_col_name).clone()) + }) .collect::>(); - // PromQL binary arithmetic only combines the shared prefix of value columns. - // Keep the output field count aligned with that zipped prefix so planning - // remains stable even when the two sides have uneven multi-field schemas. - self.ctx.field_columns = output_field_columns; - let mut field_columns = field_columns.into_iter(); + // Regular multi-field vectors combine their shared prefix. Alternative + // float/histogram lanes instead align by valid PromQL sample combinations. + self.ctx.field_columns = field_groups + .iter() + .map(|(output, _)| output.clone()) + .collect(); + let mut field_groups = field_groups.into_iter(); let join_plan = self.join_on_non_field_columns( left_input, @@ -1339,28 +1480,103 @@ impl PromPlanner { &right_context, )?; let join_plan_schema = join_plan.schema().clone(); + let promql_annotations = self.promql_annotations.clone(); + // These predicates always pass; they only evaluate otherwise-discarded pairs + // while collecting annotations. + let invalid_pair_predicates = invalid_field_pairs + .into_iter() + .filter(|_| promql_annotations.is_some()) + .map(|(left_col_name, right_col_name)| { + let left_field = join_plan_schema + .qualified_field_with_name(Some(&left_table_ref), left_col_name) + .context(DataFusionPlanningSnafu)?; + let right_field = join_plan_schema + .qualified_field_with_name(Some(&right_table_ref), right_col_name) + .context(DataFusionPlanningSnafu)?; + let left_is_histogram = + left_field.1.data_type() == &Self::native_histogram_arrow_type(); + let right_is_histogram = + right_field.1.data_type() == &Self::native_histogram_arrow_type(); + let drop_expr = Self::native_histogram_binary_expr( + *op, + DfExpr::Column(left_field.into()), + left_is_histogram, + DfExpr::Column(right_field.into()), + right_is_histogram, + true, + promql_annotations.clone(), + )? + .with_context(|| UnexpectedPlanExprSnafu { + desc: "invalid native histogram pair produced no drop expression", + })?; + Ok(DfExpr::Not(Box::new(drop_expr))) + }) + .collect::>>()?; + let join_plan = if let Some(predicate) = conjunction(invalid_pair_predicates) { + LogicalPlanBuilder::from(join_plan) + .filter(predicate) + .context(DataFusionPlanningSnafu)? + .build() + .context(DataFusionPlanningSnafu)? + } else { + join_plan + }; let bin_expr_builder = |_: &String| { - let (left_col_name, right_col_name) = field_columns.next().unwrap(); - let left_col = join_plan_schema - .qualified_field_with_name(Some(&left_table_ref), left_col_name) - .context(DataFusionPlanningSnafu)? - .into(); - let right_col = join_plan_schema - .qualified_field_with_name(Some(&right_table_ref), right_col_name) - .context(DataFusionPlanningSnafu)? - .into(); + let (_, field_pairs) = + field_groups + .next() + .with_context(|| UnexpectedPlanExprSnafu { + desc: "missing binary field group", + })?; + let binary_exprs = field_pairs + .into_iter() + .map(|(left_col_name, right_col_name)| { + let left_field = join_plan_schema + .qualified_field_with_name(Some(&left_table_ref), left_col_name) + .context(DataFusionPlanningSnafu)?; + let right_field = join_plan_schema + .qualified_field_with_name(Some(&right_table_ref), right_col_name) + .context(DataFusionPlanningSnafu)?; + let left_is_histogram = + left_field.1.data_type() == &Self::native_histogram_arrow_type(); + let right_is_histogram = + right_field.1.data_type() == &Self::native_histogram_arrow_type(); + let left_col = left_field.into(); + let right_col = right_field.into(); - let binary_expr_builder = Self::prom_token_to_binary_expr_builder(*op)?; - let mut binary_expr = - binary_expr_builder(DfExpr::Column(left_col), DfExpr::Column(right_col))?; - if is_comparison_op && should_return_bool { - binary_expr = DfExpr::Cast(Cast { - expr: Box::new(binary_expr), - data_type: ArrowDataType::Float64, - }); + let binary_expr_builder = Self::prom_token_to_binary_expr_builder(*op)?; + let lhs = DfExpr::Column(left_col); + let rhs = DfExpr::Column(right_col); + let mut binary_expr = match Self::native_histogram_binary_expr( + *op, + lhs.clone(), + left_is_histogram, + rhs.clone(), + right_is_histogram, + is_comparison_op && !should_return_bool, + promql_annotations.clone(), + )? { + Some(expr) => expr, + None => binary_expr_builder(lhs, rhs)?, + }; + if is_comparison_op && should_return_bool { + binary_expr = DfExpr::Cast(Cast { + expr: Box::new(binary_expr), + data_type: ArrowDataType::Float64, + }); + } + Ok(binary_expr) + }) + .collect::>>()?; + if let [binary_expr] = binary_exprs.as_slice() { + Ok(binary_expr.clone()) + } else { + Ok(DfExpr::ScalarFunction(ScalarFunction { + func: coalesce(), + args: binary_exprs, + })) } - Ok(binary_expr) }; if is_comparison_op && !should_return_bool { // PromQL comparison operators without `bool` are filters: @@ -1386,12 +1602,90 @@ impl PromPlanner { project_context.field_columns = project_field_columns; self.project_binary_join_side(filtered, project_table_ref, &project_context) } else { - self.projection_for_each_field_column(join_plan, bin_expr_builder) + let projected = + self.projection_for_each_field_column(join_plan, bin_expr_builder)?; + let preserve_any_value = Self::field_columns_are_alternative_samples( + projected.schema(), + &self.ctx.field_columns, + ); + let retain_field_columns = vec![true; self.ctx.field_columns.len()]; + self.filter_binary_projection( + projected, + has_native_histogram, + preserve_any_value, + retain_field_columns, + ) } } } } + fn filter_binary_projection( + &mut self, + input: LogicalPlan, + has_native_histogram: bool, + preserve_any_value: bool, + retain_field_columns: Vec, + ) -> Result { + if !has_native_histogram { + return Ok(input); + } + + ensure!( + retain_field_columns.len() == self.ctx.field_columns.len(), + UnexpectedPlanExprSnafu { + desc: "binary output field count changed unexpectedly", + } + ); + + let filtered = LogicalPlanBuilder::from(input) + .filter(self.create_empty_values_filter_expr(preserve_any_value)?) + .context(DataFusionPlanningSnafu)? + .build() + .context(DataFusionPlanningSnafu)?; + if retain_field_columns.iter().all(|retain| *retain) { + return Ok(filtered); + } + + let retained = self + .ctx + .field_columns + .iter() + .zip(retain_field_columns) + .filter(|(_, retain)| *retain) + .map(|(field, _)| field.clone()) + .collect::>(); + if retained.is_empty() { + return Ok(filtered); + } + self.ctx.field_columns = retained; + + let mut output_columns = self + .ctx + .field_columns + .iter() + .chain(&self.ctx.tag_columns) + .cloned() + .collect::>(); + output_columns.extend(self.ctx.time_index_column.iter().cloned()); + if self.ctx.use_tsid { + output_columns.insert(DATA_SCHEMA_TSID_COLUMN_NAME.to_string()); + } + let project_exprs = filtered + .schema() + .iter() + .filter(|(_, field)| output_columns.contains(field.name())) + .map(|(qualifier, field)| { + DfExpr::Column(Column::new(qualifier.cloned(), field.name().clone())) + }) + .collect::>(); + LogicalPlanBuilder::from(filtered) + .project(project_exprs) + .context(DataFusionPlanningSnafu)? + .build() + .context(DataFusionPlanningSnafu) + } + fn project_binary_join_side( &mut self, input: LogicalPlan, @@ -4416,6 +4710,69 @@ impl PromPlanner { } } + fn native_histogram_binary_expr( + token: TokenType, + lhs: DfExpr, + lhs_is_histogram: bool, + rhs: DfExpr, + rhs_is_histogram: bool, + filter_context: bool, + promql_annotations: Option, + ) -> Result> { + if !lhs_is_histogram && !rhs_is_histogram { + return Ok(None); + } + + let scalar_fn = |func: ScalarUdfDef, args| { + DfExpr::ScalarFunction(ScalarFunction { + func: Arc::new(func), + args, + }) + }; + let invalid_expr = || { + let message = format!( + "{}: dropped native histogram samples because this binary operation is not supported for native histograms", + token + ); + let func = if filter_context { + NativeHistogramDrop::bool_false_udf(message, promql_annotations.clone()) + } else { + NativeHistogramDrop::float_null_udf(message, promql_annotations.clone()) + }; + let args = vec![lhs.clone(), rhs.clone()]; + scalar_fn(func, args) + }; + + let expr = match (token.id(), lhs_is_histogram, rhs_is_histogram) { + (token::T_ADD, true, true) => scalar_fn( + NativeHistogramAdd::scalar_udf_with_collector(promql_annotations.clone()), + vec![lhs, rhs], + ), + (token::T_SUB, true, true) => scalar_fn( + NativeHistogramSub::scalar_udf_with_collector(promql_annotations.clone()), + vec![lhs, rhs], + ), + (token::T_MUL, true, false) => { + scalar_fn(NativeHistogramMulScalar::scalar_udf(), vec![lhs, rhs]) + } + (token::T_MUL, false, true) => { + scalar_fn(NativeHistogramScalarMul::scalar_udf(), vec![lhs, rhs]) + } + (token::T_DIV, true, false) => { + scalar_fn(NativeHistogramDivScalar::scalar_udf(), vec![lhs, rhs]) + } + (token::T_EQLC, true, true) => { + scalar_fn(NativeHistogramEq::scalar_udf(), vec![lhs, rhs]) + } + (token::T_NEQ, true, true) => { + scalar_fn(NativeHistogramNotEq::scalar_udf(), vec![lhs, rhs]) + } + _ => invalid_expr(), + }; + + Ok(Some(expr)) + } + /// Return a lambda to build binary expression from token. /// Because some binary operator are function in DataFusion like `atan2` or `^`. #[allow(clippy::type_complexity)] @@ -4501,18 +4858,117 @@ impl PromPlanner { } fn align_binary_field_columns<'a>( + left_schema: &DFSchemaRef, + right_schema: &DFSchemaRef, left_field_columns: &'a [String], right_field_columns: &'a [String], - ) -> (Vec, Vec<(&'a String, &'a String)>) { - let field_pairs = left_field_columns - .iter() - .zip(right_field_columns.iter()) - .collect::>(); - let output_field_columns = field_pairs - .iter() - .map(|(left_col_name, _)| (*left_col_name).clone()) - .collect(); - (output_field_columns, field_pairs) + op: TokenType, + left_is_scalar: bool, + right_is_scalar: bool, + ) -> ( + Vec<(String, Vec>)>, + Vec>, + ) { + // Mixed vectors store mutually exclusive float and histogram samples in two columns. + // Retain each valid sample combination and group expressions by their output lane. + let left_alternative = Self::alternative_sample_columns(left_schema, left_field_columns); + let right_alternative = Self::alternative_sample_columns(right_schema, right_field_columns); + let alternative_alignment = match (left_alternative, right_alternative) { + (Some(output_names), Some(_)) => Some(( + output_names, + left_field_columns + .iter() + .flat_map(|left| right_field_columns.iter().map(move |right| (left, right))) + .collect::>(), + )), + (Some(output_names), None) if right_field_columns.len() == 1 => Some(( + output_names, + left_field_columns + .iter() + .map(|left| (left, &right_field_columns[0])) + .collect::>(), + )), + (None, Some(output_names)) if left_field_columns.len() == 1 => Some(( + output_names, + right_field_columns + .iter() + .map(|right| (&left_field_columns[0], right)) + .collect::>(), + )), + _ => None, + }; + let mut invalid_pairs = Vec::new(); + if let Some(((float_output, histogram_output), field_pairs)) = alternative_alignment { + let mut float_pairs = Vec::new(); + let mut histogram_pairs = Vec::new(); + for (left, right) in field_pairs { + let left_is_histogram = Self::field_column_is_native_histogram(left_schema, left); + let right_is_histogram = + Self::field_column_is_native_histogram(right_schema, right); + match Self::binary_result_is_histogram(op, left_is_histogram, right_is_histogram) { + Some(false) => float_pairs.push((left, right)), + Some(true) => histogram_pairs.push((left, right)), + None => invalid_pairs.push((left, right)), + } + } + if !float_pairs.is_empty() || !histogram_pairs.is_empty() { + return ( + [ + (!float_pairs.is_empty()).then(|| (float_output.to_string(), float_pairs)), + (!histogram_pairs.is_empty()) + .then(|| (histogram_output.to_string(), histogram_pairs)), + ] + .into_iter() + .flatten() + .collect(), + invalid_pairs, + ); + } + } + + if left_is_scalar && !right_is_scalar && left_field_columns.len() == 1 { + return ( + right_field_columns + .iter() + .map(|right| (right.clone(), vec![(&left_field_columns[0], right)])) + .collect(), + invalid_pairs, + ); + } + if right_is_scalar && !left_is_scalar && right_field_columns.len() == 1 { + return ( + left_field_columns + .iter() + .map(|left| (left.clone(), vec![(left, &right_field_columns[0])])) + .collect(), + invalid_pairs, + ); + } + + ( + left_field_columns + .iter() + .zip(right_field_columns.iter()) + .map(|(left, right)| (left.clone(), vec![(left, right)])) + .collect(), + invalid_pairs, + ) + } + + fn binary_result_is_histogram( + token: TokenType, + lhs_is_histogram: bool, + rhs_is_histogram: bool, + ) -> Option { + match (token.id(), lhs_is_histogram, rhs_is_histogram) { + (_, false, false) => Some(false), + (token::T_ADD | token::T_SUB, true, true) + | (token::T_MUL, true, false) + | (token::T_MUL, false, true) + | (token::T_DIV, true, false) => Some(true), + (token::T_EQLC | token::T_NEQ, true, true) => Some(false), + _ => None, + } } fn plan_has_tsid_column(plan: &LogicalPlan) -> bool { @@ -4540,6 +4996,15 @@ impl PromPlanner { .is_some_and(|data_type| data_type == &Self::native_histogram_arrow_type()) } + fn field_columns_contain_native_histogram( + schema: &DFSchemaRef, + field_columns: &[String], + ) -> bool { + field_columns + .iter() + .any(|field| Self::field_column_is_native_histogram(schema, field)) + } + fn field_column_is_float_range(schema: &DFSchemaRef, field_column: &str) -> bool { Self::field_column_type(schema, field_column).is_some_and(|data_type| { matches!( @@ -4940,9 +5405,13 @@ impl PromPlanner { } ensure!( - left_context.field_columns.len() == 1, + left_context.field_columns.len() == 1 + || Self::field_columns_are_alternative_samples( + left.schema(), + &left_context.field_columns, + ), MultiFieldsNotSupportedSnafu { - operator: "AND operator" + operator: "AND/UNLESS operator" } ); // Generate join plan. @@ -5072,24 +5541,28 @@ impl PromPlanner { (false, false) => {} } - // checks ensure!( - left_context.field_columns.len() == right_context.field_columns.len(), - CombineTableColumnMismatchSnafu { - left: left_context.field_columns.clone(), - right: right_context.field_columns.clone() + !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_context.field_columns.len() == 1 || left_has_alternative_samples, MultiFieldsNotSupportedSnafu { operator: "OR operator" } ); ensure!( - !left.schema().fields().is_empty() && !right.schema().fields().is_empty(), - UnexpectedPlanExprSnafu { - desc: "OR operator input has zero columns", + right_context.field_columns.len() == 1 || right_has_alternative_samples, + MultiFieldsNotSupportedSnafu { + operator: "OR operator" } ); @@ -5122,62 +5595,125 @@ impl PromPlanner { .with_context(|| TimeIndexNotFoundSnafu { table: right_qualifier_string.clone(), })?; - // Take the name of first field column. The length is checked above. - let left_field_col = left_context.field_columns.first().unwrap(); - let right_field_col = right_context.field_columns.first().unwrap(); - let left_field = left - .schema() + 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() - .find(|(_, field)| field.name() == left_field_col) - .map(|(qualifier, field)| (qualifier.cloned(), field.data_type().clone())) - .with_context(|| ColumnNotFoundSnafu { - col: left_field_col.clone(), - })?; - let right_field = right - .schema() + .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::>>()?; + let right_fields = right_context + .field_columns .iter() - .find(|(_, field)| field.name() == right_field_col) - .map(|(qualifier, field)| (qualifier.cloned(), field.data_type().clone())) - .with_context(|| ColumnNotFoundSnafu { - col: right_field_col.clone(), - })?; - let target_field_type = if left_field.1 == right_field.1 { - left_field.1.clone() - } else if matches!( - left_field.1, - ArrowDataType::Int8 - | ArrowDataType::Int16 - | ArrowDataType::Int32 - | ArrowDataType::Int64 - | ArrowDataType::UInt8 - | ArrowDataType::UInt16 - | ArrowDataType::UInt32 - | ArrowDataType::UInt64 - | ArrowDataType::Float32 - | ArrowDataType::Float64 - ) && matches!( - right_field.1, - ArrowDataType::Int8 - | ArrowDataType::Int16 - | ArrowDataType::Int32 - | ArrowDataType::Int64 - | ArrowDataType::UInt8 - | ArrowDataType::UInt16 - | ArrowDataType::UInt32 - | ArrowDataType::UInt64 - | ArrowDataType::Float32 - | ArrowDataType::Float64 - ) { + .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::>>()?; + 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, 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::>(), + right_fields + .iter() + .map(|(_, _, data_type)| data_type) + .collect::>() + ), + } + .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.1, right_field.1 + 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::>(); + 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| { + 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| { @@ -5244,8 +5780,15 @@ impl PromPlanner { // remove time index column all_columns_set.remove(&left_time_index_column); all_columns_set.remove(&right_time_index_column); - // remove field column in the right - if left_field_col != right_field_col { + 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::>(); @@ -5284,11 +5827,52 @@ impl PromPlanner { .alias(col.clone()) } }; + let null_histogram = + ScalarValue::try_new_null(&native_histogram_type).context(DataFusionPlanningSnafu)?; + let mixed_value_expr = |fields: &[(String, Option, 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 { + expr: Box::new(expr), + data_type: 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 col == left_field_col && left_field.1 != target_field_type { + 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 { expr: Box::new(DfExpr::Column(Column::new( - left_field.0.clone(), + left_field.1.clone(), left_field_col, ))), data_type: target_field_type.clone(), @@ -5310,9 +5894,13 @@ impl PromPlanner { // `skip(1)` to skip the time index column let right_proj_exprs_without_time_index = all_columns.iter().skip(1).map(|col| { // expr - if col == left_field_col { - let expr = DfExpr::Column(Column::new(right_field.0.clone(), right_field_col)); - if right_field.1 != target_field_type { + 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 { expr: Box::new(expr), data_type: target_field_type.clone(), @@ -5529,7 +6117,11 @@ impl PromPlanner { visible_tags.sort_unstable(); output_context.time_index_column = Some(left_time_index_column); output_context.tag_columns = visible_tags; - output_context.field_columns = vec![output_field_col]; + 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; @@ -5551,6 +6143,10 @@ impl PromPlanner { where F: FnMut(&String) -> Result, { + // Keep the generated float/histogram lane names while an element-wise operation + // preserves both sample types, so downstream operators still recognize the pair. + let preserve_field_names = + Self::field_columns_are_alternative_samples(input.schema(), &self.ctx.field_columns); let table_ref = self.ctx.table_name.clone().map(TableReference::bare); let non_field_columns_iter = self .ctx @@ -5572,10 +6168,12 @@ impl PromPlanner { .collect::>>()?; // alias the computation exprs to remove qualifier - self.ctx.field_columns = result_field_columns - .iter() - .map(|expr| expr.schema_name().to_string()) - .collect(); + if !preserve_field_names { + self.ctx.field_columns = result_field_columns + .iter() + .map(|expr| expr.schema_name().to_string()) + .collect(); + } let field_columns_iter = result_field_columns .into_iter() .zip(self.ctx.field_columns.iter()) @@ -5594,24 +6192,32 @@ impl PromPlanner { .context(DataFusionPlanningSnafu) } - /// Build a filter plan that filter on value column. Notice that only one value column - /// is expected. - fn filter_on_field_column( - &self, - input: LogicalPlan, - mut name_to_expr: F, - ) -> Result + /// Build a filter plan on one value column or a float/histogram alternative pair. + fn filter_on_field_column(&self, input: LogicalPlan, name_to_expr: F) -> Result where F: FnMut(&String) -> Result, { ensure!( - self.ctx.field_columns.len() == 1, + self.ctx.field_columns.len() == 1 + || Self::field_columns_are_alternative_samples( + input.schema(), + &self.ctx.field_columns, + ), UnsupportedExprSnafu { name: "filter on multi-value input" } ); - let field_column_filter = name_to_expr(&self.ctx.field_columns[0])?; + let field_column_filters = self + .ctx + .field_columns + .iter() + .map(name_to_expr) + .collect::>>()?; + let field_column_filter = + disjunction(field_column_filters).context(UnsupportedExprSnafu { + name: "filter on empty input", + })?; LogicalPlanBuilder::from(input) .filter(field_column_filter) @@ -5743,7 +6349,9 @@ mod test { CounterResetHint, NativeHistogram, build_histogram_array, }; use common_query::prelude::{greptime_native_histogram, greptime_timestamp, greptime_value}; + use common_query::prometheus::PROMETHEUS_STALE_NAN_BITS; use common_query::test_util::DummyDecoder; + use common_recordbatch::RecordBatch as GreptimeRecordBatch; use datafusion::arrow::array::{ Array, Float64Array, Int64Array, StringArray, TimestampMillisecondArray, }; @@ -5761,8 +6369,9 @@ mod test { use promql_parser::parser; use session::context::QueryContext; use substrait::{DFLogicalSubstraitConvertor, SubstraitPlan}; - use table::metadata::{TableInfoBuilder, TableMetaBuilder}; - use table::test_util::EmptyTable; + use table::Table; + use table::metadata::{FilterPushDownType, TableInfoBuilder, TableMetaBuilder}; + use table::test_util::{EmptyTable, MemTable as GreptimeMemTable}; use super::*; use crate::QueryEngineContext; @@ -5896,6 +6505,7 @@ mod test { enum DirectOrValue { Float64(f64), Int64(i64), + NativeHistogram(NativeHistogram), Utf8(&'static str), } @@ -5904,6 +6514,7 @@ mod test { match self { Self::Float64(_) => ArrowDataType::Float64, Self::Int64(_) => ArrowDataType::Int64, + Self::NativeHistogram(_) => native_histogram_value_type().as_arrow_type(), Self::Utf8(_) => ArrowDataType::Utf8, } } @@ -5911,6 +6522,7 @@ mod test { match self { Self::Float64(v) => Arc::new(Float64Array::from(vec![*v])), Self::Int64(v) => Arc::new(Int64Array::from(vec![*v])), + Self::NativeHistogram(v) => build_histogram_array(&[Some(v.clone())]), Self::Utf8(v) => Arc::new(StringArray::from(vec![*v])), } } @@ -5933,6 +6545,121 @@ mod test { } } + fn operator_metric_table( + name: &str, + table_id: u32, + tag: &str, + value: DirectOrValue, + ) -> table::TableRef { + let value_type = match &value { + DirectOrValue::Float64(_) => ConcreteDataType::float64_datatype(), + DirectOrValue::Int64(_) => ConcreteDataType::int64_datatype(), + DirectOrValue::NativeHistogram(_) => native_histogram_value_type().clone(), + DirectOrValue::Utf8(_) => ConcreteDataType::string_datatype(), + }; + let schema = Arc::new(Schema::new(vec![ + ColumnSchema::new( + "tag".to_string(), + ConcreteDataType::string_datatype(), + false, + ), + ColumnSchema::new( + "ts".to_string(), + ConcreteDataType::timestamp_millisecond_datatype(), + false, + ) + .with_time_index(true), + ColumnSchema::new("v".to_string(), value_type, true), + ])); + let batch = RecordBatch::try_new( + schema.arrow_schema().clone(), + vec![ + Arc::new(StringArray::from(vec![tag])) as Arc, + Arc::new(TimestampMillisecondArray::from(vec![1_000])), + value.array(), + ], + ) + .unwrap(); + let backing = GreptimeMemTable::new_with_catalog( + name, + GreptimeRecordBatch::from_df_record_batch(schema.clone(), batch), + table_id, + DEFAULT_CATALOG_NAME.to_string(), + DEFAULT_SCHEMA_NAME.to_string(), + ); + let meta = TableMetaBuilder::empty() + .schema(schema) + .primary_key_indices(vec![0]) + .value_indices(vec![2]) + .next_column_id(3) + .build() + .unwrap(); + let info = Arc::new( + TableInfoBuilder::default() + .table_id(table_id) + .name(name) + .meta(meta) + .build() + .unwrap(), + ); + Arc::new(Table::new( + info, + FilterPushDownType::Unsupported, + backing.data_source(), + )) + } + + fn operator_table_provider() -> DfTableSourceProvider { + let catalog = MemoryCatalogManager::with_default_setup(); + let tables = [ + operator_metric_table("lf", 2_001, "a", DirectOrValue::Float64(2.0)), + operator_metric_table( + "lh", + 2_002, + "b", + DirectOrValue::NativeHistogram(direct_or_histogram()), + ), + operator_metric_table("rf", 2_003, "b", DirectOrValue::Float64(3.0)), + operator_metric_table( + "rh", + 2_004, + "a", + DirectOrValue::NativeHistogram(direct_or_histogram()), + ), + operator_metric_table("fallback", 2_005, "c", DirectOrValue::Float64(7.0)), + ]; + for table in tables { + let info = table.table_info(); + catalog + .register_table_sync(RegisterTableRequest { + catalog: DEFAULT_CATALOG_NAME.to_string(), + schema: DEFAULT_SCHEMA_NAME.to_string(), + table_name: info.name.clone(), + table_id: info.ident.table_id, + table, + }) + .unwrap(); + } + DfTableSourceProvider::new( + catalog, + false, + QueryContext::arc(), + DummyDecoder::arc(), + false, + ) + } + + fn operator_eval_stmt(expr: &str) -> EvalStmt { + let time = UNIX_EPOCH.checked_add(Duration::from_secs(1)).unwrap(); + EvalStmt { + expr: parser::parse(expr).unwrap(), + start: time, + end: time, + interval: Duration::from_secs(1), + lookback_delta: Duration::from_secs(5), + } + } + struct DirectOrSource { name: &'static str, empty: bool, @@ -6093,6 +6820,66 @@ mod test { execute(plan, &build_query_engine_state()).await } + async fn mixed_direct_or(histogram_on_left: bool) -> (PromPlanner, LogicalPlan) { + let sample = |histogram: bool| { + if histogram { + DirectOrValue::NativeHistogram(direct_or_histogram()) + } else { + DirectOrValue::Float64(1.25) + } + }; + let left = tagged_source( + "lhs", + false, + ( + "k", + Some(if histogram_on_left { + "histogram" + } else { + "float" + }), + ), + sample(histogram_on_left), + ); + let right = tagged_source( + "rhs", + false, + ( + "k", + Some(if histogram_on_left { + "float" + } else { + "histogram" + }), + ), + sample(!histogram_on_left), + ); + let table_provider = build_test_table_provider_with_fields( + &[(DEFAULT_SCHEMA_NAME.to_string(), "dummy".to_string())], + &[], + ) + .await; + let mut planner = PromPlanner { + table_provider, + ctx: PromPlannerContext::default(), + promql_annotations: None, + }; + let left_context = direct_or_context("lhs", &["job", "k"], "v"); + let right_context = direct_or_context("rhs", &["job", "k"], "v"); + let plan = planner + .or_operator( + scan(&left), + scan(&right), + left_context.tag_columns.iter().cloned().collect(), + right_context.tag_columns.iter().cloned().collect(), + left_context, + right_context, + &or_modifier("lhs or on(k) rhs"), + ) + .unwrap(); + (planner, plan) + } + fn assert_no_internal_or_keys(schema: &DFSchema) { assert!( schema @@ -6119,6 +6906,23 @@ mod test { .collect() } + fn histograms(batches: &[RecordBatch], column: &str) -> Vec { + batches + .iter() + .flat_map(|batch| { + let values = batch + .column_by_name(column) + .unwrap() + .as_any() + .downcast_ref::() + .unwrap(); + (0..values.len()).filter_map(|row| { + common_query::native_histogram::read_histogram(values, row).unwrap() + }) + }) + .collect() + } + fn rows(batches: &[RecordBatch]) -> Vec<(f64, Option)> { let mut rows = batches .iter() @@ -8833,6 +9637,110 @@ mod test { assert!(!plan.contains("PromHistogramFold"), "{plan}"); } + #[tokio::test] + async fn timestamp_filters_native_histogram_stale_marker_before_projection() { + let mut stale = direct_or_histogram(); + stale.sum = f64::from_bits(PROMETHEUS_STALE_NAN_BITS); + let table = operator_metric_table( + "stale_histogram", + 2_100, + "a", + DirectOrValue::NativeHistogram(stale), + ); + let catalog = MemoryCatalogManager::with_default_setup(); + catalog + .register_table_sync(RegisterTableRequest { + catalog: DEFAULT_CATALOG_NAME.to_string(), + schema: DEFAULT_SCHEMA_NAME.to_string(), + table_name: "stale_histogram".to_string(), + table_id: 2_100, + table, + }) + .unwrap(); + let provider = DfTableSourceProvider::new( + catalog, + false, + QueryContext::arc(), + DummyDecoder::arc(), + false, + ); + let state = build_query_engine_state(); + let plan = PromPlanner::stmt_to_plan( + provider, + &operator_eval_stmt("timestamp(stale_histogram)"), + &state, + ) + .await + .unwrap(); + let plan_text = plan.display_indent_schema().to_string(); + assert!(plan_text.contains(TIMESTAMP_VALUE_PREFIX), "{plan_text}"); + + let (_, batches) = execute(plan, &state).await; + assert_eq!(batches.iter().map(RecordBatch::num_rows).sum::(), 0); + } + + #[tokio::test] + async fn timestamp_filters_stale_marker_from_mixed_sample_companion() { + let histograms = build_histogram_array(&[None]); + let schema = Arc::new(ArrowSchema::new(vec![ + Field::new( + "timestamp", + ArrowDataType::Timestamp(ArrowTimeUnit::Millisecond, None), + false, + ), + Field::new( + greptime_native_histogram(), + histograms.data_type().clone(), + true, + ), + Field::new(greptime_value(), ArrowDataType::Float64, true), + ])); + let batch = RecordBatch::try_new( + schema.clone(), + vec![ + Arc::new(TimestampMillisecondArray::from(vec![1_000])), + histograms, + Arc::new(Float64Array::from(vec![f64::from_bits( + PROMETHEUS_STALE_NAN_BITS, + )])), + ], + ) + .unwrap(); + let table = Arc::new(MemTable::try_new(schema, vec![vec![batch]]).unwrap()); + let input = LogicalPlanBuilder::scan("mixed", provider_as_source(table), None) + .unwrap() + .build() + .unwrap(); + let input = LogicalPlan::Extension(Extension { + node: Arc::new(SeriesDivide::new( + Vec::new(), + "timestamp".to_string(), + input, + )), + }); + let input = LogicalPlan::Extension(Extension { + node: Arc::new(InstantManipulate::new( + 1_000, + 1_000, + 5_000, + 1_000, + "timestamp".to_string(), + Vec::new(), + Some(greptime_native_histogram().to_string()), + input, + )), + }); + // Match timestamp()'s parent projection, which otherwise prunes the companion lane. + let plan = LogicalPlanBuilder::from(input) + .project([col("timestamp")]) + .unwrap() + .build() + .unwrap(); + + let (_, batches) = execute(plan, &build_query_engine_state()).await; + assert_eq!(batches.iter().map(RecordBatch::num_rows).sum::(), 0); + } + #[tokio::test] async fn native_histogram_rate_can_feed_count() { let plan = native_histogram_plan("histogram_count(rate(some_metric[5m]))").await; @@ -10654,7 +11562,6 @@ Projection: count(prometheus_tsdb_head_series.greptime_value) AS my_series, prom .contains("OR value fields have incompatible types") ); } - #[tokio::test] async fn test_or_with_histogram_quantile_missing_le_column() { let case = r#"histogram_quantile(0.99, non_existent_histogram_bucket) or normal_metric"#; @@ -10871,4 +11778,572 @@ Projection: count(prometheus_tsdb_head_series.greptime_value) AS my_series, prom "{plan:?}" ); } + + #[tokio::test] + async fn test_direct_or_preserves_float_and_native_histogram_samples() { + for histogram_on_left in [false, true] { + let (planner, plan) = mixed_direct_or(histogram_on_left).await; + + let float_field = &planner.ctx.field_columns[0]; + let histogram_field = &planner.ctx.field_columns[1]; + assert!(float_field.starts_with(OR_FLOAT_FIELD_PREFIX)); + assert!(histogram_field.starts_with(OR_HISTOGRAM_FIELD_PREFIX)); + assert_eq!( + plan.schema() + .field_with_name(None, float_field) + .unwrap() + .data_type(), + &ArrowDataType::Float64 + ); + assert_eq!( + plan.schema() + .field_with_name(None, histogram_field) + .unwrap() + .data_type(), + &native_histogram_value_type().as_arrow_type() + ); + + let (optimized, batches) = execute(plan, &build_query_engine_state()).await; + assert_no_internal_or_keys(optimized.schema()); + let mut sample_kinds = batches + .iter() + .flat_map(|batch| { + let values = batch.column_by_name(float_field).unwrap(); + let histograms = batch.column_by_name(histogram_field).unwrap(); + (0..batch.num_rows()) + .map(|row| (values.is_valid(row), histograms.is_valid(row))) + }) + .collect::>(); + sample_kinds.sort_unstable(); + assert_eq!(sample_kinds, vec![(false, true), (true, false)]); + } + } + + #[tokio::test] + async fn test_mixed_binary_operator_aligns_both_alternative_inputs() { + let state = build_query_engine_state(); + let plan = PromPlanner::stmt_to_plan( + operator_table_provider(), + &operator_eval_stmt("(lf or on(tag) lh) * on(tag) (rf or on(tag) rh)"), + &state, + ) + .await + .unwrap(); + let plan_text = plan.display_indent_schema().to_string(); + assert!( + plan_text.contains("prom_native_histogram_mul_scalar"), + "{plan_text}" + ); + assert!( + plan_text.contains("prom_native_histogram_scalar_mul"), + "{plan_text}" + ); + let float_field = plan + .schema() + .fields() + .iter() + .find(|field| field.name().starts_with(OR_FLOAT_FIELD_PREFIX)) + .unwrap() + .name() + .clone(); + let histogram_field = plan + .schema() + .fields() + .iter() + .find(|field| field.name().starts_with(OR_HISTOGRAM_FIELD_PREFIX)) + .unwrap() + .name() + .clone(); + + let (_, batches) = execute(plan, &state).await; + assert_eq!(batches.iter().map(RecordBatch::num_rows).sum::(), 2); + assert!(values(&batches, &float_field).is_empty()); + let mut sums = histograms(&batches, &histogram_field) + .into_iter() + .map(|histogram| histogram.sum) + .collect::>(); + sums.sort_by(f64::total_cmp); + assert_eq!(sums, vec![2.0, 3.0]); + } + + #[tokio::test] + async fn test_mixed_binary_operator_reports_only_dropped_samples() { + for (query, expected_rows, expected_infos) in [ + ("(lf or on(tag) lh) + on(tag) (rf or on(tag) rh)", 0, 1), + ("(lf or on(tag) lh) + on(tag) (lf or on(tag) lh)", 2, 0), + ("(lf or on(tag) lh) % on(tag) lh", 0, 1), + ] { + let state = build_query_engine_state(); + let annotations = PromqlAnnotationCollector::default(); + let plan = PromPlanner::stmt_to_plan_with_annotations( + operator_table_provider(), + &operator_eval_stmt(query), + &state, + Some(annotations.clone()), + ) + .await + .unwrap(); + + let (_, batches) = execute(plan, &state).await; + assert_eq!( + batches.iter().map(RecordBatch::num_rows).sum::(), + expected_rows, + "{query}" + ); + let mut warnings = vec![]; + let mut infos = vec![]; + annotations.append_to(&mut warnings, &mut infos); + assert!(warnings.is_empty(), "{query}: {warnings:?}"); + assert_eq!(infos.len(), expected_infos, "{query}: {infos:?}"); + } + } + + #[tokio::test] + async fn test_mixed_or_can_feed_another_or() { + let state = build_query_engine_state(); + let plan = PromPlanner::stmt_to_plan( + operator_table_provider(), + &operator_eval_stmt("lf or on(tag) lh or on(tag) fallback"), + &state, + ) + .await + .unwrap(); + let float_field = plan + .schema() + .fields() + .iter() + .find(|field| field.name().starts_with(OR_FLOAT_FIELD_PREFIX)) + .unwrap() + .name() + .clone(); + let histogram_field = plan + .schema() + .fields() + .iter() + .find(|field| field.name().starts_with(OR_HISTOGRAM_FIELD_PREFIX)) + .unwrap() + .name() + .clone(); + + let (_, batches) = execute(plan, &state).await; + assert_eq!(batches.iter().map(RecordBatch::num_rows).sum::(), 3); + let mut float_values = values(&batches, &float_field); + float_values.sort_by(f64::total_cmp); + assert_eq!(float_values, vec![2.0, 7.0]); + assert_eq!(histograms(&batches, &histogram_field).len(), 1); + } + + #[tokio::test] + async fn test_mixed_fields_align_with_single_float_vector() { + let (planner, mixed) = mixed_direct_or(false).await; + let scale = tagged_source( + "scale", + false, + ("k", Some("float")), + DirectOrValue::Float64(2.0), + ); + let scale = scan(&scale); + let scale_fields = vec!["v".to_string()]; + let PromExpr::Binary(binary) = parser::parse("lhs * rhs").unwrap() else { + unreachable!() + }; + + let (groups, invalid_pairs) = PromPlanner::align_binary_field_columns( + mixed.schema(), + scale.schema(), + &planner.ctx.field_columns, + &scale_fields, + binary.op, + false, + false, + ); + assert!(invalid_pairs.is_empty()); + assert_eq!( + groups + .iter() + .map(|(output, _)| output.clone()) + .collect::>(), + planner.ctx.field_columns + ); + assert_eq!(groups.len(), 2); + assert!( + groups + .iter() + .flat_map(|(_, pairs)| pairs) + .all(|(_, right)| *right == &scale_fields[0]) + ); + + let (groups, invalid_pairs) = PromPlanner::align_binary_field_columns( + scale.schema(), + mixed.schema(), + &scale_fields, + &planner.ctx.field_columns, + binary.op, + false, + false, + ); + assert!(invalid_pairs.is_empty()); + assert_eq!( + groups + .iter() + .map(|(output, _)| output.clone()) + .collect::>(), + planner.ctx.field_columns + ); + assert_eq!(groups.len(), 2); + assert!( + groups + .iter() + .flat_map(|(_, pairs)| pairs) + .all(|(left, _)| *left == &scale_fields[0]) + ); + } + + #[tokio::test] + async fn test_non_bool_comparison_filters_mixed_sample_lanes() { + let (planner, input) = mixed_direct_or(false).await; + let input_schema = input.schema().clone(); + let plan = planner + .filter_on_field_column(input, |field| { + if PromPlanner::field_column_is_native_histogram(&input_schema, field) { + Ok(lit(false)) + } else { + Ok(col(field).gt(lit(0.0))) + } + }) + .unwrap(); + let float_field = planner.ctx.field_columns[0].clone(); + + let (_, batches) = execute(plan, &build_query_engine_state()).await; + assert_eq!(values(&batches, &float_field), vec![1.25]); + assert_eq!(batches.iter().map(RecordBatch::num_rows).sum::(), 1); + } + + #[tokio::test] + async fn test_mixed_left_and_unless_preserve_sample_lanes() { + for (expression, expected_sample_kind) in [ + ("lhs and on(k) mask", (false, true)), + ("lhs unless on(k) mask", (true, false)), + ] { + let (mut planner, left) = mixed_direct_or(false).await; + let left_context = planner.ctx.clone(); + let float_field = left_context.field_columns[0].clone(); + let histogram_field = left_context.field_columns[1].clone(); + let mask = tagged_source( + "mask", + false, + ("k", Some("histogram")), + DirectOrValue::Float64(1.0), + ); + let PromExpr::Binary(binary) = parser::parse(expression).unwrap() else { + unreachable!() + }; + let plan = planner + .set_op_on_non_field_columns( + left, + scan(&mask), + left_context, + direct_or_context("mask", &["job", "k"], "v"), + binary.op, + &binary.modifier, + ) + .unwrap(); + + let (_, batches) = execute(plan, &build_query_engine_state()).await; + let sample_kinds = batches + .iter() + .flat_map(|batch| { + let floats = batch.column_by_name(&float_field).unwrap(); + let histograms = batch.column_by_name(&histogram_field).unwrap(); + (0..batch.num_rows()) + .map(|row| (floats.is_valid(row), histograms.is_valid(row))) + }) + .collect::>(); + assert_eq!(sample_kinds, vec![expected_sample_kind], "{expression}"); + } + } + + #[tokio::test] + async fn test_mixed_fields_arithmetic_broadcasts_computed_scalar() { + let plan = PromPlanner::stmt_to_plan( + build_test_mixed_native_histogram_table_provider("some_metric").await, + &build_eval_stmt("some_metric * scalar(vector(2))"), + &build_query_engine_state(), + ) + .await + .unwrap(); + let schema = plan.schema(); + assert_eq!( + schema + .field_with_unqualified_name(greptime_value()) + .unwrap() + .data_type(), + &ArrowDataType::Float64 + ); + assert_eq!( + schema + .field_with_unqualified_name(greptime_native_histogram()) + .unwrap() + .data_type(), + &native_histogram_value_type().as_arrow_type() + ); + assert!( + plan.display_indent_schema() + .to_string() + .contains("prom_native_histogram_mul_scalar"), + "{plan:?}" + ); + } + + #[tokio::test] + async fn test_unsupported_histogram_binary_does_not_block_or_fallback() { + let state = build_query_engine_state(); + let plan = PromPlanner::stmt_to_plan( + operator_table_provider(), + &operator_eval_stmt("((lf or on(tag) lh) % 2) or on(tag) lh"), + &state, + ) + .await + .unwrap(); + let float_field = plan + .schema() + .fields() + .iter() + .find(|field| field.data_type() == &ArrowDataType::Float64) + .unwrap() + .name() + .clone(); + let histogram_field = plan + .schema() + .fields() + .iter() + .find(|field| field.data_type() == &native_histogram_value_type().as_arrow_type()) + .unwrap() + .name() + .clone(); + + let (_, batches) = execute(plan, &state).await; + assert_eq!(values(&batches, &float_field), vec![0.0]); + assert_eq!(histograms(&batches, &histogram_field).len(), 1); + } + + #[tokio::test] + async fn test_unary_negates_mixed_float_and_native_histogram_samples() { + for histogram_on_left in [false, true] { + let (mut planner, input) = mixed_direct_or(histogram_on_left).await; + let plan = planner.negate_field_columns(input).unwrap(); + assert!(PromPlanner::field_columns_are_alternative_samples( + plan.schema(), + &planner.ctx.field_columns + )); + let float_field = planner + .ctx + .field_columns + .iter() + .find(|field| field.starts_with(OR_FLOAT_FIELD_PREFIX)) + .unwrap(); + let histogram_field = planner + .ctx + .field_columns + .iter() + .find(|field| field.starts_with(OR_HISTOGRAM_FIELD_PREFIX)) + .unwrap(); + + let (_, batches) = execute(plan, &build_query_engine_state()).await; + assert_eq!(values(&batches, float_field), vec![-1.25]); + let histogram = batches + .iter() + .find_map(|batch| { + let values = batch + .column_by_name(histogram_field) + .unwrap() + .as_any() + .downcast_ref::() + .unwrap(); + (0..values.len()).find_map(|row| { + common_query::native_histogram::read_histogram(values, row).unwrap() + }) + }) + .unwrap(); + assert_eq!(histogram.count, -1.0); + assert_eq!(histogram.sum, -1.0); + assert_eq!(histogram.reset_hint, CounterResetHint::Gauge); + } + } + + #[tokio::test] + async fn test_mixed_or_value_aliases_do_not_replace_labels() { + let left = source( + "lhs", + false, + 1, + vec![("job", Some("job")), ("k", Some("float"))], + DirectOrValue::Float64(1.0), + ); + let right = source( + "rhs", + false, + 1, + vec![ + ("job", Some("job")), + ("k", Some("histogram")), + (greptime_value(), Some("value-label")), + ], + DirectOrValue::NativeHistogram(direct_or_histogram()), + ); + let table_provider = build_test_table_provider_with_fields( + &[(DEFAULT_SCHEMA_NAME.to_string(), "dummy".to_string())], + &[], + ) + .await; + let mut planner = PromPlanner { + table_provider, + ctx: PromPlannerContext::default(), + promql_annotations: None, + }; + let left = LogicalPlanBuilder::from(scan(&left)) + .project(vec![ + col("ts"), + col("job"), + col("k"), + col("v").alias(greptime_value()), + ]) + .unwrap() + .build() + .unwrap(); + let left_context = direct_or_context("lhs", &["job", "k"], greptime_value()); + let right_context = direct_or_context("rhs", &["job", "k", greptime_value()], "v"); + let plan = planner + .or_operator( + left, + scan(&right), + left_context.tag_columns.iter().cloned().collect(), + right_context.tag_columns.iter().cloned().collect(), + left_context, + right_context, + &or_modifier("lhs or on(k) rhs"), + ) + .unwrap(); + + assert_eq!( + plan.schema() + .field_with_name(None, greptime_value()) + .unwrap() + .data_type(), + &ArrowDataType::Utf8 + ); + assert!( + planner + .ctx + .field_columns + .iter() + .all(|field| { field != greptime_value() && field != greptime_native_histogram() }) + ); + assert!(PromPlanner::field_columns_are_alternative_samples( + plan.schema(), + &planner.ctx.field_columns + )); + let (_, batches) = execute(plan, &build_query_engine_state()).await; + assert_eq!(batches.iter().map(RecordBatch::num_rows).sum::(), 2); + let labels = batches + .iter() + .flat_map(|batch| { + batch + .column_by_name(greptime_value()) + .unwrap() + .as_any() + .downcast_ref::() + .unwrap() + .iter() + .flatten() + }) + .collect::>(); + assert_eq!(labels, vec!["value-label"]); + } + + #[tokio::test] + async fn test_mixed_or_routes_float_histogram_and_label_functions() { + for (function, expected) in [("abs", 1.25), ("round", 1.0), ("histogram_count", 1.0)] { + let (mut planner, input) = mixed_direct_or(false).await; + let preserve_any_value = PromPlanner::field_columns_are_alternative_samples( + input.schema(), + &planner.ctx.field_columns, + ); + let PromExpr::Call(call) = parser::parse(&format!("{function}(lhs)")).unwrap() else { + unreachable!() + }; + let state = build_query_engine_state(); + let (mut exprs, _) = planner + .create_function_expr(&call.func, vec![], input.schema(), &state) + .unwrap(); + exprs.insert(0, planner.create_time_index_column_expr().unwrap()); + exprs.extend(planner.create_tag_column_exprs().unwrap()); + let plan = LogicalPlanBuilder::from(input) + .project(exprs) + .unwrap() + .filter( + planner + .create_empty_values_filter_expr(preserve_any_value) + .unwrap(), + ) + .unwrap() + .build() + .unwrap(); + let (_, batches) = execute(plan, &state).await; + let values = batches + .iter() + .flat_map(|batch| { + batch + .schema() + .fields() + .iter() + .position(|field| field.data_type() == &ArrowDataType::Float64) + .map(|index| { + batch + .column(index) + .as_any() + .downcast_ref::() + .unwrap() + .iter() + .flatten() + }) + .into_iter() + .flatten() + }) + .collect::>(); + assert_eq!(values, vec![expected], "{function}"); + } + + let (mut planner, input) = mixed_direct_or(false).await; + let preserve_any_value = PromPlanner::field_columns_are_alternative_samples( + input.schema(), + &planner.ctx.field_columns, + ); + let PromExpr::Call(call) = + parser::parse(r#"label_replace(lhs, "copy", "$1", "k", "(.*)")"#).unwrap() + else { + unreachable!() + }; + let args = planner.create_function_args(&call.args.args).unwrap(); + let state = build_query_engine_state(); + let (mut exprs, _) = planner + .create_function_expr(&call.func, args.literals, input.schema(), &state) + .unwrap(); + exprs.insert(0, planner.create_time_index_column_expr().unwrap()); + exprs.extend(planner.create_tag_column_exprs().unwrap()); + let plan = LogicalPlanBuilder::from(input) + .project(exprs) + .unwrap() + .filter( + planner + .create_empty_values_filter_expr(preserve_any_value) + .unwrap(), + ) + .unwrap() + .build() + .unwrap(); + let (_, batches) = execute(plan, &state).await; + let sample_count = batches.iter().map(RecordBatch::num_rows).sum::(); + assert_eq!(sample_count, 2); + } }