Skip to main content

query/promql/planner/
island.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
15//! The binary-island fast path of the PromQL planner.
16
17use std::collections::{BTreeSet, HashMap};
18
19use common_query::prelude::OTLP_AGGREGATION_TEMPORALITY_LABEL;
20use datafusion::common::DFSchemaRef;
21use datafusion::logical_expr::expr::Alias;
22use datafusion::logical_expr::{LogicalPlan, LogicalPlanBuilder};
23use datafusion::prelude::{Column, Expr as DfExpr, JoinType};
24use datafusion_common::{NullEquality, TableReference};
25use datafusion_expr::lit;
26use promql_parser::label::{METRIC_NAME, MatchOp, Matcher};
27use promql_parser::parser::token::{self, TokenType};
28use promql_parser::parser::{
29    BinaryExpr as PromBinaryExpr, Expr as PromExpr, Offset, ParenExpr, UnaryExpr,
30    VectorMatchCardinality, VectorSelector,
31};
32use snafu::ResultExt;
33
34use crate::promql::error::{DataFusionPlanningSnafu, Result};
35use crate::promql::planner::{PromPlanner, PromPlannerContext};
36
37/// Prefix for generated binary island leaf aliases.
38const BINARY_ISLAND_LEAF_ALIAS_PREFIX: &str = "__prom_v";
39
40#[derive(Debug, Clone, PartialEq, Eq, Hash)]
41struct VectorLeafKey {
42    metric_name: String,
43    matchers: Vec<(String, String, String)>,
44    or_matchers: Vec<Vec<(String, String, String)>>,
45    offset_ms: i128,
46    at: String,
47}
48
49#[derive(Debug, Clone)]
50struct IslandLeaf {
51    selector: VectorSelector,
52    display_table: String,
53}
54
55#[derive(Debug, Clone)]
56enum IslandExpr {
57    VectorLeaf(usize),
58    Scalar(DfExpr),
59    Unary {
60        input: Box<IslandExpr>,
61    },
62    Binary {
63        op: TokenType,
64        lhs: Box<IslandExpr>,
65        rhs: Box<IslandExpr>,
66    },
67}
68
69impl IslandExpr {
70    fn try_new(expr: &PromExpr, env: &mut IslandCollectEnv) -> Option<Self> {
71        if let Some(expr) = PromPlanner::try_build_literal_expr(expr) {
72            return Some(Self::Scalar(expr));
73        }
74
75        match expr {
76            PromExpr::Paren(ParenExpr { expr }) => Self::try_new(expr, env),
77            PromExpr::VectorSelector(selector) => {
78                let leaf = env.intern_leaf(selector)?;
79                Some(Self::VectorLeaf(leaf))
80            }
81            PromExpr::Unary(UnaryExpr { expr }) => {
82                let input = Self::try_new(expr, env)?;
83                Some(Self::Unary {
84                    input: Box::new(input),
85                })
86            }
87            PromExpr::Binary(PromBinaryExpr {
88                lhs,
89                rhs,
90                op,
91                modifier,
92            }) if matches!(
93                op.id(),
94                token::T_ADD
95                    | token::T_SUB
96                    | token::T_MUL
97                    | token::T_DIV
98                    | token::T_MOD
99                    | token::T_POW
100                    | token::T_ATAN2
101            ) && modifier.as_ref().is_none_or(|modifier| {
102                !modifier.return_bool
103                    && modifier.matching.is_none()
104                    && matches!(modifier.card, VectorMatchCardinality::OneToOne)
105                    && modifier.fill_values.lhs.is_none()
106                    && modifier.fill_values.rhs.is_none()
107            }) =>
108            {
109                let lhs = Self::try_new(lhs, env)?;
110                let rhs = Self::try_new(rhs, env)?;
111                Some(Self::Binary {
112                    op: *op,
113                    lhs: Box::new(lhs),
114                    rhs: Box::new(rhs),
115                })
116            }
117            _ => None,
118        }
119    }
120}
121
122#[derive(Debug, Default)]
123struct IslandCollectEnv {
124    leaf_by_key: HashMap<VectorLeafKey, usize>,
125    leaves: Vec<IslandLeaf>,
126    vector_occurrences: usize,
127}
128
129#[derive(Debug)]
130struct PlannedIslandLeaf {
131    plan: LogicalPlan,
132    ctx: PromPlannerContext,
133    alias: TableReference,
134    display_table: String,
135}
136
137#[derive(Debug)]
138struct IslandFieldExprs {
139    exprs: Vec<DfExpr>,
140    names: Vec<String>,
141    scalar: bool,
142}
143
144impl VectorLeafKey {
145    fn from_selector(selector: &VectorSelector) -> Option<Self> {
146        let mut metric_name = selector.name.clone();
147        let mut matchers = Vec::with_capacity(selector.matchers.matchers.len());
148        let matcher_key = |matcher: &Matcher| {
149            (
150                matcher.name.clone(),
151                matcher.op.to_string(),
152                matcher.value.clone(),
153            )
154        };
155
156        for matcher in &selector.matchers.matchers {
157            if matcher.name == METRIC_NAME {
158                if matcher.op != MatchOp::Equal || metric_name.is_some() {
159                    return None;
160                }
161                metric_name = Some(matcher.value.clone());
162            } else {
163                matchers.push(matcher_key(matcher));
164            }
165        }
166        matchers.sort();
167
168        let mut or_matchers = selector
169            .matchers
170            .or_matchers
171            .iter()
172            .map(|group| {
173                let mut group = group.iter().map(matcher_key).collect::<Vec<_>>();
174                group.sort();
175                group
176            })
177            .collect::<Vec<_>>();
178        or_matchers.sort();
179
180        Some(Self {
181            metric_name: metric_name?,
182            matchers,
183            or_matchers,
184            offset_ms: match &selector.offset {
185                Some(Offset::Pos(duration)) => duration.as_millis() as i128,
186                Some(Offset::Neg(duration)) => -(duration.as_millis() as i128),
187                None => 0,
188            },
189            at: format!("{:?}", selector.at),
190        })
191    }
192}
193
194impl IslandCollectEnv {
195    fn intern_leaf(&mut self, selector: &VectorSelector) -> Option<usize> {
196        self.vector_occurrences += 1;
197        let key = VectorLeafKey::from_selector(selector)?;
198        if let Some(id) = self.leaf_by_key.get(&key) {
199            return Some(*id);
200        }
201
202        let id = self.leaves.len();
203        self.leaves.push(IslandLeaf {
204            selector: selector.clone(),
205            display_table: key.metric_name.clone(),
206        });
207        self.leaf_by_key.insert(key, id);
208        Some(id)
209    }
210}
211
212impl PromPlanner {
213    pub(super) async fn try_plan_binary_island(
214        &mut self,
215        binary_expr: &PromBinaryExpr,
216    ) -> Result<Option<LogicalPlan>> {
217        let original_ctx = self.ctx.clone();
218        let mut collect_env = IslandCollectEnv::default();
219        let Some(island_expr) =
220            IslandExpr::try_new(&PromExpr::Binary(binary_expr.clone()), &mut collect_env)
221        else {
222            return Ok(None);
223        };
224
225        if collect_env.leaves.is_empty()
226            || collect_env.vector_occurrences <= collect_env.leaves.len()
227        {
228            return Ok(None);
229        }
230
231        let mut planned_leaves = Vec::with_capacity(collect_env.leaves.len());
232        for (idx, leaf) in collect_env.leaves.iter().enumerate() {
233            let plan = self
234                .prom_vector_selector_to_plan(&leaf.selector, false)
235                .await?;
236            let ctx = self.ctx.clone();
237            let alias = TableReference::bare(format!("{BINARY_ISLAND_LEAF_ALIAS_PREFIX}{idx}"));
238            let plan = LogicalPlanBuilder::from(plan)
239                .alias(alias.clone())
240                .context(DataFusionPlanningSnafu)?
241                .build()
242                .context(DataFusionPlanningSnafu)?;
243            planned_leaves.push(PlannedIslandLeaf {
244                plan,
245                ctx,
246                alias,
247                display_table: leaf.display_table.clone(),
248            });
249        }
250
251        if planned_leaves.iter().any(|leaf| {
252            Self::field_columns_contain_native_histogram(
253                leaf.plan.schema(),
254                &leaf.ctx.field_columns,
255            )
256        }) {
257            self.ctx = original_ctx;
258            return Ok(None);
259        }
260
261        if !Self::binary_island_join_contexts_supported(&planned_leaves) {
262            self.ctx = original_ctx;
263            return Ok(None);
264        }
265
266        let mut input = planned_leaves[0].plan.clone();
267        for right_idx in 1..planned_leaves.len() {
268            input = self.join_binary_island_leaf(
269                input,
270                &planned_leaves[0],
271                &planned_leaves[right_idx],
272            )?;
273        }
274
275        let field_exprs =
276            Self::build_binary_island_field_exprs(&island_expr, &planned_leaves, input.schema())?;
277        if field_exprs.scalar || field_exprs.exprs.is_empty() {
278            self.ctx = original_ctx;
279            return Ok(None);
280        }
281
282        let plan = self.project_binary_island(
283            input,
284            &planned_leaves[0].alias,
285            &planned_leaves[0].ctx,
286            field_exprs,
287        )?;
288        Ok(Some(plan))
289    }
290
291    fn binary_island_join_contexts_supported(leaves: &[PlannedIslandLeaf]) -> bool {
292        if leaves
293            .iter()
294            .any(|leaf| leaf.ctx.time_index_column.is_none())
295        {
296            return false;
297        }
298
299        if leaves.len() <= 1 {
300            return true;
301        }
302
303        let first_tags = leaves[0].ctx.tag_columns.iter().collect::<BTreeSet<_>>();
304
305        leaves.iter().skip(1).all(|leaf| {
306            (Self::plan_has_tsid_column(&leaves[0].plan) && Self::plan_has_tsid_column(&leaf.plan))
307                || leaf.ctx.tag_columns.iter().collect::<BTreeSet<_>>() == first_tags
308        })
309    }
310
311    fn join_binary_island_leaf(
312        &self,
313        left: LogicalPlan,
314        first_leaf: &PlannedIslandLeaf,
315        right_leaf: &PlannedIslandLeaf,
316    ) -> Result<LogicalPlan> {
317        let only_join_time_index = (first_leaf.ctx.tag_columns.is_empty()
318            || right_leaf.ctx.tag_columns.is_empty())
319            && !first_leaf
320                .ctx
321                .tag_columns
322                .iter()
323                .chain(&right_leaf.ctx.tag_columns)
324                .any(|tag| tag == OTLP_AGGREGATION_TEMPORALITY_LABEL);
325        let (mut left_keys, mut right_keys, force_empty_join) = self.binary_join_key_columns(
326            left.schema(),
327            right_leaf.plan.schema(),
328            &first_leaf.ctx,
329            &right_leaf.ctx,
330            only_join_time_index,
331            &None,
332        )?;
333
334        if let (Some(left_time_index_column), Some(right_time_index_column)) = (
335            first_leaf.ctx.time_index_column.clone(),
336            right_leaf.ctx.time_index_column.clone(),
337        ) {
338            left_keys.insert(left_time_index_column);
339            right_keys.insert(right_time_index_column);
340        }
341
342        LogicalPlanBuilder::from(left)
343            .join_detailed(
344                right_leaf.plan.clone(),
345                JoinType::Inner,
346                (
347                    left_keys
348                        .into_iter()
349                        .map(|name| Column::new(Some(first_leaf.alias.clone()), name))
350                        .collect::<Vec<_>>(),
351                    right_keys
352                        .into_iter()
353                        .map(|name| Column::new(Some(right_leaf.alias.clone()), name))
354                        .collect::<Vec<_>>(),
355                ),
356                force_empty_join.then_some(lit(false)),
357                NullEquality::NullEqualsNull,
358            )
359            .context(DataFusionPlanningSnafu)?
360            .build()
361            .context(DataFusionPlanningSnafu)
362    }
363
364    fn build_binary_island_field_exprs(
365        expr: &IslandExpr,
366        leaves: &[PlannedIslandLeaf],
367        schema: &DFSchemaRef,
368    ) -> Result<IslandFieldExprs> {
369        match expr {
370            IslandExpr::VectorLeaf(id) => {
371                let leaf = &leaves[*id];
372                let exprs = leaf
373                    .ctx
374                    .field_columns
375                    .iter()
376                    .map(|field| {
377                        schema
378                            .qualified_field_with_name(Some(&leaf.alias), field)
379                            .context(DataFusionPlanningSnafu)
380                            .map(|field| DfExpr::Column(field.into()))
381                    })
382                    .collect::<Result<Vec<_>>>()?;
383                let names = leaf
384                    .ctx
385                    .field_columns
386                    .iter()
387                    .map(|field| format!("{}.{}", leaf.display_table, field))
388                    .collect();
389                Ok(IslandFieldExprs {
390                    exprs,
391                    names,
392                    scalar: false,
393                })
394            }
395            IslandExpr::Scalar(expr) => Ok(IslandFieldExprs {
396                exprs: vec![expr.clone()],
397                names: vec![expr.schema_name().to_string()],
398                scalar: true,
399            }),
400            IslandExpr::Unary { input } => {
401                let input = Self::build_binary_island_field_exprs(input, leaves, schema)?;
402                let mut exprs = Vec::with_capacity(input.exprs.len());
403                let mut names = Vec::with_capacity(input.names.len());
404                for (expr, name) in input.exprs.into_iter().zip(input.names) {
405                    exprs.push(DfExpr::Negative(Box::new(expr)));
406                    names.push(format!("-{name}"));
407                }
408                Ok(IslandFieldExprs {
409                    exprs,
410                    names,
411                    scalar: input.scalar,
412                })
413            }
414            IslandExpr::Binary { op, lhs, rhs } => {
415                let same_leaf = match (&**lhs, &**rhs) {
416                    (IslandExpr::VectorLeaf(left), IslandExpr::VectorLeaf(right))
417                        if left == right =>
418                    {
419                        Some(*left)
420                    }
421                    _ => None,
422                };
423                let lhs = Self::build_binary_island_field_exprs(lhs, leaves, schema)?;
424                let rhs = Self::build_binary_island_field_exprs(rhs, leaves, schema)?;
425                let expr_builder = Self::prom_token_to_binary_expr_builder(*op)?;
426                let scalar = lhs.scalar && rhs.scalar;
427                let op = op.to_string();
428
429                let (exprs, names) = match (lhs.scalar, rhs.scalar) {
430                    (true, true) => {
431                        let expr = expr_builder(lhs.exprs[0].clone(), rhs.exprs[0].clone())?;
432                        let name = format!("{} {op} {}", lhs.names[0], rhs.names[0]);
433                        (vec![expr], vec![name])
434                    }
435                    (true, false) => {
436                        let mut exprs = Vec::with_capacity(rhs.exprs.len());
437                        let mut names = Vec::with_capacity(rhs.names.len());
438                        for (rhs_expr, rhs_name) in rhs.exprs.into_iter().zip(rhs.names) {
439                            exprs.push(expr_builder(lhs.exprs[0].clone(), rhs_expr)?);
440                            names.push(format!("{} {op} {rhs_name}", lhs.names[0]));
441                        }
442                        (exprs, names)
443                    }
444                    (false, true) => {
445                        let mut exprs = Vec::with_capacity(lhs.exprs.len());
446                        let mut names = Vec::with_capacity(lhs.names.len());
447                        for (lhs_expr, lhs_name) in lhs.exprs.into_iter().zip(lhs.names) {
448                            exprs.push(expr_builder(lhs_expr, rhs.exprs[0].clone())?);
449                            names.push(format!("{lhs_name} {op} {}", rhs.names[0]));
450                        }
451                        (exprs, names)
452                    }
453                    (false, false) => {
454                        let mut exprs = Vec::new();
455                        let mut names = Vec::new();
456                        for (idx, ((lhs_expr, rhs_expr), (mut lhs_name, mut rhs_name))) in lhs
457                            .exprs
458                            .into_iter()
459                            .zip(rhs.exprs)
460                            .zip(lhs.names.into_iter().zip(rhs.names))
461                            .enumerate()
462                        {
463                            if let Some(leaf) = same_leaf {
464                                let field = leaves[leaf]
465                                    .ctx
466                                    .field_columns
467                                    .get(idx)
468                                    .cloned()
469                                    .unwrap_or_else(|| lhs_name.clone());
470                                lhs_name = format!("lhs.{field}");
471                                rhs_name = format!("rhs.{field}");
472                            }
473                            exprs.push(expr_builder(lhs_expr, rhs_expr)?);
474                            names.push(format!("{lhs_name} {op} {rhs_name}"));
475                        }
476                        (exprs, names)
477                    }
478                };
479
480                Ok(IslandFieldExprs {
481                    exprs,
482                    names,
483                    scalar,
484                })
485            }
486        }
487    }
488
489    fn project_binary_island(
490        &mut self,
491        input: LogicalPlan,
492        base_alias: &TableReference,
493        base_ctx: &PromPlannerContext,
494        field_exprs: IslandFieldExprs,
495    ) -> Result<LogicalPlan> {
496        self.ctx = base_ctx.clone();
497
498        let schema = input.schema();
499        let non_field_exprs = base_ctx
500            .tag_columns
501            .iter()
502            .chain(base_ctx.time_index_column.iter())
503            .map(|column| {
504                schema
505                    .qualified_field_with_name(Some(base_alias), column)
506                    .context(DataFusionPlanningSnafu)
507                    .map(|field| DfExpr::Column(field.into()))
508            });
509        let tsid_expr = Self::optional_tsid_projection(schema, Some(base_alias), base_ctx.use_tsid)
510            .into_iter()
511            .map(Ok);
512
513        self.ctx.field_columns = field_exprs.names;
514        let field_exprs = field_exprs
515            .exprs
516            .into_iter()
517            .zip(self.ctx.field_columns.iter())
518            .map(|(expr, name)| Ok(DfExpr::Alias(Alias::new(expr, None::<String>, name))));
519
520        let project_exprs = non_field_exprs
521            .chain(tsid_expr)
522            .chain(field_exprs)
523            .collect::<Result<Vec<_>>>()?;
524
525        let plan = LogicalPlanBuilder::from(input)
526            .project(project_exprs)
527            .context(DataFusionPlanningSnafu)?
528            .build()
529            .context(DataFusionPlanningSnafu)?;
530
531        self.ctx.table_name = None;
532        self.ctx.schema_name = None;
533
534        Ok(plan)
535    }
536}