Skip to main content

petaplot_core/storage/
arrow_layout.rs

1use arrow_array::{Array, Float32Array, Float64Array, RecordBatch};
2use crate::error::{Result, TeraError};
3
4/// Abstracción zero-copy para una columna de serie temporal contigua en memoria.
5pub enum TimeSeriesColumn<'a> {
6    F32(&'a [f32]),
7    F64(&'a [f64]),
8    OwningF32(Vec<f32>),
9}
10
11impl<'a> TimeSeriesColumn<'a> {
12    /// Retorna la cantidad total de muestras en la columna.
13    pub fn len(&self) -> usize {
14        match self {
15            Self::F32(slice) => slice.len(),
16            Self::F64(slice) => slice.len(),
17            Self::OwningF32(vec) => vec.len(),
18        }
19    }
20
21    /// Indica si la columna está vacía.
22    pub fn is_empty(&self) -> bool {
23        self.len() == 0
24    }
25
26    /// Intenta extraer un slice contiguo `&[f32]`. Si los datos originales son `f64`, realiza una conversión eficiente.
27    pub fn to_f32_slice(&'a self) -> &'a [f32] {
28        match self {
29            Self::F32(slice) => slice,
30            Self::OwningF32(vec) => vec.as_slice(),
31            Self::F64(_) => panic!("Llama a convert_f64_to_f32() para obtener una representación f32"),
32        }
33    }
34}
35
36/// Extrae de forma zero-copy una columna numérica de un `RecordBatch` de Apache Arrow.
37pub fn extract_float_column<'a>(
38    batch: &'a RecordBatch,
39    column_index: usize,
40) -> Result<TimeSeriesColumn<'a>> {
41    if column_index >= batch.num_columns() {
42        return Err(TeraError::Arrow(format!(
43            "Índice de columna {} fuera de rango (columnas disponibles: {})",
44            column_index,
45            batch.num_columns()
46        )));
47    }
48
49    let column = batch.column(column_index);
50
51    if let Some(arr) = column.as_any().downcast_ref::<Float32Array>() {
52        let values_slice: &[f32] = arr.values().as_ref();
53        Ok(TimeSeriesColumn::F32(values_slice))
54    } else if let Some(arr) = column.as_any().downcast_ref::<Float64Array>() {
55        let values_slice: &[f64] = arr.values().as_ref();
56        Ok(TimeSeriesColumn::F64(values_slice))
57    } else {
58        Err(TeraError::InvalidLayout(format!(
59            "La columna en índice {} no es del tipo Float32 o Float64",
60            column_index
61        )))
62    }
63}
64
65/// Convierte directamente un slice de bytes alineados en un `&[f32]` sin copia (Zero-Copy Transmute).
66pub fn cast_bytes_to_f32(bytes: &[u8]) -> Result<&[f32]> {
67    if bytes.len() % std::mem::size_of::<f32>() != 0 {
68        return Err(TeraError::InvalidLayout(format!(
69            "El tamaño de los bytes ({}) no es múltiplo de size_of::<f32>() (4 bytes)",
70            bytes.len()
71        )));
72    }
73
74    let (head, body, tail) = unsafe { bytes.align_to::<f32>() };
75    if !head.is_empty() || !tail.is_empty() {
76        return Err(TeraError::InvalidLayout(
77            "El buffer de bytes no tiene la alineación de memoria requerida para f32".to_string(),
78        ));
79    }
80
81    Ok(body)
82}
83
84#[cfg(test)]
85mod tests {
86    use super::*;
87    use std::sync::Arc;
88    use arrow_schema::{DataType, Field, Schema};
89
90    #[test]
91    fn test_cast_bytes_to_f32() -> Result<()> {
92        let original: Vec<f32> = vec![1.0, 2.5, -3.14, 42.0, 0.0];
93        let bytes: &[u8] = bytemuck::cast_slice(&original);
94
95        let recovered = cast_bytes_to_f32(bytes)?;
96        assert_eq!(recovered, original.as_slice());
97        Ok(())
98    }
99
100    #[test]
101    fn test_arrow_float32_extraction() -> Result<()> {
102        let array = Float32Array::from(vec![10.0, 20.0, 30.0]);
103        let schema = Arc::new(Schema::new(vec![
104            Field::new("signal", DataType::Float32, false),
105        ]));
106
107        let batch = RecordBatch::try_new(schema, vec![Arc::new(array)])
108            .map_err(|e| TeraError::Arrow(e.to_string()))?;
109
110        let col = extract_float_column(&batch, 0)?;
111        assert_eq!(col.len(), 3);
112        assert_eq!(col.to_f32_slice(), &[10.0, 20.0, 30.0]);
113        Ok(())
114    }
115}