Skip to main content

promql/
extension_plan.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15mod absent;
16mod empty_metric;
17mod histogram_fold;
18mod instant_manipulate;
19mod normalize;
20mod planner;
21mod range_manipulate;
22mod scalar_calculate;
23mod series_divide;
24#[cfg(test)]
25mod test_util;
26mod union_distinct_on;
27
28pub use absent::{Absent, AbsentExec, AbsentStream};
29use common_query::native_histogram::{SUM_FIELD, native_histogram_value_type};
30use common_query::prometheus::is_prometheus_stale_nan;
31use datafusion::arrow::array::{Array, Float64Array, StructArray};
32use datafusion::arrow::datatypes::{
33    ArrowPrimitiveType, DataType, TimeUnit, TimestampMillisecondType,
34};
35use datafusion::common::{Column, DFSchemaRef};
36use datafusion::error::{DataFusionError, Result as DataFusionResult};
37use datafusion::logical_expr::{Expr, Extension, LogicalPlan};
38use datatypes::data_type::DataType as _;
39pub use empty_metric::{EmptyMetric, EmptyMetricExec, EmptyMetricStream, build_special_time_expr};
40pub use histogram_fold::{
41    HistogramFold, HistogramFoldExec, HistogramFoldOperation, HistogramFoldStream,
42};
43pub use instant_manipulate::{InstantManipulate, InstantManipulateExec, InstantManipulateStream};
44pub use normalize::{SeriesNormalize, SeriesNormalizeExec, SeriesNormalizeStream};
45pub use planner::PromExtensionPlanner;
46pub use range_manipulate::{RangeManipulate, RangeManipulateExec, RangeManipulateStream};
47pub use scalar_calculate::ScalarCalculate;
48pub use series_divide::{SeriesDivide, SeriesDivideExec, SeriesDivideStream};
49pub use union_distinct_on::{UnionDistinctOn, UnionDistinctOnExec, UnionDistinctOnStream};
50
51pub type Millisecond = <TimestampMillisecondType as ArrowPrimitiveType>::Native;
52
53pub(crate) fn timestamp_unit(data_type: &DataType) -> datafusion::error::Result<TimeUnit> {
54    match data_type {
55        DataType::Timestamp(unit, _) => Ok(*unit),
56        _ => Err(datafusion::error::DataFusionError::Execution(
57            "Time index column is not a timestamp".into(),
58        )),
59    }
60}
61
62pub(crate) fn nanoseconds_per_native_tick(unit: TimeUnit) -> i128 {
63    match unit {
64        TimeUnit::Second => 1_000_000_000,
65        TimeUnit::Millisecond => 1_000_000,
66        TimeUnit::Microsecond => 1_000,
67        TimeUnit::Nanosecond => 1,
68    }
69}
70
71/// Recovers the offset serialized only by an immediately underlying normalize node.
72///
73/// This is decode-only recovery for manipulators whose wire messages have no offset field.
74/// Follow identity projections (as used by `timestamp()`), but stop at other nodes or changed
75/// time columns to avoid applying an inner selector's offset again to an outer subquery.
76pub(crate) fn local_offset(plan: &LogicalPlan, time_index: &str) -> Millisecond {
77    let Some(index) = plan.schema().index_of_column_by_name(None, time_index) else {
78        return 0;
79    };
80    let (qualifier, field) = plan.schema().qualified_field(index);
81    let mut time_index = Column::new(qualifier.cloned(), field.name().clone());
82    let mut plan = plan;
83
84    loop {
85        match plan {
86            LogicalPlan::Extension(Extension { node }) => {
87                return node
88                    .as_any()
89                    .downcast_ref::<SeriesNormalize>()
90                    .and_then(|normalize| normalize.offset_for_time_index(&time_index))
91                    .unwrap_or_default();
92            }
93            LogicalPlan::Projection(projection) => {
94                let Some(output_index) = projection.schema.maybe_index_of_column(&time_index)
95                else {
96                    return 0;
97                };
98                let expr = &projection.expr[output_index];
99                let source = match expr {
100                    Expr::Column(column) => column,
101                    Expr::Alias(alias) => {
102                        let Expr::Column(column) = alias.expr.as_ref() else {
103                            return 0;
104                        };
105                        if alias.name != column.name {
106                            return 0;
107                        }
108                        column
109                    }
110                    _ => return 0,
111                };
112                let Some(input_index) = projection.input.schema().maybe_index_of_column(source)
113                else {
114                    return 0;
115                };
116                let (qualifier, field) = projection.input.schema().qualified_field(input_index);
117                time_index = Column::new(qualifier.cloned(), field.name().clone());
118                plan = projection.input.as_ref();
119            }
120            _ => return 0,
121        }
122    }
123}
124
125const METRIC_NUM_SERIES: &str = "num_series";
126
127fn prometheus_stale_sample_column(column: &dyn Array) -> Option<(&dyn Array, &Float64Array)> {
128    let values = if let Some(values) = column.as_any().downcast_ref::<Float64Array>() {
129        values
130    } else {
131        let histograms = column.as_any().downcast_ref::<StructArray>()?;
132        if histograms.data_type() != &native_histogram_value_type().as_arrow_type() {
133            return None;
134        }
135        histograms
136            .column_by_name(SUM_FIELD)?
137            .as_any()
138            .downcast_ref::<Float64Array>()?
139    };
140    Some((column, values))
141}
142
143fn is_prometheus_stale_sample((column, values): (&dyn Array, &Float64Array), row: usize) -> bool {
144    column.is_valid(row) && values.is_valid(row) && is_prometheus_stale_nan(values.value(row))
145}
146
147/// Utilities for handling unfix logic in extension plans
148/// Convert column name to index for serialization
149pub fn serialize_column_index(schema: &DFSchemaRef, column_name: &str) -> u64 {
150    schema
151        .index_of_column_by_name(None, column_name)
152        .map(|idx| idx as u64)
153        .unwrap_or(u64::MAX) // make sure if not found, it will report error in deserialization
154}
155
156/// Convert index back to column name for deserialization
157pub fn resolve_column_name(
158    index: u64,
159    schema: &DFSchemaRef,
160    context: &str,
161    column_type: &str,
162) -> DataFusionResult<String> {
163    let columns = schema.columns();
164    columns
165        .get(index as usize)
166        .ok_or_else(|| {
167            DataFusionError::Internal(format!(
168                "Failed to get {} column at idx {} during unfixing {} with columns:{:?}",
169                column_type, index, context, columns
170            ))
171        })
172        .map(|field| field.name().to_string())
173}
174
175/// Batch process multiple column indices
176pub fn resolve_column_names(
177    indices: &[u64],
178    schema: &DFSchemaRef,
179    context: &str,
180    column_type: &str,
181) -> DataFusionResult<Vec<String>> {
182    indices
183        .iter()
184        .map(|idx| resolve_column_name(*idx, schema, context, column_type))
185        .collect()
186}
187
188#[cfg(test)]
189mod tests {
190    use std::sync::Arc;
191
192    use datafusion::arrow::datatypes::{DataType, Field, Schema, TimeUnit};
193    use datafusion::common::ToDFSchema;
194    use datafusion::logical_expr::{EmptyRelation, Extension, LogicalPlan, Projection};
195    use datafusion_expr::col;
196
197    use super::*;
198
199    fn input() -> LogicalPlan {
200        LogicalPlan::EmptyRelation(EmptyRelation {
201            produce_one_row: false,
202            schema: Arc::new(Schema::new(vec![
203                Field::new(
204                    "timestamp",
205                    DataType::Timestamp(TimeUnit::Millisecond, None),
206                    false,
207                ),
208                Field::new(
209                    "other_ts",
210                    DataType::Timestamp(TimeUnit::Millisecond, None),
211                    false,
212                ),
213                Field::new("value", DataType::Float64, true),
214            ]))
215            .to_dfschema_ref()
216            .unwrap(),
217        })
218    }
219
220    fn normalized() -> LogicalPlan {
221        LogicalPlan::Extension(Extension {
222            node: Arc::new(SeriesNormalize::new(
223                1_000,
224                "timestamp",
225                false,
226                Vec::new(),
227                input(),
228            )),
229        })
230    }
231
232    #[test]
233    fn local_offset_tracks_identity_preserving_projections() {
234        let projection =
235            Projection::try_new(vec![col("timestamp"), col("value")], Arc::new(normalized()))
236                .unwrap();
237        let projection = Projection::try_new(
238            vec![col("timestamp").alias("timestamp"), col("value")],
239            Arc::new(LogicalPlan::Projection(projection)),
240        )
241        .unwrap();
242
243        assert_eq!(
244            1_000,
245            local_offset(&LogicalPlan::Projection(projection), "timestamp")
246        );
247    }
248
249    #[test]
250    fn local_offset_rejects_a_different_timestamp_or_manipulator() {
251        let renamed = Projection::try_new(
252            vec![col("other_ts").alias("timestamp"), col("value")],
253            Arc::new(normalized()),
254        )
255        .unwrap();
256        assert_eq!(
257            0,
258            local_offset(&LogicalPlan::Projection(renamed), "timestamp")
259        );
260
261        let divide = LogicalPlan::Extension(Extension {
262            node: Arc::new(SeriesDivide::new(
263                Vec::new(),
264                "timestamp".to_string(),
265                normalized(),
266            )),
267        });
268        assert_eq!(0, local_offset(&divide, "timestamp"));
269    }
270}