1use std::collections::{BTreeSet, HashMap, HashSet};
18use std::sync::Arc;
19
20use datafusion::logical_expr::{Cast, Extension, LogicalPlan, LogicalPlanBuilder};
21use datafusion::prelude::{Column, Expr as DfExpr, JoinType};
22use datafusion::scalar::ScalarValue;
23use datafusion_common::{NullEquality, TableReference};
24use datatypes::arrow::datatypes::DataType as ArrowDataType;
25use promql::extension_plan::UnionDistinctOn;
26use promql_parser::parser::token::{self, TokenType};
27use promql_parser::parser::{BinModifier, LabelModifier, VectorMatchCardinality};
28use snafu::{OptionExt, ResultExt, ensure};
29use store_api::metric_engine_consts::DATA_SCHEMA_TSID_COLUMN_NAME;
30
31use crate::promql::error::{
32 ColumnNotFoundSnafu, CombineTableColumnMismatchSnafu, DataFusionPlanningSnafu,
33 MultiFieldsNotSupportedSnafu, Result, TimeIndexNotFoundSnafu, UnexpectedPlanExprSnafu,
34 UnexpectedTokenSnafu, UnsupportedVectorMatchSnafu,
35};
36use crate::promql::planner::{
37 OR_FLOAT_FIELD_PREFIX, OR_HISTOGRAM_FIELD_PREFIX, PromPlanner, PromPlannerContext,
38};
39
40impl PromPlanner {
41 #[allow(clippy::too_many_arguments)]
43 pub(super) fn or_operator(
44 &mut self,
45 left: LogicalPlan,
46 right: LogicalPlan,
47 left_tag_cols_set: HashSet<String>,
48 right_tag_cols_set: HashSet<String>,
49 left_context: PromPlannerContext,
50 right_context: PromPlannerContext,
51 modifier: &Option<BinModifier>,
52 ) -> Result<LogicalPlan> {
53 let left_is_empty = Self::is_zero_row_empty_relation(&left);
54 let right_is_empty = Self::is_zero_row_empty_relation(&right);
55 match (left_is_empty, right_is_empty) {
56 (true, false) => {
57 self.ctx = right_context;
58 return Ok(right);
59 }
60 (false, true) => {
61 self.ctx = left_context;
62 return Ok(left);
63 }
64 (true, true) => {
65 self.ctx = left_context;
66 return Ok(left);
67 }
68 (false, false) => {}
69 }
70
71 ensure!(
72 !left.schema().fields().is_empty() && !right.schema().fields().is_empty(),
73 UnexpectedPlanExprSnafu {
74 desc: "OR operator input has zero columns",
75 }
76 );
77 let left_has_alternative_samples =
78 Self::field_columns_are_alternative_samples(left.schema(), &left_context.field_columns);
79 let right_has_alternative_samples = Self::field_columns_are_alternative_samples(
80 right.schema(),
81 &right_context.field_columns,
82 );
83 ensure!(
84 left_context.field_columns.len() == 1 || left_has_alternative_samples,
85 MultiFieldsNotSupportedSnafu {
86 operator: "OR operator"
87 }
88 );
89 ensure!(
90 right_context.field_columns.len() == 1 || right_has_alternative_samples,
91 MultiFieldsNotSupportedSnafu {
92 operator: "OR operator"
93 }
94 );
95
96 let all_tags = left_tag_cols_set
98 .union(&right_tag_cols_set)
99 .cloned()
100 .collect::<HashSet<_>>();
101 let left_qualifier = left.schema().qualified_field(0).0.cloned();
102 let right_qualifier = right.schema().qualified_field(0).0.cloned();
103 let left_qualifier_string = left_qualifier
104 .as_ref()
105 .map(|l| l.to_string())
106 .unwrap_or_default();
107 let right_qualifier_string = right_qualifier
108 .as_ref()
109 .map(|r| r.to_string())
110 .unwrap_or_default();
111 let left_time_index_column =
112 left_context
113 .time_index_column
114 .clone()
115 .with_context(|| TimeIndexNotFoundSnafu {
116 table: left_qualifier_string.clone(),
117 })?;
118 let right_time_index_column =
119 right_context
120 .time_index_column
121 .clone()
122 .with_context(|| TimeIndexNotFoundSnafu {
123 table: right_qualifier_string.clone(),
124 })?;
125 let native_histogram_type = Self::native_histogram_arrow_type();
126 let is_numeric = |data_type: &ArrowDataType| {
127 matches!(
128 data_type,
129 ArrowDataType::Int8
130 | ArrowDataType::Int16
131 | ArrowDataType::Int32
132 | ArrowDataType::Int64
133 | ArrowDataType::UInt8
134 | ArrowDataType::UInt16
135 | ArrowDataType::UInt32
136 | ArrowDataType::UInt64
137 | ArrowDataType::Float32
138 | ArrowDataType::Float64
139 )
140 };
141 let left_fields = left_context
142 .field_columns
143 .iter()
144 .map(|name| {
145 left.schema()
146 .iter()
147 .find(|(_, field)| field.name() == name)
148 .map(|(qualifier, field)| {
149 (name.clone(), qualifier.cloned(), field.data_type().clone())
150 })
151 .with_context(|| ColumnNotFoundSnafu { col: name.clone() })
152 })
153 .collect::<Result<Vec<_>>>()?;
154 let right_fields = right_context
155 .field_columns
156 .iter()
157 .map(|name| {
158 right
159 .schema()
160 .iter()
161 .find(|(_, field)| field.name() == name)
162 .map(|(qualifier, field)| {
163 (name.clone(), qualifier.cloned(), field.data_type().clone())
164 })
165 .with_context(|| ColumnNotFoundSnafu { col: name.clone() })
166 })
167 .collect::<Result<Vec<_>>>()?;
168 let left_field = &left_fields[0];
169 let right_field = &right_fields[0];
170 let left_field_col = &left_field.0;
171 let right_field_col = &right_field.0;
172 let fields_are_samples = |fields: &[(String, Option<TableReference>, ArrowDataType)]| {
173 fields.iter().all(|(_, _, data_type)| {
174 is_numeric(data_type) || data_type == &native_histogram_type
175 })
176 };
177 let mixed_sample_types = if left_has_alternative_samples || right_has_alternative_samples {
178 if !fields_are_samples(&left_fields) || !fields_are_samples(&right_fields) {
179 return UnexpectedPlanExprSnafu {
180 desc: format!(
181 "OR value fields have incompatible types: {:?} and {:?}",
182 left_fields
183 .iter()
184 .map(|(_, _, data_type)| data_type)
185 .collect::<Vec<_>>(),
186 right_fields
187 .iter()
188 .map(|(_, _, data_type)| data_type)
189 .collect::<Vec<_>>()
190 ),
191 }
192 .fail();
193 }
194 true
195 } else {
196 (left_field.2 == native_histogram_type && is_numeric(&right_field.2))
197 || (right_field.2 == native_histogram_type && is_numeric(&left_field.2))
198 };
199 let target_field_type = if mixed_sample_types {
200 ArrowDataType::Float64
203 } else if left_field.2 == right_field.2 {
204 left_field.2.clone()
205 } else if is_numeric(&left_field.2) && is_numeric(&right_field.2) {
206 ArrowDataType::Float64
207 } else {
208 return UnexpectedPlanExprSnafu {
209 desc: format!(
210 "OR value fields have incompatible types: {:?} and {:?}",
211 left_field.2, right_field.2
212 ),
213 }
214 .fail();
215 };
216 let (mixed_float_field_col, mixed_histogram_field_col) = if mixed_sample_types {
217 let mut reserved_names = left
218 .schema()
219 .fields()
220 .iter()
221 .chain(right.schema().fields().iter())
222 .map(|field| field.name().clone())
223 .collect::<HashSet<_>>();
224 for (name, _, _) in left_fields.iter().chain(&right_fields) {
225 reserved_names.remove(name);
226 }
227 reserved_names.extend(all_tags.iter().cloned());
228 let unique_name = |prefix: &str, reserved_names: &mut HashSet<String>| {
229 let mut index = 0;
230 loop {
231 let name = format!("{prefix}{index}");
232 index += 1;
233 if reserved_names.insert(name.clone()) {
234 break name;
235 }
236 }
237 };
238 let float_field = unique_name(OR_FLOAT_FIELD_PREFIX, &mut reserved_names);
239 let histogram_field = unique_name(OR_HISTOGRAM_FIELD_PREFIX, &mut reserved_names);
240 (float_field, histogram_field)
241 } else {
242 (left_field_col.clone(), String::new())
243 };
244 let left_tag_types = left_tag_cols_set
245 .iter()
246 .map(|label| {
247 left.schema()
248 .fields()
249 .iter()
250 .find(|field| field.name() == label)
251 .map(|field| (label.clone(), field.data_type().clone()))
252 .with_context(|| ColumnNotFoundSnafu { col: label.clone() })
253 })
254 .collect::<Result<HashMap<_, _>>>()?;
255 let right_tag_types = right_tag_cols_set
256 .iter()
257 .map(|label| {
258 right
259 .schema()
260 .fields()
261 .iter()
262 .find(|field| field.name() == label)
263 .map(|field| (label.clone(), field.data_type().clone()))
264 .with_context(|| ColumnNotFoundSnafu { col: label.clone() })
265 })
266 .collect::<Result<HashMap<_, _>>>()?;
267 let mut target_tag_types = HashMap::with_capacity(all_tags.len());
268 for label in &all_tags {
269 let Some(data_type) =
270 Self::common_label_data_type(left_tag_types.get(label), right_tag_types.get(label))
271 else {
272 return UnexpectedPlanExprSnafu {
273 desc: format!(
274 "OR label {label} has incompatible types: {:?} and {:?}",
275 left_tag_types.get(label),
276 right_tag_types.get(label)
277 ),
278 }
279 .fail();
280 };
281 target_tag_types.insert(label.clone(), data_type);
282 }
283 let left_has_tsid = left
284 .schema()
285 .fields()
286 .iter()
287 .any(|field| field.name() == DATA_SCHEMA_TSID_COLUMN_NAME);
288 let right_has_tsid = right
289 .schema()
290 .fields()
291 .iter()
292 .any(|field| field.name() == DATA_SCHEMA_TSID_COLUMN_NAME);
293
294 let mut all_columns_set = left
296 .schema()
297 .fields()
298 .iter()
299 .chain(right.schema().fields().iter())
300 .map(|field| field.name().clone())
301 .collect::<HashSet<_>>();
302 if !(left_has_tsid && right_has_tsid) {
305 all_columns_set.remove(DATA_SCHEMA_TSID_COLUMN_NAME);
306 }
307 all_columns_set.remove(&left_time_index_column);
309 all_columns_set.remove(&right_time_index_column);
310 if mixed_sample_types {
311 for (name, _, _) in left_fields.iter().chain(&right_fields) {
312 all_columns_set.remove(name);
313 }
314 all_columns_set.extend(all_tags.iter().cloned());
315 all_columns_set.insert(mixed_float_field_col.clone());
316 all_columns_set.insert(mixed_histogram_field_col.clone());
317 } else if left_field_col != right_field_col {
318 all_columns_set.remove(right_field_col);
320 }
321 let mut all_columns = all_columns_set.into_iter().collect::<Vec<_>>();
322 all_columns.sort_unstable();
324 all_columns.insert(0, left_time_index_column.clone());
326 let mut occupied_column_names = left
327 .schema()
328 .fields()
329 .iter()
330 .chain(right.schema().fields().iter())
331 .map(|field| field.name().clone())
332 .collect::<HashSet<_>>();
333
334 let aligned_label_expr = |col: &String, source_types: &HashMap<String, ArrowDataType>| {
336 let target_type = &target_tag_types[col];
337 if let Some(source_type) = source_types.get(col) {
338 let expr = DfExpr::Column(Column::new(None::<String>, col));
339 if source_type == target_type {
340 expr
341 } else {
342 DfExpr::Cast(Cast::new(Box::new(expr), target_type.clone())).alias(col.clone())
343 }
344 } else {
345 DfExpr::Literal(
346 Self::string_scalar_value(target_type, None)
347 .expect("target label type is a string"),
348 None,
349 )
350 .alias(col.clone())
351 }
352 };
353 let null_histogram =
354 ScalarValue::try_new_null(&native_histogram_type).context(DataFusionPlanningSnafu)?;
355 let mixed_value_expr = |fields: &[(String, Option<TableReference>, ArrowDataType)],
356 output_col: &String| {
357 if output_col == &mixed_float_field_col {
358 if let Some((name, qualifier, data_type)) = fields
359 .iter()
360 .find(|(_, _, data_type)| is_numeric(data_type))
361 {
362 let expr = DfExpr::Column(Column::new(qualifier.clone(), name));
363 if data_type == &ArrowDataType::Float64 {
364 expr.alias(output_col)
365 } else {
366 DfExpr::Cast(Cast::new(Box::new(expr), ArrowDataType::Float64))
367 .alias(output_col)
368 }
369 } else {
370 DfExpr::Literal(ScalarValue::Float64(None), None).alias(output_col)
371 }
372 } else {
373 fields
374 .iter()
375 .find(|(_, _, data_type)| data_type == &native_histogram_type)
376 .map(|(name, qualifier, _)| {
377 DfExpr::Column(Column::new(qualifier.clone(), name)).alias(output_col)
378 })
379 .unwrap_or_else(|| {
380 DfExpr::Literal(null_histogram.clone(), None).alias(output_col)
381 })
382 }
383 };
384 let left_proj_exprs = all_columns.iter().map(|col| {
385 if mixed_sample_types
386 && (col == &mixed_float_field_col || col == &mixed_histogram_field_col)
387 {
388 mixed_value_expr(&left_fields, col)
389 } else if !mixed_sample_types
390 && col == left_field_col
391 && left_field.2 != target_field_type
392 {
393 DfExpr::Cast(Cast::new(
394 Box::new(DfExpr::Column(Column::new(
395 left_field.1.clone(),
396 left_field_col,
397 ))),
398 target_field_type.clone(),
399 ))
400 .alias(left_field_col.clone())
401 } else if target_tag_types.contains_key(col) {
402 aligned_label_expr(col, &left_tag_types)
403 } else {
404 DfExpr::Column(Column::new(None::<String>, col))
405 }
406 });
407 let right_time_index_expr = DfExpr::Column(Column::new(
408 right_qualifier.clone(),
409 right_time_index_column,
410 ))
411 .alias(left_time_index_column.clone());
412 let right_proj_exprs_without_time_index = all_columns.iter().skip(1).map(|col| {
416 if mixed_sample_types
418 && (col == &mixed_float_field_col || col == &mixed_histogram_field_col)
419 {
420 mixed_value_expr(&right_fields, col)
421 } else if !mixed_sample_types && col == left_field_col {
422 let expr = DfExpr::Column(Column::new(right_field.1.clone(), right_field_col));
423 if right_field.2 != target_field_type {
424 DfExpr::Cast(Cast::new(Box::new(expr), target_field_type.clone()))
425 .alias(left_field_col.clone())
426 } else if left_field_col != right_field_col {
427 expr.alias(left_field_col.clone())
428 } else {
429 expr
430 }
431 } else if target_tag_types.contains_key(col) {
432 aligned_label_expr(col, &right_tag_types)
433 } else {
434 DfExpr::Column(Column::new(None::<String>, col))
435 }
436 });
437 let right_proj_exprs = [right_time_index_expr]
438 .into_iter()
439 .chain(right_proj_exprs_without_time_index);
440
441 let left_projected = LogicalPlanBuilder::from(left)
442 .project(left_proj_exprs)
443 .context(DataFusionPlanningSnafu)?
444 .alias(left_qualifier_string.clone())
445 .context(DataFusionPlanningSnafu)?
446 .build()
447 .context(DataFusionPlanningSnafu)?;
448 let right_projected = LogicalPlanBuilder::from(right)
449 .project(right_proj_exprs)
450 .context(DataFusionPlanningSnafu)?
451 .alias(right_qualifier_string.clone())
452 .context(DataFusionPlanningSnafu)?
453 .build()
454 .context(DataFusionPlanningSnafu)?;
455
456 let mut match_columns = if let Some(modifier) = modifier
458 && let Some(matching) = &modifier.matching
459 {
460 match matching {
461 LabelModifier::Include(on) => on.labels.clone(),
463 LabelModifier::Exclude(ignoring) => {
465 let ignoring = ignoring.labels.iter().cloned().collect::<HashSet<_>>();
466 all_tags.difference(&ignoring).cloned().collect()
467 }
468 }
469 } else {
470 all_tags.iter().cloned().collect()
471 };
472 match_columns.sort_unstable();
474 match_columns.dedup();
475 occupied_column_names.extend(
476 left_projected
477 .schema()
478 .fields()
479 .iter()
480 .chain(right_projected.schema().fields().iter())
481 .map(|field| field.name().clone()),
482 );
483
484 let visible_schema = left_projected.schema().clone();
485 let visible_left_exprs = left_projected
486 .schema()
487 .iter()
488 .map(|(qualifier, field)| {
489 DfExpr::Column(Column::new(qualifier.cloned(), field.name().clone()))
490 })
491 .collect::<Vec<_>>();
492 let visible_right_exprs = right_projected
493 .schema()
494 .iter()
495 .map(|(qualifier, field)| {
496 DfExpr::Column(Column::new(qualifier.cloned(), field.name().clone()))
497 })
498 .collect::<Vec<_>>();
499 let mut left_match_exprs = Vec::with_capacity(match_columns.len());
500 let mut right_match_exprs = Vec::with_capacity(match_columns.len());
501 let mut next_internal_column = 0;
502
503 for label in &match_columns {
504 let left_field = if left_tag_cols_set.contains(label) {
505 Some(
506 left_projected
507 .schema()
508 .iter()
509 .find(|(_, field)| field.name() == label)
510 .map(|(qualifier, field)| (qualifier.cloned(), field.data_type().clone()))
511 .with_context(|| ColumnNotFoundSnafu { col: label.clone() })?,
512 )
513 } else {
514 None
515 };
516 let right_field = if right_tag_cols_set.contains(label) {
517 Some(
518 right_projected
519 .schema()
520 .iter()
521 .find(|(_, field)| field.name() == label)
522 .map(|(qualifier, field)| (qualifier.cloned(), field.data_type().clone()))
523 .with_context(|| ColumnNotFoundSnafu { col: label.clone() })?,
524 )
525 } else {
526 None
527 };
528 let data_type = match (left_field.as_ref(), right_field.as_ref()) {
529 (Some((_, left_type)), Some((_, right_type))) if left_type == right_type => {
530 left_type.clone()
531 }
532 (Some((_, left_type)), Some((_, right_type))) => {
533 return UnexpectedPlanExprSnafu {
534 desc: format!(
535 "OR match label {label} has incompatible types: {left_type:?} and {right_type:?}"
536 ),
537 }
538 .fail();
539 }
540 (Some((_, data_type)), None) | (None, Some((_, data_type))) => data_type.clone(),
541 (None, None) => ArrowDataType::Utf8,
542 };
543 let Some(value_type) = Self::string_value_data_type(&data_type).cloned() else {
544 return UnexpectedPlanExprSnafu {
545 desc: format!("OR match label {label} must be a string"),
546 }
547 .fail();
548 };
549 let internal_name = loop {
550 let name = format!("__promql_or_match_{next_internal_column}");
551 next_internal_column += 1;
552 if occupied_column_names.insert(name.clone()) {
553 break name;
554 }
555 };
556 left_match_exprs.push(Self::normalized_match_key_expr(
557 label,
558 left_field,
559 &value_type,
560 &internal_name,
561 ));
562 right_match_exprs.push(Self::normalized_match_key_expr(
563 label,
564 right_field,
565 &value_type,
566 &internal_name,
567 ));
568 }
569
570 let left_augmented = LogicalPlanBuilder::from(left_projected)
571 .project(visible_left_exprs.into_iter().chain(left_match_exprs))
572 .context(DataFusionPlanningSnafu)?
573 .build()
574 .context(DataFusionPlanningSnafu)?;
575 let right_augmented = LogicalPlanBuilder::from(right_projected)
576 .project(visible_right_exprs.into_iter().chain(right_match_exprs))
577 .context(DataFusionPlanningSnafu)?
578 .build()
579 .context(DataFusionPlanningSnafu)?;
580
581 let visible_field_count = visible_schema.fields().len();
583 let compare_key_indices =
584 (visible_field_count..visible_field_count + match_columns.len()).collect::<Vec<_>>();
585 let (time_qualifier, _) = visible_schema
586 .iter()
587 .find(|(_, field)| field.name() == &left_time_index_column)
588 .with_context(|| TimeIndexNotFoundSnafu {
589 table: left_qualifier_string.clone(),
590 })?;
591 let ts_col_idx = left_augmented
592 .schema()
593 .iter()
594 .position(|(qualifier, field)| {
595 qualifier == time_qualifier && field.name() == &left_time_index_column
596 })
597 .with_context(|| TimeIndexNotFoundSnafu {
598 table: left_qualifier_string.clone(),
599 })?;
600 let union_distinct_on = UnionDistinctOn::try_new(
601 left_augmented,
602 right_augmented,
603 compare_key_indices,
604 ts_col_idx,
605 )
606 .context(DataFusionPlanningSnafu)?;
607 let augmented_result = LogicalPlan::Extension(Extension {
608 node: Arc::new(union_distinct_on),
609 });
610 let result = LogicalPlanBuilder::from(augmented_result)
611 .project(visible_schema.iter().map(|(qualifier, field)| {
612 DfExpr::Column(Column::new(qualifier.cloned(), field.name().clone()))
613 }))
614 .context(DataFusionPlanningSnafu)?
615 .build()
616 .context(DataFusionPlanningSnafu)?;
617
618 let output_field_col = left_field_col.clone();
620 let mut output_context = left_context;
621 let mut visible_tags = all_tags.into_iter().collect::<Vec<_>>();
622 visible_tags.sort_unstable();
623 output_context.time_index_column = Some(left_time_index_column);
624 output_context.tag_columns = visible_tags;
625 output_context.field_columns = if mixed_sample_types {
626 vec![mixed_float_field_col, mixed_histogram_field_col]
627 } else {
628 vec![output_field_col]
629 };
630 output_context.use_tsid = left_has_tsid && right_has_tsid;
631 self.ctx = output_context;
632
633 Ok(result)
634 }
635
636 pub(super) fn set_op_on_non_field_columns(
638 &mut self,
639 mut left: LogicalPlan,
640 mut right: LogicalPlan,
641 left_context: PromPlannerContext,
642 right_context: PromPlannerContext,
643 op: TokenType,
644 modifier: &Option<BinModifier>,
645 ) -> Result<LogicalPlan> {
646 let left_tag_col_set = left_context
647 .tag_columns
648 .iter()
649 .cloned()
650 .collect::<HashSet<_>>();
651 let right_tag_col_set = right_context
652 .tag_columns
653 .iter()
654 .cloned()
655 .collect::<HashSet<_>>();
656
657 if matches!(op.id(), token::T_LOR) {
658 return self.or_operator(
659 left,
660 right,
661 left_tag_col_set,
662 right_tag_col_set,
663 left_context,
664 right_context,
665 modifier,
666 );
667 }
668
669 if let Some(modifier) = modifier {
670 ensure!(
671 matches!(
672 modifier.card,
673 VectorMatchCardinality::OneToOne | VectorMatchCardinality::ManyToMany
674 ),
675 UnsupportedVectorMatchSnafu {
676 name: modifier.card.clone(),
677 },
678 );
679 }
680
681 let output_context = left_context.clone();
682 let visible_left_schema = left.schema().clone();
683 let mut left_context = left_context;
684 let mut right_context = right_context;
685 let added_marker_to_left = if Self::only_temporality_match_label_mismatches(
686 &left_context,
687 &right_context,
688 modifier,
689 ) {
690 let aligned = Self::align_temporality_match_column(
691 left,
692 right,
693 &mut left_context,
694 &mut right_context,
695 )?;
696 left = aligned.0;
697 right = aligned.1;
698 aligned.2
699 } else {
700 false
701 };
702
703 let mut left_tag_col_set = left_context
704 .tag_columns
705 .iter()
706 .cloned()
707 .collect::<BTreeSet<_>>();
708 let mut right_tag_col_set = right_context
709 .tag_columns
710 .iter()
711 .cloned()
712 .collect::<BTreeSet<_>>();
713 if let Some(matching) = modifier
714 .as_ref()
715 .and_then(|modifier| modifier.matching.as_ref())
716 {
717 match matching {
718 LabelModifier::Include(on) => {
719 let mask = on.labels.iter().cloned().collect::<BTreeSet<_>>();
720 left_tag_col_set = left_tag_col_set.intersection(&mask).cloned().collect();
721 right_tag_col_set = right_tag_col_set.intersection(&mask).cloned().collect();
722 }
723 LabelModifier::Exclude(ignoring) => {
724 for label in &ignoring.labels {
725 let _ = left_tag_col_set.remove(label);
726 let _ = right_tag_col_set.remove(label);
727 }
728 }
729 }
730 }
731 ensure!(
732 left_tag_col_set == right_tag_col_set,
733 CombineTableColumnMismatchSnafu {
734 left: left_tag_col_set.iter().cloned().collect::<Vec<_>>(),
735 right: right_tag_col_set.iter().cloned().collect::<Vec<_>>(),
736 }
737 );
738
739 let left_time_index = left_context.time_index_column.clone().unwrap();
740 let right_time_index = right_context.time_index_column.clone().unwrap();
741
742 if left_context.time_index_column != right_context.time_index_column {
744 let right_project_exprs = right
745 .schema()
746 .fields()
747 .iter()
748 .map(|field| {
749 if field.name() == &right_time_index {
750 DfExpr::Column(Column::from_name(&right_time_index)).alias(&left_time_index)
751 } else {
752 DfExpr::Column(Column::from_name(field.name()))
753 }
754 })
755 .collect::<Vec<_>>();
756
757 right = LogicalPlanBuilder::from(right)
758 .project(right_project_exprs)
759 .context(DataFusionPlanningSnafu)?
760 .build()
761 .context(DataFusionPlanningSnafu)?;
762 }
763
764 let join_keys = left_tag_col_set
765 .into_iter()
766 .chain([left_time_index])
767 .map(Column::from_name)
768 .collect::<Vec<_>>();
769
770 ensure!(
771 left_context.field_columns.len() == 1
772 || Self::field_columns_are_alternative_samples(
773 left.schema(),
774 &left_context.field_columns,
775 ),
776 MultiFieldsNotSupportedSnafu {
777 operator: "AND/UNLESS operator"
778 }
779 );
780 let result = match op.id() {
783 token::T_LAND => LogicalPlanBuilder::from(left)
784 .distinct()
785 .context(DataFusionPlanningSnafu)?
786 .join_detailed(
787 right,
788 JoinType::LeftSemi,
789 (join_keys.clone(), join_keys),
790 None,
791 NullEquality::NullEqualsNull,
792 )
793 .context(DataFusionPlanningSnafu)?
794 .build()
795 .context(DataFusionPlanningSnafu),
796 token::T_LUNLESS => LogicalPlanBuilder::from(left)
797 .distinct()
798 .context(DataFusionPlanningSnafu)?
799 .join_detailed(
800 right,
801 JoinType::LeftAnti,
802 (join_keys.clone(), join_keys),
803 None,
804 NullEquality::NullEqualsNull,
805 )
806 .context(DataFusionPlanningSnafu)?
807 .build()
808 .context(DataFusionPlanningSnafu),
809 token::T_LOR => {
810 unreachable!()
813 }
814 _ => UnexpectedTokenSnafu { token: op }.fail(),
815 }?;
816 let result = if added_marker_to_left {
817 LogicalPlanBuilder::from(result)
818 .project(visible_left_schema.iter().map(|(qualifier, field)| {
819 DfExpr::Column(Column::new(qualifier.cloned(), field.name().clone()))
820 }))
821 .context(DataFusionPlanningSnafu)?
822 .build()
823 .context(DataFusionPlanningSnafu)?
824 } else {
825 result
826 };
827
828 self.ctx = output_context;
831 Ok(result)
832 }
833}