Files
sientia-dataops-model-manager/model_manager/sientia/models.py

556 lines
21 KiB
Python

from typing import Any
import numpy as np
import pandas as pd
from sientia_do.operations.df_preprocessor import create_features, limit_dataset, treat_nan
from sientia_do.timeseries.analyzer import TimeSeriesDiscontinuityAnalyzer
from sklearn.base import BaseEstimator, TransformerMixin
from sklearn.linear_model import LinearRegression
from sklearn.preprocessing import StandardScaler
DISCONTINUITY_TREATMENT = 'Discontinuity Treatment'
LAG_SELECTION = 'Lag Selection'
STATIC_WINDOW_REMOVAL = 'Static Window Removal'
DEFINE_VARIABLES_LIMITS = 'Define Variables Limits'
NORMALIZATION = 'Normalization'
class LinearRegressionModel(BaseEstimator, TransformerMixin):
"""
Linear Regression Model for Time Series Analysis.
Thread-safety: This class is NOT thread-safe during fit() operations.
Do not call fit() on the same instance from multiple threads simultaneously.
After fitting, predict() is thread-safe for read-only operations.
For multi-threaded environments:
- Fit the model in a single thread
- Share the fitted instance across threads for prediction only
- Or create separate instances per thread
"""
def __init__(
self,
target_variable: str = '',
variable_columns: list[str] | None = None,
model_params: dict[str, Any] | None = None,
clipping: dict[str, float] | None = None,
weights: dict[str, float] | None = None,
):
"""
Linear Regression Model for Time Series Analysis
Args:
target_variable (str): The target variable name
variable_columns (list): The input columns names in a list
model_params (dict): The parameters used for training the model \\
clipping (dict): The lower and upper limits for the target variable to be clipped \\
*Format: {'min': min_value, 'max': max_value}*
weights (dict): The weights for the Linear Regression model \\
*Format: {'variable_name': weight}*
Returns:
LinearRegressionModel: The prediction model object
"""
self.target_variable: str = target_variable
self.variable_columns: list[str] | None = variable_columns
self.model_params: dict[str, Any] | None = model_params
self.clipping: dict[str, float] | None = clipping
self.regr = LinearRegression()
self.q1_target: float | None = None
self.q3_target: float | None = None
self.weights: dict[str, float] | None = weights
def fit(self, input_data: pd.DataFrame) -> 'LinearRegressionModel':
"""
Function to fit the model
Args:
input_data (pandas.DataFrame): The data used to fit the Linear Regression model
Returns:
LinearRegressionModel: The prediction model object
"""
assert self.variable_columns is not None, 'variable_columns must be set before fitting'
X_train = input_data[self.variable_columns]
y_train = input_data[self.target_variable]
self.q1_target = y_train.quantile(0.25)
self.q3_target = y_train.quantile(0.75)
# Fit the model
self.regr.fit(X_train, y_train)
# Get the weights
round_coef = np.round(self.regr.coef_, 3)
round_intercept = np.round(self.regr.intercept_, 3)
# Save the weights
weights = dict(zip(self.variable_columns, [float(c) for c in round_coef], strict=True))
weights = dict(sorted(weights.items(), key=lambda item: abs(item[1]), reverse=True))
weights = {'Bias': float(round_intercept), **weights}
self.weights = weights
return self
def predict(self, input_data: pd.DataFrame) -> np.ndarray:
"""
Function to predict the target variable.
If clipping is True, the predictions are clipped based on the target variable quartiles.
Args:
input_data (pandas.DataFrame): The data used to predict the target variable
Returns:
numpy.ndarray: The predicted target variable
"""
X_test = input_data[self.variable_columns]
y_pred = self.regr.predict(X_test)
if self.clipping:
for i in range(len(y_pred)):
if y_pred[i] > self.clipping['max']:
y_pred[i] = self.q3_target
elif y_pred[i] < self.clipping['min']:
y_pred[i] = self.q1_target
return y_pred
class DataPreprocessor(BaseEstimator, TransformerMixin):
"""
Data Preprocessor for Time Series Analysis.
Thread-safety: This class is NOT thread-safe during fit() operations.
Do not call fit() on the same instance from multiple threads simultaneously.
After fitting, transform() is thread-safe for read-only operations IF the
input DataFrames are not shared between threads.
For multi-threaded environments:
- Fit the preprocessor in a single thread
- Share the fitted instance across threads for transform() only
- Ensure each thread passes its own DataFrame copy to transform()
- Or create separate instances per thread
"""
def __init__(
self,
date_column: str = '',
target_variable: str = '',
input_columns: list[str] | None = None,
nan_treatment: str | None = None,
lag_train: dict[str, int] | None = None,
lag_transform: dict[str, int] | None = None,
static_threshold: int | None = None,
low_lim: dict[str, float] | None = None,
upp_lim: dict[str, float] | None = None,
window: int | None = None,
scaler_name: str | None = None,
scaler_params: dict[str, Any] | None = None,
ar_var: str | None = None,
self_operations: list[str] | None = None,
cross_operations: list[str] | None = None,
created_lags: dict[str, int] | None = None,
steps_order: list[str] | None = None,
):
"""
Data Preprocessor for Time Series Analysis
Args:
date_column (str): The column name of the date in the dataset
target_variable (str): The target variable name
input_columns (list): The input columns names in a list
nan_treatment (str): The treatment for missing values \\
*Options: 'drop', 'fill linear'*
lag_train (dict): The lags for each variable to be applyed during training \\
*Format: {'variable_name': lag}*
lag_transform (dict): The lags for each variable to be applyed during transformation \\
*Format: {'variable_name': lag}*
static_threshold (int): The number of repeated values to be considered as static
low_lim (dict): The lower limits for each variable \\
*Format: {'variable_name': limit}*
upp_lim (dict): The upper limits for each variable \\
*Format: {'variable_name': limit}*
window (int): The window size for rolling window. **Not implemented yet**
scaler_name (str): The scaler name. If no scaler is used, it is 'None' \\
*Options: 'None', 'Standard Scaler'*
scaler_params (dict): The parameters for the scaler object, if it is used \\
*Format for Standard Scaler: {'variable_name': {'mean': mean, 'variance': variance}}*
ar_var (str): The autoregressive variable name. If None, it is not created
self_operations (list): The operations for feature creation using the same variable \\
*Format: ['{variable_name}\\_{operation}\\_{scalar}']* \\
*Operations: 'exp', 'pow', 'log', 'root'*
cross_operations (list): The operations for feature creation using two variables \\
*Format: ['{variable_name1}\\_{operation}\\_{variable_name2}']* \\
*Operations: '\\*', '/'*
created_lags (dict): Variables created by lagging existing ones \\
*Format: {'original_variable_name': lag}*
steps_order (list): The order of the steps to be executed in the pipeline \\
*Options for list: 'Discontinuity Treatment',
'Lag Selection',
'Static Window Removal',
'Define Variables Limits',
'Normalization',
'Feature Creation',
'Lag Creation'*
Returns:
DataPreprocessor: The data preprocessor object
"""
self.date_column = date_column
self.target_variable = target_variable
self.input_columns = input_columns
self.nan_treatment = nan_treatment
self.lag_train = lag_train if lag_train else {}
self.lag_transform = lag_transform if lag_transform else {}
self.ar_var = ar_var
self.self_operations = self_operations
self.cross_operations = cross_operations
self.created_lags = created_lags
self.static_threshold = static_threshold
self.low_lim = low_lim
self.upp_lim = upp_lim
# self.window = window # Not implemented yet
self.scaler_name = scaler_name
self.scaler_params = scaler_params
self.feature_names_order: list[str] = [] # Initialize to avoid AttributeError
if self.scaler_name == 'Standard Scaler':
self.scaler = StandardScaler()
elif self.scaler_name == 'None':
self.scaler = None
else:
self.scaler = None
# Filter steps for preprocessor class
possible_steps = [
DISCONTINUITY_TREATMENT,
LAG_SELECTION,
STATIC_WINDOW_REMOVAL,
DEFINE_VARIABLES_LIMITS,
NORMALIZATION,
'Feature Creation',
'Lag Creation',
]
self.steps_order = steps_order or possible_steps
for step in possible_steps:
if step not in self.steps_order:
self.steps_order.append(step)
def get_required_columns(self, existing_columns: list) -> list:
"""
Get the required columns to generate the input columns
Args:
existing_columns (list): The existing columns in the data
Returns:
list: The required columns
"""
required_columns: list[str] = []
# Columns for feature creation
if self.self_operations is not None:
for name in self.self_operations:
var, operation, scalar = name.split('}_{')
var = var.split('{')[1]
operation = operation.split('}')[0]
scalar = scalar.split('}')[0]
required_columns.append(var)
if self.cross_operations is not None:
for name in self.cross_operations:
var1, operation, var2 = name.split('}_{')
var1 = var1.split('{')[1]
operation = operation.split('}')[0]
var2 = var2.split('}')[0]
required_columns.append(var1)
required_columns.append(var2)
# Columns for lag creation
if self.created_lags is not None:
for var in self.created_lags.keys():
required_columns.append(var)
# Check if any column in required_columns is not in existing_columns
required_columns = list(set(required_columns))
_to_remove: list[str] = []
for column in required_columns:
# If column was already in self_operations list, remove it
if (
self.self_operations is not None
and column not in existing_columns
and column in self.self_operations
):
_to_remove.append(column)
# If column was already in cross_operations list, remove it
if (
self.cross_operations is not None
and column not in existing_columns
and column in self.cross_operations
):
_to_remove.append(column)
# If column was already in created_lags list, remove it
if (
self.created_lags is not None
and column not in existing_columns
and column in self.created_lags
):
_to_remove.append(column)
for column in set(_to_remove):
required_columns.remove(column)
return required_columns
def get_scaler(self) -> Any:
"""
Get the scaler object
Returns:
Scaler: The scaler object
"""
return self.scaler
def treat_discontinuities(self, input_data: pd.DataFrame) -> pd.DataFrame:
"""
Treat the discontinuities in the data
Args:
input_data (pandas.DataFrame): The input data
Returns:
pandas.DataFrame: The treated data
"""
if self.nan_treatment:
input_data = treat_nan(input_data, self.nan_treatment)
return input_data
def lag_selection(self, input_data: pd.DataFrame, lag_dict: dict) -> pd.DataFrame:
"""
Select the lags for the variables
Args:
input_data (pandas.DataFrame): The input data
lag_dict (dict): The lags for each variable \\
*Format: {'variable_name': lag}*
Returns:
pandas.DataFrame: The treated data
Note:
This method modifies input_data in-place. Ensure the caller passes
a copy if the original DataFrame needs to be preserved.
"""
if lag_dict:
for var, lag in lag_dict.items():
if lag > 0:
input_data[var] = input_data[var].shift(lag)
# WARNING: Modifies DataFrame in-place
input_data.dropna(inplace=True)
return input_data
def treat_static_windows(self, input_data: pd.DataFrame) -> pd.DataFrame:
"""
Treat the static windows in the data
Args:
input_data (pandas.DataFrame): The input data
Returns:
pandas.DataFrame: The treated data
"""
if self.static_threshold:
ts_analyzer = TimeSeriesDiscontinuityAnalyzer(input_data)
ts_analyzer.infer_frequency()
for col in input_data.columns:
ts_analyzer.identify_static_windows(column=col, threshold=self.static_threshold)
ts_analyzer.treat_static_windows(
column=col, remove_window=True, threshold=self.static_threshold
)
ts_analyzer.update_total_discontinuities(col)
input_data = ts_analyzer.get_treated_data()
return input_data
def adjust_limits(self, input_data: pd.DataFrame) -> pd.DataFrame:
"""
Adjust the limits for the variables
Args:
input_data (pandas.DataFrame): The input data
Returns:
pandas.DataFrame: The treated data
"""
input_data, self.low_lim, self.upp_lim = limit_dataset(
input_data, self.low_lim, self.upp_lim
)
return input_data
def create_features(self, input_data: pd.DataFrame) -> pd.DataFrame:
"""
Create features in the data
Args:
input_data (pandas.DataFrame): The input data
Returns:
pandas.DataFrame: The treated data
"""
input_data = create_features(input_data, self.self_operations, self.cross_operations)
return input_data
def create_ar(self, input_data: pd.DataFrame) -> pd.DataFrame:
"""
Create the autoregressive variable in the data
Args:
input_data (pandas.DataFrame): The input data
Returns:
pandas.DataFrame: The treated data
Note:
This method modifies input_data in-place.
"""
if self.ar_var:
input_data[self.ar_var] = input_data[self.target_variable].shift(1)
# WARNING: Modifies DataFrame in-place
input_data.dropna(inplace=True)
return input_data
def create_lags(self, input_data: pd.DataFrame) -> pd.DataFrame:
"""
Create additional lags in the data
Args:
input_data (pandas.DataFrame): The input data
Returns:
pandas.DataFrame: The treated data
Note:
This method modifies input_data in-place.
"""
if self.created_lags:
for var, lag in self.created_lags.items():
if lag > 0 and var in input_data.columns:
new_col = f'{var}_lag{lag}'
input_data[new_col] = input_data[var].shift(lag)
# WARNING: Modifies DataFrame in-place
input_data.dropna(inplace=True)
return input_data
def fit(self, x: pd.DataFrame, y: None | pd.Series = None) -> 'DataPreprocessor':
"""
Function to preprocess the data and split it into training and testing sets
Args:
x (pandas.DataFrame): The input data
y (pandas.Series): The target variable
Returns:
DataPreprocessor: The data preprocessor object
"""
if x is not None and y is not None:
data_treat = pd.concat([x.copy(), y.copy()], axis=1)
elif x is not None:
data_treat = x.copy()
else:
raise ValueError('No data was provided')
assert self.input_columns is not None, 'input_columns must be set'
existing_columns = [col for col in data_treat.columns if col in self.input_columns]
data_treat = data_treat[existing_columns + [self.target_variable]]
for step in self.steps_order:
# Discontinuity Treatment
if step == DISCONTINUITY_TREATMENT:
data_treat = self.treat_discontinuities(data_treat)
# Lag for Model Training
if step == LAG_SELECTION:
data_treat = self.lag_selection(data_treat, self.lag_train)
# Static Window Treatment
if step == STATIC_WINDOW_REMOVAL:
data_treat = self.treat_static_windows(data_treat)
# Adjust limits
if step == DEFINE_VARIABLES_LIMITS:
data_treat = self.adjust_limits(data_treat)
# Normalization
if step == NORMALIZATION and self.scaler:
self.scaler = self.scaler.fit(data_treat[existing_columns])
self.feature_names_order = list(data_treat[existing_columns].columns)
data_treat[existing_columns] = self.scaler.transform(data_treat[existing_columns])
# Save scaler parameters
assert self.scaler_params is not None, 'scaler_params must be initialized'
for index, column in enumerate(list(existing_columns)):
mean = self.scaler.mean_[index]
variance = self.scaler.var_[index]
self.scaler_params[column] = {
'mean': round(mean, 3),
'variance': round(variance, 3),
}
return self
def transform(self, x: pd.DataFrame) -> pd.DataFrame:
"""
Function to preprocess the data
Args:
x (pandas.DataFrame): The input data
Returns:
pandas.DataFrame: The treated data
"""
if 'timestamp' in x.columns:
data_treat = x.drop(columns='timestamp')
else:
data_treat = x.copy()
assert self.input_columns is not None, 'input_columns must be set'
existing_columns = [col for col in data_treat.columns if col in self.input_columns]
required_columns = self.get_required_columns(existing_columns)
all_cols = required_columns + existing_columns + [self.target_variable]
all_cols = list(set(all_cols))
data_treat = data_treat[all_cols]
for step in self.steps_order:
# Discontinuity Treatment
if step == DISCONTINUITY_TREATMENT:
data_treat = self.treat_discontinuities(data_treat)
# Lag for Model Training
if step == LAG_SELECTION:
data_treat = self.lag_selection(data_treat, self.lag_transform)
# Static Window Treatment
if step == STATIC_WINDOW_REMOVAL:
data_treat = self.treat_static_windows(data_treat)
# Adjust limits
if step == DEFINE_VARIABLES_LIMITS:
data_treat = self.adjust_limits(data_treat)
# Normalization
if step == NORMALIZATION and self.scaler:
# Only transform feature columns, preserve target and any other required columns
feature_cols = self.feature_names_order
data_treat[feature_cols] = self.scaler.transform(data_treat[feature_cols])
# Feature Creation
if step == 'Feature Creation':
data_treat = self.create_features(data_treat)
# Lag Creation
if step == 'Lag Creation':
# Autoregressive Variable
if self.input_columns is not None and self.ar_var in self.input_columns:
data_treat = self.create_ar(data_treat)
# Additonal Lags
data_treat = self.create_lags(data_treat)
return data_treat