use crate::column::ColumnTrait;
use crate::error::{Error, Result};
use crate::ml::preprocessing::{MinMaxScaler, StandardScaler};
use crate::optimized::OptimizedDataFrame;
use std::collections::HashMap;
pub trait Transformer: std::fmt::Debug {
fn transform(&self, df: &OptimizedDataFrame) -> Result<OptimizedDataFrame>;
fn fit_transform(&mut self, df: &OptimizedDataFrame) -> Result<OptimizedDataFrame>;
fn fit(&mut self, df: &OptimizedDataFrame) -> Result<()>;
}
fn read_f64_preserving_nulls(col: &crate::column::Float64Column) -> (Vec<f64>, Vec<bool>) {
let n = col.len();
let mut values = Vec::with_capacity(n);
let mut nulls = Vec::with_capacity(n);
for i in 0..n {
match col.get(i) {
Ok(Some(v)) => {
values.push(v);
nulls.push(false);
}
_ => {
values.push(f64::NAN);
nulls.push(true);
}
}
}
(values, nulls)
}
fn read_string_preserving_nulls(col: &crate::column::StringColumn) -> (Vec<String>, Vec<bool>) {
let n = col.len();
let mut values = Vec::with_capacity(n);
let mut nulls = Vec::with_capacity(n);
for i in 0..n {
match col.get(i) {
Ok(Some(v)) => {
values.push(v.to_string());
nulls.push(false);
}
_ => {
values.push(String::new());
nulls.push(true);
}
}
}
(values, nulls)
}
fn make_f64_column(values: Vec<f64>, nulls: Vec<bool>, name: &str) -> crate::column::Float64Column {
if nulls.iter().any(|&n| n) {
crate::column::Float64Column::with_nulls(values, nulls)
} else {
crate::column::Float64Column::with_name(values, name.to_string())
}
}
fn make_string_column(
values: Vec<String>,
nulls: Vec<bool>,
name: &str,
) -> crate::column::StringColumn {
if nulls.iter().any(|&n| n) {
crate::column::StringColumn::with_nulls(values, nulls)
} else {
crate::column::StringColumn::with_name(values, name.to_string())
}
}
fn copy_column_passthrough(
result: &mut OptimizedDataFrame,
col_name: &str,
column_view: &crate::optimized::ColumnView,
) -> Result<()> {
if let Some(float_values) = column_view.as_float64() {
let (values, nulls) = read_f64_preserving_nulls(float_values);
let col = make_f64_column(values, nulls, col_name);
result.add_column(col_name.to_string(), crate::column::Column::Float64(col))?;
} else if let Some(string_values) = column_view.as_string() {
let (values, nulls) = read_string_preserving_nulls(string_values);
let col = make_string_column(values, nulls, col_name);
result.add_column(col_name.to_string(), crate::column::Column::String(col))?;
}
Ok(())
}
impl Transformer for StandardScaler {
fn transform(&self, df: &OptimizedDataFrame) -> Result<OptimizedDataFrame> {
let means = self.means.as_ref().ok_or_else(|| {
Error::InvalidOperation(
"StandardScaler::transform called before fit() (or fit_transform()); \
no fitted means available"
.into(),
)
})?;
let stds = self.stds.as_ref().ok_or_else(|| {
Error::InvalidOperation(
"StandardScaler::transform called before fit() (or fit_transform()); \
no fitted standard deviations available"
.into(),
)
})?;
let mut result = OptimizedDataFrame::new();
let column_names: Vec<String> = df.column_names().to_vec();
for col_name in &column_names {
let should_scale = if let Some(columns) = &self.columns {
columns.contains(col_name)
} else {
true };
let column_view = match df.column(col_name) {
Ok(v) => v,
Err(_) => continue,
};
if should_scale && column_view.as_float64().is_some() {
let float_values = column_view.as_float64().ok_or_else(|| {
Error::TypeMismatch("column type check failed for Float64".into())
})?;
let mean = *means.get(col_name.as_str()).ok_or_else(|| {
Error::InvalidOperation(format!(
"StandardScaler: no fitted mean for column '{}'; it was not \
present (or was entirely null) when fit() ran",
col_name
))
})?;
let std_dev = *stds.get(col_name.as_str()).ok_or_else(|| {
Error::InvalidOperation(format!(
"StandardScaler: no fitted standard deviation for column '{}'; \
it was not present (or was entirely null) when fit() ran",
col_name
))
})?;
let (raw_values, nulls) = read_f64_preserving_nulls(float_values);
let scaled_values: Vec<f64> = raw_values
.iter()
.map(|&x| {
if std_dev > 1e-10 {
(x - mean) / std_dev
} else {
0.0
}
})
.collect();
let scaled_column = make_f64_column(scaled_values, nulls, col_name);
result.add_column(
col_name.to_string(),
crate::column::Column::Float64(scaled_column),
)?;
} else {
copy_column_passthrough(&mut result, col_name, &column_view)?;
}
}
Ok(result)
}
fn fit_transform(&mut self, df: &OptimizedDataFrame) -> Result<OptimizedDataFrame> {
<Self as Transformer>::fit(self, df)?;
<Self as Transformer>::transform(self, df)
}
fn fit(&mut self, df: &OptimizedDataFrame) -> Result<()> {
let column_names: Vec<String> = match &self.columns {
Some(cols) => cols.clone(),
None => df.column_names().to_vec(),
};
let mut means = HashMap::new();
let mut stds = HashMap::new();
for col_name in column_names {
if df.column(&col_name).is_err() {
continue;
}
let column_view = df.column(&col_name)?;
if let Some(float_values) = column_view.as_float64() {
let mut values: Vec<f64> = Vec::new();
for i in 0..float_values.len() {
if let Ok(Some(val)) = float_values.get(i) {
values.push(val);
}
}
if values.is_empty() {
continue;
}
let sum: f64 = values.iter().sum();
let mean = sum / values.len() as f64;
means.insert(col_name.clone(), mean);
let var_sum: f64 = values.iter().map(|&x| (x - mean).powi(2)).sum();
let variance = var_sum / values.len() as f64;
let std_dev = variance.sqrt();
stds.insert(col_name.clone(), std_dev);
}
}
self.means = Some(means);
self.stds = Some(stds);
Ok(())
}
}
impl Transformer for MinMaxScaler {
fn transform(&self, df: &OptimizedDataFrame) -> Result<OptimizedDataFrame> {
let mins = self.min_values.as_ref().ok_or_else(|| {
Error::InvalidOperation(
"MinMaxScaler::transform called before fit() (or fit_transform()); \
no fitted minimums available"
.into(),
)
})?;
let maxs = self.max_values.as_ref().ok_or_else(|| {
Error::InvalidOperation(
"MinMaxScaler::transform called before fit() (or fit_transform()); \
no fitted maximums available"
.into(),
)
})?;
let mut result = OptimizedDataFrame::new();
let column_names: Vec<String> = df.column_names().to_vec();
for col_name in &column_names {
let should_scale = if let Some(columns) = &self.columns {
columns.contains(col_name)
} else {
true };
let column_view = match df.column(col_name) {
Ok(v) => v,
Err(_) => continue,
};
if should_scale && column_view.as_float64().is_some() {
let float_values = column_view.as_float64().ok_or_else(|| {
Error::TypeMismatch("column type check failed for Float64".into())
})?;
let min_val = *mins.get(col_name.as_str()).ok_or_else(|| {
Error::InvalidOperation(format!(
"MinMaxScaler: no fitted minimum for column '{}'; it was not \
present (or was entirely null) when fit() ran",
col_name
))
})?;
let max_val = *maxs.get(col_name.as_str()).ok_or_else(|| {
Error::InvalidOperation(format!(
"MinMaxScaler: no fitted maximum for column '{}'; it was not \
present (or was entirely null) when fit() ran",
col_name
))
})?;
let (feature_min, feature_max) = self.feature_range;
let (raw_values, nulls) = read_f64_preserving_nulls(float_values);
let scaled_values: Vec<f64> = if (max_val - min_val).abs() > 1e-10 {
raw_values
.iter()
.map(|&x| {
let scaled = (x - min_val) / (max_val - min_val);
scaled * (feature_max - feature_min) + feature_min
})
.collect()
} else {
vec![feature_min; raw_values.len()]
};
let scaled_column = make_f64_column(scaled_values, nulls, col_name);
result.add_column(
col_name.to_string(),
crate::column::Column::Float64(scaled_column),
)?;
} else {
copy_column_passthrough(&mut result, col_name, &column_view)?;
}
}
Ok(result)
}
fn fit_transform(&mut self, df: &OptimizedDataFrame) -> Result<OptimizedDataFrame> {
<Self as Transformer>::fit(self, df)?;
<Self as Transformer>::transform(self, df)
}
fn fit(&mut self, df: &OptimizedDataFrame) -> Result<()> {
let column_names: Vec<String> = match &self.columns {
Some(cols) => cols.clone(),
None => df.column_names().to_vec(),
};
let mut min_values = HashMap::new();
let mut max_values = HashMap::new();
for col_name in column_names {
if df.column(&col_name).is_err() {
continue;
}
let column_view = df.column(&col_name)?;
if let Some(float_values) = column_view.as_float64() {
let mut values: Vec<f64> = Vec::new();
for i in 0..float_values.len() {
if let Ok(Some(val)) = float_values.get(i) {
values.push(val);
}
}
if values.is_empty() {
continue;
}
let min_val = values
.iter()
.min_by(|a, b| a.partial_cmp(b).unwrap_or(std::cmp::Ordering::Equal))
.copied()
.ok_or_else(|| {
Error::InvalidOperation("Cannot compute min of empty values".into())
})?;
let max_val = values
.iter()
.max_by(|a, b| a.partial_cmp(b).unwrap_or(std::cmp::Ordering::Equal))
.copied()
.ok_or_else(|| {
Error::InvalidOperation("Cannot compute max of empty values".into())
})?;
min_values.insert(col_name.clone(), min_val);
max_values.insert(col_name.clone(), max_val);
}
}
self.min_values = Some(min_values);
self.max_values = Some(max_values);
Ok(())
}
}
pub use self::Transformer as PipelineTransformer;
pub struct Pipeline {
pub stages: Vec<Box<dyn Transformer>>,
}
impl std::fmt::Debug for Pipeline {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("Pipeline")
.field("stages_count", &self.stages.len())
.finish()
}
}
impl Pipeline {
pub fn new() -> Self {
Pipeline { stages: Vec::new() }
}
pub fn add_stage<T: Transformer + 'static>(&mut self, stage: T) -> &mut Self {
self.stages.push(Box::new(stage));
self
}
pub fn fit(&mut self, df: &OptimizedDataFrame) -> Result<()> {
let mut current_df = df.clone();
for stage in &mut self.stages {
stage.fit(¤t_df)?;
current_df = stage.transform(¤t_df)?;
}
Ok(())
}
pub fn transform(&self, df: &OptimizedDataFrame) -> Result<OptimizedDataFrame> {
let mut current_df = df.clone();
for stage in &self.stages {
current_df = stage.transform(¤t_df)?;
}
Ok(current_df)
}
pub fn fit_transform(&mut self, df: &OptimizedDataFrame) -> Result<OptimizedDataFrame> {
self.fit(df)?;
self.transform(df)
}
}