uqa_execution/
columnar_batch.rs1use std::collections::BTreeMap;
10
11use uqa_core::Value;
12use uqa_sql::ResultRow;
13
14use crate::{Batch, RowSchema};
15
16#[derive(Debug, Clone, PartialEq)]
22pub struct ColumnVector {
23 pub name: String,
24 pub values: Vec<Value>,
25}
26
27#[derive(Debug, Clone, PartialEq)]
31pub struct ColumnarBatch {
32 columns: Vec<ColumnVector>,
33 row_count: usize,
34}
35
36impl ColumnarBatch {
37 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 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 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 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}