mirror of
https://github.com/webclinic017/drift.git
synced 2026-07-27 18:57:55 +00:00
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 <szemereydaniel@gmail.com>
This commit is contained in:
committed by
GitHub
parent
4aefba33ea
commit
797d45d036
+9
-6
@@ -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
|
||||
|
||||
|
||||
@@ -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():
|
||||
|
||||
@@ -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)
|
||||
|
||||
|
||||
+4
-4
@@ -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)
|
||||
run_inference(preload_models=True, get_config=get_lightweight_ensemble_config)
|
||||
+4
-4
@@ -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
|
||||
|
||||
+34
-43
@@ -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
|
||||
|
||||
@@ -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)
|
||||
|
||||
+17
-14
@@ -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,
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user