From 797d45d036fa11aba78f4373e5222c5927371a8e Mon Sep 17 00:00:00 2001 From: Mark Aron Szulyovszky Date: Fri, 14 Jan 2022 10:34:28 +0100 Subject: [PATCH] feat(Inference): pipeline wired up (#171) * feat: Basic pipeline extended. * feat: Added conversion of model list to existing structure (model_name, model_in_time). Fixed loading of previous models and dicts. * fix: Had an unfinished function. * fix: Inference wasn't getting model_over_time. Now transformations are not getting it either yet. Co-authored-by: Daniel Szemerey --- models/saving.py | 15 +++++--- reporting/types.py | 8 ++++ run_evaluate_robustness.py | 2 +- run_inference.py | 8 ++-- run_pipeline.py | 8 ++-- training/inference.py | 77 +++++++++++++++++--------------------- training/meta_labeling.py | 7 +++- training/primary_model.py | 31 ++++++++------- training/training_steps.py | 10 +++-- 9 files changed, 89 insertions(+), 77 deletions(-) diff --git a/models/saving.py b/models/saving.py index ceba6f9..67148c4 100644 --- a/models/saving.py +++ b/models/saving.py @@ -7,10 +7,11 @@ from reporting.types import Reporting -def save_models(all_models_for_all_assets: list[Reporting.Asset], data_config:dict, training_config:dict) -> None: +def save_models(all_models_for_all_assets: list[Reporting.Asset], data_config:dict, training_config:dict, model_config:dict) -> None: dict_for_pickle = dict() dict_for_pickle['training_config'] = training_config dict_for_pickle['data_config'] = data_config + dict_for_pickle['model_config'] = model_config dict_for_pickle['all_models_for_all_assets'] = all_models_for_all_assets date_string = datetime.datetime.now().strftime("%Y-%m-%d-%H-%M") @@ -22,7 +23,7 @@ def save_models(all_models_for_all_assets: list[Reporting.Asset], data_config:di pickle.dump( dict_for_pickle, open( "output/models/{}.p".format(date_string), "wb" ) ) -def load_models(file_name:Union[str, None]) -> tuple[dict, dict, dict]: +def load_models(file_name:Union[str, None]) -> tuple[list[Reporting.Asset], dict, dict, dict]: if file_name is None: warnings.warn("No file name provided, will load latest models and configurations.") @@ -31,10 +32,12 @@ def load_models(file_name:Union[str, None]) -> tuple[dict, dict, dict]: assert len(files_in_directory) > 0, "No models found in output/models." file_name = sorted(files_in_directory)[-1] - all_models_for_all_assets = pickle.load( open( "output/models/{}".format(file_name), "rb" ) ) + packacked_dict = pickle.load( open( "output/models/{}".format(file_name), "rb" ) ) - data_config = all_models_for_all_assets['data_config'] - training_config = all_models_for_all_assets['training_config'] + data_config = packacked_dict.pop("data_config", None) + training_config = packacked_dict.pop("training_config", None) + model_config = packacked_dict.pop("model_config", None) + all_models_for_all_assets = packacked_dict.pop("all_models_for_all_assets", None) - return all_models_for_all_assets, data_config, training_config + return all_models_for_all_assets, data_config, training_config, model_config diff --git a/reporting/types.py b/reporting/types.py index 0b1a8ef..a801d8b 100644 --- a/reporting/types.py +++ b/reporting/types.py @@ -25,6 +25,14 @@ class Reporting: self.level: str = level self.base: list[Reporting.Single_Model] = [] self.metalabeling: list[list[Reporting.Single_Model]] = [] + + def convert_step_to_tuple(self, step:str)->list[tuple[str, list[Model]]]: + if step == 'base': + return [(x.model_name, x.model_over_time) for x in self.base ] + elif step == 'metalabeling': + return [(x.model_name, x.model_over_time) for sub in self.metalabeling for x in sub] + else: + raise ValueError('Unknown step: {}'.format(step)) class Asset(): diff --git a/run_evaluate_robustness.py b/run_evaluate_robustness.py index e2f320d..d3fad4c 100644 --- a/run_evaluate_robustness.py +++ b/run_evaluate_robustness.py @@ -5,7 +5,7 @@ import pandas as pd all_results = [] all_predictions = [] for index in range(6): - _, _, _, results_1, predictions_1, _ = run_pipeline(project_name='price-prediction', with_wandb = False, sweep = False, get_config= get_default_ensemble_config) + _, _, _, _, results_1, predictions_1, _ = run_pipeline(project_name='price-prediction', with_wandb = False, sweep = False, get_config= get_default_ensemble_config) all_results.append(results_1) all_predictions.append(predictions_1) diff --git a/run_inference.py b/run_inference.py index 2678bf1..05d32b4 100644 --- a/run_inference.py +++ b/run_inference.py @@ -9,12 +9,12 @@ from typing import Callable def run_inference(preload_models:bool, get_config:Callable): if preload_models: - all_models_all_assets, data_config, training_config = load_models(None) + all_models_all_assets, data_config, training_config, model_config = load_models(None) else: - all_models_all_assets, data_config, training_config, _, _, _ = run_pipeline(project_name='price-prediction', with_wandb = False, sweep = False, get_config=get_config) + all_models_all_assets, data_config, training_config, model_config, _, _, _ = run_pipeline(project_name='price-prediction', with_wandb = False, sweep = False, get_config=get_config) - run_inference_pipeline(data_config, training_config, all_models_all_assets) + run_inference_pipeline(data_config, training_config, model_config, all_models_all_assets) if __name__ == '__main__': - run_inference(preload_models=False, get_config=get_lightweight_ensemble_config) \ No newline at end of file + run_inference(preload_models=True, get_config=get_lightweight_ensemble_config) \ No newline at end of file diff --git a/run_pipeline.py b/run_pipeline.py index 5e59937..8b34c3f 100644 --- a/run_pipeline.py +++ b/run_pipeline.py @@ -20,14 +20,14 @@ import ray ray.init() -def run_pipeline(project_name:str, with_wandb: bool, sweep: bool, get_config: Callable) -> tuple[list[Reporting.Asset], dict, dict, pd.DataFrame, pd.DataFrame, pd.DataFrame]: +def run_pipeline(project_name:str, with_wandb: bool, sweep: bool, get_config: Callable) -> tuple[list[Reporting.Asset], dict, dict, dict, pd.DataFrame, pd.DataFrame, pd.DataFrame]: wandb, model_config, training_config, data_config = __setup_config(project_name, with_wandb, sweep, get_config) reporting = __run_training(model_config, training_config, data_config) results, all_predictions, all_probabilities, all_models_all_assets = reporting.get_results() report_results(results, all_predictions, model_config, wandb, sweep, project_name) - save_models(all_models_all_assets, data_config, training_config) + save_models(all_models_all_assets, data_config, training_config, model_config) - return all_models_all_assets, data_config, training_config, results, all_predictions, all_probabilities + return all_models_all_assets, data_config, training_config, model_config, results, all_predictions, all_probabilities def __setup_config(project_name:str, with_wandb: bool, sweep: bool, get_config: Callable) -> tuple[Optional[object], dict, dict, dict]: @@ -55,7 +55,7 @@ def __run_training(model_config:dict, training_config:dict, data_config:dict): # 1. Load data, check for validity and process data (feature selection, dimensionality reduction, etc.) X, y, target_returns = load_data(**configs['data_config']) - if check_data(X, y, training_config) is False: continue + if check_data(X, y, configs['training_config']) is False: continue X, original_X = process_data(X, y, configs) # 2. Train a Primary model with optional metalabeling for each asset diff --git a/training/inference.py b/training/inference.py index bd389b6..92805fc 100644 --- a/training/inference.py +++ b/training/inference.py @@ -1,39 +1,51 @@ import pandas as pd -from data_loader.load_data import load_data from typing import Optional, Union import warnings -from reporting.types import Reporting +from data_loader.load_data import load_data +from data_loader.process_data import process_data, check_data -def run_inference_pipeline(data_config:dict, training_config:dict, all_models_all_assets:list[Reporting.Asset]): - data_params = data_config.copy() - data_params['target_asset'] = data_params['assets'][0] +from reporting.types import Reporting +from training.training_steps import primary_step, secondary_step + +def run_inference_pipeline(data_config:dict, training_config:dict, model_config:dict, all_models_all_assets:list[Reporting.Asset]): + configs = dict(model_config=model_config, training_config=training_config, data_config=data_config) + configs['data_config']['target_asset'] = data_config['assets'][0] - X, y, _ = load_data(**data_params) - input_features = __select_data(X, training_config) + primary_models, secondary_models = __select_models(configs, all_models_all_assets) - primary_step, secondary_step = __select_models(data_params, all_models_all_assets) - - result = __inference(input_features, primary_step, secondary_step) + result = __inference(configs, primary_models, secondary_models) return result -def __inference(data:pd.DataFrame, primary_step:Union[Reporting.Training_Step,None], secondary_step:Union[Reporting.Training_Step,None]) -> pd.DataFrame: - assert primary_step is not None, "No primary models found. Cancelling Inference." - - data = __primary_models(data, primary_step) +def __inference(configs:dict, primary_models:Union[Reporting.Training_Step,None], secondary_models:Union[Reporting.Training_Step, None]): + reporting = Reporting() + asset = configs['data_config']['target_asset'] + # 1. Load data, truncate it, check for validity and process data (feature selection, dimensionality reduction, etc.) + X, y, target_returns = load_data(**configs['data_config']) + X, y = __select_data(X, y, configs['training_config']) + assert check_data(X, y, configs['training_config']) == False, "Data is not valid. Cancelling Inference." + X, original_X = process_data(X, y, configs) + + # 2. Train a Primary model with optional metalabeling for each asset + training_step_primary, current_predictions = primary_step(X, y, original_X, asset, target_returns, configs, reporting, primary_models) + + # 3. Train an Ensemble model with optional metalabeling for each asset if secondary_step is not None: warnings.warn("Secondary models are not specified.") - data = __secondary_models(data, secondary_step) - - return data + training_step_secondary = secondary_step(X, y, original_X, current_predictions, asset, target_returns, configs, reporting, secondary_models) + + # 4. Save the models + reporting.all_assets.append(Reporting.Asset(ticker=asset, primary=training_step_primary, secondary=training_step_secondary)) + + return reporting -def __select_models( data_params:dict, all_models_all_assets:list[Reporting.Asset])-> tuple[Union[Reporting.Training_Step,None], Union[Reporting.Training_Step, None]]: - target_asset_name = data_params['target_asset'][1] +def __select_models( configs:dict, all_models_all_assets:list[Reporting.Asset])-> tuple[Union[Reporting.Training_Step,None], Union[Reporting.Training_Step, None]]: + target_asset_name = configs['data_config']['target_asset'][1] primary_step, secondary_step, = None, None target_asset_models = next((x for x in all_models_all_assets if x.name == target_asset_name), None) @@ -51,34 +63,13 @@ def __select_models( data_params:dict, all_models_all_assets:list[Reporting.Asse return primary_step, secondary_step -def __select_data(X:pd.DataFrame, training_config:dict)-> pd.DataFrame: +def __select_data(X:pd.DataFrame, y:pd.Series, training_config:dict)-> tuple[pd.DataFrame, pd.Series]: window_size = training_config['sliding_window_size_primary'] num_rows = X.shape[0] if num_rows <= window_size: - return X.copy() + return X.copy(), y.copy() else: - return X.truncate(before=int(num_rows-window_size), after=num_rows, copy=True) + return X.truncate(before=int(num_rows-window_size), after=num_rows, copy=True), y.truncate(before=int(num_rows-window_size), after=num_rows, copy=True) -def __primary_models(data:pd.DataFrame, models:dict)-> pd.DataFrame: - for k, model in models: - last_model = model[-1] - prediction = last_model.predict(data.to_numpy()) - - # result = evaluate_predictions( - # model_name = model_name, - # target_returns = target_returns, - # y_pred = preds, - # y_true = y, - # method = method, - # no_of_classes=no_of_classes, - # print_results = print_results, - # discretize=True - # ) - - return data - -def __secondary_models(data:pd.DataFrame, model:dict)-> pd.DataFrame: - - return data diff --git a/training/meta_labeling.py b/training/meta_labeling.py index a76afea..88fbd91 100644 --- a/training/meta_labeling.py +++ b/training/meta_labeling.py @@ -6,6 +6,7 @@ import pandas as pd from models.model_map import default_feature_selector_regression, default_feature_selector_classification from models.base import Model from reporting.types import Reporting +from typing import Union def train_meta_labeling_model( @@ -18,7 +19,8 @@ def train_meta_labeling_model( data_config: dict, model_config: dict, training_config: dict, - model_suffix: str + model_suffix: str, + preloaded_models: Union[list[Reporting.Single_Model], None] = None ) -> tuple[pd.Series, pd.Series, pd.DataFrame, list[Reporting.Single_Model]]: @@ -48,7 +50,8 @@ def train_meta_labeling_model( scaler = training_config['scaler'], no_of_classes = 'two', level = 'meta_labeling', - print_results = False + print_results = False, + preloaded_models = preloaded_models ) if len(models) > 1: meta_preds = meta_preds.mean(axis = 1) diff --git a/training/primary_model.py b/training/primary_model.py index f6d0091..7d9dea1 100644 --- a/training/primary_model.py +++ b/training/primary_model.py @@ -1,5 +1,5 @@ import pandas as pd -from typing import Literal +from typing import Literal, Union from training.walk_forward import walk_forward_train, walk_forward_inference from utils.evaluate import evaluate_predictions from models.base import Model @@ -22,6 +22,7 @@ def train_primary_model( no_of_classes: Literal['two', 'three-balanced', 'three-imbalanced'], level: str, print_results: bool, + preloaded_models: Union[list[Reporting.Single_Model], None] = None ) -> tuple[pd.DataFrame, pd.DataFrame, pd.DataFrame, list[Reporting.Single_Model]]: results = pd.DataFrame() @@ -29,23 +30,25 @@ def train_primary_model( probabilities = pd.DataFrame(index=y.index) all_models_single_asset: list[Reporting.Single_Model] = [] - + if preloaded_models is not None: + models = preloaded_models for model_name, model in models: - model_over_time, transformations_over_time = walk_forward_train( - model_name=model_name, - model = model, - X = X if model.feature_selection == 'on' else original_X, - y = y, - target_returns = target_returns, - expanding_window = expanding_window, - window_size = sliding_window_size, - retrain_every = retrain_every, - transformations= [get_scaler(scaler)], - ) + if preloaded_models is None: + model_over_time, transformations_over_time = walk_forward_train( + model_name=model_name, + model = model, + X = X if model.feature_selection == 'on' else original_X, + y = y, + target_returns = target_returns, + expanding_window = expanding_window, + window_size = sliding_window_size, + retrain_every = retrain_every, + transformations= [get_scaler(scaler)], + ) preds, probs = walk_forward_inference( model_name = model_name, - model_over_time= model_over_time, + model_over_time= model_over_time if preloaded_models is None else pd.Series(model), transformations_over_time = transformations_over_time, X = X if model.feature_selection == 'on' else original_X, expanding_window = expanding_window, diff --git a/training/training_steps.py b/training/training_steps.py index ff2bdd5..be507ff 100644 --- a/training/training_steps.py +++ b/training/training_steps.py @@ -5,6 +5,7 @@ from training.primary_model import train_primary_model from training.meta_labeling import train_meta_labeling_model from reporting.types import Reporting +from typing import Union def primary_step( @@ -14,7 +15,8 @@ def primary_step( asset:list, target_returns:pd.Series, configs: dict, - reporting: Reporting + reporting: Reporting, + preloaded_training_step: Union[Reporting.Training_Step, None] = None ) -> tuple[Reporting.Training_Step, pd.DataFrame]: training_step = Reporting.Training_Step(level='primary') model_config, training_config, data_config = itemgetter('model_config', 'training_config', 'data_config')(configs) @@ -34,7 +36,8 @@ def primary_step( scaler = training_config['scaler'], no_of_classes = data_config['no_of_classes'], level = 'primary', - print_results= True + print_results= True, + preloaded_models = preloaded_training_step.convert_step_to_tuple('base') if preloaded_training_step is not None else None ) training_step.base = all_models_for_single_asset @@ -53,7 +56,8 @@ def primary_step( data_config= data_config, model_config= model_config, training_config= training_config, - model_suffix = 'meta' + model_suffix = 'meta', + preloaded_models =preloaded_training_step.convert_step_to_tuple('metalabeling') if preloaded_training_step is not None else None ) current_result[model_name] = primary_meta_result current_predictions[model_name] = primary_meta_preds