Skip to main content

akar_processor/processor/
chunk_helpers.rs

1use 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}