1use 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
37const 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}