diff --git a/data_loader/process_data.py b/data_loader/process_data.py new file mode 100644 index 0000000..e5c32e0 --- /dev/null +++ b/data_loader/process_data.py @@ -0,0 +1,41 @@ +import pandas as pd +from operator import itemgetter + +from utils.helpers import has_enough_samples_to_train +from feature_selection.dim_reduction import reduce_dimensionality +from models.model_map import default_feature_selector_regression, default_feature_selector_classification +from feature_selection.feature_selection import select_features +import warnings + + + + +def process_data(X:pd.DataFrame, y:pd.Series, configs: dict) -> tuple[pd.DataFrame,pd.DataFrame,pd.DataFrame]: + + model_config, training_config, data_config = itemgetter('model_config', 'training_config', 'data_config')(configs) + + original_X = X.copy() + + # 2a. Dimensionality Reduction (optional) + if training_config['dimensionality_reduction']: + X_pca = reduce_dimensionality(X, int(len(X.columns) / 2)) + X = X_pca.copy() + else: + X_pca = X.copy() + + # 2b. Feature Selection + print("Feature Selection started") + # TODO: this needs to be done per model! + backup_model = default_feature_selector_regression if data_config['method'] == 'regression' else default_feature_selector_classification + X = select_features(X = X, y = y, model = model_config['primary_models'][0][1], n_features_to_select = training_config['n_features_to_select'], backup_model = backup_model, scaling = training_config['scaler']) + + return X, original_X, X_pca + +def check_data(X:pd.DataFrame, y:pd.Series, training_config:dict): + """ Returns True if data is valid, else returns False.""" + + if has_enough_samples_to_train(X, y, training_config) == False: + warnings.warn("Not enough samples to train") + return False + + return True \ No newline at end of file diff --git a/models/base.py b/models/base.py index c042cc8..3ea6343 100644 --- a/models/base.py +++ b/models/base.py @@ -3,7 +3,7 @@ from typing import Literal, Optional, Union from abc import ABC, abstractmethod import numpy as np -import numpy as np +# import numpy as np class Model(ABC): diff --git a/models/saving.py b/models/saving.py index 1f70a3b..9aa07f6 100644 --- a/models/saving.py +++ b/models/saving.py @@ -3,12 +3,15 @@ import datetime from typing import Optional, Union import os import warnings +from utils.encapsulation import Asset -def save_models(all_models_for_all_assets: dict, data_config:dict, training_config:dict) -> None: - all_models_for_all_assets['training_config'] = training_config - all_models_for_all_assets['data_config'] = data_config +def save_models(all_models_for_all_assets: list[Asset], data_config:dict, training_config:dict) -> None: + dict_for_pickle = dict() + dict_for_pickle['training_config'] = training_config + dict_for_pickle['data_config'] = data_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") @@ -16,7 +19,7 @@ def save_models(all_models_for_all_assets: dict, data_config:dict, training_conf warnings.warn("No folder exists, creating one.") os.makedirs('output/models') - pickle.dump( all_models_for_all_assets, open( "output/models/{}.p".format(date_string), "wb" ) ) + 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]: diff --git a/run_pipeline.py b/run_pipeline.py index 00e9497..8d336e6 100644 --- a/run_pipeline.py +++ b/run_pipeline.py @@ -1,26 +1,30 @@ -from config.hashing import hash_data_config -from data_loader.load_data import load_data import pandas as pd -from training.primary_model import train_primary_model +from typing import Callable, Optional +from operator import itemgetter + +from data_loader.load_data import load_data +from data_loader.process_data import process_data, check_data + from reporting.wandb import launch_wandb, register_config_with_wandb -from models.model_map import default_feature_selector_regression, default_feature_selector_classification +from reporting.reporting import report_results + from models.saving import save_models -from utils.helpers import has_enough_samples_to_train + from config.config import get_default_ensemble_config from config.preprocess import validate_config, preprocess_config -from feature_selection.feature_selection import select_features -from feature_selection.dim_reduction import reduce_dimensionality -from training.meta_labeling import train_meta_labeling_model -from reporting.reporting import report_results -from typing import Callable, Optional +from training.training_steps import primary_step, secondary_step + +from utils.encapsulation import Reporting, Asset, Training_Step + import ray ray.init() -def run_pipeline(project_name:str, with_wandb: bool, sweep: bool, get_config: Callable) -> tuple[dict, dict, dict, pd.DataFrame, pd.DataFrame, pd.DataFrame]: +def run_pipeline(project_name:str, with_wandb: bool, sweep: bool, get_config: Callable) -> tuple[list[Asset], dict, dict, pd.DataFrame, pd.DataFrame, pd.DataFrame]: wandb, model_config, training_config, data_config = __setup_config(project_name, with_wandb, sweep, get_config) - results, all_predictions, all_probabilities, all_models_all_assets = __run_training(model_config, training_config, data_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) @@ -36,134 +40,35 @@ def __setup_config(project_name:str, with_wandb: bool, sweep: bool, get_config: model_config, training_config, data_config = preprocess_config(model_config, training_config, data_config) return wandb, model_config, training_config, data_config - + + def __run_training(model_config:dict, training_config:dict, data_config:dict): - results = pd.DataFrame() - all_predictions = pd.DataFrame() - all_probabilities = pd.DataFrame() - all_models_for_all_assets = dict() - validate_config(model_config, training_config, data_config) + validate_config(model_config, training_config, data_config) + configs = dict(model_config=model_config, training_config=training_config, data_config=data_config) + reporting = Reporting() + for asset in data_config['assets']: print('--------\nPredicting: ', asset[1]) - # 1. Load data - data_params = data_config.copy() - data_params['target_asset'] = asset + configs['data_config']['target_asset'] = asset - X, y, target_returns = load_data(**data_params) - original_X = X.copy() - if has_enough_samples_to_train(X, y, training_config) == False: - print("Not enough samples to train") - continue + # 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 + X, original_X, X_pca = process_data(X, y, configs) - # 2a. Dimensionality Reduction (optional) - if training_config['dimensionality_reduction']: - X_pca = reduce_dimensionality(X, int(len(X.columns) / 2)) - X = X_pca.copy() - else: - X_pca = X.copy() - - # 2b. Feature Selection - print("Feature Selection started") - # TODO: this needs to be done per model! - backup_model = default_feature_selector_regression if data_config['method'] == 'regression' else default_feature_selector_classification - X = select_features(X = X, y = y, model = model_config['primary_models'][0][1], n_features_to_select = training_config['n_features_to_select'], backup_model = backup_model, scaling = training_config['scaler']) - - # 3. Train Primary models - current_result, current_predictions, current_probabilities, all_models_for_single_asset = train_primary_model( - ticker_to_predict = asset[1], - original_X = original_X, - X = X, - y = y, - target_returns = target_returns, - models = model_config['primary_models'], - method = data_config['method'], - expanding_window = training_config['expanding_window_primary'], - sliding_window_size = training_config['sliding_window_size_primary'], - retrain_every = training_config['retrain_every'], - scaler = training_config['scaler'], - no_of_classes = data_config['no_of_classes'], - level = 'primary', - print_results= True - ) + # 2. Train a Primary model with optional metalabeling for each asset + training_step_primary, current_predictions = primary_step(X, y, original_X, X_pca, asset, target_returns, configs, reporting) - all_models_for_all_assets[asset[1]] = dict( - name=asset[1], - primary_models = all_models_for_single_asset - ) + # 3. Train an Ensemble model with optional metalabeling for each asset + training_step_secondary = secondary_step(X, y, original_X, X_pca, current_predictions, asset, target_returns, configs, reporting) - # 4. Train a Meta-Labeling model for each Primary model and replace their predictions with the meta-labeling predictions - if training_config['primary_models_meta_labeling'] == True: - for model_name in current_result.columns: - primary_model_predictions = current_predictions[model_name] - primary_meta_result, primary_meta_preds, primary_meta_probabilities, meta_labeling_models = train_meta_labeling_model( - target_asset=asset[1], - X_pca = X_pca, - input_predictions= primary_model_predictions, - y = y, - target_returns = target_returns, - models = model_config['meta_labeling_models'], - data_config= data_config, - model_config= model_config, - training_config= training_config, - model_suffix = 'meta' - ) - current_result[model_name] = primary_meta_result - current_predictions[model_name] = primary_meta_preds - - all_models_for_all_assets[asset[1]]['primary_models'][model_name]['meta_labeling'] = meta_labeling_models - - results = pd.concat([results, current_result], axis=1) - # With static models, because of the lag in the indicator, the first prediction is NA, so we fill it with zero. - all_predictions = pd.concat([all_predictions, current_predictions], axis=1).fillna(0.) - all_probabilities = pd.concat([all_probabilities, current_probabilities], axis=1).fillna(0.) - - # 5. Ensemble primary model predictions (If Ensemble model is present) - if model_config['ensemble_model'] is not None: - - ensemble_result, ensemble_predictions, _, ensemble_models_one_asset = train_primary_model( - ticker_to_predict = asset[1], - original_X = current_predictions, - X = current_predictions, - y = y, - target_returns = target_returns, - models = [model_config['ensemble_model']], - method = data_config['method'], - expanding_window = False, - sliding_window_size = 1, - retrain_every = training_config['retrain_every'], - scaler = training_config['scaler'], - no_of_classes = data_config['no_of_classes'], - level = 'ensemble', - print_results= True, - ) - ensemble_result, ensemble_predictions = ensemble_result.iloc[:,0], ensemble_predictions.iloc[:,0] - all_models_for_all_assets[asset[1]]['secondary_model'] = ensemble_models_one_asset - - if len(model_config['meta_labeling_models']) > 0: - - # 3. Train a Meta-labeling model on the averaged level-1 model predictions - ensemble_meta_result, ensemble_meta_predictions, ensemble_meta_probabilities, ensemble_meta_labeling_models = train_meta_labeling_model( - target_asset=asset[1], - X_pca = X_pca, - input_predictions= ensemble_predictions, - y = y, - target_returns = target_returns, - models = model_config['meta_labeling_models'], - data_config= data_config, - model_config= model_config, - training_config= training_config, - model_suffix = 'ensemble' - ) - - all_models_for_all_assets[asset[1]]['secondary_model'][model_config['ensemble_model']] = dict(meta_labeling=ensemble_meta_labeling_models) - results = pd.concat([results, ensemble_meta_result], axis=1) - all_predictions = pd.concat([all_predictions, ensemble_meta_predictions], axis=1) - all_probabilities = pd.concat([all_probabilities, ensemble_meta_probabilities], axis=1).fillna(0.) - - return results, all_predictions, all_probabilities, all_models_for_all_assets + # 4. Save the models + reporting.all_assets.append(Asset(ticker=asset[1], primary=training_step_primary, secondary=training_step_secondary)) + + return reporting diff --git a/training/inference.py b/training/inference.py index efcdbdd..39e4efa 100644 --- a/training/inference.py +++ b/training/inference.py @@ -4,46 +4,51 @@ from data_loader.load_data import load_data from typing import Optional, Union import warnings +from utils.encapsulation import Asset, Single_Model, Training_Step -def run_inference_pipeline(data_config:dict, training_config:dict, all_models_all_assets:dict): +def run_inference_pipeline(data_config:dict, training_config:dict, all_models_all_assets:list[Asset]): data_params = data_config.copy() data_params['target_asset'] = data_params['assets'][0] X, y, _ = load_data(**data_params) input_features = __select_data(X, training_config) - primary_models, secondary_models = __select_models(data_params, all_models_all_assets) + primary_step, secondary_step = __select_models(data_params, all_models_all_assets) - result = __inference(input_features, primary_models, secondary_models) + result = __inference(input_features, primary_step, secondary_step) return result -def __inference(data:pd.DataFrame, primary_models:Union[dict,None], secondary_models:Union[dict,None]) -> pd.DataFrame: - assert primary_models is not None, "No primary models found. Cancelling Inference." +def __inference(data:pd.DataFrame, primary_step:Union[Training_Step,None], secondary_step:Union[Training_Step,None]) -> pd.DataFrame: + assert primary_step is not None, "No primary models found. Cancelling Inference." - data = __primary_models(data, primary_models) + data = __primary_models(data, primary_step) - if secondary_models is not None: + if secondary_step is not None: warnings.warn("Secondary models are not specified.") - data = __secondary_models(data, secondary_models) + data = __secondary_models(data, secondary_step) return data -def __select_models( data_params:dict, all_models_all_assets:dict)-> tuple[Optional[dict],Optional[dict]]: +def __select_models( data_params:dict, all_models_all_assets:list[Asset])-> tuple[Union[Training_Step,None], Union[Training_Step, None]]: target_asset_name = data_params['target_asset'][1] - primary_models, secondary_models = None, None + 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) - if 'primary_models' in all_models_all_assets[target_asset_name]: - primary_models = all_models_all_assets[target_asset_name]['primary_models'] + if target_asset_models is not None: + if len(target_asset_models.primary.base)>0: + primary_step = target_asset_models.primary + else: warnings.warn("No primary models found for {}.".format(target_asset_name)) + + if len(target_asset_models.secondary.base)>0: + secondary_step = target_asset_models.secondary + else: warnings.warn("No secondary models found for {}.".format(target_asset_name)) else: - assert("No primary models found for asset: " + target_asset_name) + assert("No models found for asset: " + target_asset_name) - if 'secondary_model' in all_models_all_assets[target_asset_name]: - secondary_models = all_models_all_assets[target_asset_name]['secondary_model'] - - return primary_models, secondary_models + return primary_step, secondary_step def __select_data(X:pd.DataFrame, training_config:dict)-> pd.DataFrame: diff --git a/training/meta_labeling.py b/training/meta_labeling.py index f9ec879..e59a0c8 100644 --- a/training/meta_labeling.py +++ b/training/meta_labeling.py @@ -5,6 +5,7 @@ from feature_selection.feature_selection import select_features import pandas as pd from models.model_map import default_feature_selector_regression, default_feature_selector_classification from models.base import Model +from utils.encapsulation import Single_Model def train_meta_labeling_model( @@ -18,8 +19,9 @@ def train_meta_labeling_model( model_config: dict, training_config: dict, model_suffix: str - ) -> tuple[pd.Series, pd.Series, pd.DataFrame, dict]: + ) -> tuple[pd.Series, pd.Series, pd.DataFrame, list[Single_Model]]: + discretize = discretize_threeway_threshold(0.33) discretized_predictions = input_predictions.apply(discretize) meta_y: pd.Series = pd.concat([discretized_predictions, y], axis=1).apply(equal_except_nan, axis = 1) @@ -67,5 +69,6 @@ def train_meta_labeling_model( discretize=False ) meta_result.rename("model_" + target_asset + "_" + model_suffix, inplace=True) + return meta_result, avg_predictions_with_sizing, meta_probabilities, all_models_single_asset diff --git a/training/primary_model.py b/training/primary_model.py index f04a778..993aecb 100644 --- a/training/primary_model.py +++ b/training/primary_model.py @@ -5,6 +5,7 @@ from utils.evaluate import evaluate_predictions from models.base import Model from utils.scaler import get_scaler from utils.types import ScalerTypes +from utils.encapsulation import Training_Step, Single_Model, Asset def train_primary_model( ticker_to_predict: str, @@ -20,15 +21,17 @@ def train_primary_model( scaler: ScalerTypes, no_of_classes: Literal['two', 'three-balanced', 'three-imbalanced'], level: str, - print_results: bool - ) -> tuple[pd.DataFrame, pd.DataFrame, pd.DataFrame, dict]: + print_results: bool, + ) -> tuple[pd.DataFrame, pd.DataFrame, pd.DataFrame, list[Single_Model]]: scaler = get_scaler(scaler) results = pd.DataFrame() - all_models_single_asset = dict() predictions = pd.DataFrame(index=y.index) probabilities = pd.DataFrame(index=y.index) + all_models_single_asset:list[Single_Model] = [] + + for model_name, model in models: model_over_time, scaler_over_time = walk_forward_train( @@ -65,15 +68,16 @@ def train_primary_model( levelname=("_" + level) if level=='metalabeling' else "" column_name = "model_" + model_name + "_" + ticker_to_predict + levelname results[column_name] = result - all_models_single_asset[column_name]=dict() - all_models_single_asset[column_name][level] = model_over_time.tolist() - # all_models_single_asset[model_name]=dict() - # all_models_single_asset[model_name][level] = model_over_time.tolist() + + + all_models_single_asset.append(Single_Model(model_name=column_name, model_over_time=model_over_time.tolist())) + # column names for model outputs should be different, so we can differentiate between original data and model predictions later, where necessary predictions[column_name] = preds probs_column_name = "probs_" + ticker_to_predict + "_" + model_name + "_" + level probs.columns = [probs_column_name + "_" + c for c in probs.columns] probabilities = pd.concat([probabilities, probs], axis=1) - + + return results, predictions, probabilities, all_models_single_asset \ No newline at end of file diff --git a/training/training_steps.py b/training/training_steps.py new file mode 100644 index 0000000..b633675 --- /dev/null +++ b/training/training_steps.py @@ -0,0 +1,117 @@ +import pandas as pd +from operator import itemgetter + +from training.primary_model import train_primary_model +from training.meta_labeling import train_meta_labeling_model + +from utils.encapsulation import Reporting, Asset, Single_Model, Training_Step + + +def primary_step(X: pd.DataFrame, y:pd.Series, original_X:pd.DataFrame, X_pca:pd.DataFrame, asset:list, target_returns:pd.Series, configs: dict, reporting: Reporting) -> tuple[Training_Step, pd.DataFrame]: + training_step = Training_Step(level='primary') + model_config, training_config, data_config = itemgetter('model_config', 'training_config', 'data_config')(configs) + + # 3. Train Primary models + current_result, current_predictions, current_probabilities, all_models_for_single_asset = train_primary_model( + ticker_to_predict = asset[1], + original_X = original_X, + X = X, + y = y, + target_returns = target_returns, + models = model_config['primary_models'], + method = data_config['method'], + expanding_window = training_config['expanding_window_primary'], + sliding_window_size = training_config['sliding_window_size_primary'], + retrain_every = training_config['retrain_every'], + scaler = training_config['scaler'], + no_of_classes = data_config['no_of_classes'], + level = 'primary', + print_results= True + ) + + training_step.base = all_models_for_single_asset + + # 4. Train a Meta-Labeling model for each Primary model and replace their predictions with the meta-labeling predictions + if training_config['primary_models_meta_labeling'] == True: + for model_name in current_result.columns: + primary_model_predictions = current_predictions[model_name] + primary_meta_result, primary_meta_preds, primary_meta_probabilities, meta_labeling_models = train_meta_labeling_model( + target_asset=asset[1], + X_pca = X_pca, + input_predictions= primary_model_predictions, + y = y, + target_returns = target_returns, + models = model_config['meta_labeling_models'], + data_config= data_config, + model_config= model_config, + training_config= training_config, + model_suffix = 'meta' + ) + current_result[model_name] = primary_meta_result + current_predictions[model_name] = primary_meta_preds + + training_step.metalabeling.append(meta_labeling_models) + + reporting.results = pd.concat([reporting.results, current_result], axis=1) + # With static models, because of the lag in the indicator, the first prediction is NA, so we fill it with zero. + reporting.all_predictions = pd.concat([reporting.all_predictions, current_predictions], axis=1).fillna(0.) + reporting.all_probabilities = pd.concat([reporting.all_probabilities, current_probabilities], axis=1).fillna(0.) + + return training_step, current_predictions + + +def secondary_step(X:pd.DataFrame, y:pd.Series, original_X:pd.DataFrame, X_pca:pd.DataFrame, current_predictions:pd.DataFrame, asset:list, target_returns:pd.Series, configs: dict, reporting: Reporting) -> Training_Step: + training_step = Training_Step(level='secondary') + model_config, training_config, data_config = itemgetter('model_config', 'training_config', 'data_config')(configs) + + # 5. Ensemble primary model predictions (If Ensemble model is present) + if model_config['ensemble_model'] is not None: + ensemble_result, ensemble_predictions, _, ensemble_models_one_asset = train_primary_model( + ticker_to_predict = asset[1], + original_X = current_predictions, + X = current_predictions, + y = y, + target_returns = target_returns, + models = [model_config['ensemble_model']], + method = data_config['method'], + expanding_window = False, + sliding_window_size = 1, + retrain_every = training_config['retrain_every'], + scaler = training_config['scaler'], + no_of_classes = data_config['no_of_classes'], + level = 'ensemble', + print_results= True, + ) + ensemble_result, ensemble_predictions = ensemble_result.iloc[:,0], ensemble_predictions.iloc[:,0] + + training_step.base = ensemble_models_one_asset + + reporting.results = pd.concat([reporting.results, ensemble_result], axis=1) + reporting.all_predictions = pd.concat([reporting.all_predictions, ensemble_predictions], axis=1) + + + + if len(model_config['meta_labeling_models']) > 0: + + # 3. Train a Meta-labeling model on the averaged level-1 model predictions + ensemble_meta_result, ensemble_meta_predictions, ensemble_meta_probabilities, ensemble_meta_labeling_models = train_meta_labeling_model( + target_asset=asset[1], + X_pca = X_pca, + input_predictions= ensemble_predictions, + y = y, + target_returns = target_returns, + models = model_config['meta_labeling_models'], + data_config= data_config, + model_config= model_config, + training_config= training_config, + model_suffix = 'ensemble' + ) + + training_step.metalabeling.append(ensemble_meta_labeling_models) + + + reporting.results = pd.concat([reporting.results, ensemble_meta_result], axis=1) + reporting.all_predictions = pd.concat([reporting.all_predictions, ensemble_meta_predictions], axis=1) + reporting.all_probabilities = pd.concat([reporting.all_probabilities, ensemble_meta_probabilities], axis=1).fillna(0.) + + return training_step diff --git a/utils/encapsulation.py b/utils/encapsulation.py new file mode 100644 index 0000000..5d5f50b --- /dev/null +++ b/utils/encapsulation.py @@ -0,0 +1,38 @@ +import pandas as pd +from models.base import Model + +# | Reporting +# | + + +class Single_Model: + def __init__(self, model_name:str, model_over_time:list[Model]): + self.model_name: str = model_name + self.model_over_time: list[Model] = model_over_time + + +class Training_Step: + def __init__(self, level:str): + self.level:str = level + self.base: list[Single_Model] = [] + self.metalabeling: list[list[Single_Model]] = [] + + +class Asset(): + def __init__(self, ticker:str, primary: Training_Step, secondary: Training_Step): + self.name:str = ticker + self.primary:Training_Step = primary + self.secondary:Training_Step = secondary + + +class Reporting: + def __init__(self): + self.results:pd.DataFrame = pd.DataFrame() + self.all_predictions:pd.DataFrame = pd.DataFrame() + self.all_probabilities:pd.DataFrame = pd.DataFrame() + self.all_assets:list[Asset] = [] + + def get_results(self)->tuple[pd.DataFrame, pd.DataFrame, pd.DataFrame, list[Asset]]: + return self.results, self.all_predictions, self.all_probabilities, self.all_assets + +