From e80fffdb65ca404baaedce80844ed6c232996ea1 Mon Sep 17 00:00:00 2001 From: Mark Aron Szulyovszky Date: Sun, 23 Jan 2022 18:37:43 +0100 Subject: [PATCH] refactor(Config): use a Config object instead of dictionary of dictionaries! (#184) * refactor(Config): use a Config object instead of dictionary of dictionaries! * fix(Config): use default_ensemble_config * fix(Portfolio): fixed portfolio construction --- config/config.py | 141 ++++++++++++++++++++----------- config/hashing.py | 1 - config/preprocess.py | 30 +++---- config/sweep_ensemble.yaml | 4 +- config/sweep_meta_labeling.yaml | 4 +- config/sweep_primary_models.yaml | 4 +- data_loader/load_data.py | 23 +---- data_loader/process_data.py | 6 +- environment.yml | 1 + exploration.ipynb | 2 +- reporting/reporting.py | 8 +- reporting/saving.py | 15 ++-- reporting/wandb.py | 26 +++--- run_inference.py | 22 ++--- run_pipeline.py | 56 +++++++----- training/meta_labeling.py | 14 ++- training/training_steps.py | 54 ++++++------ utils/helpers.py | 6 +- utils/types.py | 2 +- 19 files changed, 223 insertions(+), 196 deletions(-) diff --git a/config/config.py b/config/config.py index fd303b8..6e76421 100644 --- a/config/config.py +++ b/config/config.py @@ -1,19 +1,87 @@ -def get_dev_config() -> tuple[dict, dict, dict]: +from pydantic import BaseModel +from typing import Literal, Optional +from models.base import Model +from utils.types import DataCollection, DataSource, FeatureExtractor + +# RawConfig is needed to ensure we can declare config presets here with static typing, we then convert it to Config +class RawConfig(BaseModel): + primary_models_meta_labeling: bool + dimensionality_reduction: bool + n_features_to_select: int + expanding_window_base: bool + expanding_window_meta_labeling: bool + sliding_window_size_base: int + sliding_window_size_meta_labeling: int + retrain_every: int + scaler: Literal['normalize', 'minmax', 'standardize'] + + assets: list[str] + target_asset: str + other_assets: list[str] + exogenous_data: list[str] + load_non_target_asset: bool + log_returns: bool + forecasting_horizon: int + own_features: list[str] + other_features: list[str] + exogenous_features: list[str] + no_of_classes: Literal['two', 'three-balanced', 'three-imbalanced'] + index_column: Literal['date', 'int'] + + primary_models: list[str] + meta_labeling_models: list[str] + ensemble_model: Optional[str] + + +class Config(BaseModel): + primary_models_meta_labeling: bool + dimensionality_reduction: bool + n_features_to_select: int + expanding_window_base: bool + expanding_window_meta_labeling: bool + sliding_window_size_base: int + sliding_window_size_meta_labeling: int + retrain_every: int + scaler: Literal['normalize', 'minmax', 'standardize'] + + assets: DataCollection + target_asset: DataSource + other_assets: DataCollection + exogenous_data: DataCollection + load_non_target_asset: bool + log_returns: bool + forecasting_horizon: int + own_features: list[tuple[str, FeatureExtractor, list[int]]] + other_features: list[tuple[str, FeatureExtractor, list[int]]] + exogenous_features: list[tuple[str, FeatureExtractor, list[int]]] + no_of_classes: Literal['two', 'three-balanced', 'three-imbalanced'] + index_column: Literal['date', 'int'] + + primary_models: list[tuple[str, Model]] + meta_labeling_models: list[tuple[str, Model]] + ensemble_model: Optional[tuple[str, Model]] + + class Config: + arbitrary_types_allowed = True + + +def get_dev_config() -> RawConfig: - training_config = dict( + regression_models = ["Lasso"] + classification_models = ["LogisticRegression_two_class"] + + return RawConfig( primary_models_meta_labeling = False, dimensionality_reduction = False, n_features_to_select = 30, - expanding_window_primary = False, + expanding_window_base = False, expanding_window_meta_labeling = False, - sliding_window_size_primary = 380, + sliding_window_size_base = 380, sliding_window_size_meta_labeling = 1, retrain_every = 20, scaler = 'minmax', # 'normalize' 'minmax' 'standardize' - ) - data_config = dict( assets = ['daily_only_btc'], target_asset = 'BTC_USD', other_assets = [], @@ -26,37 +94,31 @@ def get_dev_config() -> tuple[dict, dict, dict]: exogenous_features = ['z_score'], index_column= 'int', no_of_classes= 'two', - narrow_format = False, - ) - regression_models = ["Lasso"] - classification_models = ["LogisticRegression_two_class"] - - model_config = dict( primary_models = classification_models, meta_labeling_models = [], ensemble_model = None ) - - return model_config, training_config, data_config - -def get_default_ensemble_config() -> tuple[dict, dict, dict]: +def get_default_ensemble_config() -> RawConfig: - training_config = dict( + regression_models = ["Lasso", "KNN", "RFR"] + classification_models = ["LogisticRegression_two_class", "LDA", "NB", "RFC", "XGB_two_class", "LGBM", "StaticMom"] + meta_labeling_models = ['LogisticRegression_two_class', 'LGBM'] + ensemble_model = 'Average' + + return RawConfig( primary_models_meta_labeling = True, dimensionality_reduction = False, n_features_to_select = 30, - expanding_window_primary = False, + expanding_window_base = False, expanding_window_meta_labeling = True, - sliding_window_size_primary = 380, + sliding_window_size_base = 380, sliding_window_size_meta_labeling = 240, retrain_every = 10, scaler = 'minmax', # 'normalize' 'minmax' 'standardize' - ) - data_config = dict( assets = ['daily_crypto'], target_asset = 'BTC_USD', other_assets = ['daily_etf'], @@ -69,41 +131,34 @@ def get_default_ensemble_config() -> tuple[dict, dict, dict]: exogenous_features = ['z_score'], index_column= 'int', no_of_classes= 'two', - narrow_format = False, - ) - regression_models = ["Lasso", "KNN", "RFR"] - classification_models = ["LogisticRegression_two_class", "LDA", "NB", "RFC", "XGB_two_class", "LGBM", "StaticMom"] - meta_labeling_models = ['LogisticRegression_two_class', 'LGBM'] - ensemble_model = 'Average' - - model_config = dict( primary_models = classification_models, meta_labeling_models = meta_labeling_models, ensemble_model = ensemble_model ) - - return model_config, training_config, data_config -def get_lightweight_ensemble_config() -> tuple[dict, dict, dict]: +def get_lightweight_ensemble_config() -> RawConfig: - training_config = dict( + regression_models = ["Lasso", "KNN"] + classification_models = ['LogisticRegression_two_class', 'SVC'] + meta_labeling_models = ['LogisticRegression_two_class', 'LGBM'] + ensemble_model = 'Average' + + return RawConfig( primary_models_meta_labeling = True, dimensionality_reduction = True, n_features_to_select = 30, - expanding_window_primary = False, + expanding_window_base = False, expanding_window_meta_labeling = True, - sliding_window_size_primary = 380, + sliding_window_size_base = 380, sliding_window_size_meta_labeling = 240, retrain_every = 40, scaler = 'minmax', # 'normalize' 'minmax' 'standardize' - ) - data_config = dict( assets = ['daily_crypto_lightweight'], - target_asset = 'BTC_USD', + target_asset = 'BCH_USD', other_assets = ['daily_etf'], exogenous_data = ['daily_glassnode'], load_non_target_asset= True, @@ -114,20 +169,10 @@ def get_lightweight_ensemble_config() -> tuple[dict, dict, dict]: exogenous_features = ['z_score'], index_column= 'int', no_of_classes= 'two', - narrow_format = False, - ) - regression_models = ["Lasso", "KNN"] - classification_models = ['LogisticRegression_two_class', 'SVC'] - meta_labeling_models = ['LogisticRegression_two_class', 'LGBM'] - ensemble_model = 'Average' - - model_config = dict( primary_models = classification_models, meta_labeling_models = meta_labeling_models, ensemble_model = ensemble_model ) - - return model_config, training_config, data_config diff --git a/config/hashing.py b/config/hashing.py index 9842ad8..07bcec8 100644 --- a/config/hashing.py +++ b/config/hashing.py @@ -19,5 +19,4 @@ def hash_data_config(data_config: dict) -> str: hash_feature_extractors(data_config['exogenous_features']), data_config['index_column'], data_config['no_of_classes'], - data_config['narrow_format'] ])) diff --git a/config/preprocess.py b/config/preprocess.py index 185f301..7778bbd 100644 --- a/config/preprocess.py +++ b/config/preprocess.py @@ -1,16 +1,19 @@ - +from config.config import Config, RawConfig from utils.helpers import flatten from feature_extractors.feature_extractor_presets import presets as feature_extractor_presets from models.model_map import get_model_map from data_loader.collections import data_collections -def preprocess_config(model_config:dict, training_config:dict, data_config:dict) -> tuple[dict, dict, dict]: - model_config = __preprocess_model_config(model_config) - data_config = __preprocess_feature_extractors_config(data_config) - data_config = __preprocess_data_collections_config(data_config) - validate_config(model_config, training_config, data_config) - return model_config, training_config, data_config +def preprocess_config(raw_config: RawConfig) -> Config: + config_dict = vars(raw_config) + config_dict = __preprocess_model_config(config_dict) + config_dict = __preprocess_feature_extractors_config(config_dict) + config_dict = __preprocess_data_collections_config(config_dict) + + config = Config(**config_dict) + validate_config(config) + return config def __preprocess_feature_extractors_config(data_dict: dict) -> dict: data_dict = data_dict.copy() @@ -42,18 +45,11 @@ def __preprocess_data_collections_config(data_dict: dict) -> dict: return data_dict -def validate_config(model_config:dict, training_config:dict, data_config:dict): +def validate_config(config: Config): # We need to make sure there's only one output from the pipeline # If level-2 model is there, we need more than one level-1 models to train - if len(model_config["meta_labeling_models"]) > 1: assert len(model_config["primary_models"]) > 0 + if len(config.meta_labeling_models) > 1: assert len(config.primary_models) > 0 # If there's no level-2 model, we need to have only one level-1 model - if len(model_config["meta_labeling_models"]) == 0: assert len(model_config["primary_models"]) == 1 + if len(config.meta_labeling_models) == 0: assert len(config.primary_models) == 1 -def get_model_name(model_config:dict) -> str: - if len(model_config["meta_labeling_models"]) > 0: - return model_config["meta_labeling_models"][0][0] - elif len(model_config["primary_models"]) == 1: - return model_config["primary_models"][0][0] - else: - raise Exception("No model name found") diff --git a/config/sweep_ensemble.yaml b/config/sweep_ensemble.yaml index d328bc9..94255f2 100644 --- a/config/sweep_ensemble.yaml +++ b/config/sweep_ensemble.yaml @@ -14,11 +14,11 @@ parameters: value: ['daily_etf'] exogenous_data: value: ['daily_glassnode'] - expanding_window_primary: + expanding_window_base: value: True expanding_window_meta_labeling: value: True - sliding_window_size_primary: + sliding_window_size_base: value: 380 sliding_window_size_meta_labeling: values: [250, 300, 380] diff --git a/config/sweep_meta_labeling.yaml b/config/sweep_meta_labeling.yaml index 2e5117f..b55b6e2 100644 --- a/config/sweep_meta_labeling.yaml +++ b/config/sweep_meta_labeling.yaml @@ -14,7 +14,7 @@ parameters: value: ['daily_etf'] exogenous_data: value: ['daily_glassnode'] - expanding_window_primary: + expanding_window_base: values: [True, False] distribution: categorical expanding_window_meta_labeling: @@ -25,7 +25,7 @@ parameters: distribution: categorical dimensionality_reduction: value: True - sliding_window_size_primary: + sliding_window_size_base: values: [180, 280, 380] distribution: categorical sliding_window_size_meta_labeling: diff --git a/config/sweep_primary_models.yaml b/config/sweep_primary_models.yaml index bb733ce..24a8d81 100644 --- a/config/sweep_primary_models.yaml +++ b/config/sweep_primary_models.yaml @@ -14,7 +14,7 @@ parameters: value: ['daily_etf'] exogenous_data: value: ['daily_glassnode'] - expanding_window_primary: + expanding_window_base: values: [True, False] distribution: categorical expanding_window_meta_labeling: @@ -23,7 +23,7 @@ parameters: value: 50 dimensionality_reduction: value: True - sliding_window_size_primary: + sliding_window_size_base: value: 380 sliding_window_size_meta_labeling: value: 380 diff --git a/data_loader/load_data.py b/data_loader/load_data.py index 451a24a..2945194 100644 --- a/data_loader/load_data.py +++ b/data_loader/load_data.py @@ -31,7 +31,6 @@ def __load_data(assets: DataCollection, exogenous_features: list[tuple[str, FeatureExtractor, list[int]]], index_column: Literal['date', 'int'], no_of_classes: Literal['two', 'three-balanced', 'three-imbalanced'], - narrow_format: bool = False ) -> tuple[pd.DataFrame, pd.Series, pd.Series]: """ Loads asset data from the specified path. @@ -51,7 +50,6 @@ def __load_data(assets: DataCollection, prefix=data_source[1], returns='log_returns' if log_returns else 'returns', feature_extractors=own_features, - narrow_format=narrow_format, ) for data_source in target_file] target_asset_df = ray.get(target_asset_future) @@ -60,7 +58,6 @@ def __load_data(assets: DataCollection, prefix=target_file[0][1], returns='returns', feature_extractors=[], - narrow_format=narrow_format, ) df_target_asset_only_returns = ray.get(target_asset_only_returns_future) @@ -69,7 +66,6 @@ def __load_data(assets: DataCollection, prefix=data_source[1], returns='log_returns' if log_returns else 'returns', feature_extractors=other_features, - narrow_format=narrow_format, ) for data_source in files] asset_dfs = ray.get(asset_futures) @@ -78,17 +74,13 @@ def __load_data(assets: DataCollection, prefix=data_source[1], returns='none', feature_extractors=exogenous_features, - narrow_format=narrow_format, ) for data_source in exogenous_data] exogenous_dfs = ray.get(exogenous_futures) dfs = target_asset_df + asset_dfs + exogenous_dfs dfs = [deduplicate_indexes(df) for df in dfs] target_df = dfs[0] - if narrow_format: - dfs = pd.concat([df.sort_index().reindex(target_df.index) for df in dfs], axis=0).fillna(0.) - else: - dfs = pd.concat([df.sort_index().reindex(target_df.index) for df in dfs], axis=1).fillna(0.) + dfs = pd.concat([df.sort_index().reindex(target_df.index) for df in dfs], axis=1).fillna(0.) dfs.index = pd.DatetimeIndex(dfs.index) @@ -96,10 +88,6 @@ def __load_data(assets: DataCollection, dfs.reset_index(drop=True, inplace=True) df_target_asset_only_returns.reset_index(drop=True, inplace=True) - if narrow_format: - dfs = dfs.drop(index=dfs.index[0], axis=0) - df_target_asset_only_returns = df_target_asset_only_returns.drop(index=dfs.index[0], axis=0) - ## Create target target_col = 'target' returns_col = target_asset[1] + '_returns' @@ -119,8 +107,7 @@ def __load_data(assets: DataCollection, def __load_df(data_source: DataSource, prefix: str, returns: Literal['none', 'price', 'returns', 'log_returns'], - feature_extractors: list[tuple[str, FeatureExtractor, list[int]]], - narrow_format: bool = False) -> pd.DataFrame: + feature_extractors: list[tuple[str, FeatureExtractor, list[int]]]) -> pd.DataFrame: df = pd.read_csv(os.path.join(data_source[0], data_source[1] + '.csv'), header=0, index_col=0).fillna(0) if returns == 'log_returns': @@ -135,10 +122,7 @@ def __load_df(data_source: DataSource, df = df.replace([np.inf, -np.inf], 0.) df = drop_columns_if_exist(df, ['open', 'high', 'low', 'close', 'volume']) - if narrow_format: - df["ticker"] = np.repeat(prefix, df.shape[0]) - else: - df.columns = [prefix + "_" + c if 'date' not in c else c for c in df.columns] + df.columns = [prefix + "_" + c if 'date' not in c else c for c in df.columns] return df @@ -239,7 +223,6 @@ def load_only_returns(assets: DataCollection, index_column: Literal['date', 'int prefix=data_source[1], returns=returns, feature_extractors=[], - narrow_format=False, ) for data_source in assets] target_asset_df = ray.get(assets_future) diff --git a/data_loader/process_data.py b/data_loader/process_data.py index 931d6b8..5271082 100644 --- a/data_loader/process_data.py +++ b/data_loader/process_data.py @@ -1,12 +1,12 @@ import pandas as pd - +from config.config import Config from utils.helpers import has_enough_samples_to_train import warnings -def check_data(X:pd.DataFrame, y:pd.Series, training_config:dict): +def check_data(X:pd.DataFrame, y:pd.Series, config: Config): """ Returns True if data is valid, else returns False.""" - if has_enough_samples_to_train(X, y, training_config) == False: + if has_enough_samples_to_train(X, y, config) == False: warnings.warn("Not enough samples to train") return False diff --git a/environment.yml b/environment.yml index 60368d8..cec6f6b 100644 --- a/environment.yml +++ b/environment.yml @@ -33,4 +33,5 @@ dependencies: - lightgbm - alphalens-reloaded - vectorbt + - pydantic prefix: /usr/local/anaconda3/envs/quant diff --git a/exploration.ipynb b/exploration.ipynb index 1cf1574..836a396 100644 --- a/exploration.ipynb +++ b/exploration.ipynb @@ -34,7 +34,7 @@ "model_config, training_config, data_config = get_default_ensemble_config()\n", "model_config, training_config, data_config = preprocess_config(model_config, training_config, data_config)\n", "\n", - "data_config['target_asset'] = data_config['assets'][0]\n", + "config.target_asset'] = config.assets'][0]\n", "X, y, target_returns = load_data(**data_config)" ] }, diff --git a/reporting/reporting.py b/reporting/reporting.py index 003ace6..0cefd56 100644 --- a/reporting/reporting.py +++ b/reporting/reporting.py @@ -1,16 +1,16 @@ from reporting.wandb import send_report_to_wandb import pandas as pd -from config.preprocess import get_model_name from utils.helpers import weighted_average +from config.config import Config -def report_results(results:pd.DataFrame, all_predictions:pd.DataFrame, model_config:dict, wandb, sweep: bool, project_name:str): +def report_results(results:pd.DataFrame, all_predictions:pd.DataFrame, config: Config, wandb, sweep: bool, project_name:str): primary_results = results[[column for column in results.columns if 'ensemble' not in column]] ensemble_results = results[[column for column in results.columns if 'ensemble' in column]] # Only send the results of the final model to wandb results_to_send = ensemble_results if ensemble_results.shape[1] > 0 else primary_results - send_report_to_wandb(results_to_send, wandb, project_name, get_model_name(model_config)) + send_report_to_wandb(results_to_send, wandb) results.to_csv('output/results.csv') primary_weights = all_predictions[[column for column in all_predictions.columns if 'ensemble' not in column]] @@ -29,7 +29,7 @@ def report_results(results:pd.DataFrame, all_predictions:pd.DataFrame, model_con print("Mean Sharpe ratio for Level-1 models: ", round(primary_avg_results.loc['sharpe'], 3)) print("Mean Probabilistic Sharpe ratio for Level-1 models: ", round(primary_avg_results.loc['prob_sharpe'].mean(), 3)) - if len(model_config['meta_labeling_models']) > 0: + if len(config.meta_labeling_models) > 0: print("Level-2 (Ensemble): Number of samples evaluated: ", ensemble_results.loc['no_of_samples'].sum()) print("Mean Sharpe ratio for Level-2 (Ensemble) models: ", round(ensemble_avg_results.loc['sharpe'].mean(), 3)) print("Mean Probabilistic Sharpe ratio for Level-2 (Ensemble) models: ", round(ensemble_avg_results.loc['prob_sharpe'].mean(), 3)) diff --git a/reporting/saving.py b/reporting/saving.py index b8725a1..b20765c 100644 --- a/reporting/saving.py +++ b/reporting/saving.py @@ -1,5 +1,6 @@ import pickle import datetime +from config.config import Config from typing import Optional, Union import os import warnings @@ -7,11 +8,9 @@ from reporting.types import Reporting -def save_models(all_models: Reporting.Asset, data_config:dict, training_config:dict, model_config:dict) -> None: +def save_models(all_models: Reporting.Asset, config: Config) -> 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['config'] = config dict_for_pickle['all_models'] = all_models date_string = datetime.datetime.now().strftime("%Y-%m-%d-%H-%M") @@ -23,7 +22,7 @@ def save_models(all_models: Reporting.Asset, data_config:dict, training_config:d pickle.dump( dict_for_pickle, open( "output/models/{}.p".format(date_string), "wb" ) ) -def load_models(file_name:Union[str, None]) -> tuple[Reporting.Asset, dict, dict, dict]: +def load_models(file_name:Union[str, None]) -> tuple[Reporting.Asset, Config]: if file_name is None: warnings.warn("No file name provided, will load latest models and configurations.") @@ -34,10 +33,8 @@ def load_models(file_name:Union[str, None]) -> tuple[Reporting.Asset, dict, dict packacked_dict = pickle.load( open( "output/models/{}".format(file_name), "rb" ) ) - data_config = packacked_dict.pop("data_config", None) - training_config = packacked_dict.pop("training_config", None) - model_config = packacked_dict.pop("model_config", None) + config = packacked_dict.pop("config", None) all_models = packacked_dict.pop("all_models", None) - return all_models, data_config, training_config, model_config + return all_models, config diff --git a/reporting/wandb.py b/reporting/wandb.py index 27f453c..f58ca60 100644 --- a/reporting/wandb.py +++ b/reporting/wandb.py @@ -1,36 +1,34 @@ import pandas as pd +from config.config import RawConfig from typing import Optional from utils.helpers import weighted_average -def launch_wandb(project_name:str, default_config:dict, sweep:bool=False): +def launch_wandb(project_name:str, default_config: RawConfig, sweep:bool=False) -> Optional[object]: from wandb_setup import get_wandb wandb = get_wandb() if wandb is None: raise Exception("Wandb can not be initalized, the environment variable WANDB_API_KEY is missing (can also use .env file)") elif sweep: - wandb.init(project=project_name, config = default_config) + wandb.init(project=project_name, config = vars(default_config)) return wandb else: - wandb.init(project=project_name, config = default_config, reinit=True) + wandb.init(project=project_name, config = vars(default_config), reinit=True) return wandb -def register_config_with_wandb(wandb: Optional[object], model_config:dict, training_config:dict, data_config:dict): - if wandb is None: return model_config, training_config, data_config +def override_config_with_wandb_values(wandb: Optional[object], raw_config: RawConfig) -> RawConfig: + if wandb is None: return raw_config - config: dict = wandb.config + wandb_config: dict = wandb.config - for k in training_config: - training_config[k] = config[k] - for k in model_config: - model_config[k] = config[k] - for k in data_config: - data_config[k] = config[k] + config_dict = vars(raw_config) + for k in config_dict: + config_dict[k] = wandb_config[k] - return model_config, training_config, data_config + return RawConfig(**config_dict) -def send_report_to_wandb(results: pd.DataFrame, wandb:Optional[object], project_name: str, model_name: str): +def send_report_to_wandb(results: pd.DataFrame, wandb:Optional[object]): if wandb is None: return run = wandb.run diff --git a/run_inference.py b/run_inference.py index 0e27f2c..bac375a 100644 --- a/run_inference.py +++ b/run_inference.py @@ -3,7 +3,7 @@ from data_loader.process_data import check_data from reporting.saving import load_models from run_pipeline import run_pipeline -from config.config import get_dev_config, get_default_ensemble_config, get_lightweight_ensemble_config +from config.config import Config, get_dev_config, get_default_ensemble_config, get_lightweight_ensemble_config from typing import Callable, Optional from reporting.types import Reporting @@ -12,30 +12,30 @@ import warnings def run_inference(preload_models:bool, get_config:Callable): if preload_models: - all_models, data_config, training_config, model_config = load_models(None) + all_models, config = load_models(None) else: - all_models, data_config, training_config, model_config, _, _, _ = run_pipeline(project_name='price-prediction', with_wandb = False, sweep = False, get_config=get_config) + all_models, config, _, _, _ = run_pipeline(project_name='price-prediction', with_wandb = False, sweep = False, get_config=get_config) - configs = dict(model_config=model_config, training_config=training_config, data_config=data_config) - __inference(configs, all_models.primary, all_models.secondary) + __inference(config, all_models.primary, all_models.secondary) -def __inference(configs: dict, primary_models: Optional[Reporting.Training_Step], secondary_models: Optional[Reporting.Training_Step]): +def __inference(config: Config, primary_models: Optional[Reporting.Training_Step], secondary_models: Optional[Reporting.Training_Step]): reporting = Reporting() - asset = configs['data_config']['target_asset'] + asset = config.target_asset # 1. Load data, check for validity and process data - X, y, target_returns = load_data(**configs['data_config']) - assert check_data(X, y, configs['training_config']) == True, "Data is not valid. Cancelling Inference." + X, y, target_returns = load_data( + ) + assert check_data(X, y, config) == True, "Data is not valid. Cancelling Inference." inference_from = X.index.stop - 2 # 2. Train a Primary model with optional metalabeling for each asset - training_step_primary, current_predictions = primary_step(X, y, target_returns, configs, reporting, from_index = inference_from, preloaded_training_step = primary_models) + training_step_primary, current_predictions = primary_step(X, y, target_returns, config, reporting, from_index = inference_from, preloaded_training_step = primary_models) # 3. Train an Ensemble model with optional metalabeling for each asset if secondary_models is not None: warnings.warn("Secondary models are not specified.") - training_step_secondary = secondary_step(X, y, current_predictions, target_returns, configs, reporting, from_index = inference_from, preloaded_training_step = secondary_models) + training_step_secondary = secondary_step(X, y, current_predictions, target_returns, config, reporting, from_index = inference_from, preloaded_training_step = secondary_models) # 4. Save the models reporting.asset = Reporting.Asset(ticker=asset, primary=training_step_primary, secondary=training_step_secondary) diff --git a/run_pipeline.py b/run_pipeline.py index 438f20c..6f421b2 100644 --- a/run_pipeline.py +++ b/run_pipeline.py @@ -4,12 +4,12 @@ from typing import Callable, Optional from data_loader.load_data import load_data from data_loader.process_data import check_data -from reporting.wandb import launch_wandb, register_config_with_wandb +from reporting.wandb import launch_wandb, override_config_with_wandb_values from reporting.reporting import report_results from reporting.saving import save_models -from config.config import get_default_ensemble_config +from config.config import Config, get_default_ensemble_config, get_lightweight_ensemble_config from config.preprocess import validate_config, preprocess_config from training.training_steps import primary_step, secondary_step @@ -20,46 +20,58 @@ import ray ray.init() -def run_pipeline(project_name:str, with_wandb: bool, sweep: bool, get_config: Callable) -> tuple[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) +def run_pipeline(project_name:str, with_wandb: bool, sweep: bool, get_config: Callable) -> tuple[Reporting.Asset, Config, pd.DataFrame, pd.DataFrame, pd.DataFrame]: + wandb, config = __setup_config(project_name, with_wandb, sweep, get_config) + reporting = __run_training(config) results, all_predictions, all_probabilities, all_models = reporting.get_results() - report_results(results, all_predictions, model_config, wandb, sweep, project_name) - save_models(all_models, data_config, training_config, model_config) + report_results(results, all_predictions, config, wandb, sweep, project_name) + save_models(all_models, config) - return all_models, data_config, training_config, model_config, results, all_predictions, all_probabilities + return all_models, 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]: - model_config, training_config, data_config = get_config() +def __setup_config(project_name:str, with_wandb: bool, sweep: bool, get_config: Callable) -> tuple[Optional[object], Config]: + raw_config = get_config() wandb = None if with_wandb: - wandb = launch_wandb(project_name=project_name, default_config=dict(**model_config, **training_config, **data_config), sweep=sweep) - model_config, training_config, data_config = register_config_with_wandb(wandb, model_config, training_config, data_config) - model_config, training_config, data_config = preprocess_config(model_config, training_config, data_config) + wandb = launch_wandb(project_name=project_name, default_config=raw_config, sweep=sweep) + raw_config = override_config_with_wandb_values(wandb, raw_config) + config = preprocess_config(raw_config) - return wandb, model_config, training_config, data_config + return wandb, config -def __run_training(model_config:dict, training_config:dict, data_config:dict): +def __run_training(config: Config): - validate_config(model_config, training_config, data_config) - configs = dict(model_config=model_config, training_config=training_config, data_config=data_config) + validate_config(config) reporting = Reporting() # 1. Load data, check for validity - X, y, target_returns = load_data(**configs['data_config']) - assert check_data(X, y, configs['training_config']) == True, "Data is not valid." + X, y, target_returns = load_data( + assets = config.assets, + other_assets = config.other_assets, + exogenous_data = config.exogenous_data, + target_asset = config.target_asset, + load_non_target_asset = config.load_non_target_asset, + log_returns = config.log_returns, + forecasting_horizon = config.forecasting_horizon, + own_features = config.own_features, + other_features = config.other_features, + exogenous_features = config.exogenous_features, + index_column = config.index_column, + no_of_classes = config.no_of_classes, + ) + assert check_data(X, y, config) == True, "Data is not valid." # 2. Train a Primary model with optional metalabeling for each asset - training_step_primary, current_predictions = primary_step(X, y, target_returns, configs, reporting, from_index = None) + training_step_primary, current_predictions = primary_step(X, y, target_returns, config, reporting, from_index = None) # 3. Train an Ensemble model with optional metalabeling for each asset - training_step_secondary = secondary_step(X, y, current_predictions, target_returns, configs, reporting, from_index = None) + training_step_secondary = secondary_step(X, y, current_predictions, target_returns, config, reporting, from_index = None) # 4. Save the models - reporting.asset = Reporting.Asset(ticker= data_config['target_asset'][1], primary=training_step_primary, secondary=training_step_secondary) + reporting.asset = Reporting.Asset(ticker= config.target_asset[1], primary=training_step_primary, secondary=training_step_secondary) return reporting diff --git a/training/meta_labeling.py b/training/meta_labeling.py index dea6125..745a3b6 100644 --- a/training/meta_labeling.py +++ b/training/meta_labeling.py @@ -5,7 +5,7 @@ import pandas as pd from models.base import Model from reporting.types import Reporting from typing import Union, Optional - +from config.config import Config def train_meta_labeling_model( target_asset: str, @@ -14,9 +14,7 @@ def train_meta_labeling_model( y: pd.Series, target_returns: pd.Series, models: list[tuple[str, Model]], - data_config: dict, - model_config: dict, - training_config: dict, + config: Config, model_suffix: str, from_index: Optional[int], preloaded_models: Optional[list[tuple[str, pd.Series, list[pd.Series]]]] = None @@ -34,11 +32,11 @@ def train_meta_labeling_model( y = meta_y, target_returns = target_returns, models = models, - expanding_window = training_config['expanding_window_meta_labeling'], - sliding_window_size = training_config['sliding_window_size_meta_labeling'], - retrain_every = training_config['retrain_every'], + expanding_window = config.expanding_window_meta_labeling, + sliding_window_size = config.sliding_window_size_meta_labeling, + retrain_every = config.retrain_every, from_index = from_index, - scaler = training_config['scaler'], + scaler = config.scaler, no_of_classes = 'two', level = 'meta_labeling', print_results = False, diff --git a/training/training_steps.py b/training/training_steps.py index 7cb6ea8..b80ac92 100644 --- a/training/training_steps.py +++ b/training/training_steps.py @@ -7,33 +7,33 @@ from training.meta_labeling import train_meta_labeling_model from reporting.types import Reporting from typing import Union, Optional +from config.config import Config def primary_step( X: pd.DataFrame, y: pd.Series, target_returns: pd.Series, - configs: dict, + config: Config, reporting: Reporting, from_index: Optional[int], preloaded_training_step: Optional[Reporting.Training_Step] = 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) # 3. Train Primary models current_result, current_predictions, current_probabilities, all_models_for_single_asset = train_primary_model( - ticker_to_predict = data_config['target_asset'][1], + ticker_to_predict = config.target_asset[1], X = X, y = y, target_returns = target_returns, - models = model_config['primary_models'], - expanding_window = training_config['expanding_window_primary'], - sliding_window_size = training_config['sliding_window_size_primary'], - retrain_every = training_config['retrain_every'], + models = config.primary_models, + expanding_window = config.expanding_window_base, + sliding_window_size = config.sliding_window_size_base, + retrain_every = config.retrain_every, from_index = from_index, - scaler = training_config['scaler'], - no_of_classes = data_config['no_of_classes'], + scaler = config.scaler, + no_of_classes = config.no_of_classes, level = 'primary', print_results= True, preloaded_models = preloaded_training_step.get_base() if preloaded_training_step is not None else None @@ -42,20 +42,18 @@ def primary_step( 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: + if 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 = data_config['target_asset'][1], + target_asset = config.target_asset[1], X = X, 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', + models = config.meta_labeling_models, + config = config, from_index = from_index, preloaded_models = preloaded_training_step.get_metalabeling()[model_name] if preloaded_training_step is not None else None ) @@ -78,27 +76,27 @@ def secondary_step( y:pd.Series, current_predictions:pd.DataFrame, target_returns:pd.Series, - configs: dict, + config: Config, reporting: Reporting, from_index: Optional[int], preloaded_training_step: Optional[Reporting.Training_Step] = None, ) -> Reporting.Training_Step: training_step = Reporting.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: + if config.ensemble_model is not None: ensemble_result, ensemble_predictions, _, ensemble_models_one_asset = train_primary_model( - ticker_to_predict = data_config['target_asset'][1], + ticker_to_predict = config.target_asset[1], X = current_predictions, y = y, target_returns = target_returns, - models = [model_config['ensemble_model']], expanding_window = False, + models = [config.ensemble_model], + expanding_window = False, sliding_window_size = 1, - retrain_every = training_config['retrain_every'], + retrain_every = config.retrain_every, from_index = from_index, - scaler = training_config['scaler'], - no_of_classes = data_config['no_of_classes'], + scaler = config.scaler, + no_of_classes = config.no_of_classes, level = 'ensemble', print_results= True, preloaded_models = preloaded_training_step.get_base() if preloaded_training_step is not None else None @@ -111,19 +109,17 @@ def secondary_step( reporting.all_predictions = pd.concat([reporting.all_predictions, ensemble_predictions], axis=1) - if len(model_config['meta_labeling_models']) > 0: + if len(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 = data_config['target_asset'][1], + target_asset = config.target_asset[1], X = X, 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, + models = config.meta_labeling_models, + config = config, model_suffix = 'ensemble', from_index = from_index, preloaded_models = preloaded_training_step.get_metalabeling()[ensemble_predictions.name] if preloaded_training_step is not None else None diff --git a/utils/helpers.py b/utils/helpers.py index 4540834..4dadda0 100644 --- a/utils/helpers.py +++ b/utils/helpers.py @@ -4,6 +4,8 @@ import os import string import random from typing import Union +from config.config import Config + def get_files_from_dir(path: str) -> list[str]: return [f for f in os.listdir(path) if os.path.isfile(os.path.join(path,f)) and not f.startswith('.')] @@ -17,10 +19,10 @@ def get_first_valid_return_index(series: pd.Series) -> int: return 0 return nested_result[0] -def has_enough_samples_to_train(X: pd.DataFrame, y: pd.Series, training_config: dict) -> bool: +def has_enough_samples_to_train(X: pd.DataFrame, y: pd.Series, config: Config) -> bool: first_valid_index = get_first_valid_return_index(X.iloc[:,0]) samples_to_train = len(y) - first_valid_index - return samples_to_train > training_config['sliding_window_size_primary'] + training_config['sliding_window_size_meta_labeling'] + 100 + return samples_to_train > config.sliding_window_size_base + config.sliding_window_size_meta_labeling + 100 def flatten(list_of_lists: list) -> list: return [item for sublist in list_of_lists for item in sublist] diff --git a/utils/types.py b/utils/types.py index 0dc5699..76f3a12 100644 --- a/utils/types.py +++ b/utils/types.py @@ -8,7 +8,7 @@ Name = str FeatureExtractorConfig = tuple[Name, FeatureExtractor, list[Period]] Path = str FileName = str -DataSource = list[tuple[Path, FileName]] +DataSource = tuple[Path, FileName] DataCollection = list[DataSource] ScalerTypes = Literal['normalize', 'minmax', 'standardize', 'none'] \ No newline at end of file