use arrow::array::Array;
use cobre_core::EntityId;
use parquet::arrow::arrow_reader::ParquetRecordBatchReaderBuilder;
use std::fs::File;
use std::path::Path;
use crate::LoadError;
use crate::parquet_helpers::{
extract_optional_int32, extract_required_float64, extract_required_int32,
};
#[derive(Debug, Clone, PartialEq)]
pub struct FphaDeviationPointRow {
pub hydro_id: EntityId,
pub stage_id: Option<i32>,
pub v: f64,
pub q: f64,
pub fph_exact: f64,
pub fpha_fitted: f64,
pub deviation: f64,
pub relative: f64,
}
pub fn parse_fpha_deviation_points(path: &Path) -> Result<Vec<FphaDeviationPointRow>, LoadError> {
let file = File::open(path).map_err(|e| LoadError::io(path, e))?;
let builder = ParquetRecordBatchReaderBuilder::try_new(file)
.map_err(|e| LoadError::parse(path, e.to_string()))?;
let reader = builder
.build()
.map_err(|e| LoadError::parse(path, e.to_string()))?;
let mut rows: Vec<FphaDeviationPointRow> = Vec::new();
for batch_result in reader {
let batch = batch_result.map_err(|e| LoadError::parse(path, e.to_string()))?;
let hydro_id_col = extract_required_int32(&batch, "hydro_id", path)?;
let v_col = extract_required_float64(&batch, "v", path)?;
let q_col = extract_required_float64(&batch, "q", path)?;
let fph_exact_col = extract_required_float64(&batch, "fph_exact", path)?;
let fpha_fitted_col = extract_required_float64(&batch, "fpha_fitted", path)?;
let deviation_col = extract_required_float64(&batch, "deviation", path)?;
let relative_col = extract_required_float64(&batch, "relative", path)?;
let stage_id_col = extract_optional_int32(&batch, "stage_id", path)?;
let n = batch.num_rows();
rows.reserve(n);
for i in 0..n {
let hydro_id = EntityId::from(hydro_id_col.value(i));
let v = v_col.value(i);
let q = q_col.value(i);
let fph_exact = fph_exact_col.value(i);
let fpha_fitted = fpha_fitted_col.value(i);
let deviation = deviation_col.value(i);
let relative = relative_col.value(i);
let stage_id = stage_id_col
.filter(|col| !col.is_null(i))
.map(|col| col.value(i));
rows.push(FphaDeviationPointRow {
hydro_id,
stage_id,
v,
q,
fph_exact,
fpha_fitted,
deviation,
relative,
});
}
}
rows.sort_by(|a, b| {
a.hydro_id
.0
.cmp(&b.hydro_id.0)
.then_with(|| a.stage_id.cmp(&b.stage_id))
});
Ok(rows)
}