diff --git a/tests/test_evaluation.py b/tests/test_evaluation.py index d0fd9b0..0ce2d98 100644 --- a/tests/test_evaluation.py +++ b/tests/test_evaluation.py @@ -46,15 +46,14 @@ class EvenOddStubModel(BaseEstimator, ClassifierMixin, Model): self.window_length = window_length def fit(self, X, y): - assert len(X) == self.window_length for i in range(len(X)): assert y[i] == -1 if X[i][0] == 1 else 1 def predict(self, X): - return np.array([-1 if row[0] == 1 else 1 for row in X]) + return np.array([-1 if row[-1] == 1 else 1 for row in X]) def predict_proba(self, X): - return np.array([[row[0] + 1, 0] for row in X]) + return np.array([[row[-1] + 1, 0] for row in X]) def test_evaluation(): @@ -70,7 +69,6 @@ def test_evaluation(): X=X, y=y, forward_returns=y, - expanding_window=False, window_size=window_length, retrain_every=retrain_every, from_index=None, diff --git a/tests/test_walk_forward.py b/tests/test_walk_forward.py index b64d01d..51315cb 100644 --- a/tests/test_walk_forward.py +++ b/tests/test_walk_forward.py @@ -42,9 +42,8 @@ class IncrementingStubModel(Model, BaseEstimator, ClassifierMixin): self.window_length = window_length def fit(self, X, y): - assert len(X) == self.window_length for i in range(len(X)): - assert X[i][0] + 1 == y[i] + assert X[i][-1] + 1 == y[i] def predict(self, X): return np.array([row[0] + 1 for row in X]) @@ -66,7 +65,6 @@ def test_walk_forward_train_test(): X=X, y=y, forward_returns=y, - expanding_window=False, window_size=window_length, retrain_every=retrain_every, from_index=None, diff --git a/training/train_model.py b/training/train_model.py index ee72dde..ee1aa14 100644 --- a/training/train_model.py +++ b/training/train_model.py @@ -42,7 +42,6 @@ def train_model( X=X, y=y, forward_returns=forward_returns, - expanding_window=True, window_size=initial_window_size, retrain_every=retrain_every, from_index=from_index, diff --git a/training/walk_forward/inference_batched.py b/training/walk_forward/inference_batched.py index 3c9e46c..17899c6 100644 --- a/training/walk_forward/inference_batched.py +++ b/training/walk_forward/inference_batched.py @@ -57,8 +57,6 @@ def walk_forward_inference_batched( X, model_over_time, transformations_over_time, - expanding_window, - window_size, ) for index in tqdm(batch_indices) ] @@ -76,8 +74,6 @@ def __inference_from_window( X: XDataFrame, model_over_time: ModelOverTime, transformations_over_time: TransformationsOverTime, - expanding_window: bool, - window_size: int, ) -> list[tuple[int, float, pd.Series]]: current_model = model_over_time[X.index[index_start]] current_transformations = [ diff --git a/training/walk_forward/inference_parallel.py b/training/walk_forward/inference_parallel.py deleted file mode 100644 index b6cb264..0000000 --- a/training/walk_forward/inference_parallel.py +++ /dev/null @@ -1,117 +0,0 @@ -import pandas as pd -from models.base import Model -from training.types import ( - ModelOverTime, - TransformationsOverTime, - PredictionsSeries, - ProbabilitiesDataFrame, -) -from transformations.base import Transformation -from utils.helpers import get_first_valid_return_index -from tqdm import tqdm -from typing import Optional -from data_loader.types import XDataFrame -import ray - - -def walk_forward_inference( - model_name: str, - model_over_time: ModelOverTime, - transformations_over_time: TransformationsOverTime, - X: XDataFrame, - expanding_window: bool, - window_size: int, - retrain_every: int, - class_labels: list[int], - from_index: Optional[pd.Timestamp], -) -> tuple[PredictionsSeries, ProbabilitiesDataFrame]: - predictions = pd.Series(index=X.index, dtype="object").rename(model_name) - probabilities = pd.DataFrame( - index=X.index, columns=[str(label) for label in class_labels] - ) - - inference_from = ( - max( - get_first_valid_return_index(model_over_time), - get_first_valid_return_index(X.iloc[:, 0]), - ) - if from_index is None - else X.index.to_list().index(from_index) - ) - inference_till = X.shape[0] - first_model = model_over_time[inference_from] - - if first_model.only_column is not None: - X = X[[column for column in X.columns if first_model.only_column in column]] - - if first_model.data_transformation == "original": - transformations_over_time = [] - - batch_size = int((inference_till - inference_from) / 10) - batched_results = ray.get( - [ - __inference_from_window.remote( - index, - index + batch_size, - inference_from, - retrain_every, - X, - model_over_time, - transformations_over_time, - expanding_window, - window_size, - ) - for index in range(inference_from, inference_till) - ] - ) - for batch in batched_results: - for index, prediction, probs in batch: - predictions[X.index[index]] = prediction - probabilities.loc[X.index[index]] = probs - - return predictions, probabilities - - -@ray.remote -def __inference_from_window( - index_start: int, - index_end: int, - inference_from: int, - retrain_every: int, - X: XDataFrame, - model_over_time: ModelOverTime, - transformations_over_time: TransformationsOverTime, - expanding_window: bool, - window_size: int, -) -> list[tuple[int, float, pd.Series]]: - - results = [] - for index in range(index_start, index_end): - last_model_index = index - ((index - inference_from) % retrain_every) - train_window_start = ( - X.index[inference_from] - if expanding_window - else X.index[index - window_size - 1] - ) - - current_model = model_over_time[X.index[last_model_index]] - current_transformations = [ - transformation_over_time[X.index[last_model_index]] - for transformation_over_time in transformations_over_time - ] - - if current_model.predict_window_size == "window_size": - next_timestep = X.loc[train_window_start : X.index[index]] - else: - # we need to get a Dataframe out of it, since the transformation step always expects a 2D array, but it's equivalent to X.iloc[index] - next_timestep = X.loc[X.index[index] : X.index[index]] - - for transformation in current_transformations: - next_timestep = transformation.transform(next_timestep) - - next_timestep = next_timestep.to_numpy() - - prediction, probs = current_model.predict(next_timestep) - results.append((index, prediction, probs)) - - return results diff --git a/training/walk_forward/process_transformations_parallel.py b/training/walk_forward/process_transformations_parallel.py index 0073350..f487468 100644 --- a/training/walk_forward/process_transformations_parallel.py +++ b/training/walk_forward/process_transformations_parallel.py @@ -39,7 +39,6 @@ def walk_forward_process_transformations( preprocess_transformations_window.remote( X, y, - window_size, transformations, first_nonzero_return, index, @@ -61,7 +60,6 @@ def walk_forward_process_transformations( def preprocess_transformations_window( X: XDataFrame, y: ySeries, - window_size: int, transformations: list[Transformation], first_nonzero_return: int, index: int, diff --git a/training/walk_forward/train.py b/training/walk_forward/train.py index c72db88..dff5515 100644 --- a/training/walk_forward/train.py +++ b/training/walk_forward/train.py @@ -13,7 +13,6 @@ def walk_forward_train( X: XDataFrame, y: ySeries, forward_returns: ForwardReturnSeries, - expanding_window: bool, window_size: int, retrain_every: int, from_index: Optional[pd.Timestamp], @@ -40,11 +39,7 @@ def walk_forward_train( transformations_over_time = [] for index in tqdm(range(train_from, train_till, retrain_every)): - train_window_start = ( - X.index[first_nonzero_return] - if expanding_window - else X.index[index - window_size - 1] - ) + train_window_start = X.index[first_nonzero_return] train_window_end = X.index[index - 1] current_transformations = [ diff --git a/training/walk_forward/train_parallel.py b/training/walk_forward/train_parallel.py index da2143c..7d05684 100644 --- a/training/walk_forward/train_parallel.py +++ b/training/walk_forward/train_parallel.py @@ -15,7 +15,6 @@ def walk_forward_train( X: XDataFrame, y: ySeries, forward_returns: ForwardReturnSeries, - expanding_window: bool, window_size: int, retrain_every: int, from_index: Optional[pd.Timestamp], @@ -46,11 +45,9 @@ def walk_forward_train( train_on_window.remote( index, first_nonzero_return, - window_size, X, y, model, - expanding_window, transformations_over_time, ) for index in tqdm(range(train_from, train_till, retrain_every)) @@ -66,18 +63,12 @@ def walk_forward_train( def train_on_window( index: int, first_nonzero_return: int, - window_size: int, X: XDataFrame, y: ySeries, model: Model, - expanding_window: bool, transformations_over_time: TransformationsOverTime, ) -> tuple[int, Model]: - train_window_start = ( - X.index[first_nonzero_return] - if expanding_window - else X.index[index - window_size - 1] - ) + train_window_start = X.index[first_nonzero_return] train_window_end = X.index[index - 1] current_transformations = [