Skip to main content

uqa_execution/
columnar_batch.rs

1//
2// Unified Query Algebra
3//
4// Copyright (c) 2023-2026 Cognica, Inc.
5//
6
7//! Positional, column-oriented batches for public result transfer.
8
9use std::collections::BTreeMap;
10
11use uqa_core::Value;
12use uqa_sql::ResultRow;
13
14use crate::{Batch, RowSchema};
15
16/// A column vector preserves its position in the declared output schema.
17///
18/// Duplicate labels remain distinct because conversion from a physical batch
19/// reads each logical position directly. The named-row compatibility
20/// constructor cannot recover values that were already collapsed by a map.
21#[derive(Debug, Clone, PartialEq)]
22pub struct ColumnVector {
23    pub name: String,
24    pub values: Vec<Value>,
25}
26
27/// Column-oriented transfer batch. Physical operators may remain row-oriented
28/// internally; API consumers can process one bounded batch without retaining
29/// the complete result set or rebuilding columns themselves.
30#[derive(Debug, Clone, PartialEq)]
31pub struct ColumnarBatch {
32    columns: Vec<ColumnVector>,
33    row_count: usize,
34}
35
36impl ColumnarBatch {
37    /// Convert a positional physical batch without crossing the legacy
38    /// map-backed row boundary. Repeated output labels are matched to repeated
39    /// logical schema positions in order, so differently-valued duplicate
40    /// columns remain distinct for wire and cursor consumers.
41    pub fn from_batch(schema: &[String], batch: Batch) -> Self {
42        let row_count = batch.rows.len();
43        let mut occurrences = BTreeMap::<&str, usize>::new();
44        let positions = schema
45            .iter()
46            .map(|name| {
47                let occurrence = occurrences.entry(name.as_str()).or_default();
48                let position = batch
49                    .schema
50                    .columns()
51                    .iter()
52                    .enumerate()
53                    .filter(|(_, input)| input == &name)
54                    .nth(*occurrence)
55                    .map(|(position, _)| position);
56                *occurrence += 1;
57                position
58            })
59            .collect::<Vec<_>>();
60        let mut columns = schema
61            .iter()
62            .map(|name| ColumnVector {
63                name: name.clone(),
64                values: Vec::with_capacity(row_count),
65            })
66            .collect::<Vec<_>>();
67        for row in &batch.rows {
68            let view = batch.schema.view(row);
69            for (column, position) in columns.iter_mut().zip(&positions) {
70                column.values.push(
71                    position
72                        .and_then(|position| view.value_at(position))
73                        .cloned()
74                        .unwrap_or(Value::Null),
75                );
76            }
77        }
78        Self { columns, row_count }
79    }
80
81    /// Move map-backed rows into positional columns. A missing projected value
82    /// has SQL `NULL` semantics.
83    pub fn from_rows(schema: &[String], rows: Vec<ResultRow>) -> Self {
84        let row_count = rows.len();
85        let occurrences = schema.iter().fold(BTreeMap::new(), |mut counts, name| {
86            *counts.entry(name.as_str()).or_insert(0_usize) += 1;
87            counts
88        });
89        let mut columns = schema
90            .iter()
91            .map(|name| ColumnVector {
92                name: name.clone(),
93                values: Vec::with_capacity(row_count),
94            })
95            .collect::<Vec<_>>();
96        if occurrences.values().all(|count| *count == 1) {
97            for mut row in rows {
98                for column in &mut columns {
99                    column
100                        .values
101                        .push(row.remove(&column.name).unwrap_or(Value::Null));
102                }
103            }
104            return Self { columns, row_count };
105        }
106        for mut row in rows {
107            let mut remaining = occurrences.clone();
108            for column in &mut columns {
109                let remaining_for_name = remaining
110                    .get_mut(column.name.as_str())
111                    .expect("column occurrence was counted");
112                *remaining_for_name -= 1;
113                let value = if *remaining_for_name == 0 {
114                    row.remove(&column.name)
115                } else {
116                    row.get(&column.name).cloned()
117                };
118                column.values.push(value.unwrap_or(Value::Null));
119            }
120        }
121        Self { columns, row_count }
122    }
123
124    pub fn schema(&self) -> RowSchema {
125        RowSchema::new(
126            self.columns
127                .iter()
128                .map(|column| column.name.clone())
129                .collect(),
130        )
131    }
132
133    pub fn columns(&self) -> &[ColumnVector] {
134        &self.columns
135    }
136
137    pub fn len(&self) -> usize {
138        self.row_count
139    }
140
141    pub fn is_empty(&self) -> bool {
142        self.row_count == 0
143    }
144
145    /// Convert the column vectors to row-major positional values without
146    /// passing through a name-keyed map.
147    pub fn into_positional_rows(self) -> Vec<Vec<Value>> {
148        let mut rows = (0..self.row_count)
149            .map(|_| Vec::with_capacity(self.columns.len()))
150            .collect::<Vec<_>>();
151        for column in self.columns {
152            for (row, value) in rows.iter_mut().zip(column.values) {
153                row.push(value);
154            }
155        }
156        rows
157    }
158
159    /// Convert back to the legacy named-row representation. Duplicate output
160    /// labels necessarily collapse because [`ResultRow`] is a map.
161    pub fn into_rows(self) -> Vec<ResultRow> {
162        let mut rows = (0..self.row_count)
163            .map(|_| BTreeMap::new())
164            .collect::<Vec<ResultRow>>();
165        for column in self.columns {
166            for (row, value) in rows.iter_mut().zip(column.values) {
167                row.insert(column.name.clone(), value);
168            }
169        }
170        rows
171    }
172}
173
174#[cfg(test)]
175mod tests {
176    use super::*;
177
178    #[test]
179    fn conversion_is_positional_and_fills_missing_values_with_null() {
180        let mut first = ResultRow::new();
181        first.insert("a".into(), Value::Int(1));
182        first.insert("b".into(), Value::Str("x".into()));
183        let mut second = ResultRow::new();
184        second.insert("a".into(), Value::Int(2));
185
186        let batch = ColumnarBatch::from_rows(&["b".into(), "a".into()], vec![first, second]);
187        assert_eq!(batch.len(), 2);
188        assert_eq!(
189            batch.columns()[0].values,
190            vec![Value::Str("x".into()), Value::Null]
191        );
192        assert_eq!(
193            batch.columns()[1].values,
194            vec![Value::Int(1), Value::Int(2)]
195        );
196    }
197
198    #[test]
199    fn duplicate_schema_labels_remain_visible_as_separate_slots() {
200        let mut row = ResultRow::new();
201        row.insert("value".into(), Value::Int(7));
202        let batch = ColumnarBatch::from_rows(&["value".into(), "value".into()], vec![row]);
203        assert_eq!(batch.columns().len(), 2);
204        assert_eq!(batch.columns()[0].values, batch.columns()[1].values);
205    }
206
207    #[test]
208    fn physical_batch_preserves_different_duplicate_values() {
209        let schema = RowSchema::new(vec!["value".into(), "value".into()]);
210        let batch = Batch::from_physical_rows(
211            schema,
212            vec![crate::PhysicalRow::from_values(vec![
213                Value::Int(1),
214                Value::Int(2),
215            ])],
216        );
217        let batch = ColumnarBatch::from_batch(&["value".into(), "value".into()], batch);
218        assert_eq!(batch.columns()[0].values, [Value::Int(1)]);
219        assert_eq!(batch.columns()[1].values, [Value::Int(2)]);
220    }
221
222    #[test]
223    fn positional_rows_preserve_duplicate_values() {
224        let schema = RowSchema::new(vec!["value".into(), "value".into()]);
225        let batch = Batch::from_physical_rows(
226            schema,
227            vec![crate::PhysicalRow::from_values(vec![
228                Value::Int(1),
229                Value::Int(2),
230            ])],
231        );
232        assert_eq!(
233            ColumnarBatch::from_batch(&["value".into(), "value".into()], batch)
234                .into_positional_rows(),
235            [vec![Value::Int(1), Value::Int(2)]]
236        );
237    }
238}