Skip to main content

akar_processor/processor/
union_helpers.rs

1use crate::processor::chunk_helpers::{extract_all_rows_from_chunks, rows_to_columns};
2use akar_common::error::ProcessorError;
3use akar_common::types::Value;
4use akar_common::vector::DataChunk;
5use akar_planner::logical_operator::LogicalOperator;
6
7pub fn flatten_union_child(op: &LogicalOperator) -> Vec<LogicalOperator> {
8    match op {
9        LogicalOperator::Projection(p) if p.expressions.is_empty() => p.children.clone(),
10        other => vec![other.clone()],
11    }
12}
13
14pub fn merge_union_chunks(
15    left: Vec<DataChunk>,
16    right: Vec<DataChunk>,
17    all: bool,
18) -> Result<Vec<DataChunk>, ProcessorError> {
19    if left.is_empty() {
20        return Ok(right);
21    }
22    if right.is_empty() {
23        return Ok(left);
24    }
25
26    let num_fields = left[0].num_fields();
27    for chunk in &right {
28        if chunk.num_fields() != num_fields {
29            return Err(format!(
30                "UNION column count mismatch: left has {num_fields} columns, right has {} columns",
31                chunk.num_fields()
32            )
33            .into());
34        }
35    }
36
37    let mut left_rows = extract_all_rows_from_chunks(&left);
38    let right_rows = extract_all_rows_from_chunks(&right);
39    left_rows.extend(right_rows);
40
41    let mut deduped: Vec<Vec<Value>> = Vec::with_capacity(left_rows.len());
42    if !all {
43        for row in &left_rows {
44            if !deduped.contains(row) {
45                deduped.push(row.clone());
46            }
47        }
48    } else {
49        deduped = left_rows;
50    }
51
52    if deduped.is_empty() {
53        return Ok(vec![DataChunk::new(vec![], vec![])]);
54    }
55
56    let (fields, field_types) = rows_to_columns(&deduped);
57    let final_size = deduped.len();
58    let field_names = left.first().map(|c| c.field_names.clone()).unwrap_or_default();
59
60    Ok(vec![DataChunk {
61        fields,
62        field_types,
63        size: final_size,
64        field_names,
65        sel_vector: None,
66    }])
67}
68
69pub fn merge_optional_chunks(left: Vec<DataChunk>, right: Vec<DataChunk>) -> Result<Vec<DataChunk>, ProcessorError> {
70    if left.is_empty() {
71        return Ok(left);
72    }
73    if right.is_empty() {
74        return Ok(left);
75    }
76
77    let left_rows = extract_all_rows_from_chunks(&left);
78    let right_rows = extract_all_rows_from_chunks(&right);
79
80    let num_left_cols = left_rows.first().map(|r| r.len()).unwrap_or(0);
81    let num_right_cols = right_rows.first().map(|r| r.len()).unwrap_or(0);
82    let max_rows = left_rows.len();
83
84    let mut combined: Vec<Vec<Value>> = Vec::with_capacity(max_rows);
85    for i in 0..max_rows {
86        let mut row = Vec::with_capacity(num_left_cols + num_right_cols);
87        if i < left_rows.len() {
88            row.extend_from_slice(&left_rows[i]);
89        }
90        if i < right_rows.len() {
91            row.extend_from_slice(&right_rows[i]);
92        } else {
93            row.extend(std::iter::repeat_n(Value::Null, num_right_cols));
94        }
95        combined.push(row);
96    }
97
98    if combined.is_empty() {
99        return Ok(vec![]);
100    }
101
102    let (fields, field_types) = rows_to_columns(&combined);
103    let size = combined.len();
104    Ok(vec![DataChunk {
105        fields,
106        field_types,
107        size,
108        field_names: vec![],
109        sel_vector: None,
110    }])
111}