akar_processor/processor/
chunk_helpers.rs1use akar_common::types::{PhysicalTypeID, Value};
2use akar_common::vector::ValueVector;
3
4pub fn extract_all_rows_from_chunks(chunks: &[akar_common::vector::DataChunk]) -> Vec<Vec<Value>> {
5 let mut all_rows = Vec::new();
6 for chunk in chunks {
7 let rows = extract_all_rows(chunk);
8 all_rows.extend(rows);
9 }
10 all_rows
11}
12
13pub fn extract_all_rows(chunk: &akar_common::vector::DataChunk) -> Vec<Vec<Value>> {
14 let num_rows = chunk.size;
15 let num_cols = chunk.fields.len();
16 let mut rows: Vec<Vec<Value>> = Vec::with_capacity(num_rows);
17 for row in 0..num_rows {
18 let mut values: Vec<Value> = Vec::with_capacity(num_cols);
19 for col in 0..num_cols {
20 let val = chunk.get_value(col, row).unwrap_or(Value::Null);
21 values.push(val);
22 }
23 rows.push(values);
24 }
25 rows
26}
27
28pub fn rows_to_columns(rows: &[Vec<Value>]) -> (Vec<arrow::array::ArrayRef>, Vec<PhysicalTypeID>) {
29 if rows.is_empty() {
30 return (Vec::new(), Vec::new());
31 }
32 let num_cols = rows[0].len();
33 let num_rows = rows.len();
34
35 let mut fields = Vec::with_capacity(num_cols);
36 let mut field_types = Vec::with_capacity(num_cols);
37
38 for col in 0..num_cols {
39 let first_val = &rows[0][col];
40 let phys_type = value_to_physical_type(first_val);
41 let mut vec = ValueVector::new(phys_type, num_rows.max(1));
42 for (row_idx, row) in rows.iter().enumerate() {
43 let val = &row[col];
44 let _ = vec.set_value(row_idx, val);
45 }
46 vec.resize(num_rows);
47 fields.push(akar_common::arrow_vector::ArrowVector::from_legacy(&vec).array);
48 field_types.push(phys_type);
49 }
50 (fields, field_types)
51}
52
53pub fn value_to_physical_type(val: &Value) -> PhysicalTypeID {
54 match val {
55 Value::Null => PhysicalTypeID::Any,
56 Value::Bool(_) => PhysicalTypeID::Bool,
57 Value::Int64(_) => PhysicalTypeID::Int64,
58 Value::Int32(_) => PhysicalTypeID::Int32,
59 Value::Int16(_) => PhysicalTypeID::Int16,
60 Value::Int8(_) => PhysicalTypeID::Int8,
61 Value::UInt64(_) => PhysicalTypeID::UInt64,
62 Value::UInt32(_) => PhysicalTypeID::UInt32,
63 Value::UInt16(_) => PhysicalTypeID::UInt16,
64 Value::UInt8(_) => PhysicalTypeID::UInt8,
65 Value::Int128(_) => PhysicalTypeID::Int128,
66 Value::Double(_) => PhysicalTypeID::Double,
67 Value::Float(_) => PhysicalTypeID::Float,
68 Value::String(_) | Value::Blob(_) => PhysicalTypeID::String,
69 Value::Date(_)
70 | Value::Timestamp(_)
71 | Value::TimestampTz(_)
72 | Value::TimestampNs(_)
73 | Value::TimestampMs(_)
74 | Value::TimestampSec(_)
75 | Value::Interval(_) => PhysicalTypeID::Int64,
76 Value::InternalID(_) => PhysicalTypeID::Int64,
77 Value::UInt128(_) => PhysicalTypeID::Int128,
78 Value::Json(_) => PhysicalTypeID::String,
79 Value::DTime(_) => PhysicalTypeID::Int64,
80 Value::Union(_, _) => PhysicalTypeID::Struct,
81 Value::List(_) => PhysicalTypeID::List,
82 Value::Map(_) => PhysicalTypeID::Struct,
83 Value::Struct(_) => PhysicalTypeID::Struct,
84 }
85}